diff --git a/cmd/mesh-controller/missed_merges.go b/cmd/mesh-controller/missed_merges.go new file mode 100644 index 0000000..4a83e46 --- /dev/null +++ b/cmd/mesh-controller/missed_merges.go @@ -0,0 +1,131 @@ +package main + +import ( + "context" + "fmt" + "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: + } + err := catchUpOnMerges(ctx, time.Now(), announced, catalogued, f.SourceMoved, func(format string, args ...any) { + fmt.Printf(format+"\n", args...) + }) + // 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) + } + 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 +} diff --git a/cmd/mesh-controller/missed_merges_test.go b/cmd/mesh-controller/missed_merges_test.go new file mode 100644 index 0000000..f4df12d --- /dev/null +++ b/cmd/mesh-controller/missed_merges_test.go @@ -0,0 +1,180 @@ +package main + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// announcedList is the events stream's merges as a test gives them. +type announcedList []link.AnnouncedMerge + +func (a announcedList) AnnouncedMerges(_ context.Context, since time.Time) ([]link.AnnouncedMerge, error) { + var out []link.AnnouncedMerge + for _, m := range a { + if !m.At.Before(since) { + out = append(out, m) + } + } + return out, nil +} + +// aCatalogue is what the inventory holds, changed the way acting on a merge changes it: every module +// the merge moved is marked as looked at (inventory.SourceMoved writes source_seen = now()). +type aCatalogue struct { + entries []inventory.Entry + acted []string + fail error +} + +func (c *aCatalogue) read(context.Context) ([]inventory.Entry, map[string][]inventory.ReadRepository, error) { + return append([]inventory.Entry(nil), c.entries...), nil, nil +} + +func (c *aCatalogue) act(now func() time.Time) func(context.Context, link.SourceMoved) error { + return func(_ context.Context, m link.SourceMoved) error { + if c.fail != nil { + return c.fail + } + c.acted = append(c.acted, m.Repo+"@"+m.Commit[:8]) + for _, moved := range wouldMove(m, c.entries, nil) { + for i := range c.entries { + if c.entries[i].Manifest.Module == moved.Manifest.Module { + c.entries[i].Source.Head = m.Commit + c.entries[i].Source.Seen = now() + } + } + } + return nil + } +} + +func at(s string) time.Time { + t, err := time.Parse(time.RFC3339, s) + if err != nil { + panic(err) + } + return t +} + +func announced(repo, commit, mergedAt, onTheBus string, paths ...string) link.AnnouncedMerge { + return link.AnnouncedMerge{ + SourceMoved: link.SourceMoved{Owner: "novox", Repo: repo, Base: "main", Commit: commit, + MergedAt: mergedAt, Paths: paths, + CloneURL: "http://forge.internal:20000/novox/" + repo + ".git"}, + At: at(onTheBus), + } +} + +func built(module, repo, path, commit, seen string) inventory.Entry { + e := fromRepo(module, "http://forge.internal:20000/novox/"+repo+".git", path) + e.Source.BuiltFrom, e.Source.Head, e.Source.Seen = commit, commit, at(seen) + return e +} + +// **novox/hq issue 266, as it happened.** The forge announced a merge of the tools repository on the +// events stream; the bus never handed it to the controller, which acted on the merges around it and +// not on this one, and said nothing. Read back from the stream, it is the one merge that would still +// move something — so it is said and acted on, once, and only after the controller's own consumer +// has had its time with it. +func TestAMergeTheBusNeverHandedOverIsActedOnLate(t *testing.T) { + cat := &aCatalogue{entries: []inventory.Entry{ + built("mesh-tools", "mesh-tools", "", "8b789578aaaaaaaa", "2026-10-04T15:24:32Z"), + built("node-tools", "mesh-tools", "node-tools", "8b789578aaaaaaaa", "2026-10-04T15:24:32Z"), + // Acted on when it was announced: looked at after it was merged. + built("gitea", "mesh-catalog", "modules/gitea", "5c2157b8bbbbbbbb", "2026-10-05T22:39:21Z"), + }} + stream := announcedList{ + // Nothing the mesh holds is built from the records repository. + announced("hq", "88f7f79fcccccccc", "2026-10-05T22:43:00Z", "2026-10-05T22:43:04Z", "04-ISSUES/x.md"), + // Acted on: its module was looked at since. + announced("mesh-catalog", "78328d4adddddddd", "2026-10-05T22:39:00Z", "2026-10-05T22:39:21Z", "modules/gitea/x.ts"), + // Never handed over. + announced("mesh-tools", "9730bd89c3e48d0e", "2026-10-05T22:46:47Z", "2026-10-05T22:47:06Z", + "node-tools/internal/console/console.go"), + } + var said []string + say := func(format string, args ...any) { said = append(said, fmt.Sprintf(format, args...)) } + clock := at("2026-10-05T22:50:00Z") + now := func() time.Time { return clock } + pass := func() { + t.Helper() + if err := catchUpOnMerges(context.Background(), clock, stream, cat.read, cat.act(now), say); err != nil { + t.Fatal(err) + } + } + + pass() + if len(cat.acted) != 0 { + t.Fatalf("a merge three minutes old was taken from the controller's own consumer: %v", cat.acted) + } + + clock = at("2026-10-05T22:58:00Z") + pass() + if strings.Join(cat.acted, ",") != "mesh-tools@9730bd89" { + t.Fatalf("acted on %v, wanted the one merge never handed over", cat.acted) + } + if len(said) != 1 || !strings.Contains(said[0], "novox/mesh-tools merged into main (9730bd89)") || + !strings.Contains(said[0], "mesh-tools and node-tools are behind it") { + t.Fatalf("the missed merge was not said as one: %q", said) + } + + clock = at("2026-10-05T23:03:00Z") + pass() + if len(cat.acted) != 1 || len(said) != 1 { + t.Fatalf("a merge acted on was acted on again: %v %q", cat.acted, said) + } +} + +// A merge that changed none of the held modules' files moves nothing, so it is never "missed"; one +// that could not be acted on is said and tried again on the next pass. +func TestAMissedMergeThatCouldNotBeActedOnIsTriedAgain(t *testing.T) { + cat := &aCatalogue{ + entries: []inventory.Entry{built("gitea", "mesh-catalog", "modules/gitea", "5c2157b8bbbbbbbb", "2026-10-05T20:00:00Z")}, + fail: errors.New("the store is restarting"), + } + stream := announcedList{ + announced("mesh-catalog", "aaaaaaaa11111111", "2026-10-05T21:00:00Z", "2026-10-05T21:00:10Z", + "modules/plex/module.json", "modules/plex/x.ts"), + announced("mesh-catalog", "bbbbbbbb22222222", "2026-10-05T21:10:00Z", "2026-10-05T21:10:10Z", "modules/gitea/x.ts"), + } + var said []string + say := func(format string, args ...any) { said = append(said, fmt.Sprintf(format, args...)) } + clock := at("2026-10-05T22:00:00Z") + now := func() time.Time { return clock } + if err := catchUpOnMerges(context.Background(), clock, stream, cat.read, cat.act(now), say); err != nil { + t.Fatal(err) + } + if len(said) != 2 || !strings.Contains(said[1], "could not act on it") { + t.Fatalf("a failed catch-up was not said: %q", said) + } + cat.fail = nil + clock = at("2026-10-05T22:05:00Z") + if err := catchUpOnMerges(context.Background(), clock, stream, cat.read, cat.act(now), say); err != nil { + t.Fatal(err) + } + if strings.Join(cat.acted, ",") != "mesh-catalog@bbbbbbbb" { + t.Fatalf("acted on %v, wanted only the merge that changed a held module", cat.acted) + } +} + +// A merge older than the look-back is left to the operator: a controller that did not run this +// missed it, and acting on it days later would be a surprise rebuild. +func TestAMergeOlderThanTheLookBackIsLeftAlone(t *testing.T) { + cat := &aCatalogue{entries: []inventory.Entry{built("gitea", "mesh-catalog", "modules/gitea", "5c2157b8bbbbbbbb", "2026-10-01T00:00:00Z")}} + stream := announcedList{announced("mesh-catalog", "cccccccc33333333", "2026-10-03T00:00:00Z", "2026-10-03T00:00:05Z", "modules/gitea/x.ts")} + clock := at("2026-10-05T22:00:00Z") + if err := catchUpOnMerges(context.Background(), clock, stream, cat.read, cat.act(func() time.Time { return clock }), + func(string, ...any) {}); err != nil { + t.Fatal(err) + } + if len(cat.acted) != 0 { + t.Fatalf("a merge of three days ago was acted on: %v", cat.acted) + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index b7fb7be..b7ef31f 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -127,6 +127,9 @@ func serve(ctx context.Context) error { if err := server.Follows(following{open}); err != nil { return err } + // And the merges the bus announced and never handed over, read back on a timer and acted on late + // rather than never (novox/hq issue 266). + go catchingUpOnMerges(ctx, open, server) // 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 { diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 3fa7aab..74e433a 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -8,6 +8,7 @@ import ( "path" "regexp" "strings" + "sync" "time" "github.com/novox/mesh-controller/internal/catalogue" @@ -266,6 +267,11 @@ func notNow(err error) error { // (novox/hq 04-ISSUES/131). Nothing is pushed here: what a finished build does to the machines // running the module is the upgrade's decision, taken when the catalogue announces it. func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { + // One merge acted on at a time, whoever hands it over: the bus, or the catch-up that reads back + // what the bus did not hand over (novox/hq issue 266). Each judges against what the other wrote. + actingOnMerges.Lock() + defer actingOnMerges.Unlock() + inv := f.open.inventory entries, err := inv.Catalogued(ctx) if err != nil { @@ -276,34 +282,7 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { return notNow(err) } - // Two kinds of module are affected by one merge, and they are affected differently. - // - // A module **built from** this repository and branch has moved: the mesh records the new commit - // as what its source now has, and only what the merge actually changed is rebuilt. A module that - // only **packages source from** it has not moved — its own source is somewhere else, at the - // commit it already records — so it is rebuilt and its record left alone. Writing this commit as - // its source would make it permanently behind a repository its manifest does not come from. - var from, packaging []inventory.Entry - already := 0 - for _, e := range entries { - switch { - case sourceIs(e.Source, m): - if e.Source.BuiltFrom == m.Commit { - already++ - continue - } - // **A merge older than the last look at the source is history, not a move.** The forge - // announces what it finds merged, and an old merge surfacing late would otherwise move - // the recorded head backwards and rebuild everything built from that repository, once - // per old merge (2026-09-28). - if isHistory(m.MergedAt, e.Source.Seen) { - continue - } - from = append(from, e) - case readsFrom(read[e.Manifest.Module], m): - packaging = append(packaging, e) - } - } + from, packaging, already := mergeCandidates(m, entries, read) if len(from) == 0 && len(packaging) == 0 { // "Already built from it" and "nothing reads it" are different facts, and reading the first // as the second sends somebody looking for a broken trigger when the mesh is up to date. @@ -438,6 +417,57 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { return nil } +// actingOnMerges keeps one merge acted on at a time (novox/hq issue 266). +var actingOnMerges sync.Mutex + +// mergeCandidates is what one merge could move, judged against what the catalogue holds. +// +// Two kinds of module are affected by one merge, and they are affected differently. +// +// A module **built from** this repository and branch has moved: the mesh records the new commit as +// what its source now has, and only what the merge actually changed is rebuilt. A module that only +// **packages source from** it has not moved — its own source is somewhere else, at the commit it +// already records — so it is rebuilt and its record left alone. Writing this commit as its source +// would make it permanently behind a repository its manifest does not come from. `already` counts +// the modules built from this repository that are already built from this very commit. +func mergeCandidates(m link.SourceMoved, entries []inventory.Entry, + read map[string][]inventory.ReadRepository) (from, packaging []inventory.Entry, already int) { + for _, e := range entries { + switch { + case sourceIs(e.Source, m): + if e.Source.BuiltFrom == m.Commit { + already++ + continue + } + // **A merge older than the last look at the source is history, not a move.** The forge + // announces what it finds merged, and an old merge surfacing late would otherwise move + // the recorded head backwards and rebuild everything built from that repository, once + // per old merge (2026-09-28). + if isHistory(m.MergedAt, e.Source.Seen) { + continue + } + from = append(from, e) + case readsFrom(read[e.Manifest.Module], m): + packaging = append(packaging, e) + } + } + return from, packaging, already +} + +// wouldMove is the modules built from the merged repository that acting on this merge would mark as +// moved and rebuild — SourceMoved's judgement, made without acting (novox/hq issue 266). Empty for a +// merge already acted on: acting marks each of them as looked at, so the merge then reads as history. +// +// **Only the modules built from it, never the ones that merely package source from it.** Acting +// records nothing about those, so a merge acted on would go on reading as unacted for them, and be +// acted on again on every look. A merge that moves both is caught by the first kind, and acting on it +// rebuilds the second as well. +func wouldMove(m link.SourceMoved, entries []inventory.Entry, + read map[string][]inventory.ReadRepository) []inventory.Entry { + from, _, _ := mergeCandidates(m, entries, read) + return whatTheMergeTouched(from, entries, m) +} + // sourceIs is whether a recorded source is the repository and branch a merge announced. A source on // the git seat is recorded as its path on the forge; one elsewhere as the URL it was cloned from. // An empty recorded ref is the repository's default branch, which is what a merge into the base diff --git a/internal/link/merges.go b/internal/link/merges.go new file mode 100644 index 0000000..e018d5e --- /dev/null +++ b/internal/link/merges.go @@ -0,0 +1,79 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "sort" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// AnnouncedMerge is one merge the forge announced, as the events stream holds it: what it said, and +// when the bus took it. +type AnnouncedMerge struct { + SourceMoved + // At is when the bus took the announcement — the stream's own time, not the forge's. + At time.Time + // Seq is its place in the events stream. + Seq uint64 +} + +// MergedSubject is where the forge's merges land: the controller's own follow of them. +var MergedSubject = broker.ControllerFollows[3] + +// readQuiet is how long a read of the stream waits for one more message before it takes the stream +// as read to its end. The stream answers at once when it holds something; this is only the wait at +// the end. +const readQuiet = 2 * time.Second + +// AnnouncedMerges is every merge the forge announced since a moment, oldest first, read back from +// the events stream (novox/hq issue 266). +// +// **Read on a consumer of its own, filtered on the one subject.** The controller's durable consumer +// carries several filters, and the bus server the mesh ran when this was written (2.10) skips a +// message now and then on a consumer with more than one filter: it moves its delivered pointer past +// the message without ever handing it over, so the controller never hears of it and nothing says so. +// A consumer filtered on a single subject reads the stream another way and was not seen to skip. This +// one is ordered, ephemeral and acknowledges nothing, so reading it changes nothing on the bus. +func (s *Server) AnnouncedMerges(ctx context.Context, since time.Time) ([]AnnouncedMerge, error) { + if s.js == nil { + return nil, errors.New("this control plane is not on the bus, so it cannot read what the forge announced") + } + sub, err := s.js.Context().SubscribeSync(MergedSubject, nats.OrderedConsumer(), nats.StartTime(since)) + if err != nil { + return nil, fmt.Errorf("reading the forge's merges from the events stream: %w", err) + } + defer func() { _ = sub.Unsubscribe() }() + + var out []AnnouncedMerge + for { + wait, cancel := context.WithTimeout(ctx, readQuiet) + msg, err := sub.NextMsgWithContext(wait) + cancel() + if err != nil { + if ctx.Err() != nil { + return nil, ctx.Err() + } + // Nothing more within the quiet wait: the stream has been read to its end. + break + } + meta, err := msg.Metadata() + if err != nil { + continue + } + var moved SourceMoved + if json.Unmarshal(msg.Data, &moved) == nil && moved.Commit != "" && moved.Repo != "" { + out = append(out, AnnouncedMerge{SourceMoved: moved, At: meta.Timestamp, Seq: meta.Sequence.Stream}) + } + if meta.NumPending == 0 { + break + } + } + sort.Slice(out, func(i, j int) bool { return out[i].Seq < out[j].Seq }) + return out, nil +}