Merge pull request 'Keep each walk's phases and say a delivery over its budget (hq ADR 0282 slice 1, issue 382)' (#206) from feat/382-walk-phases into main
This commit was merged in pull request #206.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 " +
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -510,6 +510,14 @@ func advancePlans(ctx context.Context, open *stores) {
|
||||
advanceHeld(ctx, open)
|
||||
}
|
||||
|
||||
// keepWalkPhases keeps where the walks that ended lately spent their time (novox/hq ADR 0282), outside the hold
|
||||
// on the plans so it never lengthens it: measured, never acted on, and an error only said.
|
||||
func keepWalkPhases(ctx context.Context, open *stores) {
|
||||
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.
|
||||
func advanceHeld(ctx context.Context, open *stores) {
|
||||
inv := open.inventory
|
||||
@@ -864,8 +872,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 +958,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
|
||||
@@ -1193,6 +1219,7 @@ func sayUnsent(p *inventory.Plan, rollsOut func(string) bool) {
|
||||
// planTicker advances open plans on a timer, for the steps outcomes alone cannot take.
|
||||
func planTicker(ctx context.Context, open *stores) {
|
||||
advancePlans(ctx, open)
|
||||
keepWalkPhases(ctx, open)
|
||||
tick := time.NewTicker(30 * time.Second)
|
||||
defer tick.Stop()
|
||||
for {
|
||||
@@ -1201,6 +1228,7 @@ func planTicker(ctx context.Context, open *stores) {
|
||||
return
|
||||
case <-tick.C:
|
||||
advancePlans(ctx, open)
|
||||
keepWalkPhases(ctx, open)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
+1
-1
@@ -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" }
|
||||
} }
|
||||
|
||||
@@ -18,6 +18,15 @@ func deliveryVerbs() []Verb {
|
||||
Input: schema(map[string]string{"state": "one state, e.g. delivering or held",
|
||||
"repository": "owner/repository", "group": "a group's id (its branch name)",
|
||||
"all": "\"true\": the final ones of the last thirty days too"}, nil, "all")},
|
||||
// Delivery time against the delivery budgets (novox/hq ADR 0282): what probe D16 reads, optional until its
|
||||
// holder serves it.
|
||||
{Name: "times", Description: "Delivery time: from a merge to every machine of its walk running the build. " +
|
||||
"Each class's delivery budget (a leaf module five minutes, a core module ten), the last days' median and " +
|
||||
"worst, how many within and over, the newest delivery of each class, every delivery over its budget with its " +
|
||||
"longest phase, and the delivered ones whose time is not known. Given a delivery's id, its time phase by phase.",
|
||||
Input: schema(map[string]string{"days": "how many days back (default 7)", "id": "one delivery's id: its phases"}, nil),
|
||||
// Optional until mesh-delivery serves it: its holder lives in the catalogue (design 33 §7).
|
||||
Optional: true},
|
||||
{Name: "show", Description: "One delivery or group whole: its delivery plan (what it builds, what each " +
|
||||
"machine receives, what is not an ordinary send), every transition with when and why, the machine " +
|
||||
"steps of its walk, its group and its order.",
|
||||
|
||||
@@ -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 first-node 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).
|
||||
|
||||
@@ -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
|
||||
// "<class>/<phase>", and "<class>/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`,
|
||||
|
||||
@@ -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;
|
||||
@@ -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 first-node 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 first-node 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 r := []rune(said); len(r) > 200 {
|
||||
said = string(r[: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
|
||||
|
||||
@@ -0,0 +1,456 @@
|
||||
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 first-node 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 first-node 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, ok && r.At.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 of a send made at or after a walk's send to it, within
|
||||
// ApplySilentAfter of it — 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; one sent past the bound is a
|
||||
// machine that slept, listed as silent and never moving the walk's end (ADR 0282 decision 1).
|
||||
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 and started < $4
|
||||
order by started limit 1`,
|
||||
DurationApply, node, sent.Add(-appliedSlack), sent.Add(ApplySilentAfter)).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
|
||||
@@ -0,0 +1,302 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine that slept past the bound and reported later is listed silent: its late report never moves the
|
||||
// walk's end; and a report sent past the bound is not read as the walk's.
|
||||
func TestALateReportIsSilentNotTheEnd(t *testing.T) {
|
||||
sent := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC)
|
||||
p := walkOf(t, recordedWalk)
|
||||
p.Modules["mesh-delivery"].SentAt = &sent
|
||||
p.Modules["mesh-delivery"].Rest = map[string]SentDeclaration{"laptop": {Digest: "d"}, "server": {Digest: "e"}}
|
||||
w := p.WalkPhases(sent.Add(3*time.Hour), func(node string, _ time.Time) (AppliedReport, bool) {
|
||||
if node == "server" {
|
||||
return AppliedReport{At: sent.Add(20 * time.Second), Outcome: OutcomeApplied}, true
|
||||
}
|
||||
return AppliedReport{At: sent.Add(2 * time.Hour), Outcome: OutcomeApplied}, true
|
||||
})
|
||||
if len(w.Silent) != 1 || w.Silent[0] != "laptop" || w.End == nil || !w.End.Equal(sent.Add(20*time.Second)) {
|
||||
t.Fatalf("a late report: silent %v, end %v", w.Silent, w.End)
|
||||
}
|
||||
inv := ForTest(t)
|
||||
ctx := t.Context()
|
||||
if err := inv.RecordDuration(ctx, Duration{Kind: DurationApply, Subject: "laptop", Node: "laptop", Ref: "late@1",
|
||||
Started: sent.Add(time.Hour), Took: time.Second, Detail: OutcomeApplied}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, ok, err := inv.FirstAppliedAfter(ctx, "laptop", sent); err != nil || ok {
|
||||
t.Fatalf("a send past the bound read as the walk's: %v %v", ok, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user