diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index c4fa5e1..c6266f7 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, channel, link.KeyModuleBuilt, "builder", on, announced); err != nil { + if err := link.EmitEvent(publishCtx, link.OverAMQP{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 2319213..7e9c154 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, server.Channel(), ident, node, raw, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverAMQP{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, server.Channel(), ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverAMQP{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, server.Channel(), ident, s.node, body, + if err := link.Declare(ctx, link.OverAMQP{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, server.Channel(), ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, link.OverAMQP{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/link/bus.go b/internal/link/bus.go new file mode 100644 index 0000000..a353c58 --- /dev/null +++ b/internal/link/bus.go @@ -0,0 +1,121 @@ +package link + +import ( + "context" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" + "github.com/nats-io/nats.go" +) + +// 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. +// +// 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. +type Bus interface { + // PublishEvent announces something that happened, under the emitter's own name. 1:many, and + // nobody is obliged to act (ADR 0041). + PublishEvent(ctx context.Context, key, source, node string, body []byte) error + + // PublishDeclaration delivers one node what it should be. Addressed to that node alone: a + // declaration is not an event, and replaying yesterday's is actively harmful + // (design 29 §4, the *state* shape). + PublishDeclaration(ctx context.Context, node string, body []byte) error +} + +// --- AMQP, the bus the mesh runs on today ----------------------------------------------------- + +// OverAMQP is the bus as a channel. +type OverAMQP struct{ Channel *amqp.Channel } + +func (b OverAMQP) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { + id, err := eventID() + if err != nil { + return err + } + return b.Channel.PublishWithContext(ctx, EventsExchange, key, false, false, amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + MessageId: id, + Timestamp: time.Now().UTC(), + Body: body, + Headers: amqp.Table{ + "x-event-id": id, + "x-source": source, + "x-node": node, + "x-time": time.Now().UTC().Format(time.RFC3339), + "content-type": "application/json", + }, + }) +} + +func (b OverAMQP) 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. + return b.Channel.PublishWithContext(ctx, "", QueueFor(node), false, false, amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + Body: body, + }) +} + +// --- NATS, the bus being built ---------------------------------------------------------------- + +// OverNATS is the bus as a JetStream context. +type OverNATS struct{ JS nats.JetStreamContext } + +// EventSubject is where a module's event lands. Derived from the emitter, never taken from the +// caller: a source that could differ from the subject is an envelope that can lie about its +// origin, and on NATS the account's permissions make the subject the authority (design 29 §2). +func EventSubject(source, key string) string { + return "mesh.mod." + source + ".event." + key +} + +// DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so +// a node that was away gets exactly the current one and a replayed older one is refused by +// sequence — the wire-level answer to novox/hq issue 107. +func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" } + +func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { + id, err := eventID() + if err != nil { + return err + } + h := nats.Header{} + h.Set("x-event-id", id) + h.Set("x-source", source) + h.Set("x-node", node) + h.Set("x-time", time.Now().UTC().Format(time.RFC3339)) + h.Set("content-type", "application/json") + + // The id is also the publish's message id, so the server refuses a duplicate inside its + // window. That narrows the window a consumer must deduplicate in; it does not remove the + // requirement, because the window is finite (design 19, delivery). + _, err = b.JS.PublishMsg(&nats.Msg{ + Subject: EventSubject(source, key), + Header: h, + Data: body, + }, nats.MsgId(id), nats.Context(ctx)) + if err != nil { + return fmt.Errorf("emitting %s: %w", key, err) + } + return nil +} + +func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error { + _, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx)) + if err != nil { + return fmt.Errorf("declaring to %s: %w", node, err) + } + return nil +} diff --git a/internal/link/bus_test.go b/internal/link/bus_test.go new file mode 100644 index 0000000..afd5f94 --- /dev/null +++ b/internal/link/bus_test.go @@ -0,0 +1,95 @@ +package link + +import ( + "context" + "encoding/json" + "os" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// **Both implementations, one fixture.** The point of the seam is not that the transport can be +// swapped — it is that the two can be held to the same envelope while both ship, so the day the +// bus moves is a configuration change rather than a discovery. +// +// Against a real server, because what the fixture pins is what reaches the wire: +// +// docker run -d --rm --name t -p 14222:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/link/ -run TestTheNatsBus +func TestTheNatsBusEmitsTheEnvelopeTheFixturePins(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + f := loadFixture(t, "events/module-event.json") + + conn, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + js, err := conn.JetStream() + if err != nil { + t.Fatal(err) + } + if _, err := js.AddStream(&nats.StreamConfig{ + Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}, + }); err != nil && err != nats.ErrStreamNameAlreadyInUse { + t.Fatal(err) + } + + bus := OverNATS{JS: js} + body, _ := json.Marshal(f.Given.Body) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := EmitEvent(ctx, bus, f.Given.Key, f.Given.Module, f.Given.Node, f.Given.Body); err != nil { + t.Fatal(err) + } + + // Read back from the stream, not from the thing that wrote it. + raw, err := js.GetLastMsg("EVENTS", f.Wire.Subject) + if err != nil { + t.Fatalf("nothing landed on %s, which the fixture names: %v", f.Wire.Subject, err) + } + for _, h := range f.Wire.RequiredHeaders { + if raw.Header.Get(h) == "" { + t.Errorf("%s is not set on the wire, and the fixture requires it", h) + } + } + if got := raw.Header.Get("x-source"); got != f.Given.Module { + t.Errorf("x-source is %q; the subject says %q", got, f.Given.Module) + } + if string(raw.Data) != string(body) { + t.Errorf("the payload is %s, expected the body alone: %s", raw.Data, body) + } + // The envelope must not also be nested inside the payload. + var nested map[string]any + if json.Unmarshal(raw.Data, &nested) == nil { + if _, has := nested["key"]; has { + t.Error("the payload carries the envelope's own fields, which the fixture refuses") + } + } +} + +// The subject a declaration lands on is one node's, and nothing else's — the state shape. +func TestADeclarationIsAddressedToOneNode(t *testing.T) { + if got := DeclareSubject("anchor"); got != "mesh.node.anchor.declare" { + t.Fatalf("a declaration would go to %q", got) + } + if DeclareSubject("anchor") == DeclareSubject("laptop") { + t.Fatal("two nodes share a declaration subject, so each would apply the other's") + } +} + +// A module cannot emit under another's name: the subject is derived from the source, and the +// server's permissions make that subject the authority. +func TestAnEventsSubjectIsDerivedFromItsSource(t *testing.T) { + if got := EventSubject("shop", "order.placed"); got != "mesh.mod.shop.event.order.placed" { + t.Fatalf("an event would land on %q", got) + } + if EventSubject("shop", "x") == EventSubject("billing", "x") { + t.Fatal("two modules share an event subject, so neither owns its own name") + } +} diff --git a/internal/link/declare.go b/internal/link/declare.go index 9f6fd7e..665c810 100644 --- a/internal/link/declare.go +++ b/internal/link/declare.go @@ -5,8 +5,6 @@ import ( "encoding/json" "fmt" "time" - - amqp "github.com/rabbitmq/amqp091-go" ) // Signer is whatever holds the control plane's signing key. @@ -21,7 +19,7 @@ type Signer interface { // something else, and the node would refuse a declaration that was genuinely the mesh's. // // Published to the node's own queue, which its account alone may read. -func Declare(ctx context.Context, channel *amqp.Channel, signer Signer, node string, +func Declare(ctx context.Context, bus Bus, signer Signer, node string, declaration []byte, timeout time.Duration) error { if !json.Valid(declaration) { @@ -41,13 +39,5 @@ func Declare(ctx context.Context, channel *amqp.Channel, signer Signer, node str publish, cancel := context.WithTimeout(ctx, timeout) defer cancel() - // Published to the queue directly rather than through the 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. - return channel.PublishWithContext(publish, "", QueueFor(node), false, false, - amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - Body: body, - }) + return bus.PublishDeclaration(publish, node, body) } diff --git a/internal/link/events.go b/internal/link/events.go index 3d461a3..225a19a 100644 --- a/internal/link/events.go +++ b/internal/link/events.go @@ -6,9 +6,6 @@ import ( "encoding/hex" "encoding/json" "fmt" - "time" - - amqp "github.com/rabbitmq/amqp091-go" ) // Emitting a module event from Go. @@ -29,32 +26,17 @@ const ( // EmitEvent publishes one module event, in the envelope the sdk's consumers expect. // -// Persistent, because an event that a broker restart loses is not an announcement. The publish is -// not confirmed here: the caller has already done the work the event describes, and a build that -// succeeded must not be reported as failed because saying so failed. -func EmitEvent(ctx context.Context, channel *amqp.Channel, eventType, source, node string, body any) error { +// The envelope is the transport's to write (bus.go) and this is only what goes in it, which is +// what lets one conformance fixture hold both implementations to the same headers. +func EmitEvent(ctx context.Context, bus Bus, eventType, source, node string, body any) error { payload, err := json.Marshal(body) if err != nil { return fmt.Errorf("cannot serialise a %s event: %w", eventType, err) } - id, err := eventID() - if err != nil { - return err - } - return channel.PublishWithContext(ctx, EventsExchange, eventType, false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - MessageId: id, - Timestamp: time.Now().UTC(), - Body: payload, - Headers: amqp.Table{ - "x-event-id": id, - "x-source": source, - "x-node": node, - "x-time": time.Now().UTC().Format(time.RFC3339), - "content-type": "application/json", - }, - }) + // The publish is not confirmed by the caller: it has already done the work the event + // describes, and a build that succeeded must not be reported as failed because saying so + // failed. Each transport decides what "published" means for it. + return bus.PublishEvent(ctx, eventType, source, node, payload) } // eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able diff --git a/internal/link/serve.go b/internal/link/serve.go index 925c237..26ef006 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, s.channel, KeyModuleBuilt, "control-plane", "", a); err != nil { + if err := EmitEvent(ctx, OverAMQP{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",