diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index da06248..79b7556 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -222,7 +222,7 @@ func run(ctx context.Context, command string, opts options) error { return w.Flush() case "enrol", "enroll": - return enrol(opts) + return enrol(ctx, opts) case "version": fmt.Println(version) @@ -386,7 +386,7 @@ func runApply(ctx context.Context, opts options, d *declaration.Declaration, sou // the one-time secret together with a public key it generated itself. // // The mesh issues no identity. This machine arrives holding one; what it receives is being known. -func enrol(opts options) error { +func enrol(ctx context.Context, opts options) error { tokenText, name := &opts.token, &opts.nodeName if strings.TrimSpace(*tokenText) == "" { return errors.New("enrol --token : the token is carried to this machine by a " + @@ -429,13 +429,44 @@ func enrol(opts options) error { defer conn.Close() fmt.Println("\nthe broker presented the certificate this token pins") + conn.Close() + mine, err := identity.Generate(*name) if err != nil { return err } fmt.Printf("generated this node's identity: %s\n", mine.PublicBase64()) - return errors.New("the link is not built: this machine has verified the broker and made its " + - "identity, and there is nothing yet to present them to.\n" + - "Nothing has been saved, so this can be run again unchanged") + // What this machine can be asked to do, gathered before joining rather than after. The + // control plane cannot decide what a node should run without it, so it travels with the + // request instead of being asked for in a second round trip. + detected := profile.Detect(ctx, profile.Default(nil), opts.timeout) + reported := map[string]any{} + if raw, err := json.Marshal(detected); err == nil { + _ = json.Unmarshal(raw, &reported) + } + + reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret, + mine.Public, reported, opts.timeout) + if err != nil { + return err + } + + // 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 + + // 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 + // nobody — worse than not having joined, because nothing would look wrong. + if err := identity.Save(identityPath, mine); err != nil { + return fmt.Errorf( + "the mesh accepted this node as %q and its identity could not be saved: %w\n"+ + "That token is spent, so getting back needs a new one", reply.Node, err) + } + + fmt.Printf("\nenrolled as %s\n", reply.Node) + fmt.Printf(" identity %s\n", identityPath) + fmt.Printf(" queue %s\n", reply.Queue) + return nil } diff --git a/examples/substrate-first-node.lock b/examples/substrate-first-node.lock index 2b39f61..b5cde50 100644 --- a/examples/substrate-first-node.lock +++ b/examples/substrate-first-node.lock @@ -1,11 +1,14 @@ // substrate-first-node.lock — what a machine must be before a mesh exists. // -// Steps 0 to 4 of the bootstrap (novox/hq 03-DESIGN/01-to-be/07-the-substrate.md): a container +// Steps 0 to 5 of the bootstrap (novox/hq 03-DESIGN/01-to-be/07-the-substrate.md): a container // runtime, a store, a database per context, that context's schema, and the broker. // -// It stops before step 5 (a virtual host, a credential, a certificate) and step 6 (the control -// plane runs), because nothing consumes them yet. A bundle naming a control plane that serves -// nothing would be a bundle whose last step cannot be checked. +// It stops before step 6, where the control plane runs. +// +// The broker generates its OWN certificate, in its own image, into a volume it then mounts read +// only. Self-signed, because at this moment there is no mesh to issue one and no public name to +// obtain one for -- and it does not matter, because what a joining node checks is the fingerprint +// pinned in its token, not a chain or a name (novox/hq ADR 0004). The subject is decoration. // // PINNED BY DIGEST, and the digest is not decoration: a tag can be made to point at a different // image, and this file is applied on a machine with no mesh to ask about anything. These digests @@ -41,6 +44,7 @@ "POSTGRES_PASSWORD": "bootstrap", "PGDATA": "/var/lib/postgresql/data/pgdata" }, + "ports": ["127.0.0.1:5432:5432"], "volumes": ["mesh-store-data:/var/lib/postgresql/data"] }, { @@ -58,21 +62,40 @@ "verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw inventory"] }, { - "id": "inventory-schema", + "id": "identity-database", + "type": "action", + "in": "mesh-store", + "command": ["sh", "-c", "psql -U postgres -c 'CREATE DATABASE identity'"], + "verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw identity"] + }, + { + "id": "context-schemas", "type": "action", "command": ["docker", "run", "--rm", "--network", "container:mesh-store", "-e", "MESH_STORE_INVENTORY=postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable", - "192.0.2.250:5000/mesh-control@sha256:1c27a43c4431c2b580404e8e1768cd858009e265e80b6f1591eb6de1123fc411", + "-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", + "192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396", "migrate"], - "verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node"] + "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"] + }, + { + "id": "broker-certificate", + "type": "action", + "command": ["docker", "run", "--rm", "--entrypoint", "sh", "-v", "mesh-broker-tls:/tls", + "192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336", + "-c", "test -f /tls/tls.crt || (openssl req -x509 -newkey rsa:2048 -nodes -keyout /tls/tls.key -out /tls/tls.crt -days 3650 -subj '/CN=mesh-broker' >/dev/null 2>&1 && chmod 644 /tls/tls.crt && chmod 600 /tls/tls.key)"], + "verify": ["docker", "run", "--rm", "--entrypoint", "sh", "-v", "mesh-broker-tls:/tls", + "192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336", + "-c", "test -s /tls/tls.crt && openssl x509 -in /tls/tls.crt -noout"] }, { "id": "broker", "type": "container", "name": "mesh-broker", "image": "192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336", - "ports": ["5672:5672"], - "volumes": ["mesh-broker-data:/var/lib/lavinmq"] + "ports": ["5671:5671", "127.0.0.1:5672:5672", "127.0.0.1:15672:15672"], + "volumes": ["mesh-broker-data:/var/lib/lavinmq", "mesh-broker-tls:/tls:ro"], + "args": ["--amqps-port=5671", "--cert=/tls/tls.crt", "--key=/tls/tls.key"] }, { "id": "broker-ready", @@ -80,6 +103,24 @@ "in": "mesh-broker", "command": ["sh", "-c", "for i in $(seq 1 60); do lavinmqctl status >/dev/null 2>&1 && exit 0; sleep 1; done; exit 1"], "verify": ["lavinmqctl", "status"] + }, + { + "id": "control-plane", + "type": "container", + "name": "mesh-control", + "image": "192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396", + "network": "host", + "args": ["serve"], + "volumes": ["mesh-broker-tls:/broker-tls:ro"], + "env": { + "MESH_STORE_INVENTORY": "postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable", + "MESH_STORE_IDENTITY": "postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", + "MESH_BROKER_AMQP": "amqp://guest:guest@127.0.0.1:5672/", + "MESH_BROKER_MANAGEMENT": "http://guest:guest@127.0.0.1:15672", + "MESH_BROKER_ADDRESS": "192.0.2.10:5671", + "MESH_BROKER_CERTIFICATE": "/broker-tls/tls.crt" + } } + ] } diff --git a/go.mod b/go.mod index 5e7610d..8ae86a4 100644 --- a/go.mod +++ b/go.mod @@ -1,3 +1,5 @@ module github.com/novox/mesh-host go 1.24 + +require github.com/rabbitmq/amqp091-go v1.14.0 // indirect diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..c9d50b5 --- /dev/null +++ b/go.sum @@ -0,0 +1,2 @@ +github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= +github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 3d564e7..876a107 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -584,8 +584,12 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner) ( } } - args := []string{"run", "--detach", "--name", r.Name, "--restart", "unless-stopped", - "--label", specLabel + "=" + want, "--label", idLabel + "=" + r.ID} + args := []string{"run", "--detach", "--name", r.Name, "--restart", "unless-stopped"} + if r.Network != "" { + args = append(args, "--network", r.Network) + } + args = append(args, + "--label", specLabel+"="+want, "--label", idLabel+"="+r.ID) for _, k := range sortedKeys(r.Env) { args = append(args, "--env", k+"="+r.Env[k]) } diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index 486d9e4..2a504db 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -168,6 +168,13 @@ type Container struct { Ports []string `json:"ports,omitempty"` Volumes []string `json:"volumes,omitempty"` Args []string `json:"args,omitempty"` + // Network is the container's network, passed to the runtime unchanged. + // + // Needed because the control plane must reach the store and the broker on the machine it was + // raised on, before there is any mesh to arrange that. The alternative was publishing ports + // and guessing an address that works from inside a container, which is the same thing with a + // worse failure mode. + Network string `json:"network,omitempty"` } func (c *Container) Identity() string { return c.ID } diff --git a/internal/link/enrol.go b/internal/link/enrol.go new file mode 100644 index 0000000..9ef73ff --- /dev/null +++ b/internal/link/enrol.go @@ -0,0 +1,149 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net/url" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The wire format shared with the control plane, which defines it separately because this binary +// requires nothing present and does not import it. A test on each side asserts the field names. +const ( + Exchange = "mesh" + KeyEnrol = "enrol" +) + +// QueueFor is the queue this node consumes from — the only one its account may read. +func QueueFor(node string) string { return "node." + node } + +// EnrolRequest is what this node says when joining. +type EnrolRequest struct { + Node string `json:"node"` + Secret string `json:"secret"` + PublicKey []byte `json:"public_key"` + Profile map[string]any `json:"profile,omitempty"` +} + +// EnrolReply is what the mesh says back. +type EnrolReply struct { + Accepted bool `json:"accepted"` + Node string `json:"node,omitempty"` + Queue string `json:"queue,omitempty"` + Refusal string `json:"refusal,omitempty"` +} + +// ErrRefused is what a node gets when the mesh will not have it. +var ErrRefused = errors.New("the mesh refused this enrolment") + +// Enrol presents this node's key and its one-time secret, and waits to be told it is known. +// +// The broker has already authenticated this connection: the account was created when the token +// was issued and the secret is its password. So this is not how the node gets in — it is what it +// says once it is in, and the secret travels again because the control plane must not have to ask +// the broker who connected. +func Enrol(ctx context.Context, address, pin, node, secret string, public []byte, + profile map[string]any, timeout time.Duration) (EnrolReply, error) { + + config, err := PinnedConfig(pin) + if err != nil { + return EnrolReply{}, err + } + + // The account name is the node's, and the password is the token's secret. Escaped because a + // name or secret containing a colon or an at-sign would otherwise change which host this + // connects to — a credential silently redirecting a connection is the worst shape this could + // take. + dsn := fmt.Sprintf("amqps://%s:%s@%s/", + url.QueryEscape(node), url.QueryEscape(secret), address) + + conn, err := amqp.DialConfig(dsn, amqp.Config{ + TLSClientConfig: config, + Dial: amqp.DefaultDial(timeout), + }) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return EnrolReply{}, err + } + // Not quoted back: the DSN carries the one-time secret. + return EnrolReply{}, fmt.Errorf("cannot reach the broker at %s as %s: %w", address, node, err) + } + defer conn.Close() + + channel, err := conn.Channel() + if err != nil { + return EnrolReply{}, err + } + defer channel.Close() + + // This node's own queue, which its account is scoped to and nothing else may read. + queue, err := channel.QueueDeclare(QueueFor(node), true, false, false, false, nil) + if err != nil { + return EnrolReply{}, fmt.Errorf( + "cannot declare this node's queue %s: %w", QueueFor(node), err) + } + + replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil) + if err != nil { + return EnrolReply{}, err + } + + request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, Profile: profile} + body, err := json.Marshal(request) + if err != nil { + return EnrolReply{}, err + } + + correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano()) + publish, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false, + amqp.Publishing{ + ContentType: "application/json", + CorrelationId: correlation, + ReplyTo: queue.Name, + Body: body, + }); err != nil { + return EnrolReply{}, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err) + } + + // Waited for rather than assumed. A published message that nothing answers means the control + // plane is not running, and a node that carried on regardless would believe it had joined a + // mesh that has never heard of it. + deadline := time.NewTimer(timeout) + defer deadline.Stop() + closed := conn.NotifyClose(make(chan *amqp.Error, 1)) + + for { + select { + case <-ctx.Done(): + return EnrolReply{}, ctx.Err() + case reason := <-closed: + return EnrolReply{}, fmt.Errorf("the broker closed the connection: %v", reason) + case <-deadline.C: + return EnrolReply{}, fmt.Errorf( + "the broker accepted this node's connection and nothing answered within %s. The "+ + "mesh's broker is running and its control plane is not", timeout) + case delivery, ok := <-replies: + if !ok { + return EnrolReply{}, errors.New("the broker stopped delivering") + } + // Anything else on this queue is not the answer to this question. + if delivery.CorrelationId != correlation { + continue + } + var reply EnrolReply + if err := json.Unmarshal(delivery.Body, &reply); err != nil { + return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) + } + if !reply.Accepted { + return reply, fmt.Errorf("%w: %s", ErrRefused, reply.Refusal) + } + return reply, nil + } + } +}