Walk phases: a report past the silent bound never moves a walk's end; keep phases outside the hold on the plans; first-node gate in words (review of #206)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.71s)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.71s)
This commit is contained in:
@@ -508,7 +508,11 @@ 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.
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
@@ -1215,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 {
|
||||
@@ -1223,6 +1228,7 @@ func planTicker(ctx context.Context, open *stores) {
|
||||
return
|
||||
case <-tick.C:
|
||||
advancePlans(ctx, open)
|
||||
keepWalkPhases(ctx, open)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -164,7 +164,7 @@ 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 with its readings, and its " +
|
||||
"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)"},
|
||||
|
||||
@@ -243,7 +243,7 @@ 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
|
||||
// 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"`
|
||||
}
|
||||
@@ -315,7 +315,7 @@ type GateReading struct {
|
||||
Said string `json:"said,omitempty"`
|
||||
}
|
||||
|
||||
// maxReadings bounds a gate's readings: a judging that never passes reads every few seconds for ten minutes.
|
||||
// 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
|
||||
@@ -324,8 +324,8 @@ 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]
|
||||
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 {
|
||||
|
||||
@@ -29,7 +29,7 @@ const (
|
||||
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
|
||||
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
|
||||
@@ -320,7 +320,7 @@ func (p Plan) points(now time.Time, applied AppliedLookup, w *WalkPhases) []poin
|
||||
}
|
||||
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"})
|
||||
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"})
|
||||
@@ -330,7 +330,7 @@ func (p Plan) points(now time.Time, applied AppliedLookup, w *WalkPhases) []poin
|
||||
for _, s := range lastSends {
|
||||
r, ok := applied(s.node, s.at)
|
||||
switch {
|
||||
case !ok && now.Sub(s.at) > ApplySilentAfter:
|
||||
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)
|
||||
@@ -429,18 +429,19 @@ func words(d time.Duration) string {
|
||||
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.
|
||||
// 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
|
||||
order by (detail = $4) desc, started limit 1`,
|
||||
DurationApply, node, sent.Add(-appliedSlack), OutcomeApplied).Scan(&started, &ms, &outcome)
|
||||
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
|
||||
}
|
||||
|
||||
@@ -273,3 +273,30 @@ func TestAWalksMomentsAreKeptAndItsRestsReportRead(t *testing.T) {
|
||||
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