Keep each walk's phases and say a delivery over its budget (hq ADR 0282 slice 1, issue 382)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.62s)

The operator's budget (a leaf module running everywhere within five minutes
of its merge, a core module within ten) cannot be held to without knowing
where a walk's time goes. A walk now keeps when its batch's window closed,
when it was cut and its class (migration 0093), each gate reading, and what
each machine of the rest was sent; 'delivery walks' and plan-moved say its
phases from the merge to every machine of the rest reporting the build
applied, and ended walks keep them as walk-phase durations per class. Probe
D16 reads mesh-delivery's 'times' and raises delivery.<class>.over-budget.
A phase that cannot be measured is said unknown, never zero. Measurement
only: no walk is held, sent or judged differently.
This commit is contained in:
jochen
2026-10-10 20:15:01 +02:00
parent 690b75f659
commit c640341c50
16 changed files with 1224 additions and 18 deletions
+48 -1
View File
@@ -478,7 +478,8 @@ func spelledOf(m inventory.BatchedMerge) string {
// namedOf is one merge as a walk or batch names it: its commit, and its pull request and moves as announced.
func namedOf(m inventory.BatchedMerge, repository string) inventory.PlanMerge {
k := announcedOf(m)
return inventory.PlanMerge{Repository: repository, Commit: m.Commit, Number: k.Number, Title: k.Title, Moves: k.Moves}
return inventory.PlanMerge{Repository: repository, Commit: m.Commit, Number: k.Number, Title: k.Title, Moves: k.Moves,
Merged: m.Merged, Heard: m.Heard}
}
// laterMerge says a was merged after b: by the forge's merge time, then by when each was heard.
@@ -772,6 +773,8 @@ func planBatch(ctx context.Context, open *stores, batch *inventory.Plan, carry [
}
plan.Delivery.Merges = named
plan.Delivery.Alone = batch.Delivery != nil && batch.Delivery.Alone
// **Its moments and its class** (novox/hq ADR 0282 decision 6): measured, never acted on.
plan.Times = walkTimesAtCut(*batch, plan, entries, now)
if len(moved) == 0 {
plan.State = inventory.PlanDone
plan.Tiers = [][]string{}
@@ -820,6 +823,50 @@ func planBatch(ctx context.Context, open *stores, batch *inventory.Plan, carry [
return nil
}
// walkTimesAtCut is a walk's own moments as it is cut (novox/hq ADR 0282 decision 6): when its batch's window
// closed — no merge for the window's length, or its maximum, whichever came first — when it was cut, and its
// class. A merge walked alone had no window.
func walkTimesAtCut(batch, walk inventory.Plan, entries []inventory.Entry, now time.Time) *inventory.PlanTimes {
cut := now
t := &inventory.PlanTimes{Cut: &cut, Class: classOf(walk, entries)}
if batch.Delivery != nil && batch.Delivery.Batch != nil {
w := batch.Delivery.Batch
closed := w.ClosesAt
if !w.AtMost.IsZero() && w.AtMost.Before(closed) {
closed = w.AtMost
}
if !closed.IsZero() {
if closed.After(now) {
closed = now
}
t.WindowClosed = &closed
}
}
return t
}
// resolverSeat is the seat of the mesh's resolver: a module claiming it is a core module (ADR 0282 decision 1).
const resolverSeat = "mesh-dns-resolver"
// classOf is a walk's class (novox/hq ADR 0282 decision 1): core when it walks a module on the controller's own
// path or one holding the mesh's resolver, leaf otherwise.
func classOf(walk inventory.Plan, entries []inventory.Entry) string {
resolvers := map[string]bool{}
for _, e := range entries {
for _, c := range e.Manifest.Claims {
if c.Name == resolverSeat {
resolvers[e.Manifest.Module] = true
}
}
}
for m := range walk.Modules {
if _, own := onTheControllersPath[m]; own || resolvers[m] {
return inventory.ClassCore
}
}
return inventory.ClassLeaf
}
// 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.
+72
View File
@@ -596,12 +596,81 @@ func deliveryCommand(ctx context.Context, args []string) error {
if err != nil {
return err
}
// Each walk's phases, with the rest's reports (novox/hq ADR 0282 decision 6).
now := time.Now()
for i := range walks {
walks[i].Phases = walks[i].WalkPhases(now, appliedFrom(ctx, inv))
}
return answer(map[string]any{"held": deliverySeatHeld(entries), "walks": walks,
"own-path": sortedKeysOf(ownPathWords())})
}
return fmt.Errorf("delivery %s: plan, order, check, go, stop or walks", sub)
}
// appliedFrom answers a machine's first report after a send from the controller's `apply` durations; a lookup
// that fails is a report not read, which leaves the walk's end unknown rather than wrong.
func appliedFrom(ctx context.Context, inv *inventory.Inventory) inventory.AppliedLookup {
return func(node string, sent time.Time) (inventory.AppliedReport, bool) {
r, ok, err := inv.FirstAppliedAfter(ctx, node, sent)
if err != nil {
fmt.Fprintf(os.Stderr, "the report of %s after %s could not be read: %v\n", node, sent.Format(time.RFC3339), err)
return inventory.AppliedReport{}, false
}
return r, ok
}
}
// recordWalkPhases keeps, once, each phase of every walk ended lately whose end is known, as a duration of kind
// walk-phase per class (novox/hq ADR 0282 decision 6): what `durations` summarises. A walk whose end is unknown
// past ApplySilentAfter keeps its measured phases without its total.
func recordWalkPhases(ctx context.Context, inv *inventory.Inventory, now time.Time) error {
recent, err := inv.RecentPlans(ctx, 30)
if err != nil {
return err
}
for _, p := range recent {
if p.State != inventory.PlanDone || p.Release != nil || now.Sub(p.Updated) > 2*time.Hour {
continue
}
anyKept, totalKept, err := inv.WalkPhasesKept(ctx, p.ID)
if err != nil {
return err
}
if totalKept {
continue
}
ph := p.WalkPhases(now, appliedFrom(ctx, inv))
if ph != nil && ph.End == nil && anyKept {
continue // kept without its end; kept again only once its end is known
}
if ph == nil || (ph.End == nil && now.Sub(p.Updated) < inventory.ApplySilentAfter+time.Minute) {
continue
}
class := ph.Class
if class == "" {
class = "unclassed"
}
for _, x := range ph.Phases {
if x.State != inventory.PhaseMeasured || x.Start == nil {
continue
}
if err := inv.RecordDuration(ctx, inventory.Duration{Kind: inventory.DurationWalkPhase,
Subject: class + "/" + x.Name, Ref: fmt.Sprintf("%s/%s/%d", p.ID, x.Name, x.Tier), Started: *x.Start,
Took: time.Duration(x.TookMS) * time.Millisecond, Detail: p.Named()}); err != nil {
return err
}
}
if ph.End != nil && ph.From != nil {
if err := inv.RecordDuration(ctx, inventory.Duration{Kind: inventory.DurationWalkPhase,
Subject: class + "/total", Ref: p.ID + "/total", Started: *ph.From,
Took: time.Duration(ph.TotalMS) * time.Millisecond, Detail: ph.Said}); err != nil {
return err
}
}
}
return nil
}
// ownPathWords is the controller's own path as words, for an answer.
func ownPathWords() map[string]string { return onTheControllersPath }
@@ -728,6 +797,9 @@ func sayPlanMoved(ctx context.Context, bus link.Bus, p inventory.Plan) {
}
func publishPlanMoved(ctx context.Context, bus link.Bus, p inventory.Plan) {
// Its phases so far (novox/hq ADR 0282 decision 6): the rest's reports come after the walk ends, and are
// read by whoever asks for the walk (`delivery walks`).
p.Phases = p.WalkPhases(time.Now(), nil)
body, err := json.Marshal(p)
if err != nil {
return
+5
View File
@@ -116,6 +116,11 @@ var probeRegistry = []probe{
{ID: probeDeliveriesID, Asserts: "no delivery is held past its state's bound unsaid: mesh-delivery's " +
"`stalled`, each with the transition its table lets healer H2 take", From: "ADR 0239",
Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries},
// The delivery budgets (novox/hq ADR 0282 decision 7): the newest delivery of each class within its budget,
// read from mesh-delivery's `times`; a measurement said, never a delivery held.
{ID: probeBudgetsID, Asserts: "the newest delivery of a leaf module ran on every machine within five minutes of " +
"its merge, and of a core module within ten: mesh-delivery's `times`", From: "ADR 0282",
Kind: kindOverBudget, Phase: 3, run: probeBudgets},
// A client of the bus reconnecting in a loop (novox/hq issue 327), from the server's record of closed
// connections, which the bus's own module reads.
{ID: probeReconnectsID, Asserts: "no user of the bus had its connection dropped more than twelve times in the " +
+4
View File
@@ -603,6 +603,7 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
pastBound := now.Sub(*g.Since) > gateBound
switch {
case worst == healthBroken:
g.Read(now, false, g.BrokenWhy)
var judging []string
for _, m := range modules {
if reading[m] != healthBroken && !passedAlone(m) {
@@ -621,8 +622,10 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
// Waiting on a provider that is unhealthy: not a pass, and not a failure at the bound either —
// the provider's own condition says what is wrong (ADR 0240 rule 5).
g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing
g.Read(now, false, why)
case worst == healthNotYet:
g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing
g.Read(now, false, why)
if pastBound {
fail(fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why))
}
@@ -630,6 +633,7 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
// Healthy, or waiting for a person (ADR 0254): a pass, the wait carried along in the verdict.
g.Passes++
g.LastPass, g.Last, g.Failing = &now, "", nil
g.Read(now, true, "")
if g.Passes >= gatePasses && settled {
decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes,
now.Sub(*g.Since).Round(time.Second))+waitsSaid(g.Waits), now)
+1 -1
View File
@@ -311,7 +311,7 @@ func usage() {
the self-check: the last verdict, a run now, the probes, the signals' ages
healers [--days N] [--json] the healers, what they did lately, and their brake (to-be 45 §7)
durations [--kind K] [--days N] [--json]
apply, heartbeat, plan-tier and build durations, per machine or module
apply, heartbeat, plan-tier, build and walk-phase durations
collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go
builder issue <name> a broker account for a build machine, scoped to build work,
delivered as the builder module's broker secret (module add it first)
+124
View File
@@ -0,0 +1,124 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A delivery over its budget is loud (novox/hq ADR 0282 decision 7, issue 382): the self-check reads
// mesh-delivery's `times` and raises the warning `delivery.<class>.over-budget` when the newest delivery of a
// class whose delivery time is known took longer than its class's budget — five minutes for a leaf module, ten
// for a core module — naming the delivery and its longest phase. It clears when the next delivery of that class
// lands within its budget. Measurement only: the condition says, it never holds or acts on a delivery.
// probeBudgetsID is the probe that reads the delivery times.
const probeBudgetsID = "D16"
// kindOverBudget is what a delivery over its class's budget raises.
const kindOverBudget = "over-budget"
// deliveryBudgets are the budgets by class (ADR 0282 decision 1); the probe raises nothing for another class.
var deliveryBudgets = map[string]time.Duration{inventory.ClassLeaf: 5 * time.Minute, inventory.ClassCore: 10 * time.Minute}
// timesAnswer is what mesh-delivery's `times` answers, as far as the probe reads it.
type timesAnswer struct {
Classes []timesClass `json:"classes"`
}
// timesClass is one class's line of `times`.
type timesClass struct {
Class string `json:"class"`
Latest *timesLatest `json:"latest,omitempty"`
}
// timesLatest is the newest delivery of a class whose delivery time is known.
type timesLatest struct {
ID string `json:"id"`
TookMS int64 `json:"took_ms"`
Longest string `json:"longest,omitempty"`
Landed string `json:"landed,omitempty"`
}
// deliveryTimes is what the delivery's owner says of its delivery times; nothing when no holder is on record or
// none answers (D3 says that one).
func deliveryTimes(ctx context.Context, conn *nats.Conn, held bool) (*timesAnswer, error) {
if !held {
return nil, nil
}
raw, err := askDeliveryOwner(ctx, conn, "times", map[string]any{})
if errors.Is(err, link.ErrNothingServes) {
return nil, nil
}
if err != nil {
return nil, err
}
var a timesAnswer
if err := json.Unmarshal(raw, &a); err != nil {
return nil, fmt.Errorf("%s.times answered something unreadable: %w", catalogue.DeliverySeat, err)
}
return &a, nil
}
// overBudgetObservations are the conditions of the classes whose newest delivery took longer than its budget:
// strictly longer, so a delivery of exactly its budget is within it.
func overBudgetObservations(a *timesAnswer) []conditions.Observation {
if a == nil {
return nil
}
var out []conditions.Observation
for _, c := range a.Classes {
budget, ok := deliveryBudgets[c.Class]
if !ok || c.Latest == nil {
continue
}
took := time.Duration(c.Latest.TookMS) * time.Millisecond
if took <= budget {
continue
}
longest := c.Latest.Longest
if longest == "" {
longest = "not known"
}
out = append(out, conditions.Observation{Scope: conditions.ScopeDelivery, ID: c.Class, Kind: kindOverBudget,
Severity: conditions.Warning,
Summary: fmt.Sprintf("the %s delivery %s took %s from its merge to running everywhere, over its budget of %s; "+
"its longest phase: %s — `mesh-delivery.times`", c.Class, c.Latest.ID, humanDuration(took),
humanDuration(budget), longest),
Said: fmt.Sprintf("%s took %s (budget %s), longest phase %s", c.Latest.ID, took.Round(time.Second), budget, longest),
Headline: fmt.Sprintf("A %s delivery took %s, over its %s budget", c.Class, humanDuration(took),
humanDuration(budget)),
Explanation: fmt.Sprintf("The delivery %s took %s from its merge until every machine ran it; a %s module is "+
"held to %s (ADR 0282). Most of the time went to %s. Nothing was held or changed because of this: it is "+
"a measurement.", c.Latest.ID, humanDuration(took), c.Class, humanDuration(budget), longest),
Resolved: fmt.Sprintf("the next %s delivery lands within %s", c.Class, humanDuration(budget)),
})
}
return out
}
// probeBudgets is D16: the newest delivery of each class lands within its class's budget.
func probeBudgets(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
entries, err := d.open.inventory.Catalogued(ctx)
if err != nil {
return nil, err
}
var conn *nats.Conn
if d.js != nil {
conn = d.js.Conn()
}
a, err := deliveryTimes(ctx, conn, deliverySeatHeld(entries))
if err != nil {
return nil, err
}
return overBudgetObservations(a), nil
}
+111
View File
@@ -0,0 +1,111 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// A leaf delivery of 5 minutes 10 seconds raises delivery.leaf.over-budget naming its longest phase; one of
// exactly five minutes, a core one of nine and a class without a budget raise nothing (ADR 0282 decision 7).
func TestADeliveryOverItsBudgetIsLoud(t *testing.T) {
a := &timesAnswer{Classes: []timesClass{
{Class: inventory.ClassLeaf, Latest: &timesLatest{ID: "novox/mesh-catalog@abc", TookMS: (5*time.Minute + 10*time.Second).Milliseconds(),
Longest: "judgement 2m40s"}},
{Class: inventory.ClassCore, Latest: &timesLatest{ID: "novox/mesh-controller@def", TookMS: (9 * time.Minute).Milliseconds(),
Longest: "build 1m"}},
{Class: "unclassed", Latest: &timesLatest{ID: "x@y", TookMS: time.Hour.Milliseconds()}},
}}
obs := overBudgetObservations(a)
if len(obs) != 1 {
t.Fatalf("one class over its budget, got %d: %+v", len(obs), obs)
}
o := obs[0]
if o.Key() != "delivery.leaf.over-budget" || !strings.Contains(o.Summary, "judgement 2m40s") ||
!strings.Contains(o.Summary, "novox/mesh-catalog@abc") {
t.Fatalf("the condition names its delivery and longest phase: %s — %s", o.Key(), o.Summary)
}
// The next leaf delivery within its budget clears it: the probe raises nothing for the class.
a.Classes[0].Latest = &timesLatest{ID: "novox/mesh-catalog@ghi", TookMS: (5 * time.Minute).Milliseconds()}
if obs := overBudgetObservations(a); len(obs) != 0 {
t.Fatalf("a delivery of exactly its budget is within it: %+v", obs)
}
a.Classes[1].Latest.TookMS = (10*time.Minute + time.Second).Milliseconds()
if obs := overBudgetObservations(a); len(obs) != 1 || obs[0].Key() != "delivery.core.over-budget" {
t.Fatalf("a core delivery over ten minutes: %+v", obs)
}
if obs := overBudgetObservations(nil); obs != nil {
t.Fatal("no answer raises nothing")
}
// No delivery of a class with a known time yet: nothing said of it.
if obs := overBudgetObservations(&timesAnswer{Classes: []timesClass{{Class: inventory.ClassLeaf}}}); len(obs) != 0 {
t.Fatalf("a class with no delivery: %+v", obs)
}
}
// The probe's question is one the controller's grant names: a question the bus refuses checks nothing.
func TestTheControllerMayAskForTheDeliveryTimes(t *testing.T) {
found := false
for _, v := range broker.VerbsTheControllerAsksTheDeliveryOwner {
found = found || (v.Seat == catalogue.DeliverySeat && v.Verb == "times")
}
if !found {
t.Fatal("the grant does not name mesh-delivery.times")
}
}
// A walk's class is core when it walks a module of the controller's own path or one holding the mesh's resolver,
// leaf otherwise; its window closed at its batch's window, or its maximum, never after its cut.
func TestAWalksClassAndWindowAreReadAtItsCut(t *testing.T) {
entries := []inventory.Entry{
{Manifest: catalogue.Manifest{Module: "dnsmasq", Claims: []catalogue.Claim{{Name: resolverSeat}}}},
{Manifest: catalogue.Manifest{Module: "gitea"}},
}
walk := func(modules ...string) inventory.Plan {
p := inventory.Plan{Modules: map[string]*inventory.PlanModule{}}
for _, m := range modules {
p.Modules[m] = &inventory.PlanModule{}
}
return p
}
for want, w := range map[string]inventory.Plan{
inventory.ClassLeaf: walk("gitea"), inventory.ClassCore: walk("gitea", "mesh-controller"),
} {
if got := classOf(w, entries); got != want {
t.Errorf("%v: %s, want %s", w.Modules, got, want)
}
}
if got := classOf(walk("dnsmasq"), entries); got != inventory.ClassCore {
t.Errorf("the resolver's holder is core: %s", got)
}
now := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC)
batch := inventory.Plan{Delivery: &inventory.PlanDelivery{Batch: &inventory.PlanBatch{
ClosesAt: now.Add(-20 * time.Second), AtMost: now.Add(5 * time.Minute)}}}
times := walkTimesAtCut(batch, walk("gitea"), entries, now)
if times.Cut == nil || !times.Cut.Equal(now) || times.WindowClosed == nil || !times.WindowClosed.Equal(now.Add(-20*time.Second)) ||
times.Class != inventory.ClassLeaf {
t.Fatalf("times at the cut: %+v", times)
}
batch.Delivery.Batch.AtMost = now.Add(-time.Minute)
if times := walkTimesAtCut(batch, walk("gitea"), entries, now); !times.WindowClosed.Equal(now.Add(-time.Minute)) {
t.Fatalf("a window closed at its maximum: %v", times.WindowClosed)
}
if times := walkTimesAtCut(inventory.Plan{Delivery: &inventory.PlanDelivery{Alone: true}}, walk("gitea"), entries, now); times.WindowClosed != nil {
t.Fatalf("a merge walked alone had no window: %v", times.WindowClosed)
}
}
// What each machine of the rest was sent is kept only for the machines the send reached.
func TestTheRestIsKeptForTheMachinesTheSendReached(t *testing.T) {
got := restOf([]string{"ace", "g14"}, []string{"ace", "shanks"}, map[string]inventory.SentDeclaration{"ace": {Digest: "d1"}})
if len(got) != 1 || got["ace"].Digest != "d1" {
t.Fatalf("rest: %+v", got)
}
if restOf([]string{"g14"}, []string{"ace"}, nil) != nil {
t.Fatal("no machine reached: nothing kept")
}
}
+22
View File
@@ -508,6 +508,10 @@ func advancePlans(ctx context.Context, open *stores) {
fmt.Printf("plans: what waits for a gate could not be looked at: %v\n", err)
}
advanceHeld(ctx, open)
// Where the walks that ended lately spent their time (novox/hq ADR 0282): measured, never acted on.
if err := recordWalkPhases(ctx, open.inventory, time.Now()); err != nil {
fmt.Printf("plans: the phases of the walks ended lately could not be kept: %v\n", err)
}
}
// advanceHeld is advancePlans for a caller already holding the plans.
@@ -864,8 +868,12 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
strings.Join(machines, ", "), p.Tier, err)
}
now := time.Now().UTC()
// What each machine of the rest was sent, kept for when it reports it applied (novox/hq ADR 0282 decision
// 6): measured, never acted on.
sentWhat := sentNow(ctx, open.inventory, sent)
for _, m := range rest {
p.Modules[m].SentAt = &now
p.Modules[m].Rest = restOf(restTo[m], sent, sentWhat)
}
fmt.Printf("%s: tier %d built; sent %s to %s, one send each\n", p.ID, p.Tier, strings.Join(rest, ", "),
strings.Join(sent, ", "))
@@ -946,6 +954,20 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
return true, nil
}
// restOf is, for one module, every machine of its rest that the send reached and what it carried there.
func restOf(to, sent []string, what map[string]inventory.SentDeclaration) map[string]inventory.SentDeclaration {
out := map[string]inventory.SentDeclaration{}
for _, n := range to {
if slices.Contains(sent, n) {
out[n] = what[n]
}
}
if len(out) == 0 {
return nil
}
return out
}
// firstSend sends one machine, in one send, every module of the plan's tier whose first machine it is
// (novox/hq issue 281), and records the send on each: what that machine ran of it before — read once,
// before the send, so a module the same send carries is never read as already moved — the machines it