diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 4e8aefe..8eeb3ce 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -772,7 +772,12 @@ func enrol(ctx context.Context, opts options) error { // who knows its public key (novox/hq issue 083). proof := mine.Sign(link.EnrolProof(token.Secret, mine.Public, mine.Overlay.Public, sealing.Public, serving.Public)) - reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret, + // The token says where to go and which certificate that address must present. It says nothing + // about which bus is there, and does not need to: every token names the one the mesh runs on + // today until the rollout (novox/hq ADR 0116 step 5), and that is what an empty Transport is. + reply, err := link.Enrol(ctx, + link.Approach{Address: token.Broker, Fingerprint: token.Fingerprint}, + *name, token.Secret, mine.Public, mine.Overlay.Public, sealing.Public, serving.Public, reported, proof, found, opts.timeout) if err != nil { diff --git a/examples/foundation-first-node-nats.lock b/examples/foundation-first-node-nats.lock new file mode 100644 index 0000000..50fefd1 --- /dev/null +++ b/examples/foundation-first-node-nats.lock @@ -0,0 +1,191 @@ +// foundation-first-node-nats.lock — what a machine must be before a mesh exists, on the bus being +// built (novox/hq ADR 0106, design 25). +// +// The same twelve steps as foundation-first-node.lock, with one difference that matters: **the mesh +// writes its own user list, and at genesis there is no mesh yet to write it.** So this carries the +// first one — the controller's own account, at a well-known bootstrap password, exactly as the store +// is reached at `postgres:bootstrap` and the old bus at `guest:guest`. It is rotated with those, and +// from the controller's first composition onward the file is the controller's to write. +// +// The accounts file is its own file beside the server's configuration, because the server's own +// settings belong to whoever raises it and the users belong to the mesh (design 25 §4). Both live in +// one directory, of necessity: an include path is resolved relative to the including file's own +// directory, so a server given an absolute one looks for it underneath that directory and refuses to +// start. +// +// No `verify` on the TLS block, deliberately — that setting makes the server demand a *client* +// certificate, and nothing in the mesh presents one: a host pins this server's exact certificate and +// authenticates with a password (ADR 0004, design 25 §4). + +{ + "declaration": 1, + "resources": [ + { + "id": "container-runtime", + "type": "package", + "package": "docker" + }, + { + "id": "container-runtime-running", + "type": "service", + "unit": "docker.service", + "state": "running", + "boot": "enabled" + }, + // **A filter before anything listens** (novox/hq issue 054, ADR 0088). The store and the + // broker are adopted as modules later and so bind to every interface from the moment they + // start; the packet filter that governs who may reach them is a module too, installed a + // dozen steps later. Between the two, a control-node facing the network had its store and + // its bus open to anyone who could reach the machine. So the foundation carries a filter of + // its own — the same table the filter module will replace wholesale once it can derive one: + // drop by default, keep loopback, replies, ssh and the mesh's own ports (the bus a node + // enrols over, the registry a node pulls from), and let the container runtime's own + // networks through the forward chain so containers keep working. A published container port + // is forwarded, never input (issue 047), which is why the forward chain is where the store's + // and broker's ports are refused from outside — and a container on this machine dialling a + // port this machine publishes reaches it through the runtime's proxy, which IS input, which + // is why the bus and the registry are opened in both chains, exactly as the derived ruleset + // does. + { + "id": "base-filter-package", + "type": "package", + "package": "nftables" + }, + { + "id": "base-filter", + "type": "file", + "path": "/etc/nftables.conf", + "mode": "0644", + "content": "#!/usr/sbin/nft -f\n# the foundation's own filter, until the mesh derives one (novox/hq issue 054)\ntable inet mesh {}\ndelete table inet mesh\n\ntable inet mesh {\n\tchain input {\n\t\ttype filter hook input priority filter; policy drop;\n\t\tct state established,related accept\n\t\tct state invalid drop\n\t\tiif lo accept\n\t\ticmp type echo-request accept\n\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n\t\t# ssh, from anywhere — never closed\n\t\ttcp dport 22 accept\n\t\t# the mesh's own, from anywhere: the bus a node enrols over and a container on this machine reaches through the proxy, the registry a node pulls from\n\t\ttcp dport 5671 accept\n\t\ttcp dport 5000 accept\n\t}\n\tchain output {\n\t\ttype filter hook output priority filter; policy accept;\n\t}\n\tchain forward {\n\t\ttype filter hook forward priority filter; policy drop;\n\t\tct state established,related accept\n\t\tct state invalid drop\n\t\t# the container runtime's bridge networks, and the networks its compose files are given\n\t\tip saddr 172.16.0.0/12 accept\n\t\tip saddr 192.168.128.0/17 accept\n\t\t# the mesh's own: the bus a node enrols over, the registry a node pulls from\n\t\tct original proto-dst 5671 accept\n\t\tct original proto-dst 5000 accept\n\t}\n}\n" + }, + { + "id": "base-filter-loaded", + "type": "service", + "unit": "nftables.service", + "state": "running", + "boot": "enabled", + "restart-on": ["base-filter"] + }, + { + "id": "store", + "type": "container", + "name": "mesh-store", + "image": "192.0.2.250:5000/postgres@sha256:7abf537131b66ed5af448d90653abf1679b0c7e9a1f07efdd4c3108a401b259a", + "env": { + "POSTGRES_PASSWORD": "bootstrap", + "PGDATA": "/var/lib/postgresql/data/pgdata" + }, + "ports": ["5432:5432"], + "volumes": ["mesh-store-data:/var/lib/postgresql/data"] + }, + // Over TCP, not the socket. While the store initialises it runs a temporary server on the + // socket ONLY, then stops it and starts the real one — so a socket check passes, the action + // exits happy, and the verify a moment later lands in the gap and fails. The action and its + // verify must ask the same question, or the action can succeed into a state verify rejects. + { + "id": "store-ready", + "type": "action", + "in": "mesh-store", + "command": ["sh", "-c", "for i in $(seq 1 180); do pg_isready -h 127.0.0.1 -U postgres >/dev/null 2>&1 && exit 0; sleep 1; done; echo 'the store did not answer within 180s; its own last words follow'; pg_isready -h 127.0.0.1 -U postgres; tail -n 20 /var/lib/postgresql/data/log/*.log 2>/dev/null; exit 1"], +"verify": ["pg_isready", "-h", "127.0.0.1", "-U", "postgres"] + }, + { + "id": "inventory-database", + "type": "action", + "in": "mesh-store", + "command": ["sh", "-c", "psql -U postgres -c 'CREATE DATABASE inventory'"], + "verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw inventory"] + }, + { + "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"] + }, + // Each context owns its own database (novox/hq ADR 0008). A third one is a third database, + // created the same way and named the same way — which is the whole of adding a context to the + // bootstrap, and is why the count is not something the foundation has an opinion about. + { + "id": "licences-database", + "type": "action", + "in": "mesh-store", + "command": ["sh", "-c", "psql -U postgres -c 'CREATE DATABASE licences'"], + "verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw licences"] + }, + { + "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", + "-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable", + "-e", "MESH_STORE_LICENCES=postgres://postgres:bootstrap@127.0.0.1:5432/licences?sslmode=disable", + "192.0.2.250:5000/mesh-controller@sha256:c67db38439ff0aee242b467486765467bb95801f52175fc5727cc4e437338ace", + "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 && docker exec mesh-store psql -U postgres -d licences -tAc \"select to_regclass('public.licence')\" | grep -qx licence"] + }, + { + "id": "bus-certificate", + "type": "action", + "command": ["docker", "run", "--rm", "--entrypoint", "sh", "-v", "mesh-broker-tls:/tls", + "192.0.2.250:5000/nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927", + "-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' -addext 'subjectAltName=DNS:mesh-broker,IP:127.0.0.1' >/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/nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927", + "-c", "test -s /tls/tls.crt && openssl x509 -in /tls/tls.crt -noout"] + }, + { + "id": "bus-conf-dir", + "type": "directory", + "path": "/var/lib/mesh-bus-conf", + "mode": "0700" + }, + { + "id": "bus-conf", + "type": "file", + "path": "/var/lib/mesh-bus-conf/nats.conf", + "mode": "0644", + "content": "port: 4222\nhttp: 127.0.0.1:8222\n\ntls {\n cert_file: \"/tls/tls.crt\"\n key_file: \"/tls/tls.key\"\n}\n\njetstream {\n store_dir: \"/data\"\n}\n\ninclude accounts.conf\n" + }, + { + "id": "bus-accounts", + "type": "file", + "path": "/var/lib/mesh-bus-conf/accounts.conf", + "mode": "0600", + "content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.API.>\", \"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"_INBOX.enrol.>\", \"mesh.control.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\"] }\n subscribe: { allow: [\"$JS.API.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n" + }, + { + "id": "broker", + "type": "container", + "name": "mesh-broker", + "image": "192.0.2.250:5000/nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927", + "ports": ["5671:4222", "127.0.0.1:8222:8222"], + "volumes": ["mesh-broker-data:/data", "mesh-broker-tls:/tls:ro", "/var/lib/mesh-bus-conf:/etc/nats:ro"], + "args": ["-c", "/etc/nats/nats.conf", "-js"] + }, + { + "id": "broker-ready", + "type": "action", + "command": ["sh", "-c", "for i in $(seq 1 60); do docker run --rm --network host --entrypoint sh 192.0.2.250:5000/nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927 -c 'nc -z 127.0.0.1 5671' >/dev/null 2>&1 && exit 0; sleep 1; done; exit 1"], + "verify": ["sh", "-c", "docker run --rm --network host --entrypoint sh 192.0.2.250:5000/nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927 -c 'nc -z 127.0.0.1 5671'"] + }, + { + "id": "control-plane", + "type": "container", + "name": "mesh-controller", + "image": "192.0.2.250:5000/mesh-controller@sha256:c67db38439ff0aee242b467486765467bb95801f52175fc5727cc4e437338ace", + "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_STORE_LICENCES": "postgres://postgres:bootstrap@127.0.0.1:5432/licences?sslmode=disable", + "MESH_BUS_NATS": "nats://controller:bootstrap@127.0.0.1:5671", + "MESH_BROKER_ADDRESS": "192.0.2.10:5671", + "MESH_BROKER_CERTIFICATE": "/broker-tls/tls.crt" + } + } + + ] +} diff --git a/go.mod b/go.mod index c3444a4..47f3aaf 100644 --- a/go.mod +++ b/go.mod @@ -1,9 +1,13 @@ module github.com/novox/mesh-host -go 1.25.0 +go 1.26.0 require ( + github.com/klauspost/compress v1.20.0 // indirect + github.com/nats-io/nats.go v1.54.0 // indirect + github.com/nats-io/nkeys v0.4.16 // indirect + github.com/nats-io/nuid v1.0.1 // indirect github.com/rabbitmq/amqp091-go v1.14.0 // indirect - golang.org/x/crypto v0.55.0 // indirect - golang.org/x/sys v0.47.0 // indirect + golang.org/x/crypto v0.57.0 // indirect + golang.org/x/sys v0.48.0 // indirect ) diff --git a/go.sum b/go.sum index 661ae93..726180b 100644 --- a/go.sum +++ b/go.sum @@ -1,6 +1,18 @@ +github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA= +github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI= +github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA= +github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k= +github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= +github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= +golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo= +golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og= diff --git a/internal/link/asking.go b/internal/link/asking.go new file mode 100644 index 0000000..529d9ae --- /dev/null +++ b/internal/link/asking.go @@ -0,0 +1,55 @@ +package link + +import ( + "context" + "time" +) + +// The enrolment conversation, as the host's own words for it. +// +// **Its own seam rather than part of Link**, because almost nothing about it is the same. The +// credential is a one-time secret rather than this node's own; there is no declaration to hear; the +// whole exchange is a single question asked and possibly asked again. And the stakes differ: a node +// that fails here is not in the mesh at all, where a node that fails in Link has merely lost touch +// with one it belongs to. + +// Approach is how a node reaches a mesh it does not yet belong to. +// +// The three things a token carries about where to go, and nothing about who is asking: the address, +// the certificate that address must present, and which bus is at the other end. **Every token names +// the bus the mesh runs on today until the rollout** (novox/hq ADR 0116 step 5), so an empty +// Transport is the ordinary case rather than something missing. +type Approach struct { + Address string + Fingerprint string + Transport string +} + +// Asking is one open enrolment conversation. +type Asking interface { + // Ask puts the request to the mesh and waits for one answer, or says why none came. + // + // Called again, with the same bytes, while the mesh says "try again": the keys this node + // generated are the ones it keeps, so the same request is the same enrolment and the mesh holds + // the token for it (novox/hq issue 083). + Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) + + // Close lets go of the connection made with the token. + Close() +} + +// Present opens an enrolment conversation with the mesh. +// +// The connection is made before anything is sent, and the certificate is checked while it is being +// made — so a node pointed at the wrong bus finds out before its token has left the machine (ADR +// 0004). +func Present(ctx context.Context, to Approach, node, secret string, + timeout time.Duration) (Asking, error) { + + switch to.Transport { + case OnNATS: + return presentNats(ctx, to, node, secret, timeout) + default: + return presentCurrent(ctx, to, node, secret, timeout) + } +} diff --git a/internal/link/asking_current.go b/internal/link/asking_current.go new file mode 100644 index 0000000..d1fb792 --- /dev/null +++ b/internal/link/asking_current.go @@ -0,0 +1,129 @@ +package link + +import ( + "context" + "errors" + "fmt" + "net/url" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The enrolment conversation on the bus the mesh runs on today. +// +// Moved out of Enrol rather than changed. The reply travels on this node's own queue and is picked +// out by the correlation id the request carried, which is what this transport's reply field means +// and has always meant. + +type currentAsking struct { + conn *amqp.Connection + channel *amqp.Channel + queue string + node string + replies <-chan amqp.Delivery + closed chan *amqp.Error +} + +func presentCurrent(_ context.Context, to Approach, node, secret string, + timeout time.Duration) (Asking, error) { + + config, err := PinnedConfig(to.Fingerprint) + if err != nil { + return nil, 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), to.Address) + + conn, err := amqp.DialConfig(dsn, amqp.Config{ + TLSClientConfig: config, + Dial: amqp.DefaultDial(timeout), + }) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return nil, err + } + // Not quoted back: the DSN carries the one-time secret. + return nil, fmt.Errorf("cannot reach the broker at %s as %s: %w", to.Address, node, err) + } + + channel, err := conn.Channel() + if err != nil { + conn.Close() + return nil, err + } + + // 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 { + conn.Close() + return nil, 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 { + conn.Close() + return nil, err + } + + return ¤tAsking{ + conn: conn, channel: channel, queue: queue.Name, node: node, replies: replies, + closed: conn.NotifyClose(make(chan *amqp.Error, 1)), + }, nil +} + +func (a *currentAsking) Close() { + if a.channel != nil { + _ = a.channel.Close() + } + if a.conn != nil { + _ = a.conn.Close() + } +} + +func (a *currentAsking) Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) { + correlation := fmt.Sprintf("%s-%d", a.node, time.Now().UnixNano()) + publish, cancel := context.WithTimeout(ctx, wait) + defer cancel() + if err := a.channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false, + amqp.Publishing{ + ContentType: "application/json", + CorrelationId: correlation, + ReplyTo: a.queue, + Body: request, + }); err != nil { + return nil, 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(wait) + defer deadline.Stop() + + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + case reason := <-a.closed: + return nil, fmt.Errorf("the broker closed the connection: %v", reason) + case <-deadline.C: + return nil, 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", wait) + case delivery, ok := <-a.replies: + if !ok { + return nil, errors.New("the broker stopped delivering") + } + // Anything else on this queue is not the answer to this question. + if delivery.CorrelationId != correlation { + continue + } + return delivery.Body, nil + } + } +} diff --git a/internal/link/asking_nats.go b/internal/link/asking_nats.go new file mode 100644 index 0000000..b8df9e6 --- /dev/null +++ b/internal/link/asking_nats.go @@ -0,0 +1,166 @@ +package link + +import ( + "context" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "time" + + "github.com/nats-io/nats.go" +) + +// The enrolment conversation on the bus being built. +// +// **The reply address is the whole of what changes**, and it changes for a reason the transport +// forces rather than a preference. Core NATS request/reply puts the caller's inbox in the message's +// reply field and a plain responder answers it — but this request goes into a stream, and a message +// a JetStream consumer delivers has had that field claimed for the consumer's own ack address. So by +// the time the controller reads the request, the transport's reply field names where the +// *controller* must acknowledge. Verified against a running server (design 25 §2). +// +// The address therefore travels as a field of the request, and this subscribes it before publishing: +// a node that published first could miss an answer to a question nobody was listening for. + +// enrolInbox is where a node enrolling waits. +// +// Under `_INBOX.enrol..`, which is exactly what its enrolment user may subscribe and no +// wider — so an answer sealed to one machine cannot be read by another enrolling beside it. The +// random tail is this attempt's own: a reply left over from an attempt that timed out is not the +// answer to this question, which is what the correlation id does on the other transport. +func enrolInbox(node string) (string, error) { + tail := make([]byte, 8) + if _, err := rand.Read(tail); err != nil { + return "", fmt.Errorf("cannot make a reply address: %w", err) + } + return "_INBOX.enrol." + node + "." + hex.EncodeToString(tail), nil +} + +type natsAsking struct { + conn *nats.Conn + js nats.JetStreamContext + inbox string + answers *nats.Subscription + lost chan error +} + +func presentNats(_ context.Context, to Approach, node, secret string, + timeout time.Duration) (Asking, error) { + + config, err := PinnedConfig(to.Fingerprint) + if err != nil { + return nil, err + } + + inbox, err := enrolInbox(node) + if err != nil { + return nil, err + } + + lost := make(chan error, 1) + // The user is this token's own — `enrol.`, which may publish the enrolment subject and + // subscribe its own inbox and nothing else (design 25 §6). The secret is its password, the same + // string the request claims, so the server proves somebody holds the token and the request + // proves the same thing to the controller without it having to ask the server who connected. + conn, err := nats.Connect(natsURL(to.Address), + nats.Secure(config), + nats.UserInfo("enrol."+node, secret), + nats.Name("mesh-host/enrol/"+node), + nats.Timeout(timeout), + nats.NoReconnect(), + nats.DisconnectErrHandler(func(_ *nats.Conn, err error) { + select { + case lost <- fmt.Errorf("the bus closed the connection: %w", err): + default: + } + }), + ) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return nil, err + } + // Not quoted back with the credential: the secret is one-time and still a secret. + return nil, fmt.Errorf("cannot reach the bus at %s as %s: %w", to.Address, node, err) + } + js, err := conn.JetStream() + if err != nil { + conn.Close() + return nil, fmt.Errorf("the bus at %s has no JetStream: %w", to.Address, err) + } + + // Subscribed before anything is published, so an answer cannot arrive before there is anywhere + // for it to land. + answers, err := conn.SubscribeSync(inbox) + if err != nil { + conn.Close() + return nil, fmt.Errorf("this node cannot listen for the mesh's answer: %w", err) + } + if err := conn.Flush(); err != nil { + conn.Close() + return nil, fmt.Errorf("this node's reply address did not reach the bus: %w", err) + } + + return &natsAsking{conn: conn, js: js, inbox: inbox, answers: answers, lost: lost}, nil +} + +func (a *natsAsking) Close() { + if a.answers != nil { + _ = a.answers.Unsubscribe() + } + if a.conn != nil { + a.conn.Close() + } +} + +// Ask publishes the request with this attempt's reply address written into it, and waits there. +func (a *natsAsking) Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) { + addressed, err := withReplyTo(request, a.inbox) + if err != nil { + return nil, err + } + + publish, cancel := context.WithTimeout(ctx, wait) + defer cancel() + // Into the stream and awaited: an enrolment the bus never accepted must fail here rather than be + // assumed, because the node has nothing else to go on. + if _, err := a.js.Publish(EnrolSubject, addressed, nats.Context(publish)); err != nil { + return nil, fmt.Errorf("cannot ask the mesh to enrol this node: %w", err) + } + + // Waited for rather than assumed. A published message that nothing answers means the controller + // is not running, and a node that carried on regardless would believe it had joined a mesh that + // has never heard of it. + // + // **The wait may legitimately be several store-window cycles long**: the controller naks the + // request with a delay while its store is restarting, and the node is waiting on the other side + // of that — which is exactly the combination that would have delivered the answer to a caller + // who had given up, had the address travelled in the transport's field. + answered, cancelAnswer := context.WithTimeout(ctx, wait) + defer cancelAnswer() + for { + msg, err := a.answers.NextMsgWithContext(answered) + switch { + case err == nil: + return msg.Data, nil + case errors.Is(err, context.DeadlineExceeded): + select { + case reason := <-a.lost: + return nil, reason + default: + } + return nil, fmt.Errorf( + "the bus accepted this node's connection and nothing answered within %s. The mesh's "+ + "bus is running and its controller is not", wait) + case errors.Is(err, context.Canceled): + return nil, ctx.Err() + default: + select { + case reason := <-a.lost: + return nil, reason + default: + } + return nil, fmt.Errorf("waiting for the mesh's answer: %w", err) + } + } +} diff --git a/internal/link/asking_nats_test.go b/internal/link/asking_nats_test.go new file mode 100644 index 0000000..5d853b2 --- /dev/null +++ b/internal/link/asking_nats_test.go @@ -0,0 +1,135 @@ +package link + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// The enrolment round trip against a real server. +// +// **This is the test that keeps a reason from becoming folklore.** The reply address travels in the +// request's payload because a JetStream consumer's delivery has had the transport's reply field +// claimed for its own ack address — which is a fact about a server, not a rule anybody can check by +// reading. Both halves are asserted here: that the field really is eaten, and that the answer +// reaches the node anyway. +// +// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNatsAnEnrolment + +// asking is a conversation on a bus with no TLS. Built directly rather than through Present because +// the pin is what Present adds and PinnedConfig's own tests cover it; what is under test here is the +// address the answer comes back on. +func asking(t *testing.T, conn *nats.Conn, js nats.JetStreamContext, node string) *natsAsking { + t.Helper() + inbox, err := enrolInbox(node) + if err != nil { + t.Fatal(err) + } + answers, err := conn.SubscribeSync(inbox) + if err != nil { + t.Fatal(err) + } + if err := conn.Flush(); err != nil { + t.Fatal(err) + } + a := &natsAsking{conn: conn, js: js, inbox: inbox, answers: answers, lost: make(chan error, 1)} + t.Cleanup(a.Close) + return a +} + +// theMeshAnswers stands in for the controller: it consumes the enrolment off the stream, reads the +// reply address out of the payload — never from the transport field — and answers there. It reports +// what the transport field actually held, which is the claim design 25 §2 rests on. +func theMeshAnswers(t *testing.T, conn *nats.Conn, js nats.JetStreamContext, + reply EnrolReply) <-chan string { + t.Helper() + sawReplyField := make(chan string, 1) + sub, err := js.Subscribe(EnrolSubject, func(msg *nats.Msg) { + select { + case sawReplyField <- msg.Reply: + default: + } + var addressed struct { + ReplyTo string `json:"reply_to"` + } + if err := json.Unmarshal(msg.Data, &addressed); err != nil || addressed.ReplyTo == "" { + _ = msg.Ack() + return + } + body, _ := json.Marshal(reply) + // Published explicitly to the address the payload named, never msg.Respond — which would + // send it to whatever the transport's reply field holds, and that is the point. + _ = conn.Publish(addressed.ReplyTo, body) + _ = msg.Ack() + }, nats.Durable("controller-standin"), nats.ManualAck()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sub.Unsubscribe() }) + return sawReplyField +} + +// An enrolment is answered on the address the request carried, and the transport's own reply field +// held something else entirely. +func TestNatsAnEnrolmentIsAnsweredOnTheAddressInItsPayload(t *testing.T) { + conn, js := aBus(t) + const node = "joining" + + sawReplyField := theMeshAnswers(t, conn, js, EnrolReply{Accepted: true, Node: node, Password: "p"}) + a := asking(t, conn, js, node) + + request, _ := json.Marshal(EnrolRequest{Node: node, Secret: "t"}) + answer, err := a.Ask(context.Background(), request, 8*time.Second) + if err != nil { + t.Fatalf("no answer reached the node: %v", err) + } + var reply EnrolReply + if err := json.Unmarshal(answer, &reply); err != nil { + t.Fatal(err) + } + if !reply.Accepted || reply.Node != node { + t.Fatalf("the answer was not the mesh's: %+v", reply) + } + + // And the field the answer would have gone to, had it used the transport's: the consumer's own + // ack address. If a future server stopped doing this, the payload-borne address would still + // work and this line is what would say the reason had changed. + select { + case field := <-sawReplyField: + if field == a.inbox { + t.Fatalf("the transport's reply field held this node's inbox (%s), so the payload "+ + "address is no longer load-bearing — check design 25 §2 before relying on it", field) + } + if field == "" { + t.Fatal("the transport's reply field was empty rather than claimed, which is a third " + + "behaviour from the two design 25 §2 describes") + } + case <-time.After(2 * time.Second): + t.Fatal("the stand-in never saw the request") + } +} + +// A request that names no reply address is not answered, and the node says so as a mesh that is not +// running rather than hanging. The controller has nowhere to send an answer, which is the failure +// the payload field exists to make impossible — asserted so that a request built without it fails +// loudly here rather than quietly on a machine. +func TestNatsAnEnrolmentWithNoReplyAddressIsNotAnswered(t *testing.T) { + conn, js := aBus(t) + const node = "silent" + + theMeshAnswers(t, conn, js, EnrolReply{Accepted: true, Node: node}) + a := asking(t, conn, js, node) + + // Published without going through Ask, so the reply address is genuinely absent. + request, _ := json.Marshal(EnrolRequest{Node: node, Secret: "t"}) + if _, err := js.Publish(EnrolSubject, request); err != nil { + t.Fatal(err) + } + if _, err := a.Ask(context.Background(), []byte(`{"node":"`+node+`"}`), 0); err == nil { + t.Fatal("a node with no answer coming was told it had one") + } +} diff --git a/internal/link/bus.go b/internal/link/bus.go new file mode 100644 index 0000000..7d0f9b7 --- /dev/null +++ b/internal/link/bus.go @@ -0,0 +1,88 @@ +package link + +import ( + "context" + "fmt" + "time" + + "github.com/nats-io/nats.go" + amqp "github.com/rabbitmq/amqp091-go" +) + +// Bus is what a host needs of the mesh's bus, in the mesh's own words. +// +// A host says exactly two things unprompted: what it applied, and that it is here. They are not +// the same kind of statement and the difference is the whole of this interface — one must arrive +// and one must not be insisted on. +// +// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005): this is its own +// interface over its own client libraries, not a contract shared with the controller. The two +// agree because a conformance fixture holds them to one envelope, which is the only kind of +// agreement that survives being in different repositories. +type Bus interface { + // Report says what this node applied. **It must arrive.** A report that fails leaves the + // mesh believing the node never answered while the node believes it did, and the two go on + // disagreeing with nothing anywhere saying so — the shape of fault this project keeps + // finding. Returns false when it could not be delivered, so the caller can say so. + Report(ctx context.Context, node string, body []byte) error + + // Alive says this node is here, and nothing else. **Losing one is nothing**: the next is a + // minute away and the mesh reads a gap rather than counting arrivals. Insisting on delivery + // would turn a harmless miss into a logged failure every minute. + Alive(ctx context.Context, node string, body []byte) error +} + +// --- The bus the mesh runs on today ----------------------------------------------------------- + +// OverCurrent is the bus as a channel, until the rollout. +type OverCurrent struct{ Channel *amqp.Channel } + +func (b OverCurrent) Report(ctx context.Context, node string, body []byte) error { + // Mandatory: an unroutable report comes back rather than disappearing. + return b.Channel.PublishWithContext(ctx, Exchange, KeyReport, true, false, + amqp.Publishing{ContentType: "application/json", Body: body}) +} + +func (b OverCurrent) Alive(ctx context.Context, node string, body []byte) error { + return b.Channel.PublishWithContext(ctx, Exchange, KeyAlive, false, false, + amqp.Publishing{ContentType: "application/json", Body: body}) +} + +// --- NATS --------------------------------------------------------------------------------- + +// OverNATS is the bus as a connection. A report goes through JetStream because it must survive +// the controller's store restarting; a heartbeat does not, because it must not. +type OverNATS struct { + Conn *nats.Conn + JS nats.JetStreamContext +} + +// ReportSubject and AliveSubject are this node's own, and no other node's: a host's account may +// publish `mesh.control..>` and nothing wider, so the subject is the authority on +// which node a report is about. +func ReportSubject(node string) string { return "mesh.control." + node + ".report" } +func AliveSubject(node string) string { return "mesh.control." + node + ".alive" } + +func (b OverNATS) Report(ctx context.Context, node string, body []byte) error { + // Into the CONTROL stream and awaited: this is the message the store-window guarantee is + // about (novox/hq ADR 0083). The controller naks with a delay while its store is away and + // the message is redelivered; a publish the bus never accepted must fail here rather than + // be assumed. + if _, err := b.JS.Publish(ReportSubject(node), body, nats.Context(ctx)); err != nil { + return fmt.Errorf("reporting: %w", err) + } + return nil +} + +func (b OverNATS) Alive(ctx context.Context, node string, body []byte) error { + // Core, deliberately: a heartbeat in a stream is the mesh's least valuable message competing + // for retention with its most valuable, and a lost one is the next one. + if err := b.Conn.Publish(AliveSubject(node), body); err != nil { + return err + } + // Flushed rather than fired and forgotten, so "could not tell the mesh" means the write + // failed rather than that nobody has looked yet. + flush, cancel := context.WithTimeout(ctx, 2*time.Second) + defer cancel() + return b.Conn.FlushWithContext(flush) +} diff --git a/internal/link/enrol.go b/internal/link/enrol.go index 8f0aa6f..f0a9d5c 100644 --- a/internal/link/enrol.go +++ b/internal/link/enrol.go @@ -6,10 +6,7 @@ import ( "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 @@ -53,6 +50,16 @@ type EnrolRequest struct { // on a token this key already spent (novox/hq issue 083). Proof []byte `json:"proof,omitempty"` + // ReplyTo is where the mesh's answer goes, as a field of the request rather than the + // transport's own reply address (design 25 §2). Written by the transport that needs it — + // withReplyTo, once per attempt — because a request going into a stream has had the transport's + // reply field claimed for the consumer's ack address before the controller ever reads it. + // + // Empty on the bus the mesh runs on today, where the delivery carries the reply queue and the + // field means what it has always meant. Named here so both sides of the wire hold the same + // field name, which is what the shape test on each side is for. + ReplyTo string `json:"reply_to,omitempty"` + // Tunnel is the tunnel this node found and whose key it took as its overlay key (novox/hq ADR // 0105): everything about it but that key. Sent with the keys because it is one of them — // OverlayKey above IS this tunnel's public key when this is set — and the mesh composes the @@ -148,52 +155,15 @@ func answered(reply EnrolReply, asking time.Duration) (again bool, err error) { // 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, +func Enrol(ctx context.Context, to Approach, node, secret string, public []byte, overlayKey, sealingKey, servingKey string, profile map[string]any, proof []byte, tunnel *Tunnel, 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) + asking, err := Present(ctx, to, node, secret, timeout) if err != nil { return EnrolReply{}, err } + defer asking.Close() request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, OverlayKey: overlayKey, SealingKey: sealingKey, ServingKey: servingKey, Profile: profile, @@ -203,81 +173,47 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte return EnrolReply{}, err } - // Asked, and asked again with the same request while the mesh says "try again": the keys - // this node generated are the ones it keeps, so the same request is the same enrolment, and - // the mesh holds the token for it (novox/hq issue 083). - ask := func() (string, error) { - 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 "", fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err) - } - return correlation, nil - } - correlation, err := ask() - if err != nil { - return EnrolReply{}, err - } + // Asked, and asked again with the same request while the mesh says "try again": the keys this + // node generated are the ones it keeps, so the same request is the same enrolment, and the mesh + // holds the token for it (novox/hq issue 083). began := time.Now() - - // 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 { + answer, err := asking.Ask(ctx, body, timeout) + if err != nil { + return EnrolReply{}, err + } + var reply EnrolReply + if err := json.Unmarshal(answer, &reply); err != nil { + return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) + } + again, err := answered(reply, time.Since(began)) + if err != nil { + return reply, err + } + if !again { + return reply, nil + } 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) - } - again, err := answered(reply, time.Since(began)) - if err != nil { - return reply, err - } - if !again { - return reply, nil - } - select { - case <-ctx.Done(): - return EnrolReply{}, ctx.Err() - case <-time.After(AskAgainAfter): - } - if correlation, err = ask(); err != nil { - return EnrolReply{}, err - } - if !deadline.Stop() { - select { - case <-deadline.C: - default: - } - } - deadline.Reset(timeout) + case <-time.After(AskAgainAfter): } } } + +// withReplyTo writes this attempt's reply address into the request, as a field of its own. +// +// **Written into the bytes rather than carried beside them**, because the whole point is that the +// address survives a stream: a JetStream consumer's delivery has had the transport's reply field +// claimed for its own ack address, so a reply address that is not in the payload is one the +// controller cannot read (design 25 §2). Done by decoding and re-encoding rather than by setting the +// field before marshalling, so one request can be asked again with a fresh address each time without +// the caller knowing that is what happens. +func withReplyTo(request []byte, inbox string) ([]byte, error) { + var fields map[string]any + if err := json.Unmarshal(request, &fields); err != nil { + return nil, fmt.Errorf("this node's own enrolment request cannot be read back: %w", err) + } + fields["reply_to"] = inbox + return json.Marshal(fields) +} diff --git a/internal/link/hearing.go b/internal/link/hearing.go new file mode 100644 index 0000000..303d7e7 --- /dev/null +++ b/internal/link/hearing.go @@ -0,0 +1,84 @@ +package link + +import ( + "context" + "time" +) + +// What a host hears, as the host's own words for it. +// +// The outbound half went behind `Bus` (bus.go) and a node's two statements stopped naming a +// transport. This is the other half — dialling, and the declarations that arrive — and it is where +// the transport reached furthest: the run loop selected on a channel of the client library's own +// delivery type, so every part of holding a node in its mesh knew which bus it was on. +// +// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005). This is its own +// interface over its own libraries, and it agrees with the controller only because a conformance +// fixture holds both to one envelope. + +// Link is this node's live connection to its mesh: what it hears, and what it says. +// +// One interface rather than two, because **dialling is where the transport is chosen** and choosing +// it twice is how one half of a node ends up on a different bus from the other. +type Link interface { + // Bus is what this node says: what it applied, and that it is here. + Bus + + // Declarations is what the mesh tells this node to be. + Declarations() <-chan Declaration + + // Lost says the link ended, and why. + // + // **Read rather than discovered.** A node that finds out by noticing silence is a node that + // believed it was in the mesh for as long as the silence lasted, which is the one state ADR + // 0004 says must never look like being connected. + Lost() <-chan error + + // Close lets go of whatever was dialled. + Close() +} + +// Declaration is one thing the mesh told this node to be. +// +// **Handled, once — after the report is published.** A node that dies between applying and +// reporting leaves the declaration with the mesh and applies it again on return, which is safe +// because applying is reconciliation: it converges rather than repeating. +// +// There is one way of being done rather than two. A declaration set aside because a newer arrived +// with it is settled exactly as an applied one is, on both buses, and the difference between them +// is a fact the *report* carries — a second method here would be a distinction the transport does +// not make. +type Declaration interface { + // Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it + // received rather than what it re-encoded. + Body() []byte + + // Handled settles it. Called after the report for it has been published, either way. + Handled() error +} + +// Open opens this node's link to its mesh. +// +// Named Open rather than Dial because Dial is this package's raw TLS dial, which the enrolment path +// uses to see a certificate before it trusts anything. +// +// **Both transports ship and this is the one place that chooses** (novox/hq ADR 0116: nothing moves +// a node's bus before step 5). Until then every membership names the bus the mesh runs on today, +// and the rollout is this switch and the credential behind it — not a change anywhere in the loop +// that reads from what comes back. +func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + switch m.Transport { + case OnNATS: + return dialNats(ctx, m, timeout) + default: + return dialCurrent(ctx, m, timeout) + } +} + +// The buses a node can be on. Empty is the one the mesh runs on today, which is every node until +// the rollout — so a membership recorded before any of this existed reads as correct rather than as +// unset. +const ( + OnCurrent = "" + OnNATS = "nats" +) diff --git a/internal/link/hearing_current.go b/internal/link/hearing_current.go new file mode 100644 index 0000000..772e17e --- /dev/null +++ b/internal/link/hearing_current.go @@ -0,0 +1,164 @@ +package link + +import ( + "context" + "errors" + "fmt" + "net/url" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The host's link on the bus the mesh runs on today. +// +// Moved out of the run loop rather than changed: the dial, the queue, the prefetch window and the +// return handler are what they were, because the mesh is running on this and a bus nothing speaks +// yet is no reason to alter the one every node is on. + +// currentLink is this node's connection as a channel. +type currentLink struct { + conn *amqp.Connection + channel *amqp.Channel + arrived chan Declaration + lost chan error +} + +func dialCurrent(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + config, err := PinnedConfig(m.Fingerprint) + if err != nil { + return nil, 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 nil, err + } + return nil, fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err) + } + + channel, err := conn.Channel() + if err != nil { + conn.Close() + return nil, err + } + + queue := QueueFor(m.Node) + if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { + conn.Close() + return nil, fmt.Errorf("cannot declare this node's queue %s: %w", queue, err) + } + + // Applying is one at a time — two at once would race on the same filesystem — but SEEING is + // not: with a prefetch of one the host could never know that a newer declaration was already + // waiting, and so applied every one of a backlog in turn, at the better part of a minute each, + // becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged + // deliveries lets it drain to the newest; each declaration still survives a restart on the + // broker until it is acknowledged, which happens only after it is applied or set aside. + if err := channel.Qos(drainDepth, 0, false); err != nil { + conn.Close() + return nil, err + } + + deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil) + if err != nil { + conn.Close() + return nil, err + } + + l := ¤tLink{ + conn: conn, channel: channel, + arrived: make(chan Declaration, drainDepth), + lost: make(chan error, 1), + } + + // Published mandatory, so the broker hands back anything it cannot route rather than dropping + // it. Without this a report goes to an exchange with no matching binding, the publisher is told + // nothing, and the mesh believes this node never answered while the node believes it did — + // which is what happened when `report` was left unbound on the other side. + returned := channel.NotifyReturn(make(chan amqp.Return, 4)) + go func() { + for r := range returned { + select { + case l.lost <- fmt.Errorf("the broker could not route this node's %s: %s (%d %s)", + r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText): + default: + } + } + }() + + closed := conn.NotifyClose(make(chan *amqp.Error, 1)) + go func() { + select { + case reason := <-closed: + select { + case l.lost <- fmt.Errorf("the link closed: %v", reason): + default: + } + case <-ctx.Done(): + } + }() + + // One goroutine turning the library's deliveries into the mesh's words, so the run loop selects + // on one kind of thing whichever bus it is on. + go func() { + defer close(l.arrived) + for { + select { + case <-ctx.Done(): + return + case delivery, ok := <-deliveries: + if !ok { + select { + case l.lost <- errors.New("the broker stopped delivering"): + default: + } + return + } + select { + case l.arrived <- currentDeclaration{delivery}: + case <-ctx.Done(): + return + } + } + } + }() + + return l, nil +} + +func (l *currentLink) Declarations() <-chan Declaration { return l.arrived } +func (l *currentLink) Lost() <-chan error { return l.lost } + +func (l *currentLink) Close() { + if l.channel != nil { + _ = l.channel.Close() + } + if l.conn != nil { + _ = l.conn.Close() + } +} + +// Report and Alive are the outbound half, over the channel this link holds. +func (l *currentLink) Report(ctx context.Context, node string, body []byte) error { + return OverCurrent{Channel: l.channel}.Report(ctx, node, body) +} + +func (l *currentLink) Alive(ctx context.Context, node string, body []byte) error { + return OverCurrent{Channel: l.channel}.Alive(ctx, node, body) +} + +// currentDeclaration is one delivery from the bus the mesh has. +type currentDeclaration struct{ delivery amqp.Delivery } + +func (d currentDeclaration) Body() []byte { return d.delivery.Body } +func (d currentDeclaration) Handled() error { return d.delivery.Ack(false) } diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go new file mode 100644 index 0000000..d3416be --- /dev/null +++ b/internal/link/hearing_nats.go @@ -0,0 +1,187 @@ +package link + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go" +) + +// The host's link on the bus being built. +// +// Two things a host does here that it cannot do on the other bus, and one it must not try. +// +// **It declares nothing.** On the bus the mesh has, a host declares its own queue on connecting, +// because a queue that is not there means a node that hears nothing. Here the object it reads +// through is a durable consumer, and a host's account reaches no part of the JetStream API — by +// design, because the controller is the only writer of consumer definitions (design 25 §3). So the +// host **binds** to a consumer the controller made when this node enrolled, and a missing one is +// said as what it is rather than quietly created with whatever configuration this client happens to +// default to. +// +// **It gets order for free, and keeps the drain anyway.** The declaration subject is last-per-subject +// (design 29 §4), so a node that was away receives exactly the current declaration rather than a +// queue of superseded ones, and the stream's sequence orders them definitively — the wire-level +// answer to novox/hq issue 107. What the drain in run.go still answers is the live case: three +// pushes to a *connected* node are three deliveries whatever the stream later retains. + +// EnrolSubject is where a joining machine asks. One subject for every node, because a machine +// enrolling has no name the mesh has agreed to yet — which is why its authority to publish here is +// the whole of what its enrolment user may do. +const EnrolSubject = "mesh.control.enrol" + +// DeclareSubject is where this node's declaration lands. Its own, and no other node's: a host's +// account subscribes exactly this and the subject is the authority on which node a declaration is +// for. +func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" } + +// natsURL is a bus address as the client wants it. A membership records host and port, because that +// is what genesis sealed into it and what the other transport takes; the scheme is this transport's +// own business. +func natsURL(address string) string { + if strings.Contains(address, "://") { + return address + } + return "nats://" + address +} + +// natsLink is this node's connection as a JetStream subscription. +type natsLink struct { + conn *nats.Conn + js nats.JetStreamContext + sub *nats.Subscription + node string + arrived chan Declaration + lost chan error +} + +func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { + // Pinned exactly as the other transport is, and for once the Go client makes that easy: it + // takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the + // work. **The constraint recorded against the tool runtime does not apply here** — that client + // takes PEM strings with no verify hook, which is why the bus's certificate must carry a name + // matching the address *modules* dial it by. A host checks the fingerprint and nothing else. + config, err := PinnedConfig(m.Fingerprint) + if err != nil { + return nil, err + } + opts := []nats.Option{ + nats.Secure(config), + nats.UserInfo(m.Node, m.Password), + nats.Name("mesh-host/" + m.Node), + nats.Timeout(timeout), + // A node that has silently lost its route notices, rather than holding a connection the + // server forgot about and believing it is still in the mesh. + nats.PingInterval(10 * time.Second), + nats.MaxPingsOutstanding(2), + // Reconnection is the caller's: Hold already decides when to try again and how long to + // wait, and a client quietly reconnecting underneath it would make that reasoning a + // duplicate of the library's. + nats.NoReconnect(), + } + + conn, err := nats.Connect(natsURL(m.Broker), opts...) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return nil, err + } + return nil, fmt.Errorf("cannot reach the bus at %s: %w", m.Broker, err) + } + js, err := conn.JetStream() + if err != nil { + conn.Close() + return nil, fmt.Errorf("the bus at %s has no JetStream: %w", m.Broker, err) + } + + l := &natsLink{ + conn: conn, js: js, node: m.Node, + arrived: make(chan Declaration, drainDepth), + lost: make(chan error, 1), + } + + // Bound to the consumer the controller made for this node, named after the node because that is + // what the node's own ack grant allows (`$JS.ACK.NODES..>`). + feed := make(chan *nats.Msg, drainDepth) + // The subject as well as the binding: the client checks what is asked for against the + // consumer's own filter, and an empty subject is refused rather than taken to mean "whatever + // that consumer delivers". + sub, err := js.ChanSubscribe(DeclareSubject(m.Node), feed, nats.Bind("NODES", m.Node)) + if err != nil { + conn.Close() + return nil, fmt.Errorf( + "this node cannot read its declarations: %w. The mesh creates that when a node enrols, "+ + "and a host may not create one itself — so this is the mesh's to answer, not this "+ + "machine's", err) + } + l.sub = sub + + conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) { + select { + case l.lost <- fmt.Errorf("the link dropped: %w", err): + default: + } + }) + conn.SetClosedHandler(func(*nats.Conn) { + select { + case l.lost <- errors.New("the link closed"): + default: + } + }) + + go func() { + defer close(l.arrived) + for { + select { + case <-ctx.Done(): + return + case msg, ok := <-feed: + if !ok { + select { + case l.lost <- errors.New("the bus stopped delivering"): + default: + } + return + } + select { + case l.arrived <- natsDeclaration{msg}: + case <-ctx.Done(): + return + } + } + } + }() + + return l, nil +} + +func (l *natsLink) Declarations() <-chan Declaration { return l.arrived } +func (l *natsLink) Lost() <-chan error { return l.lost } + +func (l *natsLink) Close() { + if l.sub != nil { + _ = l.sub.Unsubscribe() + } + if l.conn != nil { + l.conn.Close() + } +} + +func (l *natsLink) Report(ctx context.Context, node string, body []byte) error { + return OverNATS{Conn: l.conn, JS: l.js}.Report(ctx, node, body) +} + +func (l *natsLink) Alive(ctx context.Context, node string, body []byte) error { + return OverNATS{Conn: l.conn, JS: l.js}.Alive(ctx, node, body) +} + +// natsDeclaration is one declaration off the NODES stream. +type natsDeclaration struct{ msg *nats.Msg } + +func (d natsDeclaration) Body() []byte { return d.msg.Data } + +// Handled acknowledges it. The ack goes to this node's own ack subject, which is the one thing +// besides its reports a node's account may publish. +func (d natsDeclaration) Handled() error { return d.msg.Ack() } diff --git a/internal/link/hearing_nats_test.go b/internal/link/hearing_nats_test.go new file mode 100644 index 0000000..9a33e38 --- /dev/null +++ b/internal/link/hearing_nats_test.go @@ -0,0 +1,228 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "fmt" + "os" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// The host's link against a real server, because every claim here is about one. +// +// Whether binding to a consumer the host did not create works, whether a declaration on the node's +// own subject arrives, whether acknowledging it removes it from the consumer's pending — none of +// that can be reasoned out, and the first two are the ones that would leave a node silently hearing +// nothing: +// +// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNats + +func aBus(t *testing.T) (*nats.Conn, nats.JetStreamContext) { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + conn, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(conn.Close) + js, err := conn.JetStream() + if err != nil { + t.Fatal(err) + } + // **Ensured and purged, not deleted and recreated.** Delete-then-add looked like a reset and is + // not one: a test that did that inherited the previous test's messages, and the symptom was a + // declaration counted as delivered twice — which reads as a redelivery bug in the code under + // test rather than as a dirty stream. Purge is defined to empty a stream; recreating one is a + // race with the server's own teardown. + for _, want := range []*nats.StreamConfig{ + {Name: "NODES", Subjects: []string{"mesh.node.*.declare"}, MaxMsgsPerSubject: 1}, + {Name: "CONTROL", Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"}, + Retention: nats.WorkQueuePolicy}, + } { + if _, err := js.StreamInfo(want.Name); err != nil { + if _, err := js.AddStream(want); err != nil { + t.Fatal(err) + } + } + if err := js.PurgeStream(want.Name); err != nil { + t.Fatal(err) + } + } + return conn, js +} + +// theMeshMakes is the consumer the controller creates when a node enrols. Made here by the test +// because the host may not: its account reaches no part of the JetStream API, which is the whole +// reason this binds rather than subscribes. +// +// Removed afterwards, and each test names its own node: two tests sharing a consumer name share its +// delivery count and its pending list, and the first thing that goes wrong reads as a fault in the +// host rather than in the test beside it. +func theMeshMakes(t *testing.T, js nats.JetStreamContext, node string) { + t.Helper() + t.Cleanup(func() { _ = js.DeleteConsumer("NODES", node) }) + if _, err := js.AddConsumer("NODES", &nats.ConsumerConfig{ + Durable: node, + FilterSubject: DeclareSubject(node), + AckPolicy: nats.AckExplicitPolicy, + AckWait: 300 * time.Second, + DeliverSubject: "_DELIVER." + node, + }); err != nil { + t.Fatal(err) + } +} + +func signedBy(t *testing.T, key ed25519.PrivateKey, declaration []byte) []byte { + t.Helper() + body, err := json.Marshal(Signed{ + Declaration: declaration, Signature: ed25519.Sign(key, declaration), + }) + if err != nil { + t.Fatal(err) + } + return body +} + +// A declaration on this node's own subject reaches the host, is applied, and acknowledging it +// empties the consumer — which is what tells the mesh the node has it. +func TestNatsADeclarationReachesTheHostAndIsSettled(t *testing.T) { + conn, js := aBus(t) + const node = "settling" + theMeshMakes(t, js, node) + + public, private, _ := ed25519.GenerateKey(nil) + m := Membership{Node: node, Signer: public} + + // Dialled directly rather than through Open: the test server has no TLS, and what is being + // checked is the subscription and the settling, not the pin — which PinnedConfig owns and its + // own tests cover. + l := &natsLink{conn: conn, js: js, node: node, + arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)} + feed := make(chan *nats.Msg, drainDepth) + sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node)) + if err != nil { + t.Fatalf("the host could not bind to the consumer the mesh made for it: %v", err) + } + defer func() { _ = sub.Unsubscribe() }() + go func() { + for msg := range feed { + l.arrived <- natsDeclaration{msg} + } + }() + + if _, err := js.Publish(DeclareSubject(node), + signedBy(t, private, []byte(`{"declared":"d1"}`))); err != nil { + t.Fatal(err) + } + + select { + case d := <-l.Declarations(): + report := handleBody(context.Background(), m, d.Body(), + func(context.Context, []byte, []byte) Report { + return Report{Applied: []string{"store"}} + }) + if report.Refused != "" { + t.Fatalf("a declaration the mesh signed was refused: %s", report.Refused) + } + if err := d.Handled(); err != nil { + t.Fatalf("the node could not acknowledge its own declaration: %v", err) + } + case <-time.After(8 * time.Second): + t.Fatal("no declaration reached the host") + } + + // **Nothing pending is the property**; a delivery count is not. Delivery is at-least-once by + // design, so pinning "delivered exactly once" would be asserting something the mesh does not + // rely on. What matters is that the acknowledgement landed, so the mesh can tell the node has + // it — and that no redelivery was needed to get there, which is what would say the node was + // too slow to answer for its own ack wait. + deadline := time.Now().Add(5 * time.Second) + var last string + for time.Now().Before(deadline) { + info, err := js.ConsumerInfo("NODES", node) + switch { + case err != nil: + last = err.Error() + case info.NumAckPending == 0 && info.NumRedelivered == 0: + return + default: + last = fmt.Sprintf("pending %d, redelivered %d", info.NumAckPending, info.NumRedelivered) + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("the declaration was not settled, so the mesh cannot tell the node has it: %s", last) +} + +// **A node that was away gets exactly the current declaration and nothing older.** Three pushed +// while nothing is listening leave one on the stream, and it is the newest — the wire-level answer +// to novox/hq issue 107, and the half of the drain that stops being the host's problem. +func TestNatsANodeThatWasAwayGetsOnlyTheNewest(t *testing.T) { + _, js := aBus(t) + const node = "returning" + _, private, _ := ed25519.GenerateKey(nil) + + for _, id := range []string{"d1", "d2", "d3"} { + if _, err := js.Publish(DeclareSubject(node), + signedBy(t, private, []byte(`{"declared":"`+id+`"}`))); err != nil { + t.Fatal(err) + } + } + info, err := js.StreamInfo("NODES") + if err != nil { + t.Fatal(err) + } + if info.State.Msgs != 1 { + t.Fatalf("%d declarations survived for one node; a node that was away would apply a backlog "+ + "of things nobody wants any more", info.State.Msgs) + } + + theMeshMakes(t, js, node) + feed := make(chan *nats.Msg, drainDepth) + sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node)) + if err != nil { + t.Fatal(err) + } + defer func() { _ = sub.Unsubscribe() }() + + select { + case msg := <-feed: + if declaredIn(msg.Data) != "d3" { + t.Fatalf("the node was given %q rather than the newest", declaredIn(msg.Data)) + } + case <-time.After(8 * time.Second): + t.Fatal("the node that was away was given nothing") + } +} + +// A report goes through the stream and a heartbeat does not: the one that must survive the +// controller's store restarting is kept, and the one that must not is not. +func TestNatsAReportIsKeptAndAHeartbeatIsNot(t *testing.T) { + conn, js := aBus(t) + bus := OverNATS{Conn: conn, JS: js} + ctx := context.Background() + + body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1"}) + if err := bus.Report(ctx, "anchor", body); err != nil { + t.Fatal(err) + } + beat, _ := json.Marshal(Alive{Node: "anchor"}) + if err := bus.Alive(ctx, "anchor", beat); err != nil { + t.Fatal(err) + } + + info, err := js.StreamInfo("CONTROL") + if err != nil { + t.Fatal(err) + } + if info.State.Msgs != 1 { + t.Fatalf("%d messages were kept; a report must be and a heartbeat must not", info.State.Msgs) + } +} diff --git a/internal/link/newest_test.go b/internal/link/newest_test.go index 7ad6cbb..5eab4a0 100644 --- a/internal/link/newest_test.go +++ b/internal/link/newest_test.go @@ -3,33 +3,46 @@ package link import ( "testing" "time" - - amqp "github.com/rabbitmq/amqp091-go" ) +// said is one declaration as a test hands it over, with no transport under it — which is what the +// seam bought: the drain's reasoning was reachable only through a real broker before. +type said struct { + body []byte + handled bool +} + +func (s *said) Body() []byte { return s.body } +func (s *said) Handled() error { s.handled = true; return nil } + +func arriving(bodies ...string) chan Declaration { + ch := make(chan Declaration, 8) + for _, b := range bodies { + ch <- &said{body: []byte(b)} + } + return ch +} + // A machine asked to be five things becomes the last one: what is already waiting supersedes what // arrived first, and everything set aside is named so it can be reported. func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) { - deliveries := make(chan amqp.Delivery, 8) - for _, id := range []string{"two", "three", "four"} { - deliveries <- amqp.Delivery{Body: []byte(id)} + waiting := arriving("two", "three", "four") + apply, superseded := newest(waiting, &said{body: []byte("one")}, 50*time.Millisecond) + if string(apply.Body()) != "four" { + t.Fatalf("applied %q, not the newest", apply.Body()) } - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("one")}, 50*time.Millisecond) - if string(apply.Body) != "four" { - t.Fatalf("applied %q, not the newest", apply.Body) - } - if len(superseded) != 3 || string(superseded[0].Body) != "one" || string(superseded[2].Body) != "three" { + if len(superseded) != 3 || string(superseded[0].Body()) != "one" || + string(superseded[2].Body()) != "three" { t.Fatalf("set aside %d: %v", len(superseded), superseded) } } // One declaration with nothing behind it is applied as it always was, after the window. func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) { - deliveries := make(chan amqp.Delivery, 1) began := time.Now() - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("only")}, 30*time.Millisecond) - if string(apply.Body) != "only" || len(superseded) != 0 { - t.Fatalf("got %q with %d set aside", apply.Body, len(superseded)) + apply, superseded := newest(arriving(), &said{body: []byte("only")}, 30*time.Millisecond) + if string(apply.Body()) != "only" || len(superseded) != 0 { + t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded)) } if time.Since(began) < 30*time.Millisecond { t.Fatal("did not wait the window for a straggler") @@ -38,13 +51,13 @@ func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) { // A straggler within the window is taken; one after it is the next push. func TestAStragglerWithinTheWindowIsTaken(t *testing.T) { - deliveries := make(chan amqp.Delivery, 2) + waiting := arriving() go func() { time.Sleep(20 * time.Millisecond) - deliveries <- amqp.Delivery{Body: []byte("late")} + waiting <- &said{body: []byte("late")} }() - apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("first")}, 100*time.Millisecond) - if string(apply.Body) != "late" || len(superseded) != 1 { - t.Fatalf("got %q with %d set aside", apply.Body, len(superseded)) + apply, superseded := newest(waiting, &said{body: []byte("first")}, 100*time.Millisecond) + if string(apply.Body()) != "late" || len(superseded) != 1 { + t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded)) } } diff --git a/internal/link/run.go b/internal/link/run.go index 2675065..0288a4f 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -6,10 +6,7 @@ import ( "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. @@ -33,6 +30,10 @@ type Membership struct { Fingerprint string Password string Signer ed25519.PublicKey + // Transport is which bus this node speaks (hearing.go). Empty is the one the mesh runs on + // today, which is every node until the rollout — so a membership recorded before any of this + // existed reads as correct rather than as unset. + Transport string } // Applier is what the host does with a declaration that has been proved to come from the mesh. @@ -190,108 +191,62 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout if say == nil { say = func(string) {} } - config, err := PinnedConfig(m.Fingerprint) + link, err := Open(ctx, m, timeout) if err != nil { return err } + defer link.Close() - 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) - } - - // Applying is one at a time — two at once would race on the same filesystem — but SEEING is - // not: with a prefetch of one the host could never know that a newer declaration was already - // waiting, and so applied every one of a backlog in turn, at the better part of a minute each, - // becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged - // deliveries lets it drain to the newest; each declaration still survives a restart on the - // broker until it is acknowledged, which happens only after it is applied or set aside. - if err := channel.Qos(drainDepth, 0, false); err != nil { - return err - } - - deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil) - if err != nil { - return err - } - // Said, because it is the event anybody watching actually wants. Without it a node logs - // every failure and nothing on success, so a log full of "trying again" and then silence - // reads as still broken when it means the opposite. - say("in the mesh, consuming " + queue) + // Said, because it is the event anybody watching actually wants. Without it a node logs every + // failure and nothing on success, so a log full of "trying again" and then silence reads as + // still broken when it means the opposite. + say("in the mesh, hearing what this node should be") // A word every so often, so the mesh can tell a node that is quiet from one that is gone. // Cheap on purpose: it carries a name and nothing else, because anything more would be a // report, and reports are rare where this is constant. beat := time.NewTicker(AliveEvery) defer beat.Stop() - publishAlive(ctx, channel, m, say, timeout) + publishAlive(ctx, link, m, say, timeout) - closed := conn.NotifyClose(make(chan *amqp.Error, 1)) - - // Published mandatory, so the broker hands back anything it cannot route rather than - // dropping it. Without this a report goes to an exchange with no matching binding, the - // publisher is told nothing, and the mesh believes this node never answered while the node - // believes it did — which is what happened when `report` was left unbound on the other side. - returned := channel.NotifyReturn(make(chan amqp.Return, 4)) - go func() { - for r := range returned { - say(fmt.Sprintf("the broker could not route this node's %s: %s (%d %s)", - r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText)) - } - }() + declarations := link.Declarations() for { select { case <-ctx.Done(): return nil case <-beat.C: - publishAlive(ctx, channel, m, say, timeout) + publishAlive(ctx, link, m, say, timeout) case unasked := <-outbox: // Said without having been asked: a reconcile found what an adopted node holds, or // its firewall, changed since it last said. - published := publishReport(ctx, channel, m, unasked.Report, say, timeout) + published := publishReport(ctx, link, m, unasked.Report, say, timeout) if unasked.Done != nil { unasked.Done(published) } - case reason := <-closed: - return fmt.Errorf("the link closed: %v", reason) - case delivery, ok := <-deliveries: + case reason := <-link.Lost(): + return reason + case declaration, ok := <-declarations: if !ok { - return errors.New("the broker stopped delivering") + // The link's own reason, when it has managed to say one: "stopped delivering" on + // its own says nothing about why, and why is the whole of what an operator wants. + select { + case reason := <-link.Lost(): + return reason + default: + return errors.New("the mesh stopped sending this node declarations") + } } - // Whatever else is already waiting supersedes this one. Each set-aside declaration - // is reported as such, then acknowledged unapplied. - delivery, superseded := newest(deliveries, delivery, drainWindow) + // Whatever else is already waiting supersedes this one. Each set-aside declaration is + // reported as such, then settled unapplied. + declaration, superseded := newest(declarations, declaration, drainWindow) for _, old := range superseded { say("set aside a declaration: a newer one arrived with it") - publishReport(ctx, channel, m, Report{Node: m.Node, Declared: declaredIn(old.Body), - Superseded: declaredIn(delivery.Body)}, say, timeout) - _ = old.Ack(false) + publishReport(ctx, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()), + Superseded: declaredIn(declaration.Body())}, say, timeout) + _ = old.Handled() } - report := handle(ctx, m, apply, delivery) + report := handleBody(ctx, m, declaration.Body(), apply) switch { case report.Refused != "": say("refused a declaration: " + report.Refused) @@ -300,12 +255,12 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout default: say(fmt.Sprintf("applied %d resource(s)", len(report.Applied))) } - publishReport(ctx, channel, m, report, say, 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, + publishReport(ctx, link, m, report, say, timeout) + // Settled after the report is published. A node that dies between applying and + // reporting leaves the declaration with the mesh and applies it again on return, // which is safe because applying is reconciliation — it converges rather than // repeating. - _ = delivery.Ack(false) + _ = declaration.Handled() } } } @@ -321,12 +276,24 @@ const ( // newest takes what is already waiting behind `first` and returns the last of them to apply, and // the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a // declaration in flight from the mesh arrives within that; one that does not is the next push. -func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) { +// +// **Its job narrows once declarations are state rather than messages, and does not disappear.** +// On the bus being built, a declaration is last-per-subject (novox/hq design 29 §4), so a node +// that was away receives exactly the current one instead of a queue of superseded ones — the +// catch-up half of what this does is then the stream's. And a stream sequence orders them +// definitively, where this window only infers order from arrival time, which is the wire-level +// answer to novox/hq issue 107. +// +// What remains is the live case: three pushes in quick succession to a *connected* node are +// three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not +// deleted — and saying which half goes is worth more than a note that it "can probably be +// removed", which is how a load-bearing window gets deleted by somebody in a hurry. +func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) { latest := first - var superseded []amqp.Delivery + var superseded []Declaration for { select { - case next, ok := <-deliveries: + case next, ok := <-arriving: if !ok { return latest, superseded } @@ -355,10 +322,6 @@ func declaredIn(body []byte) string { return d.Declared } -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 { @@ -381,36 +344,19 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R // report a command makes rather than the running host — a rekey (novox/hq ADR 0105). The same // account, the same pinned certificate and the same exchange as the running host's reports. func Publish(ctx context.Context, m Membership, report Report, timeout time.Duration) error { - config, err := PinnedConfig(m.Fingerprint) + link, err := Open(ctx, m, timeout) 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), - }) - 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() + defer link.Close() var said string - if !publishReport(ctx, channel, m, report, func(s string) { said = s }, timeout) { + if !publishReport(ctx, link, m, report, func(s string) { said = s }, timeout) { return errors.New(said) } return nil } -func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report, +func publishReport(ctx context.Context, bus Bus, m Membership, report Report, say Announce, timeout time.Duration) bool { report.Node = m.Node body, err := json.Marshal(report) @@ -424,8 +370,7 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep // Said rather than swallowed. A report that fails to publish leaves the mesh believing this // node never answered, while the node believes it did — and the two would go on disagreeing // with nothing anywhere saying so. That shape of fault is the one this project keeps finding. - if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false, - amqp.Publishing{ContentType: "application/json", Body: body}); err != nil { + if err := bus.Report(publish, m.Node, body); err != nil { say(fmt.Sprintf("applied, and could not tell the mesh: %v", err)) return false } @@ -433,7 +378,7 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep } // publishAlive says this node is here, and nothing else. -func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say Announce, +func publishAlive(ctx context.Context, bus Bus, m Membership, say Announce, timeout time.Duration) { body, err := json.Marshal(Alive{Node: m.Node}) if err != nil { @@ -444,8 +389,7 @@ func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say // Not mandatory, unlike a report. Losing one is nothing: the next is a minute away, and the // mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a // harmless miss into a logged failure every minute. - if err := channel.PublishWithContext(publish, Exchange, KeyAlive, false, false, - amqp.Publishing{ContentType: "application/json", Body: body}); err != nil { + if err := bus.Alive(publish, m.Node, body); err != nil { say("could not tell the mesh this node is here: " + err.Error()) } }