A catalogue asks for what it was not there to hear
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
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user