A merge on the controller's own path never shares a batch with one that waits for the delivery's word (hq ADR 0276, decided during the build)
This commit is contained in:
+113
-29
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user