diff --git a/cmd/mesh-controller/batches.go b/cmd/mesh-controller/batches.go index 4e876f0c..e947906f 100644 --- a/cmd/mesh-controller/batches.go +++ b/cmd/mesh-controller/batches.go @@ -37,6 +37,10 @@ import ( // own commit, the newest first, one walk each, until one is delivered; the older ones are then answered // by that one. // +// A merge on the controller's own path (a module whose walk waits for nobody's word) never shares a batch with +// one that waits for mesh-delivery's: each kind has a batch of its own, so no catalogue delivery skips its turn +// behind a controller merge (decided during the build, 2026-10-10). +// // The batch, its merges and their times are in the store (batched_merge, and the batch's own plan record in // the state `assembling` or `queued`), so a restarted controller resumes the window where it stood. @@ -185,7 +189,11 @@ func (f following) hearMerge(ctx context.Context, m link.SourceMoved, now time.T Heard: now, Event: event} // **A merge heard after a later merge of its branch** is answered by the walk holding that one. - if p, ok, err := answeredByALaterMerge(ctx, inv, heard, touches); err != nil { + // **A merge on the controller's own path never shares a batch with one that waits for the delivery's + // word** (decided during the build, 2026-10-10): batched together, the catalogue's deliveries would start + // with the controller's and skip their turn. Each kind has a batch of its own. + own := ownPath(touches) + if p, ok, err := answeredByALaterMerge(ctx, inv, heard, touches, own); err != nil { return notNow(err) } else if ok { heard.Plan = p.ID @@ -204,7 +212,7 @@ func (f following) hearMerge(ctx context.Context, m link.SourceMoved, now time.T return cutBatchesHeld(ctx, f.open, now) } - batch, err := theOpenBatch(ctx, inv, heard, now) + batch, err := theOpenBatch(ctx, inv, heard, now, own) if err != nil { return notNow(err) } @@ -243,7 +251,7 @@ func owedLate(ctx context.Context, inv *inventory.Inventory, m link.SourceMoved, // later merge, which contains it. A failed or stopped walk answers nothing more, and a walk folded into // another is followed to that one. False when none does: the merge joins the next batch. func answeredByALaterMerge(ctx context.Context, inv *inventory.Inventory, heard inventory.BatchedMerge, - moves []inventory.Entry) (inventory.Plan, bool, error) { + moves []inventory.Entry, own bool) (inventory.Plan, bool, error) { later, err := inv.LaterMergesOf(ctx, heard.Repository, heard.Branch, heard.Merged) if err != nil || len(later) == 0 { return inventory.Plan{}, false, err @@ -255,7 +263,10 @@ func answeredByALaterMerge(ctx context.Context, inv *inventory.Inventory, heard } switch { case p.Batch(): - return p, true, nil + // A batch of the other kind is not this merge's: it joins the batch of its own kind. + if p.OwnPath() == own { + return p, true, nil + } case p.Open() || p.State == inventory.PlanDone: if buildsEvery(p, moves) { return p, true, nil @@ -289,19 +300,31 @@ func buildsEvery(p inventory.Plan, modules []inventory.Entry) bool { return true } -// theOpenBatch is the batch a merge heard now joins: the one not yet cut, or a new one. -func theOpenBatch(ctx context.Context, inv *inventory.Inventory, first inventory.BatchedMerge, now time.Time) (inventory.Plan, error) { +// ownPath says a merge touches a module on the controller's own path, whose walk waits for nobody's word. +func ownPath(touches []inventory.Entry) bool { + for _, e := range touches { + if _, own := onTheControllersPath[e.Manifest.Module]; own { + return true + } + } + return false +} + +// theOpenBatch is the batch a merge heard now joins: the one of its kind not yet cut, or a new one. +func theOpenBatch(ctx context.Context, inv *inventory.Inventory, first inventory.BatchedMerge, now time.Time, own bool) (inventory.Plan, error) { batches, err := inv.Batches(ctx) if err != nil { return inventory.Plan{}, err } - if len(batches) > 0 { - return batches[0], nil + for _, b := range batches { + if b.OwnPath() == own { + return b, nil + } } b := inventory.Plan{ID: fmt.Sprintf("plan-%d", now.UnixNano()), Repository: first.Repository, Branch: first.Branch, Commit: first.Commit, Merged: first.Merged, Created: now, State: inventory.PlanAssembling, Tiers: [][]string{}, Modules: map[string]*inventory.PlanModule{}, - Delivery: &inventory.PlanDelivery{Batch: &inventory.PlanBatch{}}} + Delivery: &inventory.PlanDelivery{Batch: &inventory.PlanBatch{Own: own}}} // Kept before its first merge names it: a merge never names a plan the store does not hold. if err := inv.SavePlan(ctx, &b); err != nil { return inventory.Plan{}, err @@ -334,7 +357,7 @@ func keepBatch(ctx context.Context, inv *inventory.Inventory, b *inventory.Plan, } closes, latest := windowOf(merges, b.Created, window, atMost) b.Delivery.Merges = named - b.Delivery.Batch = &inventory.PlanBatch{ClosesAt: closes, AtMost: latest, Behind: behind} + b.Delivery.Batch = &inventory.PlanBatch{ClosesAt: closes, AtMost: latest, Behind: behind, Own: b.OwnPath()} b.State = inventory.PlanAssembling if behind != "" && windowClosed(*b, now) { b.State = inventory.PlanQueued @@ -534,37 +557,95 @@ func cutBatchesHeld(ctx context.Context, open *stores, now time.Time) error { if err != nil { return err } - var batch *inventory.Plan - if len(batches) > 0 { - batch = &batches[0] - } - if started != nil { - // **One walk at a time** (ADR 0276 decision 4): the batch waits behind it, whatever its window says. - if batch != nil { - return keepBatch(ctx, inv, batch, now, started.ID) + keepAll := func(behind string) error { + for i := range batches { + if err := keepBatch(ctx, inv, &batches[i], now, behind); err != nil { + return err + } } return nil } + if started != nil { + // **One walk at a time** (ADR 0276 decision 4): every batch waits behind it, whatever its window says. + return keepAll(started.ID) + } if len(alone) > 0 { if len(waiting) > 0 { - // The waiting walk is the one open; the search walks after it, and the batch waits behind both. - if batch != nil { - return keepBatch(ctx, inv, batch, now, waiting[0].ID) - } - return nil + // The waiting walk is the one open; the search walks after it, and the batches wait behind both. + return keepAll(waiting[0].ID) } return walkAlone(ctx, open, alone[0], now) } - if batch == nil { - return nil - } - if err := keepBatch(ctx, inv, batch, now, ""); err != nil { + if err := keepAll(""); err != nil { return err } - if !windowClosed(*batch, now) { + // The oldest batch whose window closed is cut; one of each kind at most is open (theOpenBatch). + for i := range batches { + b := &batches[i] + if !windowClosed(*b, now) { + continue + } + if b.OwnPath() { + // A walk on the controller's own path folds no walk waiting for its word: those merges are batched + // again, behind it (decided during the build, 2026-10-10). + if err := rebatchWaiting(ctx, inv, waiting, b.ID, now); err != nil { + return err + } + waiting = nil + } + if err := cutBatch(ctx, open, b, waiting, now); err != nil { + return err + } + // The batches left wait behind the walk just cut when it started; the next look does the same. + if b.Open() && !b.Waiting() { + if batches, err = inv.Batches(ctx); err != nil { + return err + } + return keepAll(b.ID) + } return nil } - return cutBatch(ctx, open, batch, waiting, now) + return nil +} + +// rebatchWaiting puts the merges of the walks waiting for their word into the open batch of their kind — a new +// one when none is — queued behind the walk about to be cut, and closes the walks as taken over by that batch: +// the batch keeps its id when it is cut, so the delivery's owner follows them to the walk that answers them. +func rebatchWaiting(ctx context.Context, inv *inventory.Inventory, waiting []inventory.Plan, behind string, now time.Time) error { + for i := range waiting { + w := waiting[i] + merges, err := inv.MergesOf(ctx, w.ID) + if err != nil { + return err + } + if len(merges) == 0 { + continue // nothing of the record's: left waiting + } + batch, err := theOpenBatch(ctx, inv, merges[0], now, false) + if err != nil { + return err + } + for _, m := range merges { + if err := inv.AnswerMerge(ctx, m.Repository, m.Commit, batch.ID, m.Alone); err != nil { + return err + } + } + if err := keepBatch(ctx, inv, &batch, now, behind); err != nil { + return err + } + w.State = inventory.PlanSuperseded + if w.Delivery == nil { + w.Delivery = &inventory.PlanDelivery{} + } + w.Delivery.TakenOverBy = batch.ID + w.Note = fmt.Sprintf("superseded at tier %d by %s before it started: a walk on the controller's own path "+ + "(%s) goes first, and that batch answers its merges after it", w.Tier, batch.ID, behind) + if err := inv.SavePlan(ctx, &w); err != nil { + return err + } + fmt.Printf(" %s is %s\n", w.ID, w.Note) + } + return nil } // cutBatch makes a closed batch one walk, folding in the walks that wait for their word: it keeps its id, and @@ -966,6 +1047,9 @@ func batchWords(b inventory.Plan, now time.Time) string { return "grouped: " + grouped + "; plan not yet calculated" } w := b.Delivery.Batch + if w.Own { + grouped += " (on the controller's own path)" + } if b.State == inventory.PlanQueued { return fmt.Sprintf("queued behind %s; grouped: %s; plan not yet calculated", w.Behind, grouped) } diff --git a/cmd/mesh-controller/batches_test.go b/cmd/mesh-controller/batches_test.go index b2fa1802..38af562a 100644 --- a/cmd/mesh-controller/batches_test.go +++ b/cmd/mesh-controller/batches_test.go @@ -669,3 +669,92 @@ func TestTheCatchUpKeepsWhatACutMadeHistory(t *testing.T) { t.Fatalf("the merge made before the cut and heard after it was dropped: %v %v", kept, err) } } + +// A merge on the controller's own path never shares a batch with one that waits for mesh-delivery's word +// (decided during the build, 2026-10-10): the two are batches of their own, cut one after the other; a walk on +// the own path folds no waiting catalogue walk but batches its merges again behind it; the catalogue's walk +// still waits for the word, so no catalogue delivery skips its turn. +func TestAMergeOnTheControllersPathNeverSharesABatch(t *testing.T) { + open := windowed(t) + asked := asksWithPaths(t) + ctx := t.Context() + if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: "mesh-delivery", Version: "1", + Claims: []catalogue.Claim{{Name: catalogue.DeliverySeat, Scope: catalogue.ScopeMesh}}}, + inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/mesh-delivery", Ref: "main", + BuiltFrom: "c0", Head: "c0"}); err != nil { + t.Fatal(err) + } + if _, err := open.inventory.Assign(ctx, "anchor", "mesh-delivery"); err != nil { + t.Fatal(err) + } + hear(t, open, catalogueMerge("aaaa0001", "app", t0), t0) + hear(t, open, repoMerge("mesh-controller", "k1", t0.Add(10*time.Second)), t0.Add(10*time.Second)) + hear(t, open, catalogueMerge("bbbb0002", "notes", t0.Add(20*time.Second)), t0.Add(20*time.Second)) + _, bs := walks(t, open) + var ownBatch, catalogueBatch inventory.Plan + for _, b := range bs { + if b.OwnPath() { + ownBatch = b + } else { + catalogueBatch = b + } + } + if len(bs) != 2 || ownBatch.ID == "" || catalogueBatch.ID == "" || len(catalogueBatch.Delivery.Merges) != 2 || + len(ownBatch.Delivery.Merges) != 1 { + t.Fatalf("a controller merge and two catalogue merges in one window: %+v", bs) + } + if !strings.Contains(batchWords(ownBatch, t0.Add(30*time.Second)), "own path") { + t.Fatalf("the own-path batch does not say so: %q", batchWords(ownBatch, t0.Add(30*time.Second))) + } + // The catalogue batch, the older, is cut first and waits for its word; the controller's batch is cut next: + // the waiting walk is batched again behind it, not folded into it. + cutAt(t, open, t0.Add(2*time.Minute)) + ws, _ := walks(t, open) + if len(ws) != 1 || !ws[0].Waiting() { + t.Fatalf("the catalogue batch was not cut into a waiting walk: %+v", ws) + } + catalogueWalk := ws[0] + cutAt(t, open, t0.Add(2*time.Minute+5*time.Second)) + ws, bs = walks(t, open) + var ownWalk inventory.Plan + for _, w := range ws { + if w.Open() { + ownWalk = w + } + } + if ownWalk.ID == "" || ownWalk.Waiting() || ownWalk.CommitOf("novox/mesh-controller") != "k1" { + t.Fatalf("the controller's batch was not cut into a started walk: %+v", ws) + } + for _, m := range []string{"app", "notes"} { + if _, in := ownWalk.Modules[m]; in { + t.Fatalf("the controller's walk builds %s, a catalogue module: it skipped its turn", m) + } + } + for _, a := range *asked { + if a[1] == "modules/app" || a[1] == "modules/notes" { + t.Fatalf("a catalogue module was asked without the word: %v", *asked) + } + } + folded, _ := open.inventory.PlanByID(ctx, catalogueWalk.ID) + if folded.State != inventory.PlanSuperseded || len(bs) != 1 || folded.Delivery.TakenOverBy != bs[0].ID || + bs[0].State != inventory.PlanQueued || bs[0].Delivery.Batch.Behind != ownWalk.ID || bs[0].OwnPath() || + len(bs[0].Delivery.Merges) != 2 { + t.Fatalf("the waiting catalogue walk is %s (taken over by %q); batches %+v", folded.State, + folded.Delivery.TakenOverBy, bs) + } + // The controller's walk done, the catalogue batch is cut and waits for its word with both merges. + ownWalk.State = inventory.PlanDone + if err := open.inventory.SavePlan(ctx, &ownWalk); err != nil { + t.Fatal(err) + } + cutAt(t, open, t0.Add(3*time.Minute)) + ws, bs = walks(t, open) + if len(bs) != 0 || ws[0].ID != folded.Delivery.TakenOverBy || !ws[0].Waiting() || len(ws[0].Delivery.Merges) != 2 { + t.Fatalf("the catalogue batch was not cut into a waiting walk once the controller's ended: %+v %+v", ws, bs) + } + for _, m := range []string{"app", "notes"} { + if _, in := ws[0].Modules[m]; !in { + t.Fatalf("the catalogue walk does not build %s: %v", m, ws[0].Modules) + } + } +} diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index 4dba027c..56045490 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -79,6 +79,15 @@ type PlanBatch struct { ClosesAt time.Time `json:"closes_at"` AtMost time.Time `json:"at_most"` Behind string `json:"behind,omitempty"` + // Own says the batch holds merges on the controller's own path, whose walk waits for nobody's word: such + // a merge never shares a batch with one that waits for mesh-delivery's (ADR 0276, decided during the + // build, 2026-10-10). + Own bool `json:"own,omitempty"` +} + +// OwnPath says the record is a batch of merges on the controller's own path. +func (p Plan) OwnPath() bool { + return p.Delivery != nil && p.Delivery.Batch != nil && p.Delivery.Batch.Own } // CommitOf is the commit of a repository the walk carries; empty when it carries none of it.