Files
mesh-controller/cmd/mesh-controller/batches.go
T
jochen 3ced96fc6f
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 delivered
A merge on the controller's own path never shares a batch with one that waits for the delivery's word (hq ADR 0276, decided during the build)
2026-10-10 14:12:18 +02:00

1114 lines
39 KiB
Go

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.
//
// A merge on the controller's own path (a module whose walk waits for nobody's word) never shares a batch with
// one that waits for mesh-delivery's: each kind has a batch of its own, so no catalogue delivery skips its turn
// behind a controller merge (decided during the build, 2026-10-10).
//
// The batch, its merges and their times are in the store (batched_merge, and the batch's own plan record in
// the state `assembling` or `queued`), so a restarted controller resumes the window where it stood.
// 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)
}
merged, err := time.Parse(time.RFC3339Nano, m.MergedAt)
if err != nil {
merged = now
}
// **What it touches, history or not** (review of this change): a merge heard after a later merge of its
// branch was cut reads as history for every module that cut marked seen, and was dropped from the record
// here, never named by the walk that carried it.
raw := m
raw.MergedAt = ""
touches := wouldMove(raw, entries, read)
var later []inventory.BatchedMerge
if len(touches) > 0 {
if later, err = inv.LaterMergesOf(ctx, repository, m.Base, merged.UTC()); err != nil {
return notNow(err)
}
}
owed, err := owedLate(ctx, inv, m, merged, touches)
if err != nil {
return notNow(err)
}
if from, packaging, already := mergeCandidates(m, entries, read); len(from) == 0 && len(packaging) == 0 && len(later) == 0 && !owed {
// "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()
// Again under the hold: a controller handing over may have kept it since.
if _, known, err := inv.MergeOf(ctx, repository, m.Commit); err != nil || known {
return notNow(err)
}
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.
// **A merge on the controller's own path never shares a batch with one that waits for the delivery's
// word** (decided during the build, 2026-10-10): batched together, the catalogue's deliveries would start
// with the controller's and skip their turn. Each kind has a batch of its own.
own := ownPath(touches)
if p, ok, err := answeredByALaterMerge(ctx, inv, heard, touches, own); err != nil {
return notNow(err)
} else if ok {
heard.Plan = p.ID
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, own)
if err != nil {
return notNow(err)
}
heard.Plan = batch.ID
if added, err := inv.AddMerge(ctx, heard); err != nil || !added {
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)
}
// owedLate says a merge that reads as history is still the record's to answer (review of this change): a cut
// marks every module it moved seen at that moment, so a merge of that branch the bus lost and the catch-up
// hands over later — made before the cut, heard after it — reads as history and was dropped. One that touches
// a module the mesh holds, of a branch whose merges the controller keeps, made since it began keeping them,
// is kept; one made before it began is old news, as before.
func owedLate(ctx context.Context, inv *inventory.Inventory, m link.SourceMoved, merged time.Time,
touches []inventory.Entry) (bool, error) {
if len(touches) == 0 {
return false, nil
}
kept, first, err := inv.KeptOnBranch(ctx, m.Owner+"/"+m.Repo, m.Base)
if err != nil || !kept {
return false, err
}
return !merged.Before(first.Add(-mergeGrace)), nil
}
// 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,
moves []inventory.Entry, own bool) (inventory.Plan, bool, error) {
later, err := inv.LaterMergesOf(ctx, heard.Repository, heard.Branch, heard.Merged)
if err != nil || len(later) == 0 {
return inventory.Plan{}, false, err
}
for _, l := range later {
p, err := followTakeOver(ctx, inv, l.Plan)
if err != nil {
return inventory.Plan{}, false, err
}
switch {
case p.Batch():
// A batch of the other kind is not this merge's: it joins the batch of its own kind.
if p.OwnPath() == own {
return p, true, nil
}
case p.Open() || p.State == inventory.PlanDone:
if buildsEvery(p, moves) {
return p, true, nil
}
}
}
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, and that there is one: a merge that moves nothing
// is no walk's to name.
func buildsEvery(p inventory.Plan, modules []inventory.Entry) bool {
if len(modules) == 0 {
return false
}
for _, e := range modules {
if _, in := p.Modules[e.Manifest.Module]; !in {
return false
}
}
return true
}
// ownPath says a merge touches a module on the controller's own path, whose walk waits for nobody's word.
func ownPath(touches []inventory.Entry) bool {
for _, e := range touches {
if _, own := onTheControllersPath[e.Manifest.Module]; own {
return true
}
}
return false
}
// theOpenBatch is the batch a merge heard now joins: the one of its kind not yet cut, or a new one.
func theOpenBatch(ctx context.Context, inv *inventory.Inventory, first inventory.BatchedMerge, now time.Time, own bool) (inventory.Plan, error) {
batches, err := inv.Batches(ctx)
if err != nil {
return inventory.Plan{}, err
}
for _, b := range batches {
if b.OwnPath() == own {
return b, nil
}
}
b := inventory.Plan{ID: fmt.Sprintf("plan-%d", now.UnixNano()), Repository: first.Repository, Branch: first.Branch,
Commit: first.Commit, Merged: first.Merged, Created: now, State: inventory.PlanAssembling,
Tiers: [][]string{}, Modules: map[string]*inventory.PlanModule{},
Delivery: &inventory.PlanDelivery{Batch: &inventory.PlanBatch{Own: own}}}
// Kept before its first merge names it: a merge never names a plan the store does not hold.
if err := inv.SavePlan(ctx, &b); err != nil {
return inventory.Plan{}, err
}
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, Own: b.OwnPath()}
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
}
// spelledOf is a kept merge's repository as the forge spelled it; empty when its announcement does not say.
func spelledOf(m inventory.BatchedMerge) string {
var e link.SourceMoved
if json.Unmarshal(m.Event, &e) == nil && e.Owner != "" {
return e.Owner + "/" + e.Repo
}
return ""
}
// 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
}
// The walk's own commits stand: a merge heard late is carried by the one of its repository's branch.
var named []inventory.PlanMerge
for _, m := range merges {
n := inventory.PlanMerge{Repository: m.Repository, Commit: m.Commit}
if e := spelledOf(m); e != "" {
n.Repository = e
}
if c := p.CommitOn(m.Repository, m.Branch); c != "" && c != m.Commit {
n.Carried = c
}
named = append(named, n)
}
sort.SliceStable(named, func(i, j int) bool {
return strings.ToLower(named[i].Repository) < strings.ToLower(named[j].Repository)
})
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
}
keepAll := func(behind string) error {
for i := range batches {
if err := keepBatch(ctx, inv, &batches[i], now, behind); err != nil {
return err
}
}
return nil
}
if started != nil {
// **One walk at a time** (ADR 0276 decision 4): every batch waits behind it, whatever its window says.
return keepAll(started.ID)
}
if len(alone) > 0 {
if len(waiting) > 0 {
// The waiting walk is the one open; the search walks after it, and the batches wait behind both.
return keepAll(waiting[0].ID)
}
return walkAlone(ctx, open, alone[0], now)
}
if err := keepAll(""); err != nil {
return err
}
// The oldest batch whose window closed is cut; one of each kind at most is open (theOpenBatch).
for i := range batches {
b := &batches[i]
if !windowClosed(*b, now) {
continue
}
if b.OwnPath() {
// A walk on the controller's own path folds no walk waiting for its word: those merges are batched
// again, behind it (decided during the build, 2026-10-10).
if err := rebatchWaiting(ctx, inv, waiting, b.ID, now); err != nil {
return err
}
waiting = nil
}
if err := cutBatch(ctx, open, b, waiting, now); err != nil {
return err
}
// The batches left wait behind the walk just cut when it started; the next look does the same.
if b.Open() && !b.Waiting() {
if batches, err = inv.Batches(ctx); err != nil {
return err
}
return keepAll(b.ID)
}
return nil
}
return nil
}
// rebatchWaiting puts the merges of the walks waiting for their word into the open batch of their kind — a new
// one when none is — queued behind the walk about to be cut, and closes the walks as taken over by that batch:
// the batch keeps its id when it is cut, so the delivery's owner follows them to the walk that answers them.
func rebatchWaiting(ctx context.Context, inv *inventory.Inventory, waiting []inventory.Plan, behind string, now time.Time) error {
for i := range waiting {
w := waiting[i]
merges, err := inv.MergesOf(ctx, w.ID)
if err != nil {
return err
}
if len(merges) == 0 {
continue // nothing of the record's: left waiting
}
batch, err := theOpenBatch(ctx, inv, merges[0], now, false)
if err != nil {
return err
}
for _, m := range merges {
if err := inv.AnswerMerge(ctx, m.Repository, m.Commit, batch.ID, m.Alone); err != nil {
return err
}
}
if err := keepBatch(ctx, inv, &batch, now, behind); err != nil {
return err
}
w.State = inventory.PlanSuperseded
if w.Delivery == nil {
w.Delivery = &inventory.PlanDelivery{}
}
w.Delivery.TakenOverBy = batch.ID
w.Note = fmt.Sprintf("superseded at tier %d by %s before it started: a walk on the controller's own path "+
"(%s) goes first, and that batch answers its merges after it", w.Tier, batch.ID, behind)
if err := inv.SavePlan(ctx, &w); err != nil {
return err
}
fmt.Printf(" %s is %s\n", w.ID, w.Note)
}
return nil
}
// cutBatch makes a closed batch one walk, folding in the walks that wait for their word: it keeps its id, and
// 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{Alone: true}}
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) {
// **News, whatever was seen since** (review of this change): each merge was judged when it was heard, and
// a cut, a fold or a failed walk since marks its modules seen — so a merge walked alone after a failed
// walk, or folded, would read as history and move nothing. When it was made is not asked again here.
m.MergedAt = ""
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
plan.Delivery.Alone = batch.Delivery != nil && batch.Delivery.Alone
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 and said to be a module's — on a linear trunk the latest contains the others. A file is
// removed when the last merge of the batch that changed it removed it. merges are in the order they were made.
func combinedMerges(merges []inventory.BatchedMerge) []link.SourceMoved {
type combined struct {
latest inventory.BatchedMerge
m link.SourceMoved
said bool
removed map[string]bool
paths []string
dirs []string
cut 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: true, removed: map[string]bool{}}
by[k] = c
keys = append(keys, k)
} else if laterMerge(b, c.latest) {
c.latest, c.m = b, e
}
c.said = c.said && e.ModuleDirsSaid
c.cut = c.cut || e.PathsTruncated
c.paths, c.dirs = union(c.paths, e.Paths), union(c.dirs, e.ModuleDirs)
for _, path := range e.Paths {
c.removed[path] = slices.Contains(e.Removed, path)
}
for _, path := range e.Removed {
c.removed[path] = true
}
}
sort.Strings(keys)
var out []link.SourceMoved
for _, k := range keys {
c := by[k]
m := c.m
m.Paths, m.ModuleDirs, m.PathsTruncated, m.ModuleDirsSaid = c.paths, c.dirs, c.cut, c.said
m.Removed = nil
for _, path := range c.paths {
if c.removed[path] {
m.Removed = append(m.Removed, path)
}
}
out = append(out, m)
}
return out
}
// 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 w.Own {
grouped += " (on the controller's own path)"
}
if b.State == inventory.PlanQueued {
return fmt.Sprintf("queued behind %s; grouped: %s; plan not yet calculated", w.Behind, grouped)
}
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
}