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()) } }