diff --git a/cmd/mesh-controller/batches.go b/cmd/mesh-controller/batches.go index 3ba71b21..291a6718 100644 --- a/cmd/mesh-controller/batches.go +++ b/cmd/mesh-controller/batches.go @@ -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. diff --git a/cmd/mesh-controller/delivery.go b/cmd/mesh-controller/delivery.go index 33b5bbfc..0a8a4a88 100644 --- a/cmd/mesh-controller/delivery.go +++ b/cmd/mesh-controller/delivery.go @@ -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 diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index 3f78f57d..4986d514 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -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 " + diff --git a/cmd/mesh-controller/gate.go b/cmd/mesh-controller/gate.go index 4ddb202d..27ccd9db 100644 --- a/cmd/mesh-controller/gate.go +++ b/cmd/mesh-controller/gate.go @@ -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) diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 15674cd0..3e4cebec 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -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 a broker account for a build machine, scoped to build work, delivered as the builder module's broker secret (module add it first) diff --git a/cmd/mesh-controller/over_budget.go b/cmd/mesh-controller/over_budget.go new file mode 100644 index 00000000..8fbc895c --- /dev/null +++ b/cmd/mesh-controller/over_budget.go @@ -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..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 +} diff --git a/cmd/mesh-controller/over_budget_test.go b/cmd/mesh-controller/over_budget_test.go new file mode 100644 index 00000000..483e068a --- /dev/null +++ b/cmd/mesh-controller/over_budget_test.go @@ -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 := ×Answer{Classes: []timesClass{ + {Class: inventory.ClassLeaf, Latest: ×Latest{ID: "novox/mesh-catalog@abc", TookMS: (5*time.Minute + 10*time.Second).Milliseconds(), + Longest: "judgement 2m40s"}}, + {Class: inventory.ClassCore, Latest: ×Latest{ID: "novox/mesh-controller@def", TookMS: (9 * time.Minute).Milliseconds(), + Longest: "build 1m"}}, + {Class: "unclassed", Latest: ×Latest{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 = ×Latest{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(×Answer{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") + } +} diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 337b3947..44dd8231 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -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 diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 06c41c48..09966256 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -232,8 +232,11 @@ func MaySubscribe(perms Permissions, subject string) bool { // // And, since novox/hq ADR 0259, `release` and `stop`: the controller asks the operator for them about a // delivery held past its bound, and calls them on the operator's warrant, with its why. +// +// And, since novox/hq ADR 0282, `times`: the self-check reads the delivery times to say a delivery over its budget. var VerbsTheControllerAsksTheDeliveryOwner = []SeatVerb{{Seat: "mesh-delivery", Verb: "stalled"}, - {Seat: "mesh-delivery", Verb: "close"}, {Seat: "mesh-delivery", Verb: "release"}, {Seat: "mesh-delivery", Verb: "stop"}} + {Seat: "mesh-delivery", Verb: "close"}, {Seat: "mesh-delivery", Verb: "release"}, {Seat: "mesh-delivery", Verb: "stop"}, + {Seat: "mesh-delivery", Verb: "times"}} // VerbsTheControllerActsOnAWarrant are the other seat verbs the controller calls when the operator's warrant // chooses them (novox/hq ADR 0259): a machine's service restarted, and a walk started or stopped through the diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 8f6c6e55..0919b2b8 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_asked.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-controller.tool.plans", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.mesh-delivery.tool.stop", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*", "mesh.seat.node-launcher.tool.secret.*", "mesh.seat.node-service-manager.tool.restart.*"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_asked.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-controller.tool.plans", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.mesh-delivery.tool.stop", "mesh.seat.mesh-delivery.tool.times", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*", "mesh.seat.node-launcher.tool.secret.*", "mesh.seat.node-service-manager.tool.restart.*"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built", "mesh.seat.operator-channel.event.decided.mesh-controller"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 9541db55..ac95d050 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -164,8 +164,9 @@ var ControllerVerbs = []Verb{ Input: schema(map[string]string{"plan": "the walk's id", "why": "why", "by": "who stopped the delivery"}, []string{"plan", "why"})}, {Name: "delivery-walks", Description: "The walks the controller keeps (novox/hq ADR 0239): every open one and the " + - "last ended ones, each whole — its tiers, each module's state, first machines and gate — and whether the " + - "delivery seat has a holder on record. Given a plan, that one.", + "last ended ones, each whole — its tiers, each module's state, first machines and gate with its readings, and its " + + "phases from its merge to every machine running it (novox/hq ADR 0282) — and whether the delivery seat has a " + + "holder on record. Given a plan, that one.", Input: schema(map[string]string{"plan": "one walk's id", "limit": "how many ended walks beside the open ones (default 50)"}, nil)}, {Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " + @@ -375,10 +376,11 @@ var ControllerVerbs = []Verb{ "recorded — who, why and the cause of each, and which causes repeat: each repeat is a healer the mesh lacks.", Input: schema(map[string]string{"days": "how many days back (default 14)"}, nil)}, {Name: "durations", Description: "How long things take, as the controller measured them: a send to its machine's " + - "report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build — per machine, " + - "repository or module, with median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).", + "report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build, and each phase of an " + + "ended walk per class (walk-phase, novox/hq ADR 0282) — per machine, repository, module or class/phase, with " + + "median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).", Input: schema(map[string]string{ - "kind": "one kind: apply, heartbeat-gap, plan-tier or build; every kind when absent", + "kind": "one kind: apply, heartbeat-gap, plan-tier, build or walk-phase; every kind when absent", "days": "how many days back (default 14)", }, nil)}, // What is wrong, and the self-check (novox/hq to-be 45 §2, §4). diff --git a/internal/inventory/durations.go b/internal/inventory/durations.go index 8b2a43c2..4ab25b12 100644 --- a/internal/inventory/durations.go +++ b/internal/inventory/durations.go @@ -17,10 +17,13 @@ const ( DurationHeartbeatGap = "heartbeat-gap" DurationPlanTier = "plan-tier" DurationBuild = "build" + // DurationWalkPhase is one phase of an ended walk, per class (novox/hq ADR 0282 decision 6): subject + // "/", and "/total" for its merge to every machine running it. + DurationWalkPhase = "walk-phase" ) // DurationKinds are every kind, in the order `durations` shows them. -var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild} +var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild, DurationWalkPhase} // DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it // in Phase 1's first live week. @@ -100,6 +103,14 @@ func (i *Inventory) Durations(ctx context.Context, kind string, since time.Time) return out, rows.Err() } +// WalkPhasesKept says whether any phase of a walk is kept as a walk-phase duration, and whether its total is. +func (i *Inventory) WalkPhasesKept(ctx context.Context, plan string) (anyKept, totalKept bool, err error) { + err = i.store.Pool().QueryRow(ctx, + `select count(*) > 0, count(*) filter (where ref = $2) > 0 from duration where kind = $1 and ref like $3`, + DurationWalkPhase, plan+"/total", plan+"/%").Scan(&anyKept, &totalKept) + return anyKept, totalKept, err +} + // ForgetOldDurations removes what is older than DurationsKeptFor, and says how many. func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) { tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`, diff --git a/internal/inventory/migrations/0093-a-walk-keeps-its-phases.sql b/internal/inventory/migrations/0093-a-walk-keeps-its-phases.sql new file mode 100644 index 00000000..3d0779c4 --- /dev/null +++ b/internal/inventory/migrations/0093-a-walk-keeps-its-phases.sql @@ -0,0 +1,10 @@ +-- A walk keeps its phases (novox/hq ADR 0282 decision 6, issue 382). +-- +-- A small fix took 45 minutes to reach a node on 2026-10-10, and where the time went was read back from the +-- walks by hand: the walk kept when its modules were asked, built, sent first, judged and sent to the rest, +-- but not when its batch's window closed or when it was cut, so the window and the wait behind another walk +-- could not be told apart. A walk now keeps those moments and its class (core or leaf, read at its cut). +-- Its readings, each module's send to the rest and the times of the merges it answers are kept in the plan's +-- own records (modules, delivery). Null for a walk kept before this: its phases before the build are said +-- unknown. +alter table release_plan add column times jsonb; diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index 44bfca4e..fdac17a7 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -54,6 +54,25 @@ type Plan struct { // walk can carry several. Repository and Commit above keep one of them, for a reader that knows one. Empty // for a walk kept before it was: Repository and Commit are then the whole of it. Commits []PlanCommit `json:"commits,omitempty"` + // Times are the moments of a walk no other field keeps (novox/hq ADR 0282 decision 6): when its batch's + // window closed, when it was cut, and its class. Nil for a walk kept before they were, whose phases before + // its build are said unknown. + Times *PlanTimes `json:"times,omitempty"` + // Phases is the walk's time from its merge, phase by phase (ADR 0282 decision 6): never kept, worked out + // from the walk's own moments by whoever says the walk (`delivery walks`, `plan-moved`). + Phases *WalkPhases `json:"phases,omitempty"` +} + +// PlanTimes are a walk's own moments beside its modules' (novox/hq ADR 0282 decision 6). +type PlanTimes struct { + // WindowClosed is when its batch's merge window closed: no merge for the window's length, or its maximum. + // Nil for a walk no window assembled (one merge walked alone). + WindowClosed *time.Time `json:"window_closed,omitempty"` + // Cut is when the batch became the walk. + Cut *time.Time `json:"cut,omitempty"` + // Class is the walk's module class, read at its cut: core when it moves a module on the controller's own + // path or the mesh's resolver, leaf otherwise (ADR 0282 decision 1). + Class string `json:"class,omitempty"` } // PlanCommit is one repository's commit a walk carries: the latest merge of its branch in the batch. @@ -76,6 +95,10 @@ type PlanMerge struct { Number int `json:"number,omitempty"` Title string `json:"title,omitempty"` Moves []string `json:"moves,omitempty"` + // Merged is when the forge made the merge and Heard when the controller heard it (novox/hq ADR 0282): where + // its delivery time starts. Zero for a merge named before they were kept. + Merged time.Time `json:"merged,omitzero"` + Heard time.Time `json:"heard,omitzero"` } // PlanBatch is a batch's window while it is one (novox/hq ADR 0276): when it closes unless another merge @@ -220,6 +243,9 @@ type PlanModule struct { // went there in one send (novox/hq issue 281), and one gate judges what one send moved. Empty for // the module the gate is kept on, and for a plan from before tiers were sent whole. GatedBy string `json:"gated_by,omitempty"` + // Rest is, per machine of the rest, the declaration the send after the gate carried there (novox/hq ADR + // 0282 decision 6): the machine's first report of it, applied, is when the build runs there. + Rest map[string]SentDeclaration `json:"rest,omitempty"` } // PlanGate is one module's rollout record at its gate (to-be 45 §8): the component, the first machine, @@ -276,6 +302,35 @@ type PlanGate struct { BrokenWhy string `json:"broken_why,omitempty"` // Returned names the broken modules already put back, at once, while the rest of the send is judged. Returned []string `json:"returned,omitempty"` + // Readings are the judging's readings with their times (novox/hq ADR 0282 decision 6): every reading that + // counted a pass, and the first that did not after one that did. At most maxReadings, the newest kept. + Readings []GateReading `json:"readings,omitempty"` +} + +// GateReading is one reading of a first-node gate. +type GateReading struct { + At time.Time `json:"at"` + Healthy bool `json:"healthy"` + // Said is what a reading that did not pass found wanting. + Said string `json:"said,omitempty"` +} + +// maxReadings bounds a gate's readings: a judging that never passes reads every few seconds for ten minutes. +const maxReadings = 24 + +// Read keeps one reading: a pass always, and a reading that did not pass only when the one before passed or +// there is none, so a judging waiting for a module to start keeps one line of it, not hundreds. +func (g *PlanGate) Read(at time.Time, healthy bool, said string) { + if !healthy && len(g.Readings) > 0 && !g.Readings[len(g.Readings)-1].Healthy { + return + } + if len(said) > 200 { + said = said[:200] + } + g.Readings = append(g.Readings, GateReading{At: at, Healthy: healthy, Said: said}) + if len(g.Readings) > maxReadings { + g.Readings = g.Readings[len(g.Readings)-maxReadings:] + } } // CarriedMove is one module's build moving on a machine with a gated send. @@ -342,7 +397,12 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error { if err != nil { return err } - var release, delivery, commits []byte + var release, delivery, commits, times []byte + if p.Times != nil { + if times, err = json.Marshal(p.Times); err != nil { + return err + } + } if len(p.Commits) > 0 { if commits, err = json.Marshal(p.Commits); err != nil { return err @@ -373,18 +433,18 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error { var revision int64 err = tx.QueryRow(ctx, `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, - branch, tier_entered, revision, epoch, release, delivery, merged_at, commits) - values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13, $14, $15, $16, $17) + branch, tier_entered, revision, epoch, release, delivery, merged_at, commits, times) + values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13, $14, $15, $16, $17, $18) on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch, tier_entered = excluded.tier_entered, revision = release_plan.revision + 1, epoch = excluded.epoch, release = excluded.release, delivery = excluded.delivery, repository = excluded.repository, commit_hash = excluded.commit_hash, merged_at = excluded.merged_at, commits = excluded.commits, - created = excluded.created + created = excluded.created, times = excluded.times where release_plan.revision = $12 returning revision`, p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered, - p.Revision, epoch, release, delivery, mergedAt(p.Merged), commits).Scan(&revision) + p.Revision, epoch, release, delivery, mergedAt(p.Merged), commits, times).Scan(&revision) if errors.Is(err, pgx.ErrNoRows) { // The row is there and at another revision — moved since this was read, or there already // when this one is new: either way not this writer's to overwrite. (A plan saved before plans @@ -467,7 +527,7 @@ func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) { func (i *Inventory) plans(ctx context.Context, tail string, args ...any) ([]Plan, error) { rows, err := i.store.Pool().Query(ctx, `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch, - coalesce(tier_entered, created), revision, coalesce(epoch, 0), release, delivery, merged_at, commits + coalesce(tier_entered, created), revision, coalesce(epoch, 0), release, delivery, merged_at, commits, times from release_plan `+tail, args...) if err != nil { return nil, err @@ -476,14 +536,19 @@ func (i *Inventory) plans(ctx context.Context, tail string, args ...any) ([]Plan var out []Plan for rows.Next() { var p Plan - var tiers, modules, release, delivery, commits []byte + var tiers, modules, release, delivery, commits, times []byte var epoch int64 var merged *time.Time if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State, &p.Tier, &tiers, &modules, &p.Note, &p.Branch, &p.TierEntered, &p.Revision, &epoch, &release, - &delivery, &merged, &commits); err != nil { + &delivery, &merged, &commits, ×); err != nil { return nil, err } + if len(times) > 0 { + if err := json.Unmarshal(times, &p.Times); err != nil { + return nil, err + } + } if len(commits) > 0 { if err := json.Unmarshal(commits, &p.Commits); err != nil { return nil, err diff --git a/internal/inventory/walkphases.go b/internal/inventory/walkphases.go new file mode 100644 index 00000000..0378d9a2 --- /dev/null +++ b/internal/inventory/walkphases.go @@ -0,0 +1,455 @@ +package inventory + +import ( + "context" + "errors" + "fmt" + "sort" + "strings" + "time" + + "github.com/jackc/pgx/v5" +) + +// Where a walk's time went (novox/hq ADR 0282 decision 6, issue 382): from its merge to every machine of its +// rest running its builds, phase by phase, worked out from the walk's own moments. Measurement only: nothing +// here decides anything about the walk. +// +// **A phase that cannot be measured is said unknown, never zero.** Where a moment is missing (a walk kept +// before it was recorded, a machine not yet reported), the phases around it are unknown, and the span between +// the moments on either side is counted as unknown time: the measured phases and the unknown time always add +// up to the total. A phase that did not happen (no window for a merge walked alone, no first machine for a +// module nobody runs) is said none, with no time. + +// The phases, in the order a walk passes them (ADR 0282's table). +const ( + PhaseWindow = "window" // the merge to its batch's window closing + PhaseQueued = "queued" // the window closed to the walk being cut: waiting behind another walk + PhaseWord = "word" // the cut to the delivery's word, for a walk that waits for one + PhaseBetween = "between-tiers" // one tier's end to the next tier's ask + PhaseBuild = "build" // a tier asked to its last module built + PhaseSend = "send-first" // built to sent to the first machines + PhaseJudge = "judgement" // sent first to the gate's last verdict: its readings + PhaseRest = "send-rest" // judged to sent to the rest + PhaseApply = "apply" // the last send to every machine of the rest reporting the build applied + PhaseOpen = "open" // a walk still running: since its last moment + PhaseMeasured = "measured" + PhaseUnknown = "unknown" + PhaseNone = "none" +) + +// The classes of a walk (ADR 0282 decision 1) and their delivery budgets. +const ( + ClassCore = "core" + ClassLeaf = "leaf" +) + +// ApplySilentAfter is how long after a send to the rest a machine that has said nothing is left out of the +// walk's end, as a machine not heard from (ADR 0282 decision 1: a sleeping laptop is listed, not waited for). +const ApplySilentAfter = 15 * time.Minute + +// WalkPhases is a walk's time, phase by phase. +type WalkPhases struct { + Class string `json:"class,omitempty"` + // From is the walk's earliest merge; End when its last phase ended: every machine of its rest reported the + // build applied. Nil End while that is not known. + From *time.Time `json:"from,omitempty"` + End *time.Time `json:"end,omitempty"` + // TotalMS is End − From; zero while either is unknown, Total its words. + TotalMS int64 `json:"total_ms,omitempty"` + Total string `json:"total,omitempty"` + // UnknownMS is the time within the walk no phase could be measured over. + UnknownMS int64 `json:"unknown_ms,omitempty"` + Phases []WalkPhase `json:"phases"` + Silent []string `json:"silent,omitempty"` + Said string `json:"said"` + Merges []MergeStart `json:"merges,omitempty"` +} + +// MergeStart is one merge the walk answers and when it was made: where that merge's delivery time starts. +type MergeStart struct { + Repository string `json:"repository"` + Commit string `json:"commit"` + Merged time.Time `json:"merged"` +} + +// WalkPhase is one phase of a walk. +type WalkPhase struct { + Name string `json:"name"` + // Tier is the tier a tier's phase belongs to; -1 for a phase of the walk. + Tier int `json:"tier"` + State string `json:"state"` + Start *time.Time `json:"start,omitempty"` + End *time.Time `json:"end,omitempty"` + TookMS int64 `json:"took_ms,omitempty"` + Took string `json:"took,omitempty"` + Said string `json:"said,omitempty"` +} + +// AppliedReport is a machine's first report, after a send, that it applied what it was sent. +type AppliedReport struct { + At time.Time + Outcome string +} + +// AppliedLookup answers a machine's first report after a send to it; false when it has not reported. +type AppliedLookup func(node string, sent time.Time) (AppliedReport, bool) + +// point is one moment of the walk, ending the phase named. +type point struct { + phase string + tier int + at *time.Time + none bool // the phase did not happen + said string // why it is none, or unknown +} + +// Phases is the walk's time phase by phase, at now; applied answers the rest's reports (nil: none read). +func (p Plan) WalkPhases(now time.Time, applied AppliedLookup) *WalkPhases { + if p.Release != nil || p.Batch() { + return nil + } + w := &WalkPhases{} + if p.Times != nil { + w.Class = p.Times.Class + } + from := p.firstMerge(w) + if from != nil { + f := from.Truncate(time.Millisecond) + from = &f + } + if from == nil { + w.Said = "when its merge was made is not kept: its delivery time is unknown" + } else { + w.From = from + } + points := p.points(now, applied, w) + // Walk the moments: each known moment after a known one is a measured phase; a missing moment makes the + // phases up to the next known one unknown, and their span unknown time. + last := from + var pending []int + for _, pt := range points { + if pt.none { + w.Phases = append(w.Phases, WalkPhase{Name: pt.phase, Tier: pt.tier, State: PhaseNone, Said: pt.said}) + continue + } + ph := WalkPhase{Name: pt.phase, Tier: pt.tier, Said: pt.said} + if pt.at == nil { + ph.State = PhaseUnknown + w.Phases = append(w.Phases, ph) + pending = append(pending, len(w.Phases)-1) + continue + } + at := pt.at.Truncate(time.Millisecond) + ph.End = &at + if last != nil && len(pending) == 0 { + start := *last + ph.State, ph.Start = PhaseMeasured, &start + ph.TookMS = at.Sub(start).Milliseconds() + ph.Took = words(at.Sub(start)) + } else { + ph.State = PhaseUnknown + if last != nil { + w.UnknownMS += at.Sub(*last).Milliseconds() + } + } + w.Phases = append(w.Phases, ph) + pending = nil + last = &at + } + if len(pending) > 0 { + // The walk's last moments are unknown: it has no end. + if w.Said == "" { + w.Said = "its end is unknown: " + w.Phases[pending[0]].Name + " " + orNot(w.Phases[pending[0]].Said) + } + return w + } + if last == nil || from == nil { + return w + } + if p.Open() { + since := *last + w.Phases = append(w.Phases, WalkPhase{Name: PhaseOpen, Tier: -1, State: PhaseOpen, Start: &since, + TookMS: now.Sub(since).Milliseconds(), Took: words(now.Sub(since)), Said: "the walk is " + p.State}) + w.Said = "open: " + p.State + ", " + words(now.Sub(*from)) + " since its merge" + return w + } + if p.State != PlanDone { + w.Said = "ended " + p.State + ": no delivery time" + return w + } + end := *last + w.End = &end + w.TotalMS = end.Sub(*from).Milliseconds() + w.Total = words(end.Sub(*from)) + if w.Said == "" { + if w.UnknownMS > 0 { + w.Said = fmt.Sprintf("%s from its merge to running everywhere, %s of it unknown", w.Total, + words(time.Duration(w.UnknownMS)*time.Millisecond)) + } else { + w.Said = w.Total + " from its merge to running everywhere" + } + } + return w +} + +// firstMerge is when the walk's earliest merge was made, and every merge's start kept on w. +func (p Plan) firstMerge(w *WalkPhases) *time.Time { + var first *time.Time + if p.Delivery != nil { + for _, m := range p.Delivery.Merges { + if m.Merged.IsZero() { + continue + } + w.Merges = append(w.Merges, MergeStart{Repository: m.Repository, Commit: m.Commit, Merged: m.Merged}) + if first == nil || m.Merged.Before(*first) { + at := m.Merged + first = &at + } + } + } + if first == nil { + for _, c := range p.Carried() { + if c.Merged.IsZero() { + continue + } + w.Merges = append(w.Merges, MergeStart{Repository: c.Repository, Commit: c.Commit, Merged: c.Merged}) + if first == nil || c.Merged.Before(*first) { + at := c.Merged + first = &at + } + } + } + return first +} + +// points are the walk's moments in order. +func (p Plan) points(now time.Time, applied AppliedLookup, w *WalkPhases) []point { + var out []point + // The window and the wait behind another walk. + cut := p.Created + if p.Times != nil && p.Times.Cut != nil { + cut = *p.Times.Cut + } + switch { + case p.Times == nil: + out = append(out, point{phase: PhaseWindow, tier: -1, said: "not kept for a walk made before ADR 0282"}, + point{phase: PhaseQueued, tier: -1, at: &cut, said: "the window and the wait behind another walk together"}) + case p.Times.WindowClosed == nil: + out = append(out, point{phase: PhaseWindow, tier: -1, none: true, said: "no window: walked on its own"}, + point{phase: PhaseQueued, tier: -1, at: &cut}) + default: + closed := *p.Times.WindowClosed + out = append(out, point{phase: PhaseWindow, tier: -1, at: &closed}, point{phase: PhaseQueued, tier: -1, at: &cut}) + } + if p.Delivery != nil && p.Delivery.Awaits != "" { + out = append(out, point{phase: PhaseWord, tier: -1, at: p.Delivery.Go, said: "waiting for " + p.Delivery.Awaits + "'s word"}) + } else { + out = append(out, point{phase: PhaseWord, tier: -1, none: true, said: "waits for no word"}) + } + var lastSends []restSend + for t, tier := range p.Tiers { + if t > p.Tier || (t == p.Tier && p.Tier < len(p.Tiers) && !askedAnyOf(p, tier)) { + break + } + var asked, built, first, judged, rest *time.Time + allBuilt, anyFirst, allJudged, allRest := true, false, true, true + for _, m := range tier { + s := p.Modules[m] + if s == nil { + allBuilt, allRest = false, false + continue + } + asked = earliest(asked, s.AskedAt) + if s.BuiltAt == nil { + if s.State != "deleted" { + allBuilt = false + } + } else { + built = latest(built, s.BuiltAt) + } + if s.FirstAt != nil { + anyFirst = true + first = earliest(first, s.FirstAt) + if s.Gate != nil { + if s.Gate.JudgedAt == nil { + allJudged = false + } else { + judged = latest(judged, s.Gate.JudgedAt) + } + } + } + if s.SentAt == nil { + if s.State != "deleted" { + allRest = false + } + } else { + rest = latest(rest, s.SentAt) + for node := range s.Rest { + lastSends = append(lastSends, restSend{module: m, node: node, at: *s.SentAt}) + } + } + } + if t > 0 { + out = append(out, point{phase: PhaseBetween, tier: t, at: asked}) + } + if !allBuilt { + built = nil + } + out = append(out, point{phase: PhaseBuild, tier: t, at: built}) + if anyFirst { + if !allJudged { + judged = nil + } + out = append(out, point{phase: PhaseSend, tier: t, at: first}, point{phase: PhaseJudge, tier: t, at: judged}) + } else { + out = append(out, point{phase: PhaseSend, tier: t, none: true, said: "no first machine: nothing to judge"}, + point{phase: PhaseJudge, tier: t, none: true, said: "no first machine: nothing to judge"}) + } + if !allRest { + rest = nil + } + out = append(out, point{phase: PhaseRest, tier: t, at: rest}) + } + if p.State != PlanDone { + return out + } + // Every machine of the rest running the build: its first report after the send, applied. + if p.Times == nil { + return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's sends are not kept for a walk made before ADR 0282"}) + } + if len(lastSends) == 0 { + return append(out, point{phase: PhaseApply, tier: -1, none: true, + said: "no machine beyond the first: the gate's readings were its run"}) + } + if applied == nil { + return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's reports were not read"}) + } + var end *time.Time + var waiting, failed []string + for _, s := range lastSends { + r, ok := applied(s.node, s.at) + switch { + case !ok && now.Sub(s.at) > ApplySilentAfter: + w.Silent = appendOnce(w.Silent, s.node) + case !ok: + waiting = appendOnce(waiting, s.node) + case r.Outcome != OutcomeApplied: + failed = appendOnce(failed, s.node+" ("+r.Outcome+")") + default: + at := r.At + end = latest(end, &at) + } + } + switch { + case len(failed) > 0: + return append(out, point{phase: PhaseApply, tier: -1, said: "not applied on " + strings.Join(failed, ", ")}) + case len(waiting) > 0: + return append(out, point{phase: PhaseApply, tier: -1, said: "waiting for " + strings.Join(waiting, ", ") + + " to report it applied"}) + case end == nil: + return append(out, point{phase: PhaseApply, tier: -1, said: "no machine of the rest heard from: " + + strings.Join(w.Silent, ", ")}) + } + // A machine may have reported before the last tier ended: the walk ends at whichever is later. + for i := len(out) - 1; i >= 0; i-- { + if out[i].at != nil { + if out[i].at.After(*end) { + e := *out[i].at + end = &e + } + break + } + } + said := "" + if len(w.Silent) > 0 { + sort.Strings(w.Silent) + said = "not waited for, not heard from: " + strings.Join(w.Silent, ", ") + } + return append(out, point{phase: PhaseApply, tier: -1, at: end, said: said}) +} + +type restSend struct { + module, node string + at time.Time +} + +func askedAnyOf(p Plan, tier []string) bool { + for _, m := range tier { + if s := p.Modules[m]; s != nil && s.AskedAt != nil { + return true + } + } + return false +} + +func earliest(a, b *time.Time) *time.Time { + if b == nil { + return a + } + if a == nil || b.Before(*a) { + t := *b + return &t + } + return a +} + +func latest(a, b *time.Time) *time.Time { + if b == nil { + return a + } + if a == nil || b.After(*a) { + t := *b + return &t + } + return a +} + +func appendOnce(to []string, s string) []string { + for _, x := range to { + if x == s { + return to + } + } + return append(to, s) +} + +func orNot(s string) string { + if s == "" { + return "is not known" + } + return "(" + s + ")" +} + +// words is a duration as a person reads it. +func words(d time.Duration) string { + if d < 0 { + return "-" + words(-d) + } + return d.Round(100 * time.Millisecond).String() +} + +// FirstAppliedAfter is a machine's first report, at or after a send to it, that it applied a declaration — the +// controller's `apply` durations, which measure every send to its first report (to-be 45 Phase 0). A declaration +// sent later carries the build too, so its report counts. +func (i *Inventory) FirstAppliedAfter(ctx context.Context, node string, sent time.Time) (AppliedReport, bool, error) { + var started time.Time + var ms int64 + var outcome string + err := i.store.Pool().QueryRow(ctx, + `select started, took_ms, detail from duration + where kind = $1 and subject = $2 and started >= $3 + order by (detail = $4) desc, started limit 1`, + DurationApply, node, sent.Add(-appliedSlack), OutcomeApplied).Scan(&started, &ms, &outcome) + if errors.Is(err, pgx.ErrNoRows) { + return AppliedReport{}, false, nil + } + if err != nil { + return AppliedReport{}, false, err + } + return AppliedReport{At: started.Add(time.Duration(ms) * time.Millisecond).UTC(), Outcome: outcome}, true, nil +} + +// appliedSlack is how much earlier than a walk's record of its send the machine's own record of it may be: +// the send is recorded on the machine first, then on the walk. +const appliedSlack = 2 * time.Second diff --git a/internal/inventory/walkphases_test.go b/internal/inventory/walkphases_test.go new file mode 100644 index 00000000..8850288a --- /dev/null +++ b/internal/inventory/walkphases_test.go @@ -0,0 +1,275 @@ +package inventory + +import ( + "encoding/json" + "strings" + "testing" + "time" +) + +// A walk recorded on the live mesh on 2026-10-10 (mesh-delivery, one tier, one machine), as `delivery walks` +// gave it, with the moments ADR 0282 adds: its window closed ninety seconds after its merge was heard, and it was +// cut eleven seconds later. +const recordedWalk = `{ + "id": "plan-1791654663505629616", "repository": "novox/mesh-catalog", "branch": "main", + "commit": "a14fa306f262d88bdcca83432a70df353d2eceb0", "merged": "2026-10-10T17:50:53Z", + "created": "2026-10-10T19:52:34.037755+02:00", "updated": "2026-10-10T19:55:37.974069+02:00", + "state": "done", "tier": 1, "tiers": [["mesh-delivery"]], + "modules": {"mesh-delivery": {"state": "built", + "asked_at": "2026-10-10T17:52:34.209423145Z", "built_at": "2026-10-10T17:52:50.818011314Z", + "sent_at": "2026-10-10T17:55:35.894985053Z", "first": ["novox"], "first_at": "2026-10-10T17:53:07.287136827Z", + "gate": {"machines": ["novox"], "since": "2026-10-10T17:53:07.287136827Z", "passes": 3, + "verdict": "passed", "judged_at": "2026-10-10T17:55:35.894985053Z"}}}, + "delivery": {"awaits": "", "merges": [{"repository": "novox/mesh-catalog", + "commit": "a14fa306f262d88bdcca83432a70df353d2eceb0", "number": 187, "merged": "2026-10-10T17:50:53Z"}]}, + "times": {"window_closed": "2026-10-10T17:52:23Z", "cut": "2026-10-10T17:52:34.037755Z", "class": "core"} +}` + +func walkOf(t *testing.T, raw string) Plan { + t.Helper() + var p Plan + if err := json.Unmarshal([]byte(raw), &p); err != nil { + t.Fatal(err) + } + return p +} + +// sumsUp checks the phases' measured time and the unknown time add up to the total, to the millisecond. +func sumsUp(t *testing.T, w *WalkPhases) { + t.Helper() + if w.End == nil || w.From == nil { + t.Fatalf("the walk has no end: %+v", w) + } + var sum int64 + for _, ph := range w.Phases { + if ph.State == PhaseMeasured { + if ph.Start == nil || ph.End == nil || ph.End.Sub(*ph.Start).Milliseconds() != ph.TookMS { + t.Errorf("phase %s %d says %dms between %v and %v", ph.Name, ph.Tier, ph.TookMS, ph.Start, ph.End) + } + sum += ph.TookMS + } + if ph.State == PhaseUnknown && ph.TookMS != 0 { + t.Errorf("an unknown phase %s says a time: %dms", ph.Name, ph.TookMS) + } + } + if sum+w.UnknownMS != w.TotalMS || w.End.Sub(*w.From).Milliseconds() != w.TotalMS { + t.Fatalf("the phases add up to %dms and %dms unknown, the total is %dms (%s): %+v", sum, w.UnknownMS, + w.TotalMS, w.End.Sub(*w.From), w.Phases) + } +} + +func phase(w *WalkPhases, name string, tier int) WalkPhase { + for _, ph := range w.Phases { + if ph.Name == name && ph.Tier == tier { + return ph + } + } + return WalkPhase{} +} + +// The phases of a recorded walk add up to its measured total: its merge at 17:50:53 to its verdict on its only +// machine at 17:55:35.894, 4m42.894s. +func TestARecordedWalksPhasesAddUpToItsTotal(t *testing.T) { + w := walkOf(t, recordedWalk).WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil) + sumsUp(t, w) + if w.TotalMS != 282894 || w.Class != ClassCore || w.UnknownMS != 0 { + t.Fatalf("total %dms (want 282894), class %q, unknown %d", w.TotalMS, w.Class, w.UnknownMS) + } + for name, want := range map[string]int64{PhaseWindow: 90000, PhaseQueued: 11037, PhaseBuild: 16781, + PhaseSend: 16469, PhaseJudge: 148607, PhaseRest: 0} { + tier := 0 + if name == PhaseWindow || name == PhaseQueued { + tier = -1 + } + if got := phase(w, name, tier); got.State != PhaseMeasured || got.TookMS != want { + t.Errorf("%s: %s %dms, want measured %dms", name, got.State, got.TookMS, want) + } + } + // The cut to the first ask (171ms) is a time no phase of ADR 0282's table names: it is counted in the build. + if got := phase(w, PhaseWord, -1); got.State != PhaseNone { + t.Errorf("a walk that waits for no word says its word phase %s", got.State) + } + if got := phase(w, PhaseApply, -1); got.State != PhaseNone || !strings.Contains(got.Said, "no machine beyond the first") { + t.Errorf("one machine only: apply %s (%s)", got.State, got.Said) + } +} + +// Two tiers, a machine of the rest each, both reports read: every phase measured, the walk ends at the last +// report applied, and the phases add up. +func TestAWalkOfTwoTiersEndsAtItsRestsLastReport(t *testing.T) { + at := func(s string) *time.Time { + v, err := time.Parse(time.RFC3339Nano, "2026-10-10T18:"+s+"Z") + if err != nil { + t.Fatal(err) + } + return &v + } + p := Plan{ID: "plan-2", State: PlanDone, Tier: 2, Tiers: [][]string{{"a"}, {"b"}}, Created: *at("01:40"), + Delivery: &PlanDelivery{Merges: []PlanMerge{{Repository: "novox/x", Commit: "c1", Merged: *at("00:00")}, + {Repository: "novox/x", Commit: "c0", Merged: *at("00:30")}}}, + Times: &PlanTimes{WindowClosed: at("01:30"), Cut: at("01:40"), Class: ClassLeaf}, + Modules: map[string]*PlanModule{ + "a": {State: "built", AskedAt: at("01:41"), BuiltAt: at("02:05"), FirstAt: at("02:20"), SentAt: at("05:00"), + Gate: &PlanGate{JudgedAt: at("04:50")}, Rest: map[string]SentDeclaration{"ace": {Digest: "d1"}}}, + "b": {State: "built", AskedAt: at("05:10"), BuiltAt: at("05:40"), FirstAt: at("05:50"), SentAt: at("08:00.5"), + Gate: &PlanGate{JudgedAt: at("07:59")}, Rest: map[string]SentDeclaration{"shanks": {Digest: "d2"}}}, + }} + reports := map[string]AppliedReport{"ace": {At: *at("05:12"), Outcome: OutcomeApplied}, + "shanks": {At: *at("08:15.25"), Outcome: OutcomeApplied}} + w := p.WalkPhases(*at("30:00"), func(node string, _ time.Time) (AppliedReport, bool) { + r, ok := reports[node] + return r, ok + }) + sumsUp(t, w) + if w.TotalMS != (8*time.Minute + 15250*time.Millisecond).Milliseconds() { + t.Fatalf("the walk's delivery time runs from its earliest merge to the last report: %s", w.Total) + } + if got := phase(w, PhaseBetween, 1); got.TookMS != 10000 { + t.Errorf("between the tiers: %dms", got.TookMS) + } + if got := phase(w, PhaseApply, -1); got.State != PhaseMeasured || got.TookMS != 14750 { + t.Errorf("apply: %s %dms", got.State, got.TookMS) + } + if len(w.Merges) != 2 { + t.Errorf("each merge's start is said, for its own delivery time: %+v", w.Merges) + } + // A report not yet in: the walk has no end, and its apply is unknown — never zero. + delete(reports, "shanks") + w = p.WalkPhases(*at("10:00"), func(node string, _ time.Time) (AppliedReport, bool) { + r, ok := reports[node] + return r, ok + }) + if w.End != nil || w.TotalMS != 0 { + t.Fatalf("a walk whose rest has not reported has no end: %+v", w) + } + if got := phase(w, PhaseApply, -1); got.State != PhaseUnknown || got.TookMS != 0 || !strings.Contains(got.Said, "shanks") { + t.Errorf("apply while shanks has not reported: %+v", got) + } + // Silent past the bound: listed, not waited for. + w = p.WalkPhases(at("08:00.5").Add(ApplySilentAfter+time.Second), func(node string, _ time.Time) (AppliedReport, bool) { + r, ok := reports[node] + return r, ok + }) + sumsUp(t, w) + if len(w.Silent) != 1 || w.Silent[0] != "shanks" { + t.Errorf("a machine silent past the bound is said, not waited for: %+v", w.Silent) + } + // A failed report: the walk has no end, said. + reports["shanks"] = AppliedReport{At: *at("08:10"), Outcome: OutcomeFailed} + w = p.WalkPhases(*at("30:00"), func(node string, _ time.Time) (AppliedReport, bool) { + r, ok := reports[node] + return r, ok + }) + if w.End != nil || !strings.Contains(phase(w, PhaseApply, -1).Said, "shanks (failed)") { + t.Errorf("a failed apply ends nothing: %+v", phase(w, PhaseApply, -1)) + } +} + +// A moment that is missing makes the phases around it unknown, and their span unknown time: never a zero. +func TestAMissingMomentIsUnknownNeverZero(t *testing.T) { + p := walkOf(t, recordedWalk) + p.Modules["mesh-delivery"].BuiltAt = nil + w := p.WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil) + sumsUp(t, w) + if got := phase(w, PhaseBuild, 0); got.State != PhaseUnknown || got.TookMS != 0 { + t.Errorf("a build without its moment: %+v", got) + } + if got := phase(w, PhaseSend, 0); got.State != PhaseUnknown { + t.Errorf("the send after a build without its moment: %+v", got) + } + if w.UnknownMS != 16781+16469 { + t.Errorf("the unknown time is the span between the moments around it: %dms", w.UnknownMS) + } + // A walk kept before ADR 0282: its window and its wait together, and its apply unknown. + p = walkOf(t, recordedWalk) + p.Times = nil + w = p.WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil) + if got := phase(w, PhaseWindow, -1); got.State != PhaseUnknown { + t.Errorf("a window not kept: %+v", got) + } + if got := phase(w, PhaseApply, -1); got.State != PhaseUnknown || w.End != nil { + t.Errorf("a rest not kept: %+v, end %v", got, w.End) + } +} + +// An open walk is said open since its last moment, with no end. +func TestAnOpenWalkHasNoEnd(t *testing.T) { + p := walkOf(t, recordedWalk) + p.State, p.Tier = PlanRolling, 0 + p.Modules["mesh-delivery"].SentAt = nil + p.Modules["mesh-delivery"].Gate.JudgedAt = nil + w := p.WalkPhases(time.Date(2026, 10, 10, 17, 54, 0, 0, time.UTC), nil) + if w.End != nil || !strings.HasPrefix(w.Said, "its end is unknown") && !strings.HasPrefix(w.Said, "open") { + t.Fatalf("an open walk: %+v", w) + } + if got := phase(w, PhaseBuild, 0); got.State != PhaseMeasured { + t.Errorf("an open walk's phases so far are measured: %+v", got) + } +} + +// A gate keeps every pass and the first reading after one that did not pass, at most maxReadings. +func TestAGateKeepsItsReadings(t *testing.T) { + var g PlanGate + t0 := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC) + g.Read(t0, false, "starting") + g.Read(t0.Add(time.Second), false, "starting") + g.Read(t0.Add(40*time.Second), true, "") + g.Read(t0.Add(80*time.Second), true, "") + g.Read(t0.Add(81*time.Second), false, "gone") + g.Read(t0.Add(82*time.Second), false, "gone") + if len(g.Readings) != 4 || g.Readings[0].Said != "starting" || !g.Readings[2].Healthy || g.Readings[3].Said != "gone" { + t.Fatalf("readings: %+v", g.Readings) + } + for i := 0; i < 50; i++ { + g.Read(t0.Add(time.Duration(100+i)*time.Second), true, "") + } + if len(g.Readings) != maxReadings { + t.Fatalf("readings are bounded: %d", len(g.Readings)) + } +} + +// A walk's moments are kept and read back with it; a machine's first report after a send is read from the +// controller's apply durations. +func TestAWalksMomentsAreKeptAndItsRestsReportRead(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + p := walkOf(t, recordedWalk) + p.Revision, p.Epoch = 0, 0 + if err := inv.SavePlan(ctx, &p); err != nil { + t.Fatal(err) + } + back, err := inv.PlanByID(ctx, p.ID) + if err != nil || back.Times == nil || back.Times.Class != ClassCore || back.Times.WindowClosed == nil || + !back.Times.WindowClosed.Equal(*p.Times.WindowClosed) { + t.Fatalf("the walk's moments were not kept: %v %+v", err, back.Times) + } + if back.Delivery.Merges[0].Merged.IsZero() { + t.Fatal("a merge's time was not kept") + } + sent := time.Now().UTC().Add(-time.Minute).Truncate(time.Millisecond) + if _, ok, err := inv.FirstAppliedAfter(ctx, "ace", sent); err != nil || ok { + t.Fatalf("no report yet: %v %v", ok, err) + } + for _, d := range []Duration{ + {Kind: DurationApply, Subject: "ace", Node: "ace", Ref: "old@1", Started: sent.Add(-time.Hour), Took: time.Second, Detail: OutcomeApplied}, + {Kind: DurationApply, Subject: "ace", Node: "ace", Ref: "d1@2", Started: sent.Add(-time.Second), Took: 12 * time.Second, Detail: OutcomeApplied}, + } { + if err := inv.RecordDuration(ctx, d); err != nil { + t.Fatal(err) + } + } + r, ok, err := inv.FirstAppliedAfter(ctx, "ace", sent) + if err != nil || !ok || r.Outcome != OutcomeApplied || !r.At.Equal(sent.Add(11*time.Second)) { + t.Fatalf("the report after the send: %+v %v %v", r, ok, err) + } + if anyKept, total, err := inv.WalkPhasesKept(ctx, p.ID); err != nil || anyKept || total { + t.Fatalf("nothing kept yet: %v %v %v", anyKept, total, err) + } + if err := inv.RecordDuration(ctx, Duration{Kind: DurationWalkPhase, Subject: "core/total", Ref: p.ID + "/total", + Started: sent, Took: time.Minute}); err != nil { + t.Fatal(err) + } + if anyKept, total, err := inv.WalkPhasesKept(ctx, p.ID); err != nil || !anyKept || !total { + t.Fatalf("the total kept: %v %v %v", anyKept, total, err) + } +}