diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 79b7556..b8e738f 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -224,6 +224,9 @@ func run(ctx context.Context, command string, opts options) error { case "enrol", "enroll": return enrol(ctx, opts) + case "run": + return runLink(ctx, opts) + case "version": fmt.Println(version) return nil @@ -455,6 +458,20 @@ func enrol(ctx context.Context, opts options) error { // The mesh's name for this node wins over what the machine called itself: the token was // issued for a node record, and that record is what the identity binds to. mine.Node = reply.Node + mine.Membership = identity.Membership{ + Broker: firstNonEmpty(reply.Broker, token.Broker), + Fingerprint: firstNonEmpty(reply.Fingerprint, token.Fingerprint), + Signer: firstNonEmpty2(reply.Signer, token.Signer), + Password: reply.Password, + } + if mine.Membership.Password == "" { + // The mesh did not replace the token's secret, so it is still this node's broker + // password. Said rather than silently kept: a one-time secret living on as a credential + // is worth knowing about. + mine.Membership.Password = token.Secret + fmt.Println("\nnote: the mesh issued no separate broker password, so the token's secret " + + "remains this node's credential") + } // Saved only now, and only once the mesh has said it knows this node. A node holding an // identity the mesh has never recorded would believe it had joined and be believed by @@ -465,8 +482,108 @@ func enrol(ctx context.Context, opts options) error { "That token is spent, so getting back needs a new one", reply.Node, err) } + if !mine.Membership.Joined() { + return fmt.Errorf( + "the mesh accepted this node as %q but did not say how to reach it again, so this "+ + "identity could not be used after a restart. Nothing was saved", reply.Node) + } + fmt.Printf("\nenrolled as %s\n", reply.Node) fmt.Printf(" identity %s\n", identityPath) fmt.Printf(" queue %s\n", reply.Queue) return nil } + +func firstNonEmpty(values ...string) string { + for _, v := range values { + if strings.TrimSpace(v) != "" { + return v + } + } + return "" +} + +func firstNonEmpty2(values ...[]byte) []byte { + for _, v := range values { + if len(v) > 0 { + return v + } + } + return nil +} + +// runLink holds this node's link to the mesh open, applying what arrives. +// +// One outbound connection and nothing listening. While it is up this node is enrolled; while it +// is down it is disconnected, which is an ordinary situation rather than a failure — the machine +// keeps running whatever it was last told, from its own store. +func runLink(ctx context.Context, opts options) error { + mine, err := identity.Load(identity.Path(opts.state)) + if errors.Is(err, identity.ErrNoIdentity) { + return errors.New("this machine has not joined a mesh. Enrol it first: " + + "mesh-host enrol --token --name ") + } + if err != nil { + return err + } + + fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker) + + apply := func(ctx context.Context, raw []byte) link.Report { + return applyDeclared(ctx, opts, raw) + } + return link.Run(ctx, link.Membership{ + Node: mine.Node, + Broker: mine.Membership.Broker, + Fingerprint: mine.Membership.Fingerprint, + Password: mine.Membership.Password, + Signer: mine.Membership.Signer, + }, apply, opts.timeout) +} + +// applyDeclared applies a declaration that has already been proved to come from the mesh. +// +// Signature checking happens before this is called, in the link. By the time anything here runs, +// the question "is this from the mesh I joined" is settled — which is why this can treat the +// bytes as instructions. +func applyDeclared(ctx context.Context, opts options, raw []byte) link.Report { + declared, err := declaration.Parse(raw) + if err != nil { + return link.Report{Refused: err.Error()} + } + + built, err := system.For(builtFor) + if err != nil { + return link.Report{Refused: err.Error()} + } + if err := system.Check(built, declared); err != nil { + return link.Report{Refused: err.Error()} + } + + known, err := store.Load(opts.state) + if err != nil { + return link.Report{Refused: err.Error()} + } + + if err := built.Confirm(ctx, apply.ExecRunner); err != nil { + return link.Report{Refused: err.Error()} + } + + outcome, updated, applyErr := apply.Apply(ctx, built, declared, known, apply.ExecRunner, nil) + + // Saved whichever way it went. Recording only on success would lose the footprint of a + // failed apply, and that footprint is on the machine either way. + if saveErr := store.Save(opts.state, updated); saveErr != nil { + return link.Report{Refused: "applied, and the node's state could not be saved: " + + saveErr.Error()} + } + + report := link.Report{} + for _, change := range outcome.Outcomes { + report.Applied = append(report.Applied, change.ID) + } + if applyErr != nil { + report.Failed = map[string]string{"apply": applyErr.Error()} + } + return report +} diff --git a/examples/substrate-first-node.lock b/examples/substrate-first-node.lock index b5cde50..cd85237 100644 --- a/examples/substrate-first-node.lock +++ b/examples/substrate-first-node.lock @@ -74,7 +74,7 @@ "command": ["docker", "run", "--rm", "--network", "container:mesh-store", "-e", "MESH_STORE_INVENTORY=postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable", "-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", - "192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396", + "192.0.2.250:5000/mesh-control@sha256:1dfcf6a879e16e671d4d6459271fd2630e30560afa1211508ab3994494786893", "migrate"], "verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node && docker exec mesh-store psql -U postgres -d identity -tAc \"select to_regclass('public.signing_key')\" | grep -qx signing_key"] }, @@ -108,7 +108,7 @@ "id": "control-plane", "type": "container", "name": "mesh-control", - "image": "192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396", + "image": "192.0.2.250:5000/mesh-control@sha256:1dfcf6a879e16e671d4d6459271fd2630e30560afa1211508ab3994494786893", "network": "host", "args": ["serve"], "volumes": ["mesh-broker-tls:/broker-tls:ro"], diff --git a/internal/identity/identity.go b/internal/identity/identity.go index d648e94..29c9ae2 100644 --- a/internal/identity/identity.go +++ b/internal/identity/identity.go @@ -28,14 +28,51 @@ func Path(statePath string) string { return filepath.Join(filepath.Dir(statePath), FileName) } -// Identity is this node's own keypair, and the name the mesh knows it by. +// Identity is this node's own keypair, the name the mesh knows it by, and what it needs to get +// back to that mesh without a person. +// +// The keypair is the node's own and was never anybody else's. Everything under Membership was +// learned at enrolment and is the mesh's answer rather than this machine's — kept here because a +// node that could not reconnect after a restart without a new token would make disconnection a +// crisis instead of an ordinary situation (novox/hq ADR 0004). type Identity struct { - // Node is the name in the mesh's records. Learned at enrolment, from the mesh — it is the one - // thing here the node does not decide for itself. + // Node is the name in the mesh's records. Learned at enrolment — the one thing here the node + // does not decide for itself. Node string `json:"node"` Public []byte `json:"public"` Private []byte `json:"private"` + + Membership Membership `json:"membership"` +} + +// Membership is how this node reaches the mesh it belongs to, and who it believes. +type Membership struct { + // Broker is an address, not a name: there is no resolution before the link. + Broker string `json:"broker"` + + // Fingerprint is checked before anything is sent, on every connection and not only the first. + Fingerprint string `json:"fingerprint"` + + // Signer is the control plane's public signing key. Kept because **each declaration is + // verified by its signature, every time** (novox/hq ADR 0004) — a node that only pinned the + // broker would make the control plane's authority transitive, and a compromised broker could + // then forge declarations, which is the whole machine. + Signer []byte `json:"signer"` + + // Password is this node's own broker account, issued at enrolment and belonging to it alone. + // Not the token's secret: that is spent, and a credential that lives for ever should not be + // the same string as one that was meant to be used once. + Password string `json:"password"` +} + +// Queue is where this node listens. Its account may read this and nothing else. +func (i Identity) Queue() string { return "node." + i.Node } + +// Joined reports whether this identity can reach its mesh unaided. +func (m Membership) Joined() bool { + return m.Broker != "" && m.Fingerprint != "" && len(m.Signer) == ed25519.PublicKeySize && + m.Password != "" } // ErrNoIdentity means this machine has not enrolled. @@ -92,6 +129,13 @@ func Load(path string) (Identity, error) { if strings.TrimSpace(i.Node) == "" { return Identity{}, fmt.Errorf("the identity at %s names no node", path) } + // Checked here rather than at the moment it is used, which would be while trying to + // reconnect on a machine nobody is watching. + if !i.Membership.Joined() { + return Identity{}, fmt.Errorf( + "the identity at %s does not say how to reach its mesh, so this node cannot "+ + "reconnect. It needs a new token", path) + } return i, nil } diff --git a/internal/link/enrol.go b/internal/link/enrol.go index 9ef73ff..ed18f2e 100644 --- a/internal/link/enrol.go +++ b/internal/link/enrol.go @@ -34,7 +34,16 @@ type EnrolReply struct { Accepted bool `json:"accepted"` Node string `json:"node,omitempty"` Queue string `json:"queue,omitempty"` - Refusal string `json:"refusal,omitempty"` + + // What this node keeps so it can come back on its own. Without these a restart would need a + // person with a new token, which would make disconnection a crisis rather than the ordinary + // situation novox/hq ADR 0004 says it is. + Password string `json:"password,omitempty"` + Broker string `json:"broker,omitempty"` + Fingerprint string `json:"fingerprint,omitempty"` + Signer []byte `json:"signer,omitempty"` + + Refusal string `json:"refusal,omitempty"` } // ErrRefused is what a node gets when the mesh will not have it. diff --git a/internal/link/messages.go b/internal/link/messages.go new file mode 100644 index 0000000..46e20c2 --- /dev/null +++ b/internal/link/messages.go @@ -0,0 +1,45 @@ +package link + +// The wire formats shared with the control plane, which defines them separately because this +// binary requires nothing present and does not import it. A test on each side asserts the field +// names, so a rename breaks both at once rather than on a real machine months later. + +// Routing keys a node may publish. Its broker account is scoped to this exchange and its own +// queue, so it can say these things and nothing else. +const ( + KeyReport = "report" +) + +// Signed is a declaration and the signature over it. +// +// novox/hq ADR 0004: the transport is verified once at connect, and **each declaration is +// verified by its signature, every time**. The two are different questions — a node connects to +// the broker and takes instruction from the control plane behind it, and pinning only the first +// would make the second transitive. +// +// The signature is over Declaration exactly as it arrived, bytes unchanged. Re-encoding before +// verifying would mean checking a signature over something other than what was sent, and any +// difference in key order or spacing would break it — so the raw message is what is signed and +// what is checked. +type Signed struct { + Declaration []byte `json:"declaration"` + Signature []byte `json:"signature"` +} + +// Report is what a node says after applying, and it is a statement rather than a write. +// +// A node states; the context that owns the data writes (novox/hq ADR 0006). The difference is the +// security boundary: something that can write cannot be prevented from writing anything, and +// something that can only state has its blast radius bounded by what this struct can say. +type Report struct { + Node string `json:"node"` + + // Applied is what this machine now owns, by resource id. + Applied []string `json:"applied,omitempty"` + + // Failed says what could not be applied, and why, in words for a person. + Failed map[string]string `json:"failed,omitempty"` + + // Refused is set when the declaration was rejected whole rather than applied in part. + Refused string `json:"refused,omitempty"` +} diff --git a/internal/link/messages_test.go b/internal/link/messages_test.go new file mode 100644 index 0000000..4be61eb --- /dev/null +++ b/internal/link/messages_test.go @@ -0,0 +1,136 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "testing" +) + +// verified runs what Run does to a delivery body, without a broker: unmarshal, check the +// signature, and only then apply. Isolating it keeps this test about the check rather than about +// AMQP, which is tested against a real broker in the lab. +func verified(t *testing.T, signer ed25519.PublicKey, body []byte) (Report, bool) { + t.Helper() + applied := false + report := handleBody(context.Background(), Membership{Node: "anchor", Signer: signer}, body, + func(context.Context, []byte) Report { + applied = true + return Report{Applied: []string{"something"}} + }) + return report, applied +} + +func signedBody(t *testing.T, private ed25519.PrivateKey, declaration string) []byte { + t.Helper() + raw, err := json.Marshal(Signed{ + Declaration: []byte(declaration), + Signature: ed25519.Sign(private, []byte(declaration)), + }) + if err != nil { + t.Fatal(err) + } + return raw +} + +func TestTheMeshsOwnDeclarationIsApplied(t *testing.T) { + public, private, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + report, applied := verified(t, public, signedBody(t, private, `{"declaration":1}`)) + if !applied { + t.Fatalf("a declaration the mesh signed was not applied: %s", report.Refused) + } +} + +func TestAForgedDeclarationIsNeverApplied(t *testing.T) { + // The check that stands between "the mesh changes this machine" and "anybody does". The host + // applies whatever the link delivers, so a forged declaration is the whole machine. + public, _, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + _, other, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + + report, applied := verified(t, public, signedBody(t, other, `{"declaration":1}`)) + if applied { + t.Fatal("a declaration signed by another key was applied") + } + if report.Refused != ErrForged.Error() { + t.Errorf("refused, but not as a forgery: %q", report.Refused) + } +} + +func TestATamperedDeclarationIsNeverApplied(t *testing.T) { + // A broker that changed the declaration in flight, keeping the signature. This is what makes + // pinning the transport insufficient on its own. + public, private, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + raw, err := json.Marshal(Signed{ + Declaration: []byte(`{"declaration":1,"resources":["something else entirely"]}`), + Signature: ed25519.Sign(private, []byte(`{"declaration":1}`)), + }) + if err != nil { + t.Fatal(err) + } + + report, applied := verified(t, public, raw) + if applied { + t.Fatal("a declaration altered after signing was applied") + } + if report.Refused != ErrForged.Error() { + t.Errorf("refused, but not as a forgery: %q", report.Refused) + } +} + +func TestAMalformedMessageIsToldApartFromAForgery(t *testing.T) { + // novox/hq ADR 0004 requires these to be distinguishable: one means somebody is trying, the + // other means something is broken, and they need different responses from a person. + public, _, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatal(err) + } + report, applied := verified(t, public, []byte("this is not a message")) + if applied { + t.Fatal("something unparseable was applied") + } + if report.Refused == ErrForged.Error() { + t.Error("a malformed message was reported as a forgery; those must be distinguishable") + } +} + +func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) { + // The contract with the control plane, which defines these separately. A matching test lives + // there; rename a field on either side and both fail. + for _, c := range []struct { + value any + expect []string + }{ + {Signed{Declaration: []byte("{}"), Signature: []byte("x")}, []string{"declaration", "signature"}}, + {Report{Node: "n", Applied: []string{"a"}, Failed: map[string]string{"k": "v"}, Refused: "r"}, + []string{"node", "applied", "failed", "refused"}}, + } { + raw, err := json.Marshal(c.value) + if err != nil { + t.Fatal(err) + } + var fields map[string]any + if err := json.Unmarshal(raw, &fields); err != nil { + t.Fatal(err) + } + for _, want := range c.expect { + if _, ok := fields[want]; !ok { + t.Errorf("%T has no %q field; the control plane uses that name", c.value, want) + } + } + if len(fields) != len(c.expect) { + t.Errorf("%T has %d fields, expected %d: %v", c.value, len(fields), len(c.expect), fields) + } + } +} diff --git a/internal/link/run.go b/internal/link/run.go new file mode 100644 index 0000000..f190975 --- /dev/null +++ b/internal/link/run.go @@ -0,0 +1,140 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "errors" + "fmt" + "net/url" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// ErrForged is what a node returns for a declaration whose signature is not the mesh's. +// +// Its own error, and it must never be confused with a malformed message. novox/hq ADR 0004 +// requires a host to tell *this is not from the mesh I joined* apart from *this is malformed*: +// the first means somebody is trying, the second means something is broken. +var ErrForged = errors.New("this declaration was not signed by the mesh this node joined") + +// Membership is what a node needs to reach its mesh again, held by the caller. +type Membership struct { + Node string + Broker string + Fingerprint string + Password string + Signer ed25519.PublicKey +} + +// Applier is what the host does with a declaration that has been proved to come from the mesh. +type Applier func(ctx context.Context, declaration []byte) Report + +// Run holds the link open, applying what arrives and reporting what happened. +// +// Outbound only, and nothing listens on this machine. The connection is the node's presence in +// the mesh: while it is up the node is enrolled, and while it is down the node is disconnected — +// which is an ordinary situation and not a failure, so this returns rather than panicking and +// leaves restarting to whatever supervises it. +func Run(ctx context.Context, m Membership, apply Applier, timeout time.Duration) error { + config, err := PinnedConfig(m.Fingerprint) + if err != nil { + return err + } + + dsn := fmt.Sprintf("amqps://%s:%s@%s/", + url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker) + conn, err := amqp.DialConfig(dsn, amqp.Config{ + TLSClientConfig: config, + Dial: amqp.DefaultDial(timeout), + // Kept short so a node that has silently lost its route notices, rather than holding a + // connection the broker forgot about and believing it is still in the mesh. + Heartbeat: 10 * time.Second, + }) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return err + } + return fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err) + } + defer conn.Close() + + channel, err := conn.Channel() + if err != nil { + return err + } + defer channel.Close() + + queue := QueueFor(m.Node) + if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { + return fmt.Errorf("cannot declare this node's queue %s: %w", queue, err) + } + + // One at a time. A declaration is applied to a machine, and applying two at once would race + // on the same filesystem — so the broker holds the next one until this one is finished, + // where it survives a restart. + if err := channel.Qos(1, 0, false); err != nil { + return err + } + + deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil) + if err != nil { + return err + } + closed := conn.NotifyClose(make(chan *amqp.Error, 1)) + + for { + select { + case <-ctx.Done(): + return nil + case reason := <-closed: + return fmt.Errorf("the link closed: %v", reason) + case delivery, ok := <-deliveries: + if !ok { + return errors.New("the broker stopped delivering") + } + report := handle(ctx, m, apply, delivery) + publishReport(ctx, channel, m, report, timeout) + // Acknowledged after the report is published. A node that dies between applying and + // reporting leaves the declaration on the broker and applies it again on return, + // which is safe because applying is reconciliation — it converges rather than + // repeating. + _ = delivery.Ack(false) + } + } +} + +func handle(ctx context.Context, m Membership, apply Applier, delivery amqp.Delivery) Report { + return handleBody(ctx, m, delivery.Body, apply) +} + +// handleBody is the whole of deciding whether to trust a message, separated from the broker so it +// can be tested as the security check it is rather than as message plumbing. +func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report { + var signed Signed + if err := json.Unmarshal(body, &signed); err != nil { + return Report{Node: m.Node, Refused: "this message is not a declaration: " + err.Error()} + } + + // Before anything is read out of it, let alone applied. The host applies whatever the link + // delivers, so this check is the difference between the mesh changing this machine and + // anybody changing it. + if !ed25519.Verify(m.Signer, signed.Declaration, signed.Signature) { + return Report{Node: m.Node, Refused: ErrForged.Error()} + } + return apply(ctx, signed.Declaration) +} + +func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report, + timeout time.Duration) { + report.Node = m.Node + body, err := json.Marshal(report) + if err != nil { + return + } + publish, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + _ = channel.PublishWithContext(publish, Exchange, KeyReport, false, false, + amqp.Publishing{ContentType: "application/json", Body: body}) +}