From 6c12780abe7a85f9533e4be75de1507d31dd679c Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 26 Sep 2026 23:51:00 +0200 Subject: [PATCH] Describe the bus on its own terms MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Comments framed the new bus by what it replaces — a comparison in almost every explanation, which reads as though NATS were a variant of the old thing rather than the mesh's nervous system. Removed throughout, and OverAMQP becomes OverCurrent: the seam's two sides are the bus the mesh runs on today and the one being built, not two protocols. What remains is the client library's own package name, which is its name. --- cmd/mesh-builder/main.go | 2 +- cmd/mesh-controller/push.go | 8 ++++---- internal/broker/nats.go | 22 +++++++++------------- internal/link/bus.go | 25 ++++++++++++------------- internal/link/serve.go | 2 +- 5 files changed, 27 insertions(+), 32 deletions(-) diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index c6266f7..f321436 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -268,7 +268,7 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis "manifest": json.RawMessage(result.Manifest), "against": result.Against, "made": result.Made, } - if err := link.EmitEvent(publishCtx, link.OverAMQP{Channel: channel}, link.KeyModuleBuilt, "builder", on, announced); err != nil { + if err := link.EmitEvent(publishCtx, link.OverCurrent{Channel: channel}, link.KeyModuleBuilt, "builder", on, announced); err != nil { // Said, not fatal: the build happened and was answered. A module the catalogue has not // heard of is a gap somebody can close; a build reported as failed because announcing // it failed is a lie about work that was done. diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 7e9c154..2ee6875 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -142,7 +142,7 @@ func declare(ctx context.Context, args []string) error { } defer server.Close() - if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { return err } fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) @@ -310,7 +310,7 @@ func pushCommand(ctx context.Context, args []string) error { if err != nil { return err } - if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { return err } // After it is away, not before. A digest recorded for something that failed to send would @@ -393,7 +393,7 @@ func pushCommand(ctx context.Context, args []string) error { return declarationWith(held, open, node, plan, settings, gens, Allocating) }, func(s readyNode, body []byte) error { - if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body, + if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { return err } @@ -611,7 +611,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error { if err != nil { return err } - if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { return err } record, err := inv.NodeByName(ctx, s.node) diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 71c3356..4a09ba3 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -1,19 +1,15 @@ // Composing the bus's own configuration. // -// On AMQP an account was made by calling the broker's management API (management.go). On NATS it -// is *composed*: the controller writes accounts, users and per-subject permissions into one file -// the host keeps current, and the server reloads it in place (novox/hq ADR 0106 — never through a -// management API; design 25 §4). +// An account is *composed*, never called for: the controller writes accounts, users and +// per-subject permissions into one file the host keeps current, and the server reloads it in +// place (novox/hq ADR 0106 — never through a management API; design 25 §4). // // Everything here is pure. Given the principals, it returns the file's text — so the whole of the -// mesh's authority model is testable as strings, with no server, which is what management.go's -// `modulePermissions` already did for the half of it that could be. +// mesh's authority model is testable as strings, with no server. // -// **NATS closes a gap AMQP left open.** management.go records it plainly: LavinMQ has no topic -// permissions, so an emitting module is granted the events exchange whole, and ADR 0042's origin -// reservation — a module publishes only under its own name — is "stamped by the sdk, not enforced -// here". NATS permissions are per subject, so that reservation becomes something the server -// refuses rather than something a library promises. +// **Permissions are per subject, so a module's own name is the server's to enforce.** ADR 0042 +// reserves a module's origin — it publishes only under its own name — and here that is a refusal +// rather than something a library promises. package broker import ( @@ -270,8 +266,8 @@ func consumerDurable(p Principal) string { // Server is everything the composed file needs that is not a principal. type Server struct { - // ClientPort carries TLS itself; there is no plaintext port beside it, which is where this - // differs from the AMQP broker's 5671/5672 pair. + // ClientPort carries TLS itself. There is no plaintext port beside it: a bus reachable + // without TLS is one a module can reach without TLS by mistake. ClientPort int MonitoringPort int TLSCert string diff --git a/internal/link/bus.go b/internal/link/bus.go index b222955..8b767e3 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -12,15 +12,14 @@ import ( // Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a // transport's** (novox/hq ADR 0116 step 3). // -// Until now every one of these functions took an `*amqp.Channel`, so the transport reached every -// caller and swapping it meant touching all of them. The seam is small — the controller sends -// exactly two kinds of message that expect no answer, and asks two kinds of question — which is -// why the bus could be replaced at all. +// Until now every one of these functions took the transport's own channel type, so the transport +// reached every caller and changing it meant touching all of them. The seam is small — the +// controller sends exactly two kinds of message that expect no answer, and asks two kinds of +// question — which is why the bus can be replaced at all. // -// Two implementations live below. Both ship: steps 1 to 4 leave every node on AMQP -// ([ADR 0116](novox/hq)), so the controller keeps speaking it and the NATS one is selected at the -// rollout. That is also what makes them comparable — the same caller, the same arguments, and a -// conformance fixture holding both to one envelope. +// Two implementations live below, and both ship until the rollout (ADR 0116: nothing moves a +// node's bus before step 5). Both shipping is what makes them comparable — the same caller, the +// same arguments, and one conformance fixture holding them to one envelope. type Bus interface { // PublishEvent announces something that happened, under the emitter's own name. 1:many, and // nobody is obliged to act (ADR 0041). @@ -32,12 +31,12 @@ type Bus interface { PublishDeclaration(ctx context.Context, node string, body []byte) error } -// --- AMQP, the bus the mesh runs on today ----------------------------------------------------- +// --- The bus the mesh runs on today ----------------------------------------------------- -// OverAMQP is the bus as a channel. -type OverAMQP struct{ Channel *amqp.Channel } +// OverCurrent is the bus the mesh runs on today, until the rollout. +type OverCurrent struct{ Channel *amqp.Channel } -func (b OverAMQP) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { +func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { id, err := eventID() if err != nil { return err @@ -58,7 +57,7 @@ func (b OverAMQP) PublishEvent(ctx context.Context, key, source, node string, bo }) } -func (b OverAMQP) PublishDeclaration(ctx context.Context, node string, body []byte) error { +func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error { // To the queue directly rather than through an exchange: a declaration is for one node, and // routing it by name through a shared exchange would mean a binding per node that nothing // removes when a node is retired. diff --git a/internal/link/serve.go b/internal/link/serve.go index 26ef006..b61f633 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -631,7 +631,7 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { sent := 0 for _, a := range announcements { a.Replay = true - if err := EmitEvent(ctx, OverAMQP{Channel: s.channel}, KeyModuleBuilt, "control-plane", "", a); err != nil { + if err := EmitEvent(ctx, OverCurrent{Channel: s.channel}, KeyModuleBuilt, "control-plane", "", a); err != nil { // Said and abandoned rather than retried: the catalogue asks again every time it // starts, and half a graph delivered twice is no better than half delivered once. s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",