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