diff --git a/cmd/mesh-control/push.go b/cmd/mesh-control/push.go index e4f0c93..8eb13f2 100644 --- a/cmd/mesh-control/push.go +++ b/cmd/mesh-control/push.go @@ -92,6 +92,11 @@ func serve(ctx context.Context) error { if err := server.Follows(following{open}); err != nil { return err } + // And a catalogue that has just started, asking for what it missed. The same type answers + // both: what a build meant and what the builds were are two questions about one record. + if err := server.Answers(following{open}); err != nil { + return err + } return server.Serve(ctx) } diff --git a/cmd/mesh-control/upgrades.go b/cmd/mesh-control/upgrades.go index f4ba58c..2823d0a 100644 --- a/cmd/mesh-control/upgrades.go +++ b/cmd/mesh-control/upgrades.go @@ -163,3 +163,34 @@ func sayUpgrade(module string, u inventory.Upgrade) string { return fmt.Sprintf("when %s moves, the machines running it are sent the new version one at "+ "a time, stopping at the first that fails", module) } + +// Announceable is every build this mesh recorded, in the shape the builder announces one. +// +// **The catalogue asks for this when it starts, and the answer is the graph's foundation** +// (novox/hq 04-ISSUES/050). A durable queue keeps what arrived after it existed, so a running +// catalogue misses nothing — but the modules built before it first ran were announced to a queue +// that did not exist, and on a fresh mesh those are always the same three: the shared base, the +// store the catalogue runs on, and the catalogue itself. +func (f following) Announceable(ctx context.Context) ([]link.Announcement, error) { + builds, err := f.open.inventory.Announceable(ctx) + if err != nil { + return nil, err + } + out := make([]link.Announcement, 0, len(builds)) + for _, b := range builds { + a := link.Announcement{ + Module: b.Module, Commit: b.Commit, Repository: b.Repository, + Path: b.Path, Ref: b.Ref, Against: b.Against, + } + if len(b.Manifest) > 0 { + a.Manifest = b.Manifest + } + for _, made := range b.Made { + a.Made = append(a.Made, link.MadeArtifact{ + Name: made.Name, Kind: made.Kind, Reference: made.Reference, + }) + } + out = append(out, a) + } + return out, nil +} diff --git a/internal/inventory/builds.go b/internal/inventory/builds.go index ef07fdb..037a531 100644 --- a/internal/inventory/builds.go +++ b/internal/inventory/builds.go @@ -3,6 +3,7 @@ package inventory import ( "context" "encoding/json" + "sort" "time" ) @@ -174,3 +175,63 @@ func manifestOrNil(raw []byte) any { } return raw } + +// Announceable is every build worth telling a catalogue about, oldest first. +// +// **Oldest first, because a graph is built in the order things happened.** Registering a module +// that stands on a base before the base itself would make the edge point at a version the +// catalogue has not seen, and the shape of a fresh mesh guarantees that order matters: the base is +// always first and always the one that was missed. +// +// Only builds that succeeded and know what they built. A failure produced no module-version, and +// announcing one would put something in the graph that was never made — the same rule the builder +// follows when it decides whether to announce at all. +// +// One row per module and commit: a module built twice at the same commit is one fact, and the +// latest row is the one whose artifacts are current. +func (i *Inventory) Announceable(ctx context.Context) ([]Build, error) { + rows, err := i.store.Pool().Query(ctx, + `select distinct on (module, commit_hash) + id, repository, ref, module, commit_hash, built_on, failed, made, + source_path, manifest, built_against, at + from build + where failed = '' and module is not null and module <> '' and commit_hash <> '' + order by module, commit_hash, at desc`) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []Build + for rows.Next() { + var b Build + var made []byte + var manifest, against []byte + if err := rows.Scan(&b.ID, &b.Repository, &b.Ref, &b.Module, &b.Commit, + &b.On, &b.Failed, &made, &b.Path, &manifest, &against, &b.At); err != nil { + return nil, err + } + if err := json.Unmarshal(made, &b.Made); err != nil { + return nil, err + } + // Null rather than empty is a build recorded before the mesh kept these, and saying so is + // the point of keeping them nullable: the replay carries nothing rather than an empty + // declaration for a module that certainly had one. + if len(manifest) > 0 { + b.Manifest = manifest + } + if len(against) > 0 { + if err := json.Unmarshal(against, &b.Against); err != nil { + return nil, err + } + } + out = append(out, b) + } + if err := rows.Err(); err != nil { + return nil, err + } + // Sorted here rather than in the query, because `distinct on` fixes the ordering it needs and + // the order that matters to a catalogue is a different one. + sort.Slice(out, func(a, b int) bool { return out[a].At.Before(out[b].At) }) + return out, nil +} diff --git a/internal/link/events.go b/internal/link/events.go index 7cb62fb..3d461a3 100644 --- a/internal/link/events.go +++ b/internal/link/events.go @@ -80,10 +80,60 @@ const KeyModuleBuilt = "module.builder.built" // 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"` diff --git a/internal/link/serve.go b/internal/link/serve.go index fef6712..2d402e7 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -51,6 +51,7 @@ type Server struct { recorder Recorder log *log.Logger upgrader Upgrader + replayer Replayer } // Records tells the server where to keep build results. @@ -89,6 +90,21 @@ func (s *Server) Follows(u Upgrader) error { return nil } +// Answers binds the queue a catalogue's catch-up request arrives on. +// +// **Not bound unless something is listening**, for the same reason upgrades are not: a durable +// queue with no consumer fills quietly and the first symptom is a broker out of disk. +func (s *Server) Answers(r Replayer) error { + if _, err := s.channel.QueueDeclare(CatchUpQueue, true, false, false, false, nil); err != nil { + return fmt.Errorf("cannot declare the %s queue: %w", CatchUpQueue, err) + } + if err := s.channel.QueueBind(CatchUpQueue, KeyCatchingUp, EventsExchange, false, nil); err != nil { + return fmt.Errorf("cannot bind %s to %s/%s: %w", CatchUpQueue, EventsExchange, KeyCatchingUp, err) + } + s.replayer = r + return nil +} + // Connect opens the control plane's own connection to the broker. func Connect(enroller Enroller, listener Listener) (*Server, error) { url := strings.TrimSpace(os.Getenv(AMQPVar)) @@ -187,6 +203,18 @@ func (s *Server) Serve(ctx context.Context) error { } } + // Its own queue and its own consumer, for the reason above: 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. + var catchups <-chan amqp.Delivery + if s.replayer != nil { + catchups, err = s.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup", + false, false, false, false, nil) + if err != nil { + return err + } + } + closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}", ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt) @@ -198,6 +226,14 @@ func (s *Server) Serve(ctx context.Context) error { select { case <-ctx.Done(): return nil + case delivery, ok := <-catchups: + if !ok { + if catchups != nil { + return errors.New("the broker stopped delivering catch-up requests") + } + continue + } + s.catchingUp(ctx, delivery) case delivery, ok := <-upgrades: // A nil channel blocks for ever, so this case simply never fires when nothing is // listening for upgrades. Closed is different, and means the broker stopped. @@ -379,6 +415,37 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { // none of those get better by being handed the same message again. Requeuing would put a poison // message at the head of a durable queue and stop every upgrade behind it, which turns one module // nobody can push into a mesh that stops following its own catalogue. +// catchingUp answers a catalogue that has just started and may have missed builds. +// +// Acknowledged before the work, deliberately: a replay that fails is not one that succeeds by +// being handed the same request again, and the catalogue asks every time it starts. Requeueing a +// poison request would stop every later catch-up behind it. +func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { + defer func() { _ = delivery.Ack(false) }() + if s.replayer == nil { + s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") + return + } + announcements, err := s.replayer.Announceable(ctx) + if err != nil { + s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err) + return + } + sent := 0 + for _, a := range announcements { + a.Replay = true + if err := EmitEvent(ctx, 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", + a.Module, short(a.Commit), err) + return + } + sent++ + } + s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent) +} + func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { defer func() { _ = delivery.Ack(false) }() var u Upgraded