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",