package main import ( "context" "fmt" "sync" "time" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/link" ) // **A merge the bus announced and the controller never acted on is caught up** (novox/hq issue 266). // // The forge's poll announces every merge on the events stream, and the controller acts on what its // consumer there hands it. On 2026-10-06 one merge was on the stream and never handed over: the bus // server moved the consumer past it — a fault of consumers with several filters in the server the // mesh ran — and the controller, which only acts on what it is handed, said nothing. The modules // built from that repository stayed behind and were built by hand. // // So the stream is read back on a timer, on a consumer of its own filtered on merges alone, and every // announcement older than mergeGrace is judged as SourceMoved would judge it. **No record of what // was handled is kept, because none is needed**: acting on a merge marks every module it moved as // looked at since, so an announcement already acted on reads as history and moves nothing. One that // would still move something was never acted on — it is said, and acted on now. const ( // mergeGrace is how long an announcement is left to the controller's own consumer before it is // judged missed. That consumer hands over one event at a time, and a merge waits behind a build // outcome that is being acted on; acting on a merge itself asks builds and does not wait for them. mergeGrace = 10 * time.Minute // mergeLookBack is how far back a pass reads. A merge missed longer ago than this was missed by a // controller that was not running this, and is the operator's to look at, not a surprise rebuild. mergeLookBack = 24 * time.Hour // mergeCatchUpEvery is how often the stream is read back. mergeCatchUpEvery = 5 * time.Minute ) // merges is what reads back the forge's announcements; the link server, or a test's list. type merges interface { AnnouncedMerges(ctx context.Context, since time.Time) ([]link.AnnouncedMerge, error) } // catchingUpOnMerges reads back the forge's announcements on a timer, until the context ends. func catchingUpOnMerges(ctx context.Context, open *stores, announced merges) { f := following{open} catalogued := func(ctx context.Context) ([]inventory.Entry, map[string][]inventory.ReadRepository, error) { entries, err := open.inventory.Catalogued(ctx) if err != nil { return nil, nil, err } read, err := open.inventory.ReadRepositories(ctx) return entries, read, err } failing := "" tick := time.NewTicker(mergeCatchUpEvery) defer tick.Stop() for { select { case <-ctx.Done(): return case <-tick.C: } watchedMerges.begin() err := catchUpOnMerges(ctx, time.Now(), announced, catalogued, f.SourceMoved, func(format string, args ...any) { fmt.Printf(format+"\n", args...) }) // What the pass found is what S5 says (novox/hq to-be 45 §3); a pass that could not read // says that instead, and leaves what the last one found standing. watchedMerges.end(time.Now(), err) // A pass that cannot read says so once, not every five minutes, and says when it reads again. why := "" if err != nil { why = err.Error() } if why != failing { if why != "" { fmt.Printf("merges the bus may not have handed over cannot be looked for: %s\n", why) } else { fmt.Println("merges the bus may not have handed over are looked for again") } failing = why } } } // catchUpOnMerges is one pass: every announcement older than mergeGrace that acting on would still // move something is said and acted on, oldest first. // // Judged twice: once against the catalogue as the pass found it, and again just before acting, // because acting on an earlier missed merge of the same repository may have moved what a later one // would have. func catchUpOnMerges(ctx context.Context, now time.Time, announced merges, catalogued func(context.Context) ([]inventory.Entry, map[string][]inventory.ReadRepository, error), act func(context.Context, link.SourceMoved) error, say func(string, ...any)) error { all, err := announced.AnnouncedMerges(ctx, now.Add(-mergeLookBack)) if err != nil { return err } entries, read, err := catalogued(ctx) if err != nil { return err } for _, a := range all { if now.Sub(a.At) < mergeGrace { continue } if len(wouldMove(a.SourceMoved, entries, read)) == 0 { continue } if entries, read, err = catalogued(ctx); err != nil { return err } moves := wouldMove(a.SourceMoved, entries, read) if len(moves) == 0 { continue } var names []string for _, e := range moves { names = append(names, e.Manifest.Module) } watchedMerges.found(missedMerge{Owner: a.Owner, Repo: a.Repo, Base: a.Base, Commit: a.Commit, At: a.At, Modules: names}) say("%s/%s merged into %s (%.8s), announced %s ago, and the controller never acted on it: the bus "+ "did not hand the announcement over (novox/hq issue 266). %s %s behind it; acting on it now", a.Owner, a.Repo, a.Base, a.Commit, now.Sub(a.At).Round(time.Minute), readableList(names), isAre(len(names))) if err := act(ctx, a.SourceMoved); err != nil { say("%s/%s moved to %.8s and the mesh could not act on it: %v; the next pass tries again", a.Owner, a.Repo, a.Commit, err) continue } if entries, read, err = catalogued(ctx); err != nil { return err } } return nil } // missedMerge is one merge the bus announced and never handed over, as a pass found it. type missedMerge struct { Owner, Repo, Base, Commit string // At is when the bus took the announcement. At time.Time Modules []string } // mergeWatch is what the passes found, for S5 (novox/hq to-be 45 §3): **a merge nothing read is // said**, urgent, even though the pass acts on it at once — the bus skipping a message is a fault of // the transport the mesh's every change rides on, and acting late is the repair, not the absence of // the fault. The next pass, finding it acted on, clears it. type mergeWatch struct { mu sync.Mutex passed time.Time err error finding []missedMerge missed []missedMerge } var watchedMerges = &mergeWatch{} func (w *mergeWatch) begin() { w.mu.Lock() defer w.mu.Unlock() w.finding = nil } func (w *mergeWatch) found(m missedMerge) { w.mu.Lock() defer w.mu.Unlock() w.finding = append(w.finding, m) } func (w *mergeWatch) end(at time.Time, err error) { w.mu.Lock() defer w.mu.Unlock() w.passed, w.err = at, err if err == nil { w.missed = w.finding } } // last is when the last pass ended, what the last pass that read found, and what the last pass // could not read. func (w *mergeWatch) last() (time.Time, []missedMerge, error) { w.mu.Lock() defer w.mu.Unlock() return w.passed, append([]missedMerge(nil), w.missed...), w.err }