package link import ( "context" "fmt" "time" "github.com/nats-io/nats.go" amqp "github.com/rabbitmq/amqp091-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 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, 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). 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 } // --- The bus the mesh runs on today ----------------------------------------------------- // OverCurrent is the bus the mesh runs on today, until the rollout. type OverCurrent struct{ Channel *amqp.Channel } func (b OverCurrent) 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 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. 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 }