Assemble merges in a rolling window and walk each batch once (hq ADR 0276, issue 362)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request

Every merge opened a walk and the next merge of the branch superseded it: two
catalogue merges 18 s apart left a walk no delivery held, and the operator
started it by hand 58 minutes later. A merge now joins the open batch, kept in
the store (migration 0089), which is cut into one walk when no merge came for
merge-window (90 s) or at merge-window-at-most (10 min): one commit per
repository, the latest of its branch, with every file the batch's merges
changed. One walk at a time; a started walk is never superseded, a waiting one
is folded into the next. The walk names every merge it answers on the wire
(delivery.merges, taken_over_by, batch). A failed walk walks its earlier merges
alone, newest first, until one is delivered. A delivery group's order becomes
tier edges inside the walk. plans shows the batch assembling; S18 and S19
bound its waits; S16 names the merges a waiting walk answers.
This commit is contained in:
jochen
2026-10-10 13:20:04 +02:00
parent 8a53532d89
commit 96fc4209d3
24 changed files with 2110 additions and 563 deletions
+971
View File
@@ -0,0 +1,971 @@
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"regexp"
"slices"
"sort"
"strconv"
"strings"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// **Merges are assembled in a rolling window, and each batch is delivered by one walk** (novox/hq ADR 0276,
// issue 362).
//
// Every merge opened a walk of its own, and the next merge of the branch superseded it. On 2026-10-10 two
// catalogue merges eighteen seconds apart left one walk that no delivery held, and the operator started it by
// hand 58 minutes later. A merge now joins the open **batch**. Each merge restarts the **merge window**; it
// closes when none came for `merge-window` (90 seconds) or at the latest `merge-window-at-most` (10 minutes)
// after the batch's first merge. A closed batch becomes one walk when no walk is open: one plan for every
// commit in it, with at most one commit per repository, the latest of its branch, which contains the earlier
// ones. The walk names every merge it answers.
//
// - One walk at a time. Merges that come while a walk runs join the next batch. A walk that started is never
// superseded; one still waiting for its delivery's word is folded into the next batch's walk, which
// answers every merge it answered — the only supersession left.
// - A merge heard after a later merge of its branch was put in a walk is answered by that walk: a batch
// still assembling, an open walk that builds everything it moves, or a walk done that did, which then
// names it at once. Otherwise it joins the next batch.
// - A failed walk marks nothing delivered. Each earlier merge it carried in a later commit is walked on its
// own commit, the newest first, one walk each, until one is delivered; the older ones are then answered
// by that one.
//
// 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.
// The merge window's defaults. The controller's settings `merge-window` and `merge-window-at-most` reach it in
// the file MESH_MERGE_WINDOW_FILE names (module.json), read at every look, so a change takes effect without a
// restart.
var (
mergeWindowDefault = 90 * time.Second
mergeWindowAtMostDefault = 10 * time.Minute
)
// mergeWindowFileVar names the file the controller's window settings are written to.
const mergeWindowFileVar = "MESH_MERGE_WINDOW_FILE"
// batchCutEvery is how often the controller looks whether a batch's window closed.
const batchCutEvery = 5 * time.Second
// batchLateAfter is how long past its maximum a batch may still be assembling, with no walk open, before
// the controller is said to have failed to cut it (ADR 0276 decision 9).
const batchLateAfter = time.Minute
// edgeGroupOrder is a delivery group's order made an edge of a walk's tiers (ADR 0276 decision 5): the
// modules of the member after depend on the modules of the member before, and are gated on them.
const edgeGroupOrder = "group-order"
// mergeWindowOf is the window and its maximum, from the controller's settings; a value unreadable is its
// default, and said.
func mergeWindowOf() (window, atMost time.Duration) {
window, atMost = mergeWindowDefault, mergeWindowAtMostDefault
path := os.Getenv(mergeWindowFileVar)
if path == "" {
return window, atMost
}
raw, err := os.ReadFile(path)
if err != nil {
fmt.Fprintf(os.Stderr, "the merge window's settings cannot be read (%v); %s and %s are used\n", err, window, atMost)
return window, atMost
}
var values map[string]any
if err := json.Unmarshal(raw, &values); err != nil {
fmt.Fprintf(os.Stderr, "the merge window's settings are not JSON (%v); %s and %s are used\n", err, window, atMost)
return window, atMost
}
for key, to := range map[string]*time.Duration{"merge-window": &window, "merge-window-at-most": &atMost} {
v, has := values[key]
if !has || v == "" {
continue
}
if d, ok := durationSetting(v); ok {
*to = d
} else {
fmt.Fprintf(os.Stderr, "the setting %s=%v is no duration; %s is used\n", key, v, *to)
}
}
return window, atMost
}
// durationSetting reads a duration as a setting says it: "90s", "10m", or a number of seconds.
func durationSetting(v any) (time.Duration, bool) {
switch t := v.(type) {
case string:
if d, err := time.ParseDuration(strings.TrimSpace(t)); err == nil && d >= 0 {
return d, true
}
if n, err := strconv.ParseFloat(strings.TrimSpace(t), 64); err == nil && n >= 0 {
return time.Duration(n * float64(time.Second)), true
}
case float64:
if t >= 0 {
return time.Duration(t * float64(time.Second)), true
}
}
return 0, false
}
// hearMerge puts a merge the forge announced into the batch that answers it, at the moment it was heard,
// and cuts what is due.
func (f following) hearMerge(ctx context.Context, m link.SourceMoved, now time.Time) 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).
actingOnMerges.Lock()
defer actingOnMerges.Unlock()
inv := f.open.inventory
repository := m.Owner + "/" + m.Repo
if kept, known, err := inv.MergeOf(ctx, repository, m.Commit); err != nil {
return notNow(err)
} else if known {
fmt.Printf("%s merged into %s (%.8s) was heard before; %s answers it\n", repository, m.Base, m.Commit,
orWords(kept.Plan, "a walk of its own commit, once the walks before it end"))
return nil
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return notNow(err)
}
read, err := readForPlanning(ctx, inv)
if err != nil {
return notNow(err)
}
if from, packaging, already := mergeCandidates(m, entries, read); 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.
if already > 0 {
fmt.Printf("%s merged into %s (%.8s); %d module(s) the mesh holds are already built from it\n",
repository, m.Base, m.Commit, already)
return nil
}
fmt.Printf("%s merged into %s (%.8s); nothing the mesh holds reads it\n", repository, m.Base, m.Commit)
return nil
}
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return notNow(err)
}
defer release()
merged, err := time.Parse(time.RFC3339Nano, m.MergedAt)
if err != nil {
merged = now
}
event, err := json.Marshal(m)
if err != nil {
return err
}
heard := inventory.BatchedMerge{Repository: repository, Branch: m.Base, Commit: m.Commit, Merged: merged.UTC(),
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, m, entries, read); err != nil {
return notNow(err)
} else if ok {
heard.Plan = p.ID
if _, err := inv.AddMerge(ctx, heard); err != nil {
return notNow(err)
}
if p.Batch() {
if err := keepBatch(ctx, inv, &p, now, ""); err != nil {
return notNow(err)
}
} else if err := nameMerges(ctx, inv, &p); err != nil {
return notNow(err)
}
fmt.Printf("%s merged into %s (%.8s) after a later merge of its branch was put in %s (%s): that %s "+
"answers it\n", repository, m.Base, m.Commit, p.ID, p.State, walkOrBatch(p))
return cutBatchesHeld(ctx, f.open, now)
}
batch, err := theOpenBatch(ctx, inv, heard, now)
if err != nil {
return notNow(err)
}
heard.Plan = batch.ID
if _, err := inv.AddMerge(ctx, heard); err != nil {
return notNow(err)
}
if err := keepBatch(ctx, inv, &batch, now, ""); err != nil {
return notNow(err)
}
fmt.Printf("%s merged into %s (%.8s); it joins %s — %s\n", repository, m.Base, m.Commit, batch.ID,
batchWords(batch, now))
return cutBatchesHeld(ctx, f.open, now)
}
// answeredByALaterMerge is the plan that answers a merge heard after a later merge of its branch was put in
// one (ADR 0276 decision 4): a batch not yet cut, whose walk plans every file of the repository's merges; an
// open walk, or a walk done, that builds every module this merge would move — built from the branch after the
// 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,
m link.SourceMoved, entries []inventory.Entry, read map[string][]inventory.ReadRepository) (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
}
moves := wouldMove(m, entries, read)
for _, l := range later {
p, err := followTakeOver(ctx, inv, l.Plan)
if err != nil {
return inventory.Plan{}, false, err
}
switch {
case p.Batch():
return p, true, nil
case p.Open() || p.State == inventory.PlanDone:
if buildsEvery(p, moves) {
return p, true, nil
}
}
}
return inventory.Plan{}, false, nil
}
// followTakeOver is a plan, or the walk that took it over when it was folded before it started.
func followTakeOver(ctx context.Context, inv *inventory.Inventory, id string) (inventory.Plan, error) {
p, err := inv.PlanByID(ctx, id)
for hops := 0; err == nil && hops < 16 && p.State == inventory.PlanSuperseded && p.Delivery != nil &&
p.Delivery.TakenOverBy != ""; hops++ {
p, err = inv.PlanByID(ctx, p.Delivery.TakenOverBy)
}
return p, err
}
// buildsEvery says a walk builds every one of the modules.
func buildsEvery(p inventory.Plan, modules []inventory.Entry) bool {
for _, e := range modules {
if _, in := p.Modules[e.Manifest.Module]; !in {
return false
}
}
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) {
batches, err := inv.Batches(ctx)
if err != nil {
return inventory.Plan{}, err
}
if len(batches) > 0 {
return batches[0], 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{}}}
// 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
}
return b, nil
}
// keepBatch writes a batch as its merges and the clock say it now: the commits it holds, the merges they
// answer, its window, and assembling or queued behind a walk. Saved only when that changed.
func keepBatch(ctx context.Context, inv *inventory.Inventory, b *inventory.Plan, now time.Time, behind string) error {
merges, err := inv.MergesOf(ctx, b.ID)
if err != nil {
return err
}
window, atMost := mergeWindowOf()
// The note's seconds left are no change of the batch's: it is saved when what it holds or its state moved.
unNoted := func(p inventory.Plan) string {
p.Note = ""
return planSnapshot(p)
}
before := unNoted(*b)
commits, named := carriedOf(merges)
b.Commits = commits
if len(commits) > 0 {
newest := newestCommit(commits)
b.Repository, b.Branch, b.Commit, b.Merged = newest.Repository, newest.Branch, newest.Commit, newest.Merged
}
if b.Delivery == nil {
b.Delivery = &inventory.PlanDelivery{}
}
closes, latest := windowOf(merges, b.Created, window, atMost)
b.Delivery.Merges = named
b.Delivery.Batch = &inventory.PlanBatch{ClosesAt: closes, AtMost: latest, Behind: behind}
b.State = inventory.PlanAssembling
if behind != "" && windowClosed(*b, now) {
b.State = inventory.PlanQueued
} else {
b.Delivery.Batch.Behind = ""
}
b.Note = batchWords(*b, now)
if unNoted(*b) == before {
return nil
}
return inv.SavePlan(ctx, b)
}
// windowOf is when a batch's window closes, unless another merge comes, and when at the latest: the last
// merge heard plus the window, and the batch's first merge plus the maximum.
func windowOf(merges []inventory.BatchedMerge, created time.Time, window, atMost time.Duration) (closes, latest time.Time) {
var last time.Time
for _, m := range merges {
if m.Heard.After(last) {
last = m.Heard
}
}
if last.IsZero() {
last = created
}
return last.Add(window), created.Add(atMost)
}
// windowClosed says a batch's window is closed: no merge for the window's length, or its maximum reached.
func windowClosed(b inventory.Plan, now time.Time) bool {
if b.Delivery == nil || b.Delivery.Batch == nil {
return true
}
w := b.Delivery.Batch
return !now.Before(w.ClosesAt) || !now.Before(w.AtMost)
}
// carriedOf is a set of merges as a walk carries them: the latest merge of each repository's branch, and
// every merge named with the commit that carries it.
func carriedOf(merges []inventory.BatchedMerge) ([]inventory.PlanCommit, []inventory.PlanMerge) {
latest := map[string]inventory.BatchedMerge{}
spelled := map[string]string{}
for _, m := range merges {
k := strings.ToLower(m.Repository) + "\x00" + m.Branch
if l, ok := latest[k]; !ok || laterMerge(m, l) {
latest[k] = m
}
var e link.SourceMoved
if json.Unmarshal(m.Event, &e) == nil && e.Owner != "" {
spelled[k] = e.Owner + "/" + e.Repo
}
}
var commits []inventory.PlanCommit
for k, l := range latest {
repo := spelled[k]
if repo == "" {
repo = l.Repository
}
commits = append(commits, inventory.PlanCommit{Repository: repo, Branch: l.Branch, Commit: l.Commit, Merged: l.Merged})
}
sort.Slice(commits, func(i, j int) bool {
if !strings.EqualFold(commits[i].Repository, commits[j].Repository) {
return strings.ToLower(commits[i].Repository) < strings.ToLower(commits[j].Repository)
}
return commits[i].Branch < commits[j].Branch
})
var named []inventory.PlanMerge
for _, m := range merges {
k := strings.ToLower(m.Repository) + "\x00" + m.Branch
repo := spelled[k]
if repo == "" {
repo = m.Repository
}
n := inventory.PlanMerge{Repository: repo, Commit: m.Commit}
if l := latest[k]; l.Commit != m.Commit {
n.Carried = l.Commit
}
named = append(named, n)
}
sort.SliceStable(named, func(i, j int) bool {
return strings.ToLower(named[i].Repository) < strings.ToLower(named[j].Repository)
})
return commits, named
}
// laterMerge says a was merged after b: by the forge's merge time, then by when each was heard.
func laterMerge(a, b inventory.BatchedMerge) bool {
if !a.Merged.Equal(b.Merged) {
return a.Merged.After(b.Merged)
}
return a.Heard.After(b.Heard)
}
// newestCommit is the commit merged last among those a walk carries.
func newestCommit(commits []inventory.PlanCommit) inventory.PlanCommit {
newest := commits[0]
for _, c := range commits[1:] {
if c.Merged.After(newest.Merged) {
newest = c
}
}
return newest
}
// nameMerges writes the merges a walk answers onto its record, after a merge heard late was given to it.
func nameMerges(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) error {
merges, err := inv.MergesOf(ctx, p.ID)
if err != nil {
return err
}
_, named := carriedOf(merges)
// The walk's own commits stand: a merge heard late is carried by the one of its repository.
for i := range named {
if c := p.CommitOf(named[i].Repository); c != "" && c != named[i].Commit {
named[i].Carried = c
} else if c == named[i].Commit {
named[i].Carried = ""
}
}
if p.Delivery == nil {
p.Delivery = &inventory.PlanDelivery{}
}
p.Delivery.Merges = named
return inv.SavePlan(ctx, p)
}
// batchCutter cuts what is due on a timer, until the context ends: a window closes with no merge to say so.
func batchCutter(ctx context.Context, open *stores) {
tick := time.NewTicker(batchCutEvery)
defer tick.Stop()
for {
select {
case <-ctx.Done():
return
case <-tick.C:
}
release, err := open.inventory.HoldPlans(ctx, false)
if err != nil {
continue // whoever holds the plans moves them, and cuts at the end
}
if err := cutBatchesHeld(ctx, open, time.Now().UTC()); err != nil {
fmt.Printf("batches: %v\n", err)
}
release()
}
}
// cutBatchesHeld takes the step that is due, for a caller holding the plans: after a failed walk, the search
// for the merge that brought the failure; with no walk open, the next merge of that search walked alone, or a
// closed batch cut into one walk with the waiting walks folded in; with a walk open, a closed batch said
// queued behind it.
func cutBatchesHeld(ctx context.Context, open *stores, now time.Time) error {
inv := open.inventory
if err := searchAfterFailures(ctx, inv); err != nil {
return err
}
alone, err := inv.AloneMerges(ctx)
if err != nil {
return err
}
if alone, err = coverAlone(ctx, inv, alone); err != nil {
return err
}
plans, err := inv.OpenPlans(ctx)
if err != nil {
return err
}
var started *inventory.Plan
var waiting []inventory.Plan
for i := range plans {
p := plans[i]
if p.Release != nil {
continue // the backlog's walk through the gates: no merge's
}
if p.Waiting() {
waiting = append(waiting, p)
} else if started == nil {
started = &plans[i]
}
}
batches, err := inv.Batches(ctx)
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)
}
return nil
}
if len(alone) > 0 {
if len(waiting) > 0 {
return nil // the waiting walk is the one open; the search walks after it
}
return walkAlone(ctx, open, alone[0], now)
}
if batch == nil {
return nil
}
if err := keepBatch(ctx, inv, batch, now, ""); err != nil {
return err
}
if !windowClosed(*batch, now) {
return nil
}
return cutBatch(ctx, open, batch, waiting, now)
}
// cutBatch makes a closed batch one walk, folding in the walks that wait for their word: it keeps its id, and
// the folded walks name it as the walk that took them over.
func cutBatch(ctx context.Context, open *stores, batch *inventory.Plan, waiting []inventory.Plan, now time.Time) error {
inv := open.inventory
var carry []string
for _, w := range waiting {
merges, err := inv.MergesOf(ctx, w.ID)
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
}
}
// What the folded walk was to build is built by the batch's: its merges were acted on when it was cut,
// so they read as history now.
for name := range w.Modules {
carry = append(carry, name)
}
}
if err := planBatch(ctx, open, batch, carry, now); err != nil {
return err
}
// Closed after the batch's walk is kept, never before: a controller replaced between the two leaves both
// open, which the next cut settles, rather than neither.
for i := range waiting {
w := waiting[i]
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 (%s) before it started: that walk answers its merges",
w.Tier, batch.ID, batch.Named())
if err := inv.SavePlan(ctx, &w); err != nil {
return err
}
fmt.Printf(" %s is %s\n", w.ID, w.Note)
}
return nil
}
// walkAlone walks one merge on its own commit (ADR 0276 decision 3): a plan of its own, cut at once.
func walkAlone(ctx context.Context, open *stores, m inventory.BatchedMerge, now time.Time) error {
inv := open.inventory
p := inventory.Plan{ID: fmt.Sprintf("plan-%d", now.UnixNano()), Repository: m.Repository, Branch: m.Branch,
Commit: m.Commit, Merged: m.Merged, Created: now, State: inventory.PlanQueued, Tiers: [][]string{},
Modules: map[string]*inventory.PlanModule{}, Note: "walked alone after a failed walk carried it",
Delivery: &inventory.PlanDelivery{}}
if err := inv.SavePlan(ctx, &p); err != nil {
return err
}
if err := inv.AnswerMerge(ctx, m.Repository, m.Commit, p.ID, true); err != nil {
return err
}
fmt.Printf("%s %.8s, carried by a walk that failed, is walked on its own commit as %s\n", m.Repository, m.Commit, p.ID)
return planBatch(ctx, open, &p, nil, now)
}
// planBatch plans a batch's merges as one walk, under the batch's id: per repository the latest merge of its
// branch with every file the batch's merges of it changed, the moved modules and their dependents in tiers
// along the graph and the order of any delivery group among them.
func planBatch(ctx context.Context, open *stores, batch *inventory.Plan, carry []string, now time.Time) error {
inv := open.inventory
merges, err := inv.MergesOf(ctx, batch.ID)
if err != nil {
return err
}
if len(merges) == 0 {
batch.State = inventory.PlanSuperseded
batch.Note = "a batch that holds no merge: nothing to walk"
batch.Delivery.Batch = nil
return inv.SavePlan(ctx, batch)
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
read, err := readForPlanning(ctx, inv)
if err != nil {
return err
}
edges, err := inv.Dependencies(ctx)
if err != nil {
return err
}
commits, named := carriedOf(merges)
var moved []string
var members []batchMember
for _, m := range combinedMerges(merges) {
names, err := movesOfMerge(ctx, inv, m, entries, read)
if err != nil {
return err
}
members = append(members, batchMember{merge: m, moved: names})
for _, n := range names {
if !slices.Contains(moved, n) {
moved = append(moved, n)
}
}
}
held := map[string]bool{}
for _, e := range entries {
held[e.Manifest.Module] = true
}
var also []string
for _, n := range carry {
// One the catalogue no longer holds would fail the walk's ask; it is nobody's to build.
if held[n] && !slices.Contains(moved, n) {
moved = append(moved, n)
also = append(also, n)
}
}
sort.Strings(also)
ordered := groupOrderEdges(members, entries, read, edges)
plan := planOfMoves(moved, append(append([]inventory.Edge{}, edges...), ordered...))
plan.ID, plan.Revision, plan.Created, plan.Epoch = batch.ID, batch.Revision, now, batch.Epoch
plan.Commits = commits
newest := newestCommit(commits)
plan.Repository, plan.Branch, plan.Commit, plan.Merged = newest.Repository, newest.Branch, newest.Commit, newest.Merged
// **Whether it waits for its delivery's word** (novox/hq ADR 0239): while the mesh-delivery seat has a
// holder on record, a walk that moves no module on the controller's own path is opened and waits.
plan.Delivery = awaitsFor(entries, moved)
if plan.Delivery == nil {
plan.Delivery = &inventory.PlanDelivery{}
}
plan.Delivery.Merges = named
if len(moved) == 0 {
plan.State = inventory.PlanDone
plan.Tiers = [][]string{}
plan.Note = "moved nothing the mesh holds: nothing to walk"
if err := inv.SavePlan(ctx, &plan); err != nil {
return err
}
*batch = plan
fmt.Printf("%s (%s): its merges move nothing the mesh holds\n", plan.ID, plan.Named())
return nil
}
if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
strings.Join(plan.Tiers[len(plan.Tiers)-1], ", "))
}
if len(also) > 0 {
fmt.Printf(" %s, left unbuilt by a walk folded into it, are planned here again\n", strings.Join(also, ", "))
}
var tiers []string
for i, t := range plan.Tiers {
tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", ")))
}
var answers []string
for _, n := range named {
answers = append(answers, n.Repository+"@"+short(n.Commit))
}
fmt.Printf("%s cut: %s, answering %s; %d module(s) in %d tier(s)\n %s\n", plan.ID, plan.Named(),
strings.Join(answers, ", "), len(plan.Modules), len(plan.Tiers), strings.Join(tiers, "\n "))
if plan.Waiting() {
plan.Note = waitingNote(plan)
if err := inv.SavePlan(ctx, &plan); err != nil {
return err
}
*batch = plan
fmt.Printf(" %s waits for %s's word before its first tier is asked\n", plan.ID, plan.Delivery.Awaits)
return nil
}
// Written and its first tier asked as one act on the plans (novox/hq issue 213), under the caller's hold.
if err := askTier(ctx, inv, &plan); err != nil {
return err
}
if err := inv.SavePlan(ctx, &plan); err != nil {
return err
}
*batch = plan
return nil
}
// combinedMerges is a batch's merges as one merge per repository's branch: the latest, with every file the
// merges of it changed, removed and said to be a module's — on a linear trunk the latest contains the others.
func combinedMerges(merges []inventory.BatchedMerge) []link.SourceMoved {
type combined struct {
latest inventory.BatchedMerge
m link.SourceMoved
said bool
}
by := map[string]*combined{}
var keys []string
for _, b := range merges {
var e link.SourceMoved
if err := json.Unmarshal(b.Event, &e); err != nil {
continue
}
k := strings.ToLower(b.Repository) + "\x00" + b.Branch
c, ok := by[k]
if !ok {
c = &combined{latest: b, m: e, said: e.ModuleDirsSaid}
by[k] = c
keys = append(keys, k)
} else {
paths, removed, dirs, truncated := c.m.Paths, c.m.Removed, c.m.ModuleDirs, c.m.PathsTruncated
said := c.said && e.ModuleDirsSaid
if laterMerge(b, c.latest) {
c.latest, c.m = b, e
}
c.m.Paths = union(paths, e.Paths)
c.m.Removed = union(removed, e.Removed)
c.m.ModuleDirs = union(dirs, e.ModuleDirs)
c.m.PathsTruncated = truncated || e.PathsTruncated
c.said = said
}
}
sort.Strings(keys)
var out []link.SourceMoved
for _, k := range keys {
c := by[k]
c.m.ModuleDirsSaid = c.said
// A file a later merge of the batch brought back is no longer removed.
var removed []string
for _, r := range c.m.Removed {
if !restoredLater(r, c.latest, merges) {
removed = append(removed, r)
}
}
c.m.Removed = removed
out = append(out, c.m)
}
return out
}
// restoredLater says a file removed by one merge of a batch was changed, not removed, by the latest.
func restoredLater(path string, latest inventory.BatchedMerge, _ []inventory.BatchedMerge) bool {
var e link.SourceMoved
if json.Unmarshal(latest.Event, &e) != nil {
return false
}
return slices.Contains(e.Paths, path) && !slices.Contains(e.Removed, path)
}
// union is a and b without repeats, in order.
func union(a, b []string) []string {
out := append([]string{}, a...)
for _, s := range b {
if !slices.Contains(out, s) {
out = append(out, s)
}
}
return out
}
// batchMember is one repository's merge in a batch and what it moves.
type batchMember struct {
merge link.SourceMoved
moved []string
}
// afterLines are the `after: <repository>` lines of a pull request's description, as mesh-delivery reads them.
var afterLines = regexp.MustCompile(`(?im)^\s*after:\s*([A-Za-z0-9_.\-/]+)\s*$`)
// groupOrderEdges is the order of every delivery group among a batch's merges as edges of its walk's tiers
// (ADR 0276 decision 5, ADR 0249 unchanged): merges of two repositories or more from one head branch name are
// a group; each pair its order rules give puts the modules of the member after in a later tier than the
// member before's. A pair in a contradiction orders nothing: its group was refused before it merged.
func groupOrderEdges(members []batchMember, entries []inventory.Entry, read map[string][]inventory.ReadRepository,
edges []inventory.Edge) []inventory.Edge {
byHead := map[string][]batchMember{}
for _, m := range members {
if m.merge.Head != "" {
byHead[m.merge.Head] = append(byHead[m.merge.Head], m)
}
}
var out []inventory.Edge
heads := make([]string, 0, len(byHead))
for h := range byHead {
heads = append(heads, h)
}
sort.Strings(heads)
for _, h := range heads {
group := byHead[h]
repos := map[string]bool{}
for _, m := range group {
repos[strings.ToLower(m.merge.Owner+"/"+m.merge.Repo)] = true
}
if len(repos) < 2 {
continue
}
var om []orderMember
moved := map[string][]string{}
for _, m := range group {
repository := m.merge.Owner + "/" + m.merge.Repo
id := repository + "@" + m.merge.Commit
var after []string
for _, l := range afterLines.FindAllStringSubmatch(m.merge.Body, -1) {
after = append(after, strings.TrimSuffix(l[1], ".git"))
}
om = append(om, orderMember{ID: id, Repository: repository, Base: m.merge.Base, Head: m.merge.Commit,
Number: m.merge.Number, Paths: m.merge.Paths, PathsTruncated: m.merge.PathsTruncated,
Removed: m.merge.Removed, ModuleDirs: m.merge.ModuleDirs, ModuleDirsSaid: m.merge.ModuleDirsSaid,
CloneURL: m.merge.CloneURL, After: after})
moved[id] = m.moved
}
order := orderOf(om, reachOfMembers(om, entries, read, edges), edges)
contradicted := map[string]bool{}
for _, c := range order.Contradictions {
for _, id := range c.Members {
contradicted[id] = true
}
}
for _, p := range order.Pairs {
if contradicted[p.Before] && contradicted[p.After] {
continue
}
for _, before := range moved[p.Before] {
for _, after := range moved[p.After] {
if before != after {
out = append(out, inventory.Edge{From: after, To: before, Kind: edgeGroupOrder})
}
}
}
fmt.Printf(" group %s: %s before %s (%s) inside the walk\n", h, p.Before, p.After, p.Why)
}
}
return out
}
// searchAfterFailures starts the search a failed walk leaves (ADR 0276 decision 3): every merge it answered
// in a later commit of its repository waits to be walked on its own commit. A stopped walk starts none: a
// person ended it.
func searchAfterFailures(ctx context.Context, inv *inventory.Inventory) error {
plans, err := inv.RecentPlans(ctx, 20)
if err != nil {
return err
}
for _, p := range plans {
if p.State != inventory.PlanFailed || p.Delivery == nil || p.Delivery.Stopped != "" ||
strings.HasPrefix(p.Note, "stopped by hand") || strings.HasPrefix(p.Note, "closed by hand") {
continue
}
merges, err := inv.MergesOf(ctx, p.ID)
if err != nil {
return err
}
for _, m := range merges {
if m.Alone || p.CommitOf(m.Repository) == m.Commit {
continue // its own commit was walked
}
if err := inv.AnswerMerge(ctx, m.Repository, m.Commit, "", true); err != nil {
return err
}
fmt.Printf("%s failed; %s %.8s, which it carried in %.8s, is walked on its own commit\n", p.ID,
m.Repository, m.Commit, p.CommitOf(m.Repository))
}
}
return nil
}
// coverAlone gives every merge waiting to be walked alone to a walk done that holds a later merge of its
// branch — one walked alone, or the failed walk retried by a person: the search stops at the first merge
// delivered, and that walk answers the older ones. What is left waits, the newest first.
func coverAlone(ctx context.Context, inv *inventory.Inventory, alone []inventory.BatchedMerge) ([]inventory.BatchedMerge, error) {
var left []inventory.BatchedMerge
for _, m := range alone {
later, err := inv.LaterMergesOf(ctx, m.Repository, m.Branch, m.Merged)
if err != nil {
return nil, err
}
covered := false
for _, l := range later {
p, err := inv.PlanByID(ctx, l.Plan)
if err != nil {
return nil, err
}
if p.State != inventory.PlanDone {
continue
}
if err := inv.AnswerMerge(ctx, m.Repository, m.Commit, p.ID, true); err != nil {
return nil, err
}
if err := nameMerges(ctx, inv, &p); err != nil {
return nil, err
}
fmt.Printf("%s %.8s is answered by %s, which delivered a later merge of its branch on its own commit\n",
m.Repository, m.Commit, p.ID)
covered = true
break
}
if !covered {
left = append(left, m)
}
}
return left, nil
}
// batchWords is a batch in the operator's words (ADR 0276 decision 7): "assembling: 13 s left (at the latest
// 00:31:09); grouped: novox/mesh-catalog@48bda475 (answers 553b7191), …; plan not yet calculated".
func batchWords(b inventory.Plan, now time.Time) string {
grouped := groupedWords(b)
if b.Delivery == nil || b.Delivery.Batch == nil {
return "grouped: " + grouped + "; plan not yet calculated"
}
w := b.Delivery.Batch
if b.State == inventory.PlanQueued {
return fmt.Sprintf("queued behind %s; grouped: %s; plan not yet calculated", w.Behind, grouped)
}
left := w.ClosesAt.Sub(now)
if w.AtMost.Before(w.ClosesAt) {
left = w.AtMost.Sub(now)
}
if left < 0 {
left = 0
}
return fmt.Sprintf("assembling: %s left (at the latest %s); grouped: %s; plan not yet calculated",
secondsWords(left), w.AtMost.Local().Format("15:04:05"), grouped)
}
// groupedWords is a batch's commits, each with the earlier merges it answers.
func groupedWords(b inventory.Plan) string {
var out []string
for _, c := range b.Carried() {
s := c.Repository + "@" + short(c.Commit)
var answers []string
if b.Delivery != nil {
for _, m := range b.Delivery.Merges {
if strings.EqualFold(m.Repository, c.Repository) && m.Carried == c.Commit {
answers = append(answers, short(m.Commit))
}
}
}
if len(answers) > 0 {
s += " (answers " + strings.Join(answers, ", ") + ")"
}
out = append(out, s)
}
if len(out) == 0 {
return "nothing yet"
}
return strings.Join(out, ", ")
}
// secondsWords is a short wait as a person reads it.
func secondsWords(d time.Duration) string {
if d < 2*time.Minute {
return fmt.Sprintf("%d s", int(d.Round(time.Second)/time.Second))
}
return d.Round(time.Second).String()
}
// walkOrBatch is what a plan record is, in a word.
func walkOrBatch(p inventory.Plan) string {
if p.Batch() {
return "batch"
}
return "walk"
}
// orWords is s, or the words for when it is empty.
func orWords(s, none string) string {
if s == "" {
return none
}
return s
}
+537
View File
@@ -0,0 +1,537 @@
package main
import (
"bytes"
"io"
"os"
"slices"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The suite's tests of a merge, but for the merge window's own (this file's), plan the merge as it is heard: a
// window of nothing closes at once, so a merge with no walk open is cut into its walk at once, as every merge
// was planned before novox/hq ADR 0276. The window's tests set it through the controller's settings.
func init() { mergeWindowDefault = 0 }
// windowed is a mesh whose controller's settings give a merge window of 90 seconds and at most 10 minutes, as its
// settings file says them, and
// modules built from the catalogue (app, notes) and from three repositories of their own.
func windowed(t *testing.T) *stores {
t.Helper()
open := aMesh(t)
ctx := t.Context()
inv := open.inventory
settings := t.TempDir() + "/merge-window.json"
if err := os.WriteFile(settings, []byte(`{"merge-window": "90s", "merge-window-at-most": "10m"}`), 0o600); err != nil {
t.Fatal(err)
}
t.Setenv(mergeWindowFileVar, settings)
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "mesh-controller", Version: "1"},
inventory.Source{Repository: "novox/mesh-controller", Seat: "git", Ref: "main", BuiltFrom: "c0",
Head: "c0"}); err != nil {
t.Fatal(err)
}
for _, m := range []struct{ module, repository, path string }{
{"app", "novox/mesh-catalog", "modules/app"},
{"notes", "novox/mesh-catalog", "modules/notes"},
{"one", "novox/one", ""},
{"two", "novox/two", ""},
{"three", "novox/three", ""},
} {
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: m.module, Version: "1"},
inventory.Source{Repository: m.repository, Seat: "git", Path: m.path, Ref: "main", BuiltFrom: "c0",
Head: "c0"}); err != nil {
t.Fatal(err)
}
}
return open
}
// t0 is the moment a test's clock starts: merges are dated by it, and heard after it.
var t0 = time.Now().UTC().Truncate(time.Second)
// catalogueMerge is a merge of the catalogue changing one module's directory, made at a moment.
func catalogueMerge(commit, module string, made time.Time) link.SourceMoved {
return link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: commit,
MergedAt: made.Format(time.RFC3339Nano), Paths: []string{"modules/" + module + "/module.json"},
ModuleDirs: []string{"modules/" + module}, ModuleDirsSaid: true}
}
// repoMerge is a merge of a repository of one module's own.
func repoMerge(repo, commit string, made time.Time) link.SourceMoved {
return link.SourceMoved{Owner: "novox", Repo: repo, Base: "main", Commit: commit,
MergedAt: made.Format(time.RFC3339Nano), Paths: []string{"main.go"}}
}
// hear is a merge heard at a moment.
func hear(t *testing.T, open *stores, m link.SourceMoved, now time.Time) {
t.Helper()
if err := (following{open: open}).hearMerge(t.Context(), m, now); err != nil {
t.Fatal(err)
}
}
// cutAt is the cutter looking at a moment.
func cutAt(t *testing.T, open *stores, now time.Time) {
t.Helper()
if err := cutBatchesHeld(t.Context(), open, now); err != nil {
t.Fatal(err)
}
}
// walks is every plan record that is a walk, newest first, and the batches not yet cut.
func walks(t *testing.T, open *stores) (walks, batches []inventory.Plan) {
t.Helper()
recent, err := open.inventory.RecentPlans(t.Context(), 50)
if err != nil {
t.Fatal(err)
}
for _, p := range recent {
if p.Batch() {
batches = append(batches, p)
} else {
walks = append(walks, p)
}
}
return walks, batches
}
// The replay of issue 362: two catalogue merges 18 seconds apart are one batch, cut once into one walk at the
// later commit, which names both merges — the earlier carried by the later.
func TestTwoMergesSecondsApartAreOneWalkAtTheLaterCommit(t *testing.T) {
open := windowed(t)
asked := asksWithPaths(t)
claude := catalogueMerge("553b7191claude", "app", t0)
dunst := catalogueMerge("48bda475dunst", "notes", t0.Add(17*time.Second))
hear(t, open, claude, t0.Add(time.Second))
hear(t, open, dunst, t0.Add(18*time.Second))
if ws, bs := walks(t, open); len(ws) != 0 || len(bs) != 1 {
t.Fatalf("within the window: %d walk(s), %d batch(es)", len(ws), len(bs))
}
cutAt(t, open, t0.Add(18*time.Second+89*time.Second))
if ws, _ := walks(t, open); len(ws) != 0 {
t.Fatalf("cut before the window closed: %+v", ws)
}
cutAt(t, open, t0.Add(18*time.Second+90*time.Second))
ws, bs := walks(t, open)
if len(ws) != 1 || len(bs) != 0 {
t.Fatalf("after the window: %d walk(s), %d batch(es)", len(ws), len(bs))
}
w := ws[0]
if w.Commit != dunst.Commit || len(w.Commits) != 1 || w.Commits[0].Commit != dunst.Commit {
t.Fatalf("the walk is at %s %+v, not the later commit", w.Commit, w.Commits)
}
for _, m := range []string{"app", "notes"} {
if _, in := w.Modules[m]; !in {
t.Fatalf("the walk does not build %s, which one of its merges moved: %v", m, w.Modules)
}
}
want := []inventory.PlanMerge{{Repository: "novox/mesh-catalog", Commit: claude.Commit, Carried: dunst.Commit},
{Repository: "novox/mesh-catalog", Commit: dunst.Commit}}
if w.Delivery == nil || !slices.Equal(w.Delivery.Merges, want) {
t.Fatalf("the walk answers %+v, want %+v", w.Delivery, want)
}
if len(*asked) != 2 {
t.Fatalf("the walk asked %v", *asked)
}
// Heard again, as the bus may hand it over twice: never a second merge, nor a second walk.
hear(t, open, claude, t0.Add(5*time.Minute))
if ws, bs := walks(t, open); len(ws) != 1 || len(bs) != 0 {
t.Fatalf("a merge heard twice made %d walk(s) and %d batch(es)", len(ws), len(bs))
}
}
// Merges of three repositories in one window are one walk, one commit for each.
func TestThreeRepositoriesInOneWindowAreOneWalk(t *testing.T) {
open := windowed(t)
asksRecorded(t)
for i, r := range []string{"one", "two", "three"} {
hear(t, open, repoMerge(r, "c-"+r, t0.Add(time.Duration(i)*20*time.Second)),
t0.Add(time.Duration(i)*20*time.Second+time.Second))
}
cutAt(t, open, t0.Add(41*time.Second+90*time.Second))
ws, bs := walks(t, open)
if len(ws) != 1 || len(bs) != 0 {
t.Fatalf("%d walk(s), %d batch(es)", len(ws), len(bs))
}
w := ws[0]
if len(w.Commits) != 3 || len(w.Delivery.Merges) != 3 {
t.Fatalf("the walk carries %+v and answers %+v", w.Commits, w.Delivery.Merges)
}
for _, r := range []string{"one", "two", "three"} {
if w.CommitOf("novox/"+r) != "c-"+r {
t.Fatalf("the walk carries %s at %q", r, w.CommitOf("novox/"+r))
}
if _, in := w.Modules[r]; !in {
t.Fatalf("the walk does not build %s: %v", r, w.Modules)
}
}
// Each module's ask is recorded at its own repository's commit.
for _, r := range []string{"one", "two", "three"} {
q, found, err := open.inventory.BuildRequestByID(t.Context(), w.Modules[r].Build)
if err != nil || !found || q.Commit != "c-"+r {
t.Fatalf("%s's build was not recorded at its own commit: %+v %v %v", r, q, found, err)
}
}
}
// Each merge restarts the window: a merge at second 80 keeps the batch open until second 170; the window
// closes at most ten minutes after its first merge however busy.
func TestAMergeRestartsTheWindowAndTheMaximumClosesIt(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
hear(t, open, repoMerge("one", "c1", t0), t0)
hear(t, open, repoMerge("two", "c2", t0.Add(80*time.Second)), t0.Add(80*time.Second))
cutAt(t, open, t0.Add(100*time.Second))
if ws, _ := walks(t, open); len(ws) != 0 {
t.Fatal("the window closed 90 seconds after its first merge, not after its last")
}
cutAt(t, open, t0.Add(170*time.Second))
if ws, _ := walks(t, open); len(ws) != 1 || len(ws[0].Commits) != 2 {
t.Fatalf("the window did not close 90 seconds after its last merge: %+v", ws)
}
busy := windowed(t)
asksWithPaths(t)
var last time.Time
for i := 0; i < 12; i++ {
last = t0.Add(time.Duration(i) * time.Minute)
repo := []string{"one", "two", "three"}[i%3]
hear(t, busy, repoMerge(repo, "c"+string(rune('a'+i)), last), last)
if ws, _ := walks(t, busy); len(ws) > 0 {
if i != 10 {
t.Fatalf("a busy window closed at merge %d (minute %d), not at its maximum", i, i)
}
break
}
}
ws, _ := walks(t, busy)
if len(ws) != 1 {
t.Fatalf("ten minutes of merges a minute apart were never cut: %d walk(s)", len(ws))
}
}
// `plans` lists the batch being assembled first, in the operator's words: how long is left, when at the latest,
// what is grouped, and that the plan is not yet calculated. It says so on the bus as plan-moved, too.
func TestTheAssemblingBatchIsShown(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
var said []inventory.Plan
was := inventory.PlanSaved
inventory.PlanSaved = func(p inventory.Plan) { said = append(said, p) }
t.Cleanup(func() { inventory.PlanSaved = was })
now := time.Now().UTC()
hear(t, open, catalogueMerge("553b7191claude", "app", now.Add(-20*time.Second)), now.Add(-19*time.Second))
hear(t, open, catalogueMerge("48bda475dunst", "notes", now.Add(-2*time.Second)), now.Add(-time.Second))
hear(t, open, repoMerge("one", "a6bc0931one", now), now)
out := captured(t, func() error { return plansCommand(t.Context(), nil) })
first := strings.SplitN(out, "\n", 2)[0]
for _, want := range []string{"assembling: ", " s left (at the latest ", "grouped: novox/mesh-catalog@48bda475 " +
"(answers 553b7191), novox/one@a6bc0931; plan not yet calculated"} {
if !strings.Contains(first, want) {
t.Fatalf("plans' first line %q does not say %q", first, want)
}
}
if len(said) == 0 || said[len(said)-1].State != inventory.PlanAssembling || said[len(said)-1].Delivery.Batch == nil ||
len(said[len(said)-1].Delivery.Merges) != 3 {
t.Fatalf("the batch was not said as assembling with its merges: %+v", said)
}
}
// captured is what a command printed.
func captured(t *testing.T, run func() error) string {
t.Helper()
r, w, err := os.Pipe()
if err != nil {
t.Fatal(err)
}
was := os.Stdout
os.Stdout = w
runErr := run()
os.Stdout = was
_ = w.Close()
var b bytes.Buffer
_, _ = io.Copy(&b, r)
if runErr != nil {
t.Fatal(runErr)
}
return b.String()
}
// One walk at a time: a merge heard while a walk runs joins the next batch, which is cut when the walk ends;
// the walk that started is never superseded.
func TestAMergeDuringAWalkJoinsTheNextBatch(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
hear(t, open, repoMerge("one", "c1", t0), t0)
cutAt(t, open, t0.Add(90*time.Second))
ws, _ := walks(t, open)
if len(ws) != 1 || !ws[0].Open() {
t.Fatalf("the first batch was not cut: %+v", ws)
}
first := ws[0]
hear(t, open, repoMerge("two", "c2", t0.Add(2*time.Minute)), t0.Add(2*time.Minute))
cutAt(t, open, t0.Add(5*time.Minute))
ws, bs := walks(t, open)
if len(ws) != 1 || len(bs) != 1 || bs[0].State != inventory.PlanQueued || bs[0].Delivery.Batch.Behind != first.ID {
t.Fatalf("the merge during the walk: %d walk(s), batches %+v", len(ws), bs)
}
if got, _ := open.inventory.PlanByID(t.Context(), first.ID); !got.Open() {
t.Fatalf("the started walk was %s by the next merge", got.State)
}
if !strings.HasPrefix(batchWords(bs[0], t0.Add(5*time.Minute)), "queued behind "+first.ID) {
t.Fatalf("a queued batch reads %q", batchWords(bs[0], t0.Add(5*time.Minute)))
}
first.State = inventory.PlanDone
if err := open.inventory.SavePlan(t.Context(), &first); err != nil {
t.Fatal(err)
}
cutAt(t, open, t0.Add(6*time.Minute))
ws, bs = walks(t, open)
if len(ws) != 2 || len(bs) != 0 || ws[0].Commit != "c2" {
t.Fatalf("the queued batch was not cut when the walk ended: %+v %+v", ws, bs)
}
}
// A walk waiting for its delivery's word is folded into the next batch's walk, which answers its merges; it
// names the walk that took it over.
func TestAWaitingWalkIsFoldedIntoTheNextBatch(t *testing.T) {
open := windowed(t)
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, repoMerge("one", "c1", t0), t0)
cutAt(t, open, t0.Add(90*time.Second))
ws, _ := walks(t, open)
if len(ws) != 1 || !ws[0].Waiting() {
t.Fatalf("the walk does not wait for its word: %+v", ws)
}
waiting := ws[0]
hear(t, open, repoMerge("two", "c2", t0.Add(2*time.Minute)), t0.Add(2*time.Minute))
cutAt(t, open, t0.Add(2*time.Minute+90*time.Second))
ws, bs := walks(t, open)
if len(bs) != 0 || len(ws) != 2 {
t.Fatalf("%d walk(s), %d batch(es)", len(ws), len(bs))
}
folded, err := open.inventory.PlanByID(ctx, waiting.ID)
if err != nil {
t.Fatal(err)
}
newer := ws[0]
if folded.State != inventory.PlanSuperseded || folded.Delivery.TakenOverBy != newer.ID {
t.Fatalf("the waiting walk is %s, taken over by %q", folded.State, folded.Delivery.TakenOverBy)
}
if len(newer.Delivery.Merges) != 2 || newer.CommitOf("novox/one") != "c1" || newer.CommitOf("novox/two") != "c2" {
t.Fatalf("the newer walk answers %+v", newer.Delivery.Merges)
}
if _, in := newer.Modules["one"]; !in {
t.Fatalf("what the folded walk was to build is not built: %v", newer.Modules)
}
}
// A merge heard after a later merge of its branch was walked is answered by that walk when it builds all it
// moves — named at once on a walk done — and joins the next batch when it does not.
func TestALateMergeIsAnsweredByTheWalkOfTheLaterOne(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
ctx := t.Context()
later := catalogueMerge("48bda475dunst", "notes", t0.Add(17*time.Second))
hear(t, open, later, t0.Add(18*time.Second))
cutAt(t, open, t0.Add(2*time.Minute))
ws, _ := walks(t, open)
w := ws[0]
w.State = inventory.PlanDone
if err := open.inventory.SavePlan(ctx, &w); err != nil {
t.Fatal(err)
}
// Moves only notes, which the done walk built from the branch after it: answered there, at once.
hear(t, open, catalogueMerge("553b7191early", "notes", t0), t0.Add(15*time.Minute))
got, _ := open.inventory.PlanByID(ctx, w.ID)
if got.State != inventory.PlanDone || len(got.Delivery.Merges) != 2 ||
got.Delivery.Merges[0] != (inventory.PlanMerge{Repository: "novox/mesh-catalog", Commit: "553b7191early",
Carried: later.Commit}) {
t.Fatalf("the late merge is not named by the walk that carried it: %+v", got.Delivery.Merges)
}
// Moves app, which that walk never built: the next batch walks it.
hear(t, open, catalogueMerge("1111aaaaapp", "app", t0.Add(time.Second)), t0.Add(16*time.Minute))
_, bs := walks(t, open)
if len(bs) != 1 || bs[0].Delivery.Merges[0].Commit != "1111aaaaapp" {
t.Fatalf("a late merge moving what the walk never built did not join the next batch: %+v", bs)
}
}
// A failed walk marks nothing delivered: its earlier merges are walked alone, on their own commit, the newest
// first; the first delivered answers the older ones, and the search stops.
func TestAFailedWalkWalksItsEarlierMergesAlone(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
ctx := t.Context()
oldest := catalogueMerge("aaaa0001", "app", t0)
middle := catalogueMerge("bbbb0002", "app", t0.Add(10*time.Second))
newest := catalogueMerge("cccc0003", "notes", t0.Add(20*time.Second))
for i, m := range []link.SourceMoved{oldest, middle, newest} {
hear(t, open, m, t0.Add(time.Duration(i)*10*time.Second+time.Second))
}
cutAt(t, open, t0.Add(2*time.Minute))
ws, _ := walks(t, open)
failed := ws[0]
failed.State, failed.Note = inventory.PlanFailed, "notes failed to build in tier 0"
if err := open.inventory.SavePlan(ctx, &failed); err != nil {
t.Fatal(err)
}
cutAt(t, open, t0.Add(3*time.Minute))
got, _ := open.inventory.PlanByID(ctx, failed.ID)
if got.State != inventory.PlanFailed || len(got.Delivery.Merges) != 3 {
t.Fatalf("the failed walk's record changed: %s %+v", got.State, got.Delivery.Merges)
}
ws, _ = walks(t, open)
alone := ws[0]
if alone.ID == failed.ID || alone.Commit != middle.Commit || len(alone.Delivery.Merges) != 1 ||
alone.Delivery.Merges[0] != (inventory.PlanMerge{Repository: "novox/mesh-catalog", Commit: middle.Commit}) {
t.Fatalf("the newest earlier merge was not walked alone on its own commit: %+v", alone)
}
alone.State = inventory.PlanDone
if err := open.inventory.SavePlan(ctx, &alone); err != nil {
t.Fatal(err)
}
cutAt(t, open, t0.Add(4*time.Minute))
ws, bs := walks(t, open)
if ws[0].ID != alone.ID || len(bs) != 0 {
t.Fatalf("the search went on after a merge was delivered: %+v", ws[0])
}
done, _ := open.inventory.PlanByID(ctx, alone.ID)
if len(done.Delivery.Merges) != 2 || done.Delivery.Merges[0] != (inventory.PlanMerge{Repository: "novox/mesh-catalog",
Commit: oldest.Commit, Carried: middle.Commit}) {
t.Fatalf("the oldest merge is not answered by the walk that delivered the one after it: %+v", done.Delivery.Merges)
}
// A stopped walk starts no search: a person ended it.
hear(t, open, catalogueMerge("dddd0004", "app", t0.Add(5*time.Minute)), t0.Add(5*time.Minute))
hear(t, open, catalogueMerge("eeee0005", "notes", t0.Add(5*time.Minute+time.Second)), t0.Add(5*time.Minute+time.Second))
cutAt(t, open, t0.Add(8*time.Minute))
ws, _ = walks(t, open)
if _, err := stopWalk(ctx, open.inventory, ws[0].ID, "mesh-delivery for jochen", "not now"); err != nil {
t.Fatal(err)
}
cutAt(t, open, t0.Add(9*time.Minute))
if again, _ := walks(t, open); again[0].ID != ws[0].ID {
t.Fatalf("a stopped walk's earlier merge was walked alone: %+v", again[0])
}
}
// A delivery group's order holds inside the walk (ADR 0249 unchanged): merges of two repositories from one head
// branch name, the node-engine's before the controller's, are tiered so; without the group, one tier.
func TestAGroupsOrderBecomesTiersInsideTheWalk(t *testing.T) {
for _, grouped := range []bool{true, false} {
open := windowed(t)
asksWithPaths(t)
ctx := t.Context()
for _, m := range []string{"mesh-host"} {
if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: m, Version: "1"},
inventory.Source{Repository: "novox/" + m, Seat: "git", Ref: "main", BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
}
host := repoMerge("mesh-host", "h1", t0)
ctl := repoMerge("mesh-controller", "k1", t0.Add(time.Second))
host.Head, ctl.Head = "feat/together", "feat/together"
if !grouped {
ctl.Head = "feat/alone"
}
hear(t, open, ctl, t0.Add(2*time.Second))
hear(t, open, host, t0.Add(3*time.Second))
cutAt(t, open, t0.Add(2*time.Minute))
ws, _ := walks(t, open)
if len(ws) != 1 {
t.Fatalf("%d walks", len(ws))
}
tiers := ws[0].Tiers
if grouped {
if len(tiers) != 2 || !slices.Equal(tiers[0], []string{"mesh-host"}) ||
!slices.Equal(tiers[1], []string{"mesh-controller"}) {
t.Fatalf("the group's order is not the walk's tiers: %v", tiers)
}
} else if len(tiers) != 1 {
t.Fatalf("two merges of no group were ordered: %v", tiers)
}
}
}
// A restarted controller resumes the window where it stood: the batch, its merges and their times are in the
// store, a merge handed over again is not doubled, and the batch is cut when its window closes.
func TestABatchSurvivesARestart(t *testing.T) {
open := windowed(t)
asksWithPaths(t)
hear(t, open, repoMerge("one", "c1", t0), t0)
hear(t, open, repoMerge("two", "c2", t0.Add(30*time.Second)), t0.Add(30*time.Second))
// A new controller: nothing of the first one's but the store.
again, err := openStores(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(again.Close)
hear(t, again, repoMerge("one", "c1", t0), t0.Add(60*time.Second)) // the bus hands it over again
cutAt(t, again, t0.Add(119*time.Second))
if ws, bs := walks(t, again); len(ws) != 0 || len(bs) != 1 || len(bs[0].Delivery.Merges) != 2 {
t.Fatalf("the restarted controller's batch: %d walk(s), %+v", len(ws), bs)
}
cutAt(t, again, t0.Add(120*time.Second))
ws, bs := walks(t, again)
if len(ws) != 1 || len(bs) != 0 || len(ws[0].Delivery.Merges) != 2 {
t.Fatalf("the restarted controller did not cut the batch at its window: %+v %+v", ws, bs)
}
}
// The window's settings reach the controller in its settings file, read at every look; one that says no
// duration, or no file, is the default.
func TestTheWindowIsTheControllersSetting(t *testing.T) {
file := t.TempDir() + "/merge-window.json"
t.Setenv(mergeWindowFileVar, file)
for _, c := range []struct {
content string
window, atMost time.Duration
}{
{`{"merge-window": "60s", "merge-window-at-most": "5m"}`, time.Minute, 5 * time.Minute},
{`{"merge-window": 30, "merge-window-at-most": ""}`, 30 * time.Second, mergeWindowAtMostDefault},
{`{"merge-window": "soon"}`, mergeWindowDefault, mergeWindowAtMostDefault},
} {
if err := os.WriteFile(file, []byte(c.content), 0o600); err != nil {
t.Fatal(err)
}
if w, m := mergeWindowOf(); w != c.window || m != c.atMost {
t.Errorf("%s reads as %s and %s, want %s and %s", c.content, w, m, c.window, c.atMost)
}
}
t.Setenv(mergeWindowFileVar, file+".gone")
if w, m := mergeWindowOf(); w != mergeWindowDefault || m != mergeWindowAtMostDefault {
t.Errorf("no file reads as %s and %s", w, m)
}
}
// S16 names every merge a waiting walk answers, never one commit a supersession may have ended (issue 362).
func TestAWaitingWalkNamesTheMergesItAnswers(t *testing.T) {
f := calm(t0)
f.waits = []waitFacts{{id: "plan-2", repository: "novox/mesh-catalog", commit: "48bda475dunst", awaits: "mesh-delivery",
since: t0.Add(-31 * time.Minute), merges: []inventory.PlanMerge{
{Repository: "novox/mesh-catalog", Commit: "553b7191claude", Carried: "48bda475dunst"},
{Repository: "novox/mesh-catalog", Commit: "48bda475dunst"}}}}
got := watchWaits(f)
if len(got) != 1 || !strings.Contains(got[0].Summary, "novox/mesh-catalog@553b7191") ||
!strings.Contains(got[0].Summary, "novox/mesh-catalog@48bda475") {
t.Fatalf("the waiting walk does not name both merges: %+v", got)
}
}
+103
View File
@@ -0,0 +1,103 @@
package main
import (
"context"
"fmt"
"time"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// **Every wait of a batch has a bound and a condition** (novox/hq ADR 0276 decision 9): a batch still
// assembling a minute past its maximum while no walk is open is a cut the controller failed to make (S18),
// and a batch waiting behind an open walk longer than that walk's bound waits on a walk gone wrong (S19).
// The kinds S18 and S19 raise.
const (
kindBatchNotCut = "batch-not-cut"
kindBatchBehindWalk = "batch-behind-walk"
)
// batchFacts is one batch not yet cut, as the watchdogs read it.
type batchFacts struct {
id, state, grouped string
// atMost is when its window closes at the latest; closed when it closed.
atMost, closed time.Time
// behind is the open walk it waits behind, and walkBound that walk's bound: a tier's bound for each of its
// tiers.
behind string
walkBound time.Duration
}
// gatherBatches is every batch not yet cut, with the bound of the walk it waits behind.
func gatherBatches(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]batchFacts, error) {
batches, err := inv.Batches(ctx)
if err != nil || len(batches) == 0 {
return nil, err
}
bounds, err := measuredTierBounds(ctx, inv, now)
if err != nil {
return nil, err
}
var out []batchFacts
for _, b := range batches {
f := batchFacts{id: b.ID, state: b.State, grouped: groupedWords(b)}
if w := b.Delivery; w != nil && w.Batch != nil {
f.atMost, f.closed, f.behind = w.Batch.AtMost, w.Batch.ClosesAt, w.Batch.Behind
if f.atMost.Before(f.closed) {
f.closed = f.atMost
}
}
if f.behind != "" {
if walk, err := inv.PlanByID(ctx, f.behind); err == nil {
f.walkBound = bounds.of(walk.Repository) * time.Duration(max(1, len(walk.Tiers)))
}
}
out = append(out, f)
}
return out, nil
}
// watchBatchesNotCut is S18: a batch still assembling a minute past its maximum, while no walk is open.
func watchBatchesNotCut(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, b := range f.batches {
if b.state != inventory.PlanAssembling || b.atMost.IsZero() || f.now.Sub(b.atMost) <= batchLateAfter {
continue
}
late := f.now.Sub(b.atMost)
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: b.id, Kind: kindBatchNotCut,
Severity: conditions.Warning,
Summary: fmt.Sprintf("the batch %s is still assembling %s past its latest close (%s), with no walk open: "+
"the controller did not cut it; grouped: %s", b.id, ago(late), b.atMost.UTC().Format(time.RFC3339), b.grouped),
Said: fmt.Sprintf("assembling since %s past its maximum", ago(late)),
Headline: "Merged changes are not being delivered",
Explanation: fmt.Sprintf("Merges collected for one delivery should have been planned %s ago and were not. "+
"Nothing is lost; the mesh keeps them until it plans them.", humanDuration(late)),
Resolved: "Merged changes are being delivered again"})
}
return out
}
// watchBatchesBehind is S19: a batch waiting behind an open walk longer than that walk's bound, naming it.
func watchBatchesBehind(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, b := range f.batches {
if b.state != inventory.PlanQueued || b.behind == "" || b.walkBound <= 0 || f.now.Sub(b.closed) <= b.walkBound {
continue
}
in := f.now.Sub(b.closed)
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: b.id, Kind: kindBatchBehindWalk,
Severity: conditions.Warning,
Summary: fmt.Sprintf("the batch %s has waited %s behind the walk %s, longer than that walk's bound (%s); "+
"grouped: %s; `plans %s` says where that walk stands", b.id, ago(in), b.behind, ago(b.walkBound),
b.grouped, b.behind),
Said: fmt.Sprintf("queued behind %s for %s", b.behind, ago(in)),
Headline: "Merged changes wait behind a slow delivery",
Explanation: fmt.Sprintf("Merges collected for the next delivery have waited %s for the delivery before "+
"them, which is taking longer than it should. Nothing is lost.", humanDuration(in)),
Resolved: "Merged changes no longer wait behind a slow delivery"})
}
return out
}
+6
View File
@@ -574,6 +574,12 @@ func deliveryCommand(ctx context.Context, args []string) error {
if walks, err = inv.OpenPlans(ctx); err != nil {
return err
}
// And the batch being assembled (novox/hq ADR 0276): never a walk to link or let go, shown.
batches, err := inv.Batches(ctx)
if err != nil {
return err
}
walks = append(walks, batches...)
recent, err := inv.RecentPlans(ctx, *limit)
if err != nil {
return err
+8 -5
View File
@@ -227,7 +227,9 @@ func TestAMergeWaitsForItsDeliverysWordAndThePersonsWordWorksWithoutIt(t *testin
t.Fatalf("with no holder on record the walk waited (%v) or asked %v", p.Waiting(), *asked)
}
// The holder on record: the next merge waits, asking nothing, and an advance asks nothing either.
// The holder on record: the next merge waits, asking nothing, and an advance asks nothing either. One walk
// is open at a time (novox/hq ADR 0276): the first ends before the next merge's batch is cut.
finish(p.ID)
if _, err := inv.Assign(ctx, "anchor", "mesh-delivery"); err != nil {
t.Fatal(err)
}
@@ -259,10 +261,11 @@ func TestAMergeWaitsForItsDeliverysWordAndThePersonsWordWorksWithoutIt(t *testin
t.Fatalf("the word was not kept: %+v %v", got.Delivery, err)
}
// The delivery's owner's own merge never waits for it — and takes over what the older walk had not
// built (app, folded in: ADR 0218), which goes with it on the controller's own path.
// The delivery's owner's own merge never waits for it. A walk that started is never taken over (novox/hq
// ADR 0276): the merge's batch is cut once it ended.
finish(p.ID)
p = merge("c3cccccccc", "modules/mesh-delivery/main.go")
if p.Waiting() || len(*asked) != 4 {
if p.Waiting() || len(*asked) != 3 {
t.Fatalf("mesh-delivery's own walk waited for mesh-delivery: %v %v", p.Waiting(), *asked)
}
@@ -280,7 +283,7 @@ func TestAMergeWaitsForItsDeliverysWordAndThePersonsWordWorksWithoutIt(t *testin
}
advanceHeld(ctx, open)
got, _ = inv.PlanByID(ctx, p.ID)
if got.Waiting() || !strings.HasPrefix(got.Delivery.By, "a person") || len(*asked) < 5 {
if got.Waiting() || !strings.HasPrefix(got.Delivery.By, "a person") || len(*asked) < 4 {
t.Fatalf("a person's word did not start the walk: %+v, asked %v", got.Delivery, *asked)
}
+1 -1
View File
@@ -625,7 +625,7 @@ func planStale(ctx context.Context, inv *inventory.Inventory, p inventory.Plan)
}
for _, newer := range recent {
if newer.ID == p.ID || !newer.Created.After(p.Created) || !repositoryMatches(newer.Repository, p.Repository) ||
newer.Branch != p.Branch || newer.State == inventory.PlanSuperseded {
newer.Branch != p.Branch || newer.State == inventory.PlanSuperseded || newer.Batch() {
continue
}
return inventory.PlanSuperseded, fmt.Sprintf("superseded by %s (%s %s), a newer merge of the same repository "+
-194
View File
@@ -1,194 +0,0 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// novox/hq issue 349: on 2026-10-09 the plan of mesh-catalog at a082615b (the merge of a security fix to the
// forge's module) was "superseded at tier 0 by" the plan at 8ff8197a, the merge before it, which the
// catch-up acted on late; what the later plan had not built was folded into a plan at the commit before the
// fix. Plans of one branch are ordered by when the forge made their merges, and a merge older than an open
// plan of its branch is planned at that plan's commit.
// TestTwoMergesActedOnInReverseOrderBuildTheNewerCommit replays it: the later merge acted on first, the
// earlier one second (the catch-up). One plan is left open, at the later commit, and it builds both.
func TestTwoMergesActedOnInReverseOrderBuildTheNewerCommit(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
for _, m := range []string{"gitea", "notes"} {
if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: m, Version: "1"},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + m, Ref: "main",
BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
}
at := time.Now().UTC().Add(-20 * time.Minute).Truncate(time.Second)
fix := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "a082615bfix",
MergedAt: at.Add(2 * time.Minute).Format(time.RFC3339Nano),
Paths: []string{"modules/gitea/module.json"}, ModuleDirs: []string{"modules/gitea"}, ModuleDirsSaid: true}
before := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "8ff8197abefore",
MergedAt: at.Format(time.RFC3339Nano),
Paths: []string{"modules/notes/module.json"}, ModuleDirs: []string{"modules/notes"}, ModuleDirsSaid: true}
for _, m := range []link.SourceMoved{fix, before} {
if err := (following{open: open}).SourceMoved(ctx, m); err != nil {
t.Fatal(err)
}
}
plans, err := open.inventory.OpenPlans(ctx)
if err != nil {
t.Fatal(err)
}
if len(plans) != 1 {
var said []string
for _, p := range plans {
said = append(said, p.ID+" "+p.Commit+" "+p.Note)
}
t.Fatalf("open plans: %s", strings.Join(said, "; "))
}
p := plans[0]
if p.Commit != fix.Commit {
t.Fatalf("the open plan builds %s, not the newer commit %s", p.Commit, fix.Commit)
}
for _, m := range []string{"gitea", "notes"} {
if _, has := p.Modules[m]; !has {
t.Fatalf("the open plan at the newer commit does not build %s: %v", m, p.Modules)
}
}
}
// A late older merge after the newer plan is done (review of PR 179, A1): what it moved is built from the
// newer commit, and what the newer merge already looked at is not built again at the older one.
func TestALateMergeAfterTheNewerPlanEndedBuildsTheNewerCommit(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
for _, m := range []string{"gitea", "notes"} {
if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: m, Version: "1"},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + m, Ref: "main",
BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
}
at := time.Now().UTC().Add(-20 * time.Minute).Truncate(time.Second)
fix := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "a082615bfix",
MergedAt: at.Add(2*time.Minute + 500*time.Millisecond).Format(time.RFC3339Nano),
Paths: []string{"modules/gitea/module.json"}, ModuleDirs: []string{"modules/gitea"}, ModuleDirsSaid: true}
if err := (following{open: open}).SourceMoved(ctx, fix); err != nil {
t.Fatal(err)
}
plans, err := open.inventory.OpenPlans(ctx)
if err != nil || len(plans) != 1 {
t.Fatalf("%v %v", plans, err)
}
done := plans[0]
done.State = inventory.PlanDone
if err := open.inventory.SavePlan(ctx, &done); err != nil {
t.Fatal(err)
}
if got, err := open.inventory.PlanByID(ctx, done.ID); err != nil || !got.Merged.Equal(at.Add(2*time.Minute+500*time.Millisecond)) {
t.Fatalf("the merge time kept is %s, not to the nanosecond (%v)", got.Merged, err)
}
before := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "8ff8197abefore",
MergedAt: at.Format(time.RFC3339Nano),
Paths: []string{"modules/notes/module.json", "modules/gitea/module.json"},
ModuleDirs: []string{"modules/notes", "modules/gitea"}, ModuleDirsSaid: true}
if err := (following{open: open}).SourceMoved(ctx, before); err != nil {
t.Fatal(err)
}
plans, err = open.inventory.OpenPlans(ctx)
if err != nil || len(plans) != 1 {
t.Fatalf("open plans after the late merge: %+v %v", plans, err)
}
p := plans[0]
if p.Commit != fix.Commit {
t.Fatalf("the late merge was planned at %s, not the newer commit %s", p.Commit, fix.Commit)
}
if _, has := p.Modules["notes"]; !has {
t.Fatalf("what only the late merge moved is not built: %v", p.Modules)
}
if _, has := p.Modules["gitea"]; has {
t.Fatalf("what the newer merge already built is built again: %v", p.Modules)
}
}
// Pure: the branch's order is the merges', where both plans know it; and a later merge of the branch is
// found in any state, never a release's, another repository's or another branch's.
func TestTheBranchOrderIsTheMerges(t *testing.T) {
t0 := time.Date(2026, 10, 9, 10, 0, 0, 0, time.UTC)
newerMerge := inventory.Plan{ID: "plan-1", Repository: "novox/mesh-catalog", Branch: "main", Commit: "a082615b",
Merged: t0.Add(time.Minute), Created: t0.Add(2 * time.Minute), State: inventory.PlanRolling}
olderMerge := inventory.Plan{ID: "plan-2", Repository: "novox/mesh-catalog", Branch: "main", Commit: "8ff8197a",
Merged: t0, Created: t0.Add(10 * time.Minute), State: inventory.PlanBuilding}
if earlierOnTheBranch(newerMerge, olderMerge) || !earlierOnTheBranch(olderMerge, newerMerge) {
t.Fatal("ordered by when the plans were made, not by when the merges were")
}
if _, closed := supersededBy(olderMerge, []inventory.Plan{newerMerge}, func(string) bool { return true }); len(closed) != 0 {
t.Fatalf("the plan of an older merge superseded a newer one: %s", closed[0].Note)
}
unknown := newerMerge
unknown.Merged = time.Time{}
if !earlierOnTheBranch(unknown, olderMerge) {
t.Fatal("without a merge time the plans' own order does not stand")
}
same := olderMerge
same.Merged = newerMerge.Merged
if !earlierOnTheBranch(newerMerge, same) || earlierOnTheBranch(same, newerMerge) {
t.Fatal("one merge time: the plans' own order does not stand")
}
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "8ff8197a",
MergedAt: t0.Add(500 * time.Millisecond).Format(time.RFC3339Nano)}
newerMerge.State = inventory.PlanDone
if !laterOnTheBranch(m, newerMerge) {
t.Fatal("a later merge of the branch, its plan done, was not found")
}
for name, change := range map[string]func(p *inventory.Plan, m *link.SourceMoved){
"another branch": func(_ *inventory.Plan, m *link.SourceMoved) { m.Base = "release" },
"another repository": func(p *inventory.Plan, _ *link.SourceMoved) { p.Repository = "novox/mesh-controller" },
"a release": func(p *inventory.Plan, _ *link.SourceMoved) { p.Release = &inventory.PlanRelease{} },
"no merge time": func(p *inventory.Plan, _ *link.SourceMoved) { p.Merged = time.Time{} },
"the same commit": func(p *inventory.Plan, m *link.SourceMoved) { m.Commit = p.Commit },
"an older merge": func(p *inventory.Plan, _ *link.SourceMoved) { p.Merged = t0 },
"the same moment": func(p *inventory.Plan, _ *link.SourceMoved) { p.Merged = t0.Add(500 * time.Millisecond) },
"an unreadable time": func(_ *inventory.Plan, m *link.SourceMoved) { m.MergedAt = "yesterday" },
} {
p, mm := newerMerge, m
change(&p, &mm)
if laterOnTheBranch(mm, p) {
t.Errorf("%s was taken as a later merge of the branch", name)
}
}
// To the nanosecond: two merges within a second keep their order.
m.MergedAt = t0.Add(time.Minute + 200*time.Millisecond).Format(time.RFC3339Nano)
if !laterOnTheBranch(m, inventory.Plan{Repository: "novox/mesh-catalog", Branch: "main", Commit: "x",
Merged: t0.Add(time.Minute + 700*time.Millisecond)}) {
t.Fatal("merges within one second lost their order")
}
}
// A merge time with a fraction of a second is kept whole in its plan, and two merges within one second keep
// their order through the plans and the lookup (review of PR 179).
func TestAMergeTimeKeepsItsFractionOfASecond(t *testing.T) {
at := "2026-10-09T10:57:52.123456789Z"
p := planOfMerge(link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "c1", MergedAt: at}, nil, nil)
want := time.Date(2026, 10, 9, 10, 57, 52, 123456789, time.UTC)
if !p.Merged.Equal(want) {
t.Fatalf("the plan's merge time is %s, not %s", p.Merged, want)
}
earlier := planOfMerge(link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "c0",
MergedAt: "2026-10-09T10:57:52.123456788Z"}, nil, nil)
earlier.Created = p.Created.Add(time.Second) // made after, merged before
if !earlierOnTheBranch(earlier, p) || earlierOnTheBranch(p, earlier) {
t.Fatal("two merges a nanosecond apart lost their order")
}
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "c0", MergedAt: "2026-10-09T10:57:52.123456788Z"}
if !laterOnTheBranch(m, p) {
t.Fatal("a merge a nanosecond later was not found as the later one")
}
}
+27 -1
View File
@@ -40,6 +40,32 @@ type merges interface {
AnnouncedMerges(ctx context.Context, since time.Time) ([]link.AnnouncedMerge, error)
}
// unheardMerges is the announcements of merges the controller has not heard (novox/hq ADR 0276): a merge kept
// in a batch is not acted on until its batch is cut, which can be past mergeGrace behind a long walk, and is
// not missed for that.
type unheardMerges struct {
merges
inv *inventory.Inventory
}
func (u unheardMerges) AnnouncedMerges(ctx context.Context, since time.Time) ([]link.AnnouncedMerge, error) {
all, err := u.merges.AnnouncedMerges(ctx, since)
if err != nil {
return nil, err
}
var out []link.AnnouncedMerge
for _, a := range all {
_, kept, err := u.inv.MergeOf(ctx, a.Owner+"/"+a.Repo, a.Commit)
if err != nil {
return nil, err
}
if !kept {
out = append(out, a)
}
}
return out, nil
}
// 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}
@@ -61,7 +87,7 @@ func catchingUpOnMerges(ctx context.Context, open *stores, announced merges) {
case <-tick.C:
}
watchedMerges.begin()
err := catchUpOnMerges(ctx, time.Now(), announced, catalogued, f.SourceMoved, func(format string, args ...any) {
err := catchUpOnMerges(ctx, time.Now(), unheardMerges{announced, open.inventory}, 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
+11
View File
@@ -362,6 +362,17 @@ var plainWordings = map[string]func(conditions.Observation) words{
Explanation: "Its walk across the machines has not moved for longer than usual. Nothing is lost.",
Resolved: "Resolved: the delivery moves again"}
}),
kindBatchNotCut: worded(func(o conditions.Observation) words {
return words{Headline: "Merged changes are not being delivered",
Explanation: "Merges collected for one delivery should have been planned and were not. Nothing is lost.",
Resolved: "Resolved: the merged changes are being delivered"}
}),
kindBatchBehindWalk: worded(func(o conditions.Observation) words {
return words{Headline: "Merged changes wait behind a slow delivery",
Explanation: "Merges collected for the next delivery wait for the delivery before them, which is taking " +
"longer than it should. Nothing is lost.",
Resolved: "Resolved: the merged changes no longer wait"}
}),
kindWalkWaiting: worded(func(o conditions.Observation) words {
return words{Headline: "A delivery is waiting to start",
Explanation: "A merged change is built, and mesh-delivery (the module that decides when a delivery goes " +
+13 -1
View File
@@ -164,12 +164,13 @@ func retryRefusal(p inventory.Plan, plans []inventory.Plan) error {
case p.Open():
return fmt.Errorf("%s is still %s; nothing in it failed to retry — `rebuild <module>` asks one module again", p.ID, p.State)
}
if len(failedIn(p)) > 0 {
if q, found := newerOpenPlan(p, plans); found {
return fmt.Errorf("%s supersedes it: a newer merge of %s (%s at %s) is open, and retrying %s would build "+
"what that one replaced", q.ID, q.Repository, q.ID, short(q.Commit), p.ID)
}
return nil
return oneWalkAtATime(p, plans)
}
stopped := stoppedRollouts(p)
if len(stopped) == 0 {
@@ -193,6 +194,17 @@ func retryRefusal(p inventory.Plan, plans []inventory.Plan) error {
"the older build back", m, q.ID, q.State, q.Repository, short(q.Commit), p.ID)
}
}
return oneWalkAtATime(p, plans)
}
// oneWalkAtATime refuses a retry while another walk that started is open (novox/hq ADR 0276): a walk retried
// beside it would be two walks at once.
func oneWalkAtATime(p inventory.Plan, plans []inventory.Plan) error {
for _, q := range plans {
if q.ID != p.ID && q.Open() && q.Release == nil && !q.Waiting() {
return fmt.Errorf("%s is open (%s): one walk at a time — retry %s once it ended", q.ID, q.Named(), p.ID)
}
}
return nil
}
+3
View File
@@ -182,6 +182,9 @@ func serve(ctx context.Context) (err error) {
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
// machines to report moves when they have, and a plan left by a replaced controller resumes.
go planTicker(ctx, open)
// And the merge window (novox/hq ADR 0276): a batch is cut into its walk when the window closes, which no
// merge or outcome says.
go batchCutter(ctx, open)
// And the pending assignments settled on a tick of their own (novox/hq ADR 0261): made once their module
// is registered, ended with why when its build will not register it, raised and cleared as conditions.
// Never by a read.
+11 -4
View File
@@ -135,12 +135,19 @@ func TestRetryRefusesWhatItCannotResume(t *testing.T) {
if err := retryRefusal(stopped, nil); err == nil || !strings.Contains(err.Error(), "nothing in tier 0") {
t.Errorf("a plan with nothing failed to build was retried: %v", err)
}
// Another branch's newer plan, an older one, and a failed one do not supersede it.
// A failed or done plan does not supersede it.
if err := retryRefusal(failed, []inventory.Plan{plan("plan-4", inventory.PlanFailed, at.Add(time.Hour)),
plan("plan-5", inventory.PlanDone, at.Add(time.Hour))}); err != nil {
t.Errorf("refused for a plan that does not supersede it: %v", err)
}
// Another branch's open walk, or an older one, does not supersede it, and is one walk open: one at a time
// (novox/hq ADR 0276).
other := plan("plan-3", inventory.PlanBuilding, at.Add(time.Hour))
other.Branch = "release"
if err := retryRefusal(failed, []inventory.Plan{other, plan("plan-0", inventory.PlanBuilding, at.Add(-time.Hour)),
plan("plan-4", inventory.PlanFailed, at.Add(time.Hour))}); err != nil {
t.Errorf("refused for a plan that does not supersede it: %v", err)
for _, open := range []inventory.Plan{other, plan("plan-0", inventory.PlanBuilding, at.Add(-time.Hour))} {
if err := retryRefusal(failed, []inventory.Plan{open}); err == nil || !strings.Contains(err.Error(), "one walk at a time") {
t.Errorf("retried beside the open walk %s: %v", open.ID, err)
}
}
}
+57 -108
View File
@@ -141,7 +141,8 @@ func reachableFrom(moved []string, edges []inventory.Edge) []string {
// nothing it builds, and a new controller changes nothing about the holder it orders. A
// packages edge, read from a record made before novox/hq ADR 0267, widens nothing either:
// a shared file moves each module whose build source holds it, directly.
if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf || e.Kind == inventory.EdgePackages {
if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf || e.Kind == inventory.EdgePackages ||
e.Kind == edgeGroupOrder {
continue
}
if in[e.To] && !in[e.From] {
@@ -180,115 +181,33 @@ func hasCycle(tiers [][]string, edges []inventory.Edge) bool {
return false
}
// planFor is the plan a merge produces: the moved modules and everything reachable from them,
// planOfMerge is the plan one merge produces: the moved modules and everything reachable from them,
// tiered, with the merge it answers.
func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inventory.Plan {
p := planOfMoves(moved, edges)
merged, _ := time.Parse(time.RFC3339Nano, m.MergedAt)
p.Repository, p.Branch, p.Commit, p.Merged = m.Owner+"/"+m.Repo, m.Base, m.Commit, merged.UTC()
return p
}
// planOfMoves is the walk of a set of moved modules (novox/hq ADR 0162, 0276): they and everything reachable
// from them, tiered along the graph, with no merge named yet.
func planOfMoves(moved []string, edges []inventory.Edge) inventory.Plan {
set := reachableFrom(moved, edges)
tiers := tiersOf(set, edges)
modules := map[string]*inventory.PlanModule{}
for _, name := range set {
modules[name] = &inventory.PlanModule{}
}
merged, _ := time.Parse(time.RFC3339Nano, m.MergedAt)
return inventory.Plan{
ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()),
Repository: m.Owner + "/" + m.Repo,
Branch: m.Base,
Commit: m.Commit,
Merged: merged.UTC(),
Created: time.Now().UTC(),
State: inventory.PlanBuilding,
Tiers: tiers,
Modules: modules,
ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()),
Created: time.Now().UTC(),
State: inventory.PlanBuilding,
Tiers: tiers,
Modules: modules,
}
}
// supersededBy is what a newer plan takes over from the open plans it supersedes (novox/hq issue
// 254, ADR 0218): the modules they had not finished, and those plans closed as superseded.
//
// **A merge looked at no plan but its own.** Two merges of one repository a few minutes apart were
// two open plans asking for the same modules, each sending machines what it built; and a plan that
// would never move again — waiting on a report that could not come, at 97b1b2b — stayed open for
// ever beside the newer ones, read as work in progress by everyone who looked. The newer merge is the
// newer intent for that repository and branch, so its plan takes over: every open plan of the same
// repository and branch **created before it** — by the time the plans were made, never by comparing
// commits, which have no order of their own — gives up the modules it had not built, and those are
// planned again in the newer plan beside what the newer merge moved.
//
// "Not built" is a module not yet asked, or asked and not answered; **and a module built and not
// yet sent to its machines**, where its policy rolls it out: closed, the older plan would never send
// it, and the catalogue announces no move for a rebuild (issue 189), so the newer plan builds and
// sends it. A build the older plan asked still finishes and registers as any build does — ordered by
// when it was asked (issue 219), so the newer plan's ask, made later, is the one that stands.
//
// A plan with no branch recorded is from before branches were kept, and is superseded by the next
// plan of its repository: what it had not built is folded in, so nothing is lost by it.
//
// **Newer is the branch's order, not the plans'** (novox/hq issue 349). A merge the bus did not hand
// over is acted on late, by the catch-up, so its plan is made after the plan of a merge that came after
// it — and that later-made plan of the earlier commit superseded the later merge's, and built what it
// folded in from the commit before the later merge: twice on 2026-10-09, once a security fix. So where
// both plans know when their merge was made, that decides; only where one does not do the plans' own
// times. A merge older than the newest planned merge of its branch never reaches here at its own commit:
// it is planned at that merge's commit, which contains it (laterOnTheBranch).
func supersededBy(newer inventory.Plan, open []inventory.Plan, rollsOut func(string) bool) ([]string, []inventory.Plan) {
folded := map[string]bool{}
var closed []inventory.Plan
for _, old := range open {
if old.ID == newer.ID || !old.Open() || !strings.EqualFold(old.Repository, newer.Repository) ||
(old.Branch != "" && old.Branch != newer.Branch) || !earlierOnTheBranch(old, newer) {
continue
}
var took []string
for name, s := range old.Modules {
if s != nil && s.State == planDeleted {
continue // deleted at its source: nothing to plan again
}
if s == nil || s.State != "built" || (s.SentAt == nil && rollsOut(name)) {
folded[name] = true
took = append(took, name)
}
}
sort.Strings(took)
old.State = inventory.PlanSuperseded
old.Note = fmt.Sprintf("superseded at tier %d by %s (%s at %s)", old.Tier, newer.ID, newer.Repository, short(newer.Commit))
if len(took) > 0 {
old.Note += "; " + strings.Join(took, ", ") + " planned there again"
}
closed = append(closed, old)
}
out := make([]string, 0, len(folded))
for name := range folded {
out = append(out, name)
}
sort.Strings(out)
return out, closed
}
// earlierOnTheBranch says plan a answers a merge made before b's: by when the forge made each merge where
// both are known, else by when each plan was made. A merge's own plan made again at the same commit
// (the same merge time) is ordered by when it was made. Pure.
func earlierOnTheBranch(a, b inventory.Plan) bool {
if !a.Merged.IsZero() && !b.Merged.IsZero() && !a.Merged.Equal(b.Merged) {
return a.Merged.Before(b.Merged)
}
return a.Created.Before(b.Created)
}
// laterOnTheBranch says the plan newest answers a merge into m's branch made after m: then newest's commit
// contains m's change, and m is planned there, so the newest commit of the branch is what is built (novox/hq
// issue 349). newest is the newest merge's plan of the branch in any state: a later plan already done
// built what it moved at its commit, and m planned at its own would build m's dependents back at the older
// one. Pure.
func laterOnTheBranch(m link.SourceMoved, newest inventory.Plan) bool {
merged, err := time.Parse(time.RFC3339Nano, m.MergedAt)
if err != nil || newest.Release != nil || newest.Commit == m.Commit {
return false
}
return strings.EqualFold(newest.Repository, m.Owner+"/"+m.Repo) && newest.Branch == m.Base &&
newest.Merged.After(merged)
}
// gates is what the next tier needs running from this one: a module of the tier that a later
// tier is built by — the runtime dependency — and whose policy rolls it out, must be applied by
// the machines running it before the next tier is asked. A base an image stands on need only be
@@ -316,8 +235,11 @@ func gates(p inventory.Plan, edges []inventory.Edge, rollsOut func(string) bool)
seen := map[string]bool{}
var out []string
for _, e := range edges {
if later[e.From] && inTier[e.To] && e.Kind == inventory.EdgeBuiltBy && !seen[e.To] && rollsOut(e.To) &&
!isBaseOf(e.From, e.To, edges, all) {
// A delivery group's order inside a walk (novox/hq ADR 0276 decision 5) gates as built-by does: the
// member before is running on its machines before the member after is asked.
ordered := e.Kind == edgeGroupOrder ||
e.Kind == inventory.EdgeBuiltBy && !isBaseOf(e.From, e.To, edges, all)
if later[e.From] && inTier[e.To] && ordered && !seen[e.To] && rollsOut(e.To) {
seen[e.To] = true
out = append(out, e.To)
}
@@ -417,8 +339,14 @@ func askModule(ctx context.Context, p *inventory.Plan, name string, byName map[s
p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err)
return
}
// The commit of the module's own repository the walk carries (novox/hq ADR 0276): a batch's walk carries
// one per repository.
commit := p.CommitOf(e.Source.Repository)
if commit == "" {
commit = p.Commit
}
recordAsked(ctx, inventory.BuildRequest{ID: id, Repository: e.Source.Repository, Seat: e.Source.Seat,
Path: e.Source.Path, Ref: followedBranch(e.Source.Ref), Commit: p.Commit, For: "plan"})
Path: e.Source.Path, Ref: followedBranch(e.Source.Ref), Commit: commit, For: "plan"})
state.State = "asked"
state.AskedAt = &now
state.Build = id
@@ -537,6 +465,12 @@ func advancePlans(ctx context.Context, open *stores) {
// advanceHeld is advancePlans for a caller already holding the plans.
func advanceHeld(ctx context.Context, open *stores) {
inv := open.inventory
// A walk that ended lets the next batch be cut at once, not at the cutter's next look (novox/hq ADR 0276).
defer func() {
if err := cutBatchesHeld(ctx, open, time.Now().UTC()); err != nil {
fmt.Printf("batches: %v\n", err)
}
}()
plans, err := inv.OpenPlans(ctx)
if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err)
@@ -1244,20 +1178,23 @@ func inTierSince(p inventory.Plan) time.Time {
func planLineWith(p inventory.Plan, now time.Time, pause pauseView, bound time.Duration) string {
where := fmt.Sprintf("tier %d of %d", min(p.Tier+1, len(p.Tiers)), len(p.Tiers))
switch p.State {
case inventory.PlanAssembling, inventory.PlanQueued:
// A batch not yet a walk (novox/hq ADR 0276): what it holds and how long is left.
return batchWords(p, now)
case inventory.PlanDone:
return fmt.Sprintf("%s %s done, %d tier(s)", p.Repository, short(p.Commit), len(p.Tiers))
return fmt.Sprintf("%s done, %d tier(s)", p.Named(), len(p.Tiers))
case inventory.PlanFailed:
return fmt.Sprintf("%s %s FAILED at %s: %s", p.Repository, short(p.Commit), where, p.Note)
return fmt.Sprintf("%s FAILED at %s: %s", p.Named(), where, p.Note)
case inventory.PlanSuperseded:
return fmt.Sprintf("%s %s %s", p.Repository, short(p.Commit), p.Note)
return fmt.Sprintf("%s %s", p.Named(), p.Note)
}
since := now.Sub(inTierSince(p)).Round(time.Second)
if p.Waiting() {
// Waiting for its delivery's word is no lateness of the walk's (novox/hq ADR 0239).
return fmt.Sprintf("%s %s %s, %s, for %s", p.Repository, short(p.Commit), where, waitingNote(p), since)
return fmt.Sprintf("%s %s, %s, for %s", p.Named(), where, waitingNote(p), since)
}
if waiting, paused := pausedWaiting(p, pause, now); paused {
return fmt.Sprintf("%s %s %s, %s", p.Repository, short(p.Commit), where, waiting)
return fmt.Sprintf("%s %s, %s", p.Named(), where, waiting)
}
late := ""
if since > bound {
@@ -1267,7 +1204,7 @@ func planLineWith(p inventory.Plan, now time.Time, pause pauseView, bound time.D
if p.State == inventory.PlanRolling {
what = p.Note
}
return fmt.Sprintf("%s %s %s, %s for %s%s", p.Repository, short(p.Commit), where, what, since, late)
return fmt.Sprintf("%s %s, %s for %s%s", p.Named(), where, what, since, late)
}
// planFailedBuild marks the module a failed build was for when the result names no module: by the
@@ -1534,13 +1471,25 @@ func plansCommand(ctx context.Context, args []string) error {
if err != nil {
return err
}
if len(plans) == 0 {
// **The batch being assembled first** (novox/hq ADR 0276 decision 7): what it holds, how long is left,
// and that its plan is not yet calculated.
batches, err := inv.Batches(ctx)
if err != nil {
return err
}
if len(plans) == 0 && len(batches) == 0 {
fmt.Println("no merge has produced a plan yet")
return nil
}
for _, b := range batches {
fmt.Printf("%-28s %s\n", b.ID, batchWords(b, now))
}
pause := buildSeatPause(ctx, inv, plans)
bounds := readTierBounds(ctx, inv, now)
for _, p := range plans {
if p.Batch() {
continue
}
fmt.Printf("%-28s %s\n", p.ID, planLineWith(p, now, pause, bounds.of(p.Repository)))
}
return nil
+31 -2
View File
@@ -210,6 +210,22 @@ var signalsTable = []signalRow{
newest: func(f *signalFacts) time.Time {
return newestOf(f.waits, func(w waitFacts) time.Time { return w.since })
}},
{Row: "S18", Signal: "a batch of merges is cut into its walk", Emitter: "controller's merge window",
Trigger: "each batch (novox/hq ADR 0276)",
Bound: "a minute past merge-window-at-most, while no walk is open: the controller failed to cut it",
Kind: kindBatchNotCut, Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.plansErr }, watch: watchBatchesNotCut,
newest: func(f *signalFacts) time.Time {
return newestOf(f.batches, func(b batchFacts) time.Time { return b.atMost })
}},
{Row: "S19", Signal: "a batch waiting behind an open walk is cut when it ends", Emitter: "controller's merge window",
Trigger: "each batch closed while a walk is open (novox/hq ADR 0276)",
Bound: "the walk's bound, a tier's bound for each of its tiers, from when the batch closed: naming the walk",
Kind: kindBatchBehindWalk, Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.plansErr }, watch: watchBatchesBehind,
newest: func(f *signalFacts) time.Time {
return newestOf(f.batches, func(b batchFacts) time.Time { return b.closed })
}},
{Row: "S17", Signal: "a send held for the bus's planned step is told to a person", Emitter: "controller's plan",
Trigger: "each send refused because it would replace the bus outside its step (novox/hq issue 336)",
Bound: "none: raised at the first refusal, for the operator, naming what waits, the bus build from and to, " +
@@ -372,8 +388,8 @@ func watchWaits(f *signalFacts) []conditions.Observation {
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: w.id, Kind: kindWalkWaiting,
Severity: severity,
Summary: fmt.Sprintf("the walk of %s %s has waited %s for %s's word to start: `mesh-delivery.show` for the "+
"delivery that landed as %s says why; `plans go %s --why …` starts it by hand", w.repository,
short(w.commit), ago(in), w.awaits, short(w.commit), w.id),
"deliveries that landed as %s says why; `plans go %s --why …` starts it by hand", w.repository,
short(w.commit), ago(in), w.awaits, mergesWords(w), w.id),
Said: fmt.Sprintf("waiting since %s for %s", w.since.UTC().Format(time.RFC3339), w.awaits),
Headline: deliveryName(w.modules, w.repository) + " waiting to start",
Explanation: walkWaitingWords(w, in, severity),
@@ -384,6 +400,19 @@ func watchWaits(f *signalFacts) []conditions.Observation {
return out
}
// mergesWords is the merges a waiting walk answers (novox/hq ADR 0276): every delivery it carries, not one
// commit that a supersession may already have ended (issue 362). Its own commit for a walk naming none.
func mergesWords(w waitFacts) string {
if len(w.merges) == 0 {
return short(w.commit)
}
var out []string
for _, m := range w.merges {
out = append(out, m.Repository+"@"+short(m.Commit))
}
return readableList(out)
}
func watchLoop(f *signalFacts) []conditions.Observation {
if f.loop.pending == 0 {
return nil
+22
View File
@@ -142,6 +142,28 @@ var suppressions = map[string]suppression{
since: f.now.Add(-31 * time.Minute)}}
},
},
// A batch still assembling past its maximum with no walk open (novox/hq ADR 0276): a minute's grace.
"S18": {
inside: func(f *signalFacts) {
f.batches = []batchFacts{{id: "plan-3", state: inventory.PlanAssembling, grouped: "novox/app@c0ffee11",
atMost: f.now.Add(-59 * time.Second), closed: f.now.Add(-59 * time.Second)}}
},
past: func(f *signalFacts) {
f.batches = []batchFacts{{id: "plan-3", state: inventory.PlanAssembling, grouped: "novox/app@c0ffee11",
atMost: f.now.Add(-61 * time.Second), closed: f.now.Add(-61 * time.Second)}}
},
},
// A batch queued behind an open walk past that walk's bound, named.
"S19": {
inside: func(f *signalFacts) {
f.batches = []batchFacts{{id: "plan-3", state: inventory.PlanQueued, grouped: "novox/app@c0ffee11",
closed: f.now.Add(-59 * time.Minute), behind: "plan-1", walkBound: time.Hour}}
},
past: func(f *signalFacts) {
f.batches = []batchFacts{{id: "plan-3", state: inventory.PlanQueued, grouped: "novox/app@c0ffee11",
closed: f.now.Add(-61 * time.Minute), behind: "plan-1", walkBound: time.Hour}}
},
},
// A send refused because it would replace the bus outside its planned step: said at its first refusal,
// whatever the bound (novox/hq issue 336). Inside: a new bus build waits, and no send was refused for it.
"S17": {
-99
View File
@@ -1,99 +0,0 @@
package main
import (
"reflect"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 254, ADR 0218: a newer plan takes over what the older open plans of its repository
// and branch had not built, and closes them as superseded; another repository's plan, another
// branch's, and a plan made after it are left alone.
func TestANewerPlanSupersedesTheOlderOpenPlansOfItsRepository(t *testing.T) {
at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
sent := at.Add(time.Minute)
plan := func(id, repository, branch string, created time.Time, modules map[string]*inventory.PlanModule) inventory.Plan {
return inventory.Plan{ID: id, Repository: repository, Branch: branch, Commit: id + "-commit",
Created: created, State: inventory.PlanRolling, Modules: modules}
}
older := plan("plan-1", "novox/mesh-catalog", "main", at, map[string]*inventory.PlanModule{
"gitea": {State: "built", SentAt: &sent}, // done with: stays done
"keycloak": {State: "asked"}, // asked, not answered: folded
"plex": {}, // not yet asked: folded
"agent": {State: "built"}, // built, rolls out, not sent: folded
"notes": {State: "built"}, // built, records: nothing to send
})
stuck := plan("plan-0", "Novox/Mesh-Catalog", "", at.Add(-time.Hour), map[string]*inventory.PlanModule{
"runtime": {State: "asked"},
})
other := plan("plan-2", "novox/mesh-controller", "main", at, map[string]*inventory.PlanModule{"mesh-controller": {}})
release := plan("plan-3", "novox/mesh-catalog", "release", at, map[string]*inventory.PlanModule{"lemurs": {}})
later := plan("plan-5", "novox/mesh-catalog", "main", at.Add(2*time.Hour), map[string]*inventory.PlanModule{"later": {}})
done := plan("plan-6", "novox/mesh-catalog", "main", at, map[string]*inventory.PlanModule{"finished": {}})
done.State = inventory.PlanDone
newer := plan("plan-4", "novox/mesh-catalog", "main", at.Add(time.Hour), nil)
newer.Commit = "97b1b2b0c0ffee"
rollsOut := func(m string) bool { return m != "notes" }
folded, closed := supersededBy(newer, []inventory.Plan{stuck, older, other, release, later, done, newer}, rollsOut)
if want := []string{"agent", "keycloak", "plex", "runtime"}; !reflect.DeepEqual(folded, want) {
t.Fatalf("folded %v, wanted %v", folded, want)
}
var ids []string
for _, p := range closed {
ids = append(ids, p.ID)
if p.State != inventory.PlanSuperseded || p.Open() {
t.Errorf("%s was left %s", p.ID, p.State)
}
if !strings.Contains(p.Note, "plan-4") || !strings.Contains(p.Note, "97b1b2b0") {
t.Errorf("%s does not name the plan that superseded it: %q", p.ID, p.Note)
}
}
if want := []string{"plan-0", "plan-1"}; !reflect.DeepEqual(ids, want) {
t.Fatalf("superseded %v, wanted %v — another repository, another branch, a later plan and a "+
"finished one are left alone", ids, want)
}
if other.State != inventory.PlanRolling {
t.Fatal("the plan handed in was changed in place")
}
if line := planLine(closed[1], time.Now()); !strings.Contains(line, "superseded") {
t.Fatalf("a superseded plan reads %q", line)
}
}
// novox/hq issue 254: a person closes a plan that will not move again, by its id.
func TestAPersonClosesAStuckPlan(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
stuck := inventory.Plan{ID: "plan-97b1b2b", Repository: "novox/mesh-catalog", Commit: "97b1b2b",
Created: time.Now().UTC(), State: inventory.PlanRolling, Tier: 1, Tiers: [][]string{{"a"}, {"b"}},
Modules: map[string]*inventory.PlanModule{"a": {State: "built"}, "b": {}}}
if err := open.inventory.SavePlan(ctx, &stuck); err != nil {
t.Fatal(err)
}
if err := plansCommand(ctx, []string{"close", stuck.ID}); err == nil || !strings.Contains(err.Error(), "--why") {
t.Fatalf("a plan was closed by hand without saying why: %v", err)
}
if err := plansCommand(ctx, []string{"close", stuck.ID, "--why", "its report will not come"}); err != nil {
t.Fatal(err)
}
closed, err := open.inventory.PlanByID(ctx, stuck.ID)
if err != nil {
t.Fatal(err)
}
if closed.State != inventory.PlanFailed || !strings.Contains(closed.Note, "closed by hand") ||
!strings.Contains(closed.Note, "its report will not come") {
t.Fatalf("the plan was left %s: %q", closed.State, closed.Note)
}
if err := plansCommand(ctx, []string{"close", stuck.ID, "--why", "again"}); err == nil {
t.Fatal("a plan already closed was closed again")
}
if argv, err := argvFor("plans", map[string]any{"close": stuck.ID, "why": "w"}); err != nil ||
!reflect.DeepEqual(argv, []string{"plans", "close", stuck.ID, "--why", "w"}) {
t.Fatalf("the seat's verb does not close a plan: %v %v", argv, err)
}
}
+25 -136
View File
@@ -298,42 +298,20 @@ func notNow(err error) error {
return err
}
// SourceMoved is the forge announcing a merge: every module recorded as built from that
// repository and branch is marked as moved to the merge commit, and built — bases first, so a
// module that stands on another's artifact is built after it and not against the old one
// (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.
// SourceMoved is the forge announcing a merge. It is put into the open batch (novox/hq ADR 0276), which
// becomes one walk when its merge window closes: hearMerge, and cutBatchesHeld in batches.go.
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()
return f.hearMerge(ctx, m, time.Now().UTC())
}
inv := f.open.inventory
entries, err := inv.Catalogued(ctx)
if err != nil {
return notNow(err)
}
read, err := readForPlanning(ctx, inv)
if err != nil {
return notNow(err)
}
// **A merge older than the newest planned merge of its branch is planned at that merge's commit**
// (novox/hq issue 349): the catch-up acts on a merge the bus did not hand over after the merges that
// came after it, and planned at its own commit it built what a later merge had fixed — the forge's
// security fix, on 2026-10-09 — from the commit before it, dependents and all. The later commit contains
// this merge's change. Planned there, with that merge's time, a module the later merge already looked at
// reads as history and is not built again; the rest are built from the newest commit.
newest, known, err := inv.NewestMergeOf(ctx, m.Owner+"/"+m.Repo, m.Base)
if err != nil {
return notNow(err)
}
if known && laterOnTheBranch(m, newest) {
fmt.Printf(" %s/%s %.8s was merged before %.8s (%s, %s): what it moved is built from %.8s, which "+
"contains it\n", m.Owner, m.Repo, m.Commit, newest.Commit, newest.ID, newest.State, newest.Commit)
m.Commit, m.MergedAt = newest.Commit, newest.Merged.Format(time.RFC3339Nano)
}
// movesOfMerge is what one merge moves, acted on: every module recorded as built from that repository and
// branch is marked as moved to the merge commit; a module the merge deleted is forgotten or said; a new module
// is asked for and sent nowhere. The names returned are what the walk builds, bases first along the graph
// (novox/hq 04-ISSUES/131). Nothing is built or pushed here: the walk asks its tiers. Called when a batch is
// cut, once per repository, with the latest merge of the repository's branch and every file the batch's
// merges of it changed.
func movesOfMerge(ctx context.Context, inv *inventory.Inventory, m link.SourceMoved, entries []inventory.Entry,
read map[string][]inventory.ReadRepository) ([]string, error) {
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
@@ -341,11 +319,11 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
if already > 0 {
fmt.Printf("%s/%s merged into %s (%.8s); %d module(s) the mesh holds are already built "+
"from it\n", m.Owner, m.Repo, m.Base, m.Commit, already)
return nil
return nil, nil
}
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds reads it\n",
m.Owner, m.Repo, m.Base, m.Commit)
return nil
return nil, nil
}
// Said, never silent (novox/hq 04-ISSUES/215): a module built from this repository that follows
// another branch is not part of this merge, and whoever is waiting for its change should read why.
@@ -365,7 +343,7 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
}
for _, e := range touched {
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
return notNow(err)
return nil, notNow(err)
}
}
// **A new module is built and registered, and sent nowhere** (novox/hq issue 300): the delivery plan
@@ -377,94 +355,12 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
fmt.Printf("%s/%s merged into %s (%.8s); it changed no module the mesh holds, and adds %s: "+
"built, registered when the build lands, and sent nowhere\n", m.Owner, m.Repo, m.Base, m.Commit,
strings.Join(built, ", "))
return nil
return nil, nil
}
fmt.Printf("%s/%s merged into %s (%.8s); it changed nothing any module the mesh holds is "+
"built from\n", m.Owner, m.Repo, m.Base, m.Commit)
return nil
return nil, nil
}
// A merge produces a plan the mesh keeps (novox/hq ADR 0162): what moved and everything that
// depends on it, along the catalogue's one dependency relation, sorted into tiers. The plan is
// written before any build is asked; the first tier is asked; this returns. Outcomes advance it.
edges, err := inv.Dependencies(ctx)
if err != nil {
return notNow(err)
}
var movedNames []string
for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module)
}
// Written and its first tier asked as one act on the plans (novox/hq issue 213): a timer on
// another controller reading it between the two would ask the tier again.
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return notNow(err)
}
defer release()
plan := planOfMerge(m, movedNames, edges)
// **A newer plan supersedes the older open plans of this repository and branch** (novox/hq issue
// 254, ADR 0218): what they had not built is planned here again, and they are closed, so one plan
// works a repository's modules at a time and a stuck one ends at the next merge.
working, err := inv.OpenPlans(ctx)
if err != nil {
return notNow(err)
}
rollsOut := func(module string) bool {
u, err := inv.UpgradeOf(ctx, module)
return err == nil && u.RollOut
}
folded, superseded := supersededBy(plan, working, rollsOut)
if len(folded) > 0 {
held := map[string]bool{}
for _, e := range entries {
held[e.Manifest.Module] = true
}
names := map[string]bool{}
for _, name := range movedNames {
names[name] = true
}
var also []string
for _, name := range folded {
// One the catalogue no longer holds would fail the newer plan's ask; it is not this
// merge's to build.
if held[name] && !names[name] {
names[name] = true
movedNames = append(movedNames, name)
also = append(also, name)
}
}
if len(also) > 0 {
again := planOfMerge(m, movedNames, edges)
again.ID, again.Created = plan.ID, plan.Created
plan = again
fmt.Printf(" %s, left unbuilt by an older plan of %s, are planned here again\n",
strings.Join(also, ", "), plan.Repository)
}
}
// **Whether it waits for its delivery's word** (novox/hq ADR 0239): while the mesh-delivery seat has a
// holder on record, a walk that moves no module on the controller's own path is opened and waits.
plan.Delivery = awaitsFor(entries, movedNames)
if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
strings.Join(plan.Tiers[len(plan.Tiers)-1], ", "))
}
if err := inv.SavePlan(ctx, &plan); err != nil {
return notNow(err)
}
// Closed after the newer plan is kept, never before: a controller replaced between the two leaves
// both open, which the next merge settles, rather than neither.
for _, old := range superseded {
if err := inv.SavePlan(ctx, &old); err != nil {
return notNow(err)
}
fmt.Printf(" %s (%s at %s) is %s\n", old.ID, old.Repository, short(old.Commit), old.Note)
}
var tiers []string
for i, t := range plan.Tiers {
tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", ")))
}
fmt.Printf("%s/%s merged into %s (%.8s); plan %s, %d module(s) in %d tier(s)\n %s\n",
m.Owner, m.Repo, m.Base, m.Commit, plan.ID, len(plan.Modules), len(plan.Tiers), strings.Join(tiers, "\n "))
// Why each is in it (issue 363): the files of its build source the merge changed, or why it is read whole.
for _, e := range moved {
fmt.Printf(" %s: %s\n", e.Manifest.Module, whyMoved(e, read[e.Manifest.Module], m))
@@ -473,21 +369,11 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
fmt.Printf(" %s reads %s/%s through its build context: its own source record is left where it is\n",
e.Manifest.Module, m.Owner, m.Repo)
}
if plan.Waiting() {
plan.Note = waitingNote(plan)
if err := inv.SavePlan(ctx, &plan); err != nil {
return notNow(err)
}
fmt.Printf(" %s waits for %s's word before its first tier is asked\n", plan.ID, plan.Delivery.Awaits)
return nil
var names []string
for _, e := range moved {
names = append(names, e.Manifest.Module)
}
if err := askTier(ctx, inv, &plan); err != nil {
return notNow(err)
}
if err := inv.SavePlan(ctx, &plan); err != nil {
return notNow(err)
}
return nil
return names, nil
}
// askNewModules asks the build seat for every module a merge adds to the repository — a directory holding a
@@ -893,7 +779,10 @@ func planningView(read map[string][]inventory.ReadRepository, plans []inventory.
var commits []string
for _, p := range plans {
if st, in := p.Modules[name]; in && p.Commit != "" && (p.Open() || (st != nil && st.State == "built")) {
commits = append(commits, p.Commit)
// Every commit the walk carries (novox/hq ADR 0276): a batch's walk answers one per repository.
for _, c := range p.Carried() {
commits = append(commits, c.Commit)
}
}
}
var kept []inventory.ReadRepository
+8 -1
View File
@@ -53,6 +53,8 @@ type signalFacts struct {
plansErr error
// waits are the walks waiting for their delivery's word (S16, novox/hq ADR 0239).
waits []waitFacts
// batches are the batches not yet cut (S18, S19, novox/hq ADR 0276).
batches []batchFacts
loop loopFacts
loopErr error
@@ -159,6 +161,8 @@ type waitFacts struct {
since time.Time
// modules are the modules the walk moves, as a person names the delivery.
modules []string
// merges are the merges it answers (novox/hq ADR 0276), as the deliveries to read are found.
merges []inventory.PlanMerge
}
type loopFacts struct {
@@ -326,6 +330,9 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
f.machines, f.machinesErr = w.gatherMachines(ctx, inv, now)
f.bus, f.busErr = gatherBus(ctx, inv)
f.plans, f.waits, f.plansErr = gatherPlans(ctx, inv, now, f.bus.heldByTheBus())
if f.plansErr == nil {
f.batches, f.plansErr = gatherBatches(ctx, inv, now)
}
f.loop, f.loopErr = w.gatherLoop()
f.mergesPassed, f.merges, f.mergesErr = watchedMerges.last()
if f.mergesErr == nil && !f.mergesPassed.IsZero() && now.Sub(f.mergesPassed) > 3*mergeCatchUpEvery {
@@ -480,7 +487,7 @@ func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time, b
// S16's, whatever mesh-delivery says or does not say.
if p.Waiting() {
waits = append(waits, waitFacts{id: p.ID, repository: p.Repository, commit: p.Commit,
awaits: p.Delivery.Awaits, since: p.Created, modules: planModules(p)})
awaits: p.Delivery.Awaits, since: p.Created, modules: planModules(p), merges: p.Delivery.Merges})
continue
}
_, paused := pausedWaiting(p, pause, now)