diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index c5938a6..cb9b4e7 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -70,6 +70,8 @@ func run() error { return brokerCommand(args[1:]) case "serve": return serve(ctx) + case "declare": + return declare(ctx, args[1:]) case "version": fmt.Println(version) return nil @@ -93,6 +95,7 @@ func usage() { identity show this control plane's signing key broker show where the broker is, and what to expect there serve consume what nodes say, and answer + declare send a node a signed declaration version what this binary is Each context reaches its own store through its own credential (novox/hq ADR 0008), named @@ -389,7 +392,19 @@ func serve(ctx context.Context) error { return err } - server, err := link.Connect(link.Enrolment{Inventory: inv, Identity: ident, Broker: management}) + // Where the broker is and what to expect there, so a node can be told how to come back + // without a person and a new token. + known, err := broker.FromEnvironment() + if err != nil && !errors.Is(err, broker.ErrNotConfigured) { + return err + } + if errors.Is(err, broker.ErrNotConfigured) { + fmt.Printf("no broker address configured, so enrolled nodes will not be told how to "+ + "reconnect. Set %s and %s.\n", broker.AddressVar, broker.CertificateVar) + } + + server, err := link.Connect(link.Enrolment{ + Inventory: inv, Identity: ident, Management: management, Broker: known}) if err != nil { return err } @@ -397,3 +412,50 @@ func serve(ctx context.Context) error { return server.Serve(ctx) } + +// declare sends one node a declaration, signed. +// +// Signed here rather than trusted from the broker: a node connects to the broker and takes +// instruction from the control plane behind it, and those are two identities. If a node believed +// whatever arrived on its queue, a compromised broker could forge declarations — and since the +// host applies whatever the link delivers, that is the whole machine (novox/hq ADR 0004). +func declare(ctx context.Context, args []string) error { + if len(args) != 2 { + return errors.New("declare ") + } + node, path := args[0], args[1] + + raw, err := os.ReadFile(path) + if err != nil { + return err + } + + ident, err := openIdentity(ctx) + if err != nil { + return err + } + defer ident.Close() + + // The node has to exist before it can be told anything. Publishing to a queue nobody consumes + // would sit there looking like success. + inv, err := openInventory(ctx) + if err != nil { + return err + } + defer inv.Close() + if _, err := inv.NodeByName(ctx, node); err != nil { + return err + } + + server, err := link.Connect(nil) + if err != nil { + return err + } + defer server.Close() + + if err := link.Declare(ctx, server.Channel(), ident, node, raw, 15*time.Second); err != nil { + return err + } + fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) + return nil +} diff --git a/internal/link/declare.go b/internal/link/declare.go new file mode 100644 index 0000000..9f6fd7e --- /dev/null +++ b/internal/link/declare.go @@ -0,0 +1,53 @@ +package link + +import ( + "context" + "encoding/json" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// Signer is whatever holds the control plane's signing key. +type Signer interface { + Sign(ctx context.Context, message []byte) ([]byte, error) +} + +// Declare sends a node what it should be, signed. +// +// The signature is over the declaration exactly as it is published — the same bytes the node +// verifies. Anything that re-encoded between here and there would produce a signature over +// something else, and the node would refuse a declaration that was genuinely the mesh's. +// +// Published to the node's own queue, which its account alone may read. +func Declare(ctx context.Context, channel *amqp.Channel, signer Signer, node string, + declaration []byte, timeout time.Duration) error { + + if !json.Valid(declaration) { + return fmt.Errorf("refusing to send %s something that is not a declaration", node) + } + + signature, err := signer.Sign(ctx, declaration) + if err != nil { + return fmt.Errorf("cannot sign a declaration for %s: %w", node, err) + } + + body, err := json.Marshal(Signed{Declaration: declaration, Signature: signature}) + if err != nil { + return err + } + + publish, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + + // Published to the queue directly rather than through the exchange: a declaration is for one + // node, and routing it by name through a shared exchange would mean a binding per node that + // nothing removes when a node is retired. + return channel.PublishWithContext(publish, "", QueueFor(node), false, false, + amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + Body: body, + }) +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index b58f93e..4490f24 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -3,6 +3,8 @@ package link import ( "context" "crypto/ed25519" + "crypto/rand" + "encoding/base64" "errors" "fmt" @@ -18,9 +20,10 @@ import ( // — it holds both grants and asks each for its part, which is what the process running them is // for. type Enrolment struct { - Inventory *inventory.Inventory - Identity *identity.Identity - Broker *broker.Management + Inventory *inventory.Inventory + Identity *identity.Identity + Management *broker.Management + Broker broker.Broker } // Enrol spends the token and records what the node presented. @@ -30,35 +33,73 @@ type Enrolment struct { // winner. Only then is a key recorded — because recording a key for a node whose token turned out // to be spent would leave the mesh believing a machine that never had the right to join. func (e Enrolment) Enrol(ctx context.Context, secret string, public ed25519.PublicKey, - profile map[string]any) (string, error) { + profile map[string]any) (EnrolReply, error) { if len(public) != ed25519.PublicKeySize { - return "", fmt.Errorf("a node presented a %d-byte key, and an identity is %d", + return EnrolReply{}, fmt.Errorf("a node presented a %d-byte key, and an identity is %d", len(public), ed25519.PublicKeySize) } node, err := e.Inventory.Redeem(ctx, secret) if err != nil { - return "", err + return EnrolReply{}, err } // From here the token is gone whatever happens next, so anything that fails leaves a node // record with no live key — which is visible and fixable with a new token, where a spent // token believed to be unspent is neither. if _, err := e.Identity.RecordNodeKey(ctx, node.ID, public); err != nil { - return "", fmt.Errorf("the token was spent and the key could not be recorded, so %s has "+ - "no identity and needs a new token: %w", node.Name, err) + return EnrolReply{}, fmt.Errorf( + "the token was spent and the key could not be recorded, so %s has no identity and "+ + "needs a new token: %w", node.Name, err) + } + + key, err := e.Identity.Active(ctx) + if err != nil { + return EnrolReply{}, err + } + + reply := EnrolReply{ + Accepted: true, + Node: node.Name, + Queue: QueueFor(node.Name), + Broker: e.Broker.Address, + Fingerprint: e.Broker.Fingerprint, + Signer: key.Public, + } + + // The token's secret was the broker password up to this moment, which is what let this + // connection exist at all. It is replaced now, so the one-time thing stays one-time and the + // credential the node keeps for years is not the one that was pasted into a terminal. + if e.Management != nil { + password, err := freshPassword() + if err != nil { + return EnrolReply{}, err + } + if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil { + return EnrolReply{}, fmt.Errorf( + "the token was spent and %s's broker password could not be replaced: %w", + node.Name, err) + } + reply.Password = password } if profile != nil { - if err := e.Inventory.RecordProfile(ctx, node.ID, profile); err != nil { - // Not fatal. The profile is what the control plane needs in order to decide what this - // machine should run, and it is reported again on every connection — so losing it - // here costs a decision that can be made later, not the enrolment. - return node.Name, nil - } + // Not fatal if it fails. The profile is what the control plane needs in order to decide + // what this machine should run, and it is reported again on every connection — so losing + // it here costs a decision that can be made later, not the enrolment. + _ = e.Inventory.RecordProfile(ctx, node.ID, profile) } - return node.Name, nil + return reply, nil +} + +// freshPassword is the node's own broker credential from enrolment onward. +func freshPassword() (string, error) { + raw := make([]byte, 32) + if _, err := rand.Read(raw); err != nil { + return "", fmt.Errorf("cannot generate a broker password: %w", err) + } + return base64.RawURLEncoding.EncodeToString(raw), nil } var _ Enroller = Enrolment{} diff --git a/internal/link/protocol.go b/internal/link/protocol.go index 337253a..f3f0f1b 100644 --- a/internal/link/protocol.go +++ b/internal/link/protocol.go @@ -15,7 +15,8 @@ const ControlQueue = "control" // Routing keys. A node may publish these; it may not publish anything else, because its broker // account is scoped to this exchange and its own queue. const ( - KeyEnrol = "enrol" + KeyEnrol = "enrol" + KeyReport = "report" ) // QueueFor is the queue a node consumes from — the only one it may read. @@ -44,6 +45,24 @@ type EnrolRequest struct { Profile map[string]any `json:"profile,omitempty"` } +// Signed is a declaration and the signature over it. +// +// The signature is over Declaration exactly as it will arrive, bytes unchanged — a node verifies +// what it received rather than what it re-encoded, because any difference in key order or spacing +// would break a signature over the same meaning. +type Signed struct { + Declaration []byte `json:"declaration"` + Signature []byte `json:"signature"` +} + +// Report is what a node states after applying. It states; the owning context writes. +type Report struct { + Node string `json:"node"` + Applied []string `json:"applied,omitempty"` + Failed map[string]string `json:"failed,omitempty"` + Refused string `json:"refused,omitempty"` +} + // EnrolReply is what the mesh says back. type EnrolReply struct { // Accepted says whether the node is now known. @@ -57,6 +76,17 @@ type EnrolReply struct { // Queue is where this node listens from now on. Queue string `json:"queue,omitempty"` + // Password is this node's own broker account from now on, replacing the token's secret. + // A credential that lives for as long as the node should not be the same string as one that + // was meant to be used once. + Password string `json:"password,omitempty"` + + // Fingerprint and Signer are what the node keeps so it can reconnect and keep verifying + // without a person and a new token. + Fingerprint string `json:"fingerprint,omitempty"` + Signer []byte `json:"signer,omitempty"` + Broker string `json:"broker,omitempty"` + // Refusal says why not, in words for a person. Deliberately the same for every reason a // token can fail — unknown, spent, expired — so that guessing learns nothing. Refusal string `json:"refusal,omitempty"` diff --git a/internal/link/protocol_test.go b/internal/link/protocol_test.go new file mode 100644 index 0000000..6c0ff37 --- /dev/null +++ b/internal/link/protocol_test.go @@ -0,0 +1,39 @@ +package link + +import ( + "encoding/json" + "testing" +) + +func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) { + // The contract with the host, which defines these separately because it requires nothing + // present and does not import this. A matching test lives there; rename a field on either + // side and both fail, rather than the mismatch surfacing on a real machine. + 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"}}, + {EnrolRequest{Node: "n", Secret: "s", PublicKey: []byte("k")}, + []string{"node", "secret", "public_key"}}, + } { + 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 host reads 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/serve.go b/internal/link/serve.go index 43c6afa..7a2ee86 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -24,7 +24,7 @@ const AMQPVar = "MESH_BROKER_AMQP" type Enroller interface { // Enrol spends the token, records the key, and reports the node's name. The error is // returned to the node as a refusal; it must be the same for every reason a token can fail. - Enrol(ctx context.Context, secret string, public ed25519.PublicKey, profile map[string]any) (string, error) + Enrol(ctx context.Context, secret string, public ed25519.PublicKey, profile map[string]any) (EnrolReply, error) } // Server consumes what nodes say. @@ -77,6 +77,17 @@ func Connect(enroller Enroller) (*Server, error) { log: log.New(os.Stdout, "", log.LstdFlags)}, nil } +func (s *Server) bindOrClose(channel *amqp.Channel, conn *amqp.Connection, key string) error { + if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil { + conn.Close() + return fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err) + } + return nil +} + +// Channel is the control plane's channel, for sending declarations. +func (s *Server) Channel() *amqp.Channel { return s.channel } + func (s *Server) Close() { if s.channel != nil { _ = s.channel.Close() @@ -106,7 +117,7 @@ func (s *Server) Serve(ctx context.Context) error { } closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) - s.log.Printf("consuming %s, bound to %s/%s", ControlQueue, Exchange, KeyEnrol) + s.log.Printf("consuming %s, bound to %s/{%s,%s}", ControlQueue, Exchange, KeyEnrol, KeyReport) for { select { @@ -130,6 +141,8 @@ func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { switch delivery.RoutingKey { case KeyEnrol: s.handleEnrol(ctx, delivery) + case KeyReport: + s.handleReport(delivery) default: // Rejected without requeue: a message nothing understands will not be understood on the // next attempt either, and requeuing it would spin. @@ -138,6 +151,29 @@ func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { } } +// handleReport records what a node says it did. +// +// A node states; nothing here writes anything the node claimed about itself beyond that it was +// heard from. What it applied is its own account of its own machine, and the mesh keeps the last +// one as a copy for recovery rather than as a source (novox/hq 09-the-node-lifecycle). +func (s *Server) handleReport(delivery amqp.Delivery) { + var report Report + if err := json.Unmarshal(delivery.Body, &report); err != nil { + s.log.Printf("a report could not be read: %v", err) + _ = delivery.Reject(false) + return + } + switch { + case report.Refused != "": + s.log.Printf("%s refused a declaration: %s", report.Node, report.Refused) + case len(report.Failed) > 0: + s.log.Printf("%s applied %d and failed: %v", report.Node, len(report.Applied), report.Failed) + default: + s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied)) + } + _ = delivery.Ack(false) +} + func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { reply := EnrolReply{Refusal: "that token cannot be used"} @@ -145,14 +181,14 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { if err := json.Unmarshal(delivery.Body, &request); err != nil { s.log.Printf("an enrolment request could not be read: %v", err) } else { - name, err := s.enroller.Enrol(ctx, request.Secret, request.PublicKey, request.Profile) + accepted, err := s.enroller.Enrol(ctx, request.Secret, request.PublicKey, request.Profile) if err != nil { // Logged in full here, where an operator can see it; sent back as one refusal, so // that somebody guessing learns nothing from which reason came back. s.log.Printf("refusing enrolment for %q: %v", request.Node, err) } else { - reply = EnrolReply{Accepted: true, Node: name, Queue: QueueFor(name)} - s.log.Printf("enrolled %s", name) + reply = accepted + s.log.Printf("enrolled %s", accepted.Node) } }