Files
mesh-controller/cmd/mesh-controller/batches_watch.go
T
jochen 96fc4209d3
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
Assemble merges in a rolling window and walk each batch once (hq ADR 0276, issue 362)
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.
2026-10-10 13:20:04 +02:00

104 lines
4.2 KiB
Go

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
}