package link import ( "context" "crypto/rand" "encoding/hex" "encoding/json" "fmt" "time" amqp "github.com/rabbitmq/amqp091-go" ) // Emitting a module event from Go. // // **Every event rides one topic exchange** (novox/hq ADR 0042), which is not the direct exchange // nodes and the control plane speak over. A module that announces something publishes here, and // consumers bind their own durable queue to a pattern over it. // // This exists because the builder is a module written in Go while every other emitter is // TypeScript on the sdk. The envelope is the sdk's, reproduced exactly: the body is the payload // alone and everything about the event travels as headers. A second shape would be a second thing // for consumers to handle, and they are written against the first. const ( // EventsExchange is where every event rides. Named here rather than imported from the broker // package for the same reason BuildQueueName is duplicated there — one direction of dependency. EventsExchange = "mesh.events" ) // 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 { 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", }, }) } // eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able // to tell a redelivery from a second event, and only the emitter can say which it is. func eventID() (string, error) { raw := make([]byte, 16) if _, err := rand.Read(raw); err != nil { return "", fmt.Errorf("cannot make an event id: %w", err) } return hex.EncodeToString(raw), nil } // KeyModuleBuilt is what the builder announces when it has built something. The catalogue places // it in the module graph; nothing else need care. const KeyModuleBuilt = "module.builder.built" // KeyModuleUpgraded is the catalogue saying a module's current version has moved. // // **The control plane hooks the meaning, not the build.** The builder says what it built; the // catalogue decides whether that was an upgrade — a rebuild producing the commit already current // is not one — and only this says anything the control plane can act on. Consuming the build // directly would make the control plane re-derive a decision another module already made, and the // two would eventually disagree (novox/hq ADR 0072). const KeyModuleUpgraded = "module.mesh-catalog.upgraded" // UpgradeQueue is where those land. Durable and named, not a temporary queue: an upgrade announced // while the control plane is restarting is exactly the one that must not be missed. const UpgradeQueue = "control.upgrades" // Upgraded is what the catalogue says when a module's current version moves. type Upgraded struct { Module string `json:"module"` Commit string `json:"commit"` Previous string `json:"previous"` }