Its event queue is durable, so a running catalogue misses nothing. What it cannot have is what was announced before it first ran — and on a fresh mesh that is never arbitrary: the shared base, the store the catalogue runs on, and the catalogue itself are each necessarily built BEFORE a catalogue exists to hear about them. The graph's foundation is the part it never sees. So it says it is catching up, and the control plane re-announces what it recorded, oldest first, marked as a replay. Oldest first because a graph is built in the order things happened: registering a module that stands on a base before the base would point an edge at a version nothing has seen, and the shape of a fresh mesh guarantees the base is both first and the one that was missed. The replayer hands announcements back rather than publishing them, because the wire belongs to the link package and a replay building its own events could drift from what the builder emits — the one thing it must match exactly, since the catalogue has a single handler for both. Its own queue and its own consumer: two consumers on one queue split its messages, and a catch-up request going to whichever half was not listening is a gap that looks like a working mesh. Toward novox/hq 04-ISSUES/050. Claude-Session: https://claude.ai/code/session_01D6qtiYU3P9jk3pnAXyAFyx
143 lines
6.3 KiB
Go
143 lines
6.3 KiB
Go
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"
|
|
|
|
// KeyCatchingUp is the catalogue saying it has just started and may have missed things.
|
|
//
|
|
// **A durable queue only keeps what arrived after it existed.** The catalogue's own queue is
|
|
// durable, so nothing is lost once it is running — but the modules built before it first ran were
|
|
// announced to a queue that did not exist yet, and on a fresh mesh those are, necessarily, the
|
|
// shared base, the store the catalogue runs on, and the catalogue itself. The graph's foundation
|
|
// is the part it never hears about (novox/hq 04-ISSUES/050).
|
|
//
|
|
// So it asks, and the control plane answers with what it recorded. Asking rather than being told
|
|
// because only the catalogue knows it has a gap; the control plane cannot tell a fresh catalogue
|
|
// from one that is merely quiet.
|
|
const KeyCatchingUp = "module.mesh-catalog.catching-up"
|
|
|
|
// CatchUpQueue is where that lands. Durable, for the same reason the upgrade queue is: a catalogue
|
|
// that started while the control plane was restarting is exactly the one with a gap to fill.
|
|
const CatchUpQueue = "control.catchup"
|
|
|
|
// 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"
|
|
|
|
// Replayer answers a catalogue that says it has just started.
|
|
//
|
|
// It is handed every build the mesh recorded, oldest first, and re-announces each. The catalogue
|
|
// registers them as history: a replayed build changed nothing in the world, so announcing it as an
|
|
// upgrade would have the mesh act on news that is years old.
|
|
type Replayer interface {
|
|
// Announceable is every build worth re-announcing, oldest first.
|
|
//
|
|
// It hands them back rather than publishing them: the wire belongs to this package, and a
|
|
// replay that built its own announcements could drift from what the builder emits — which is
|
|
// the one thing it must match exactly, because the catalogue has a single handler for both.
|
|
Announceable(ctx context.Context) ([]Announcement, error)
|
|
}
|
|
|
|
// Announcement is a build, in the shape the builder announces one.
|
|
//
|
|
// The field names are the wire's, not Go's, because a catalogue reads these and a rename here is
|
|
// an event nobody handles.
|
|
type Announcement struct {
|
|
Module string `json:"module"`
|
|
Commit string `json:"commit"`
|
|
Repository string `json:"repository"`
|
|
Path string `json:"path"`
|
|
Ref string `json:"ref"`
|
|
Manifest json.RawMessage `json:"manifest,omitempty"`
|
|
Against []string `json:"against,omitempty"`
|
|
Made []MadeArtifact `json:"made,omitempty"`
|
|
// Replay says this is history rather than news: it was built once, and this is the mesh
|
|
// telling a catalogue that missed it. A consumer registers it and announces nothing — an
|
|
// upgrade that happened months ago is not one anything should act on now.
|
|
Replay bool `json:"replay,omitempty"`
|
|
}
|
|
|
|
// 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"`
|
|
}
|