Merge pull request 'A host reaches either bus, and a genesis template that raises the mesh on the new one' (#33) from feat/nats-genesis into main
This commit was merged in pull request #33.
This commit is contained in:
@@ -772,7 +772,12 @@ func enrol(ctx context.Context, opts options) error {
|
|||||||
// who knows its public key (novox/hq issue 083).
|
// who knows its public key (novox/hq issue 083).
|
||||||
proof := mine.Sign(link.EnrolProof(token.Secret, mine.Public, mine.Overlay.Public,
|
proof := mine.Sign(link.EnrolProof(token.Secret, mine.Public, mine.Overlay.Public,
|
||||||
sealing.Public, serving.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,
|
mine.Public, mine.Overlay.Public, sealing.Public, serving.Public, reported, proof, found,
|
||||||
opts.timeout)
|
opts.timeout)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -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"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
]
|
||||||
|
}
|
||||||
@@ -1,9 +1,13 @@
|
|||||||
module github.com/novox/mesh-host
|
module github.com/novox/mesh-host
|
||||||
|
|
||||||
go 1.25.0
|
go 1.26.0
|
||||||
|
|
||||||
require (
|
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
|
github.com/rabbitmq/amqp091-go v1.14.0 // indirect
|
||||||
golang.org/x/crypto v0.55.0 // indirect
|
golang.org/x/crypto v0.57.0 // indirect
|
||||||
golang.org/x/sys v0.47.0 // indirect
|
golang.org/x/sys v0.48.0 // indirect
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
|
||||||
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
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 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
|
||||||
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
|
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 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
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=
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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.<node>.`, 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.<node>`, 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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.<its own node>.>` 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)
|
||||||
|
}
|
||||||
+37
-101
@@ -6,10 +6,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/url"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
amqp "github.com/rabbitmq/amqp091-go"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// The wire format shared with the control plane, which defines it separately because this binary
|
// 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).
|
// on a token this key already spent (novox/hq issue 083).
|
||||||
Proof []byte `json:"proof,omitempty"`
|
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
|
// 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 —
|
// 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
|
// 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
|
// 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
|
// says once it is in, and the secret travels again because the control plane must not have to ask
|
||||||
// the broker who connected.
|
// 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,
|
overlayKey, sealingKey, servingKey string, profile map[string]any, proof []byte,
|
||||||
tunnel *Tunnel, timeout time.Duration) (EnrolReply, error) {
|
tunnel *Tunnel, timeout time.Duration) (EnrolReply, error) {
|
||||||
|
|
||||||
config, err := PinnedConfig(pin)
|
asking, err := Present(ctx, to, node, secret, timeout)
|
||||||
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 {
|
if err != nil {
|
||||||
return EnrolReply{}, err
|
return EnrolReply{}, err
|
||||||
}
|
}
|
||||||
|
defer asking.Close()
|
||||||
|
|
||||||
request := EnrolRequest{Node: node, Secret: secret, PublicKey: public,
|
request := EnrolRequest{Node: node, Secret: secret, PublicKey: public,
|
||||||
OverlayKey: overlayKey, SealingKey: sealingKey, ServingKey: servingKey, Profile: profile,
|
OverlayKey: overlayKey, SealingKey: sealingKey, ServingKey: servingKey, Profile: profile,
|
||||||
@@ -203,57 +173,17 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte
|
|||||||
return EnrolReply{}, err
|
return EnrolReply{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Asked, and asked again with the same request while the mesh says "try again": the keys
|
// Asked, and asked again with the same request while the mesh says "try again": the keys this
|
||||||
// this node generated are the ones it keeps, so the same request is the same enrolment, and
|
// node generated are the ones it keeps, so the same request is the same enrolment, and the mesh
|
||||||
// the mesh holds the token for it (novox/hq issue 083).
|
// holds the token for it (novox/hq issue 083).
|
||||||
ask := func() (string, error) {
|
began := time.Now()
|
||||||
correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano())
|
for {
|
||||||
publish, cancel := context.WithTimeout(ctx, timeout)
|
answer, err := asking.Ask(ctx, body, 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 {
|
if err != nil {
|
||||||
return EnrolReply{}, err
|
return EnrolReply{}, err
|
||||||
}
|
}
|
||||||
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 {
|
|
||||||
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
|
var reply EnrolReply
|
||||||
if err := json.Unmarshal(delivery.Body, &reply); err != nil {
|
if err := json.Unmarshal(answer, &reply); err != nil {
|
||||||
return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err)
|
return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err)
|
||||||
}
|
}
|
||||||
again, err := answered(reply, time.Since(began))
|
again, err := answered(reply, time.Since(began))
|
||||||
@@ -268,16 +198,22 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte
|
|||||||
return EnrolReply{}, ctx.Err()
|
return EnrolReply{}, ctx.Err()
|
||||||
case <-time.After(AskAgainAfter):
|
case <-time.After(AskAgainAfter):
|
||||||
}
|
}
|
||||||
if correlation, err = ask(); err != nil {
|
|
||||||
return EnrolReply{}, err
|
|
||||||
}
|
|
||||||
if !deadline.Stop() {
|
|
||||||
select {
|
|
||||||
case <-deadline.C:
|
|
||||||
default:
|
|
||||||
}
|
|
||||||
}
|
|
||||||
deadline.Reset(timeout)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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)
|
||||||
|
}
|
||||||
|
|||||||
@@ -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"
|
||||||
|
)
|
||||||
@@ -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) }
|
||||||
@@ -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.<node>.>`).
|
||||||
|
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() }
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -3,33 +3,46 @@ package link
|
|||||||
import (
|
import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"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
|
// 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.
|
// arrived first, and everything set aside is named so it can be reported.
|
||||||
func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) {
|
func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) {
|
||||||
deliveries := make(chan amqp.Delivery, 8)
|
waiting := arriving("two", "three", "four")
|
||||||
for _, id := range []string{"two", "three", "four"} {
|
apply, superseded := newest(waiting, &said{body: []byte("one")}, 50*time.Millisecond)
|
||||||
deliveries <- amqp.Delivery{Body: []byte(id)}
|
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 len(superseded) != 3 || string(superseded[0].Body()) != "one" ||
|
||||||
if string(apply.Body) != "four" {
|
string(superseded[2].Body()) != "three" {
|
||||||
t.Fatalf("applied %q, not the newest", apply.Body)
|
|
||||||
}
|
|
||||||
if len(superseded) != 3 || string(superseded[0].Body) != "one" || string(superseded[2].Body) != "three" {
|
|
||||||
t.Fatalf("set aside %d: %v", len(superseded), superseded)
|
t.Fatalf("set aside %d: %v", len(superseded), superseded)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// One declaration with nothing behind it is applied as it always was, after the window.
|
// One declaration with nothing behind it is applied as it always was, after the window.
|
||||||
func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) {
|
func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) {
|
||||||
deliveries := make(chan amqp.Delivery, 1)
|
|
||||||
began := time.Now()
|
began := time.Now()
|
||||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("only")}, 30*time.Millisecond)
|
apply, superseded := newest(arriving(), &said{body: []byte("only")}, 30*time.Millisecond)
|
||||||
if string(apply.Body) != "only" || len(superseded) != 0 {
|
if string(apply.Body()) != "only" || len(superseded) != 0 {
|
||||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded))
|
||||||
}
|
}
|
||||||
if time.Since(began) < 30*time.Millisecond {
|
if time.Since(began) < 30*time.Millisecond {
|
||||||
t.Fatal("did not wait the window for a straggler")
|
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.
|
// A straggler within the window is taken; one after it is the next push.
|
||||||
func TestAStragglerWithinTheWindowIsTaken(t *testing.T) {
|
func TestAStragglerWithinTheWindowIsTaken(t *testing.T) {
|
||||||
deliveries := make(chan amqp.Delivery, 2)
|
waiting := arriving()
|
||||||
go func() {
|
go func() {
|
||||||
time.Sleep(20 * time.Millisecond)
|
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)
|
apply, superseded := newest(waiting, &said{body: []byte("first")}, 100*time.Millisecond)
|
||||||
if string(apply.Body) != "late" || len(superseded) != 1 {
|
if string(apply.Body()) != "late" || len(superseded) != 1 {
|
||||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+58
-114
@@ -6,10 +6,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/url"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
amqp "github.com/rabbitmq/amqp091-go"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// ErrForged is what a node returns for a declaration whose signature is not the mesh's.
|
// 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
|
Fingerprint string
|
||||||
Password string
|
Password string
|
||||||
Signer ed25519.PublicKey
|
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.
|
// 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 {
|
if say == nil {
|
||||||
say = func(string) {}
|
say = func(string) {}
|
||||||
}
|
}
|
||||||
config, err := PinnedConfig(m.Fingerprint)
|
link, err := Open(ctx, m, timeout)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
defer link.Close()
|
||||||
|
|
||||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
// Said, because it is the event anybody watching actually wants. Without it a node logs every
|
||||||
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
// failure and nothing on success, so a log full of "trying again" and then silence reads as
|
||||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
// still broken when it means the opposite.
|
||||||
TLSClientConfig: config,
|
say("in the mesh, hearing what this node should be")
|
||||||
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)
|
|
||||||
|
|
||||||
// A word every so often, so the mesh can tell a node that is quiet from one that is gone.
|
// 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
|
// 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.
|
// report, and reports are rare where this is constant.
|
||||||
beat := time.NewTicker(AliveEvery)
|
beat := time.NewTicker(AliveEvery)
|
||||||
defer beat.Stop()
|
defer beat.Stop()
|
||||||
publishAlive(ctx, channel, m, say, timeout)
|
publishAlive(ctx, link, m, say, timeout)
|
||||||
|
|
||||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
declarations := link.Declarations()
|
||||||
|
|
||||||
// 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))
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return nil
|
return nil
|
||||||
case <-beat.C:
|
case <-beat.C:
|
||||||
publishAlive(ctx, channel, m, say, timeout)
|
publishAlive(ctx, link, m, say, timeout)
|
||||||
case unasked := <-outbox:
|
case unasked := <-outbox:
|
||||||
// Said without having been asked: a reconcile found what an adopted node holds, or
|
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||||
// its firewall, changed since it last said.
|
// 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 {
|
if unasked.Done != nil {
|
||||||
unasked.Done(published)
|
unasked.Done(published)
|
||||||
}
|
}
|
||||||
case reason := <-closed:
|
case reason := <-link.Lost():
|
||||||
return fmt.Errorf("the link closed: %v", reason)
|
return reason
|
||||||
case delivery, ok := <-deliveries:
|
case declaration, ok := <-declarations:
|
||||||
if !ok {
|
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.
|
// Whatever else is already waiting supersedes this one. Each set-aside declaration is
|
||||||
delivery, superseded := newest(deliveries, delivery, drainWindow)
|
// reported as such, then settled unapplied.
|
||||||
|
declaration, superseded := newest(declarations, declaration, drainWindow)
|
||||||
for _, old := range superseded {
|
for _, old := range superseded {
|
||||||
say("set aside a declaration: a newer one arrived with it")
|
say("set aside a declaration: a newer one arrived with it")
|
||||||
publishReport(ctx, channel, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
publishReport(ctx, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()),
|
||||||
Superseded: declaredIn(delivery.Body)}, say, timeout)
|
Superseded: declaredIn(declaration.Body())}, say, timeout)
|
||||||
_ = old.Ack(false)
|
_ = old.Handled()
|
||||||
}
|
}
|
||||||
report := handle(ctx, m, apply, delivery)
|
report := handleBody(ctx, m, declaration.Body(), apply)
|
||||||
switch {
|
switch {
|
||||||
case report.Refused != "":
|
case report.Refused != "":
|
||||||
say("refused a declaration: " + 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:
|
default:
|
||||||
say(fmt.Sprintf("applied %d resource(s)", len(report.Applied)))
|
say(fmt.Sprintf("applied %d resource(s)", len(report.Applied)))
|
||||||
}
|
}
|
||||||
publishReport(ctx, channel, m, report, say, timeout)
|
publishReport(ctx, link, m, report, say, timeout)
|
||||||
// Acknowledged after the report is published. A node that dies between applying and
|
// Settled 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,
|
// reporting leaves the declaration with the mesh and applies it again on return,
|
||||||
// which is safe because applying is reconciliation — it converges rather than
|
// which is safe because applying is reconciliation — it converges rather than
|
||||||
// repeating.
|
// 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
|
// 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
|
// 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.
|
// 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
|
latest := first
|
||||||
var superseded []amqp.Delivery
|
var superseded []Declaration
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case next, ok := <-deliveries:
|
case next, ok := <-arriving:
|
||||||
if !ok {
|
if !ok {
|
||||||
return latest, superseded
|
return latest, superseded
|
||||||
}
|
}
|
||||||
@@ -355,10 +322,6 @@ func declaredIn(body []byte) string {
|
|||||||
return d.Declared
|
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
|
// 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.
|
// 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 {
|
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
|
// 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.
|
// 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 {
|
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 {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
defer link.Close()
|
||||||
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()
|
|
||||||
var said string
|
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 errors.New(said)
|
||||||
}
|
}
|
||||||
return nil
|
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 {
|
say Announce, timeout time.Duration) bool {
|
||||||
report.Node = m.Node
|
report.Node = m.Node
|
||||||
body, err := json.Marshal(report)
|
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
|
// 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
|
// 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.
|
// 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,
|
if err := bus.Report(publish, m.Node, body); err != nil {
|
||||||
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
|
||||||
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
|
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
|
||||||
return false
|
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.
|
// 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) {
|
timeout time.Duration) {
|
||||||
body, err := json.Marshal(Alive{Node: m.Node})
|
body, err := json.Marshal(Alive{Node: m.Node})
|
||||||
if err != nil {
|
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
|
// 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
|
// mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a
|
||||||
// harmless miss into a logged failure every minute.
|
// harmless miss into a logged failure every minute.
|
||||||
if err := channel.PublishWithContext(publish, Exchange, KeyAlive, false, false,
|
if err := bus.Alive(publish, m.Node, body); err != nil {
|
||||||
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
|
||||||
say("could not tell the mesh this node is here: " + err.Error())
|
say("could not tell the mesh this node is here: " + err.Error())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user