Say a walk that waits too long, and a delivery's own stalls (hq ADR 0239)
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
mesh/delivery delivered
mesh/delivery-group group feat/mesh-delivery-waits-said delivered: every member is delivered

A walk mesh-delivery never lets go waited for ever with nothing open: S16
says it at 30 minutes, urgent at 4 hours, naming plans go. Phase B: probe
D14 reads the delivery owner's stalled and raises delivery.<id>.stalled,
and H2 takes the table's transition through its close.
This commit is contained in:
jochen
2026-10-07 00:16:32 +02:00
parent 487aa040de
commit 75213b9091
11 changed files with 515 additions and 12 deletions
+182
View File
@@ -0,0 +1,182 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
// A delivery held past its bound, said by the controller and healed by H2 from mesh-delivery's own table
// (novox/hq ADR 0239 decision 9, to-be 47 Phase B). The conditions are the controller's, one owner: the
// self-check reads the delivery owner's `stalled` (probe D14) and raises `delivery.<id>.stalled`; healer
// H2 calls the owner's `close`, which takes only the transition the table names for that state, and only
// the next observation says whether it worked (ADR 0231).
// probeDeliveriesID is the probe that reads the delivery's owner.
const probeDeliveriesID = "D14"
// kindDeliveryStalled is what a delivery held past its bound raises: H2's kind, in the delivery scope.
const kindDeliveryStalled = "stalled"
// deliveryOwnerWithin is how long the owner is given to answer.
var deliveryOwnerWithin = 10 * time.Second
// stalledLine is one delivery past its bound, as mesh-delivery's `stalled` says it.
type stalledLine struct {
ID string `json:"id"`
State string `json:"state"`
For string `json:"for"`
Bound string `json:"bound"`
H2 string `json:"h2"`
Says string `json:"says"`
}
// operatorsOnly is whether the table leaves H2 nothing to do for the line: the state is the operator's.
func (l stalledLine) operatorsOnly() bool { return l.H2 == "" || strings.HasPrefix(l.H2, "none") }
// askDeliveryOwner asks the holder of the mesh-delivery seat one of the verbs the controller is granted;
// a seam a test replaces. The holder's own refusal is an error naming it.
var askDeliveryOwner = func(ctx context.Context, conn *nats.Conn, verb string, args map[string]any) (json.RawMessage, error) {
granted := false
for _, v := range broker.VerbsTheControllerAsksTheDeliveryOwner {
granted = granted || v.Verb == verb
}
if !granted {
return nil, fmt.Errorf("the controller asks %s.%s, which its grant does not name", catalogue.DeliverySeat, verb)
}
if conn == nil {
return nil, errors.New("this controller is not on the bus")
}
answer, err := link.AskMeshSeatTool(ctx, conn, catalogue.DeliverySeat, verb, args, deliveryOwnerWithin)
if err != nil {
return nil, err
}
if answer.Error != "" {
return nil, fmt.Errorf("%s.%s refused: %s", catalogue.DeliverySeat, verb, answer.Error)
}
return unwrapToolResult(answer.Result), nil
}
// unwrapToolResult is a tool's answer whatever the runtime wrapped it in: the JSON itself, or the
// protocol's content list holding it as text.
func unwrapToolResult(raw json.RawMessage) json.RawMessage {
var wrapped struct {
Content []struct {
Text string `json:"text"`
} `json:"content"`
}
if json.Unmarshal(raw, &wrapped) == nil && len(wrapped.Content) > 0 && json.Valid([]byte(wrapped.Content[0].Text)) {
return json.RawMessage(wrapped.Content[0].Text)
}
return raw
}
// deliveriesStalled is what the owner says is held past its bound; nothing when no holder is on record,
// or none answers (D3 says that one).
func deliveriesStalled(ctx context.Context, conn *nats.Conn, held bool) ([]stalledLine, error) {
if !held {
return nil, nil
}
raw, err := askDeliveryOwner(ctx, conn, "stalled", map[string]any{})
if errors.Is(err, link.ErrNothingServes) {
return nil, nil // the holder not running is D3's holder-silent, said once there
}
if err != nil {
return nil, err
}
var lines []stalledLine
if err := json.Unmarshal(raw, &lines); err != nil {
return nil, fmt.Errorf("%s.stalled answered something unreadable: %w", catalogue.DeliverySeat, err)
}
return lines, nil
}
// stalledObservations are the conditions of what is stalled.
func stalledObservations(lines []stalledLine) []conditions.Observation {
out := make([]conditions.Observation, 0, len(lines))
for _, l := range lines {
o := conditions.Observation{Scope: conditions.ScopeDelivery, ID: l.ID, Kind: kindDeliveryStalled,
Severity: conditions.Warning,
Summary: fmt.Sprintf("the delivery %s has been %s for %s, past its bound of %s (%s): healer H2 may %s — "+
"`mesh-delivery.show %s`", l.ID, l.State, l.For, l.Bound, l.Says, l.H2, l.ID),
Said: fmt.Sprintf("%s for %s", l.State, l.For)}
if l.operatorsOnly() {
o.Resolver = conditions.ResolverOperator
}
out = append(out, o)
}
return out
}
// probeDeliveries is D14: every delivery held past its state's bound is said.
func probeDeliveries(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()
}
lines, err := deliveriesStalled(ctx, conn, deliverySeatHeld(entries))
if err != nil {
return nil, err
}
return stalledObservations(lines), nil
}
// --- H2, for a delivery ---------------------------------------------------------------------------------
// connOf is the bus a healer asks over.
func (h *healing) connOf() *nats.Conn {
if h.js == nil {
return nil
}
return h.js.Conn()
}
// appliesToAStalledDelivery is H2's for a delivery: the owner still lists it, with a transition the table
// lets H2 take; never one whose state is the operator's.
func appliesToAStalledDelivery(ctx context.Context, h *healing, c conditions.Condition) (string, bool, string, error) {
lines, err := deliveriesStalled(ctx, h.connOf(), true)
if err != nil {
return "", false, "", err
}
for _, l := range lines {
if l.ID != c.Subject.ID {
continue
}
if l.operatorsOnly() {
return "", false, "the delivery is " + l.State + ", a state the table leaves to the operator", nil
}
return c.Key, true, "", nil
}
return "", false, "the delivery's owner no longer lists it as stalled: its condition clears on the next look", nil
}
// repairDelivery is H2 for a delivery: the owner's `close`, which reads the walk again and takes only the
// transition its table names. A refusal is no repair, said; the next observation says whether it worked.
func repairDelivery(ctx context.Context, h *healing, c conditions.Condition) (string, string, error) {
raw, err := askDeliveryOwner(ctx, h.connOf(), "close", map[string]any{"id": c.Subject.ID, "why": c.Key})
if err != nil {
if errors.Is(err, link.ErrNothingServes) {
return "", "", err
}
return "closed nothing", oneLine(err.Error()), nil
}
var said string
if json.Unmarshal(raw, &said) != nil {
said = string(raw)
}
return "asked " + catalogue.DeliverySeat + " to close " + c.Subject.ID, said, nil
}
@@ -0,0 +1,171 @@
package main
import (
"context"
"encoding/json"
"errors"
"slices"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"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"
"github.com/novox/mesh-controller/internal/testbus"
)
// novox/hq ADR 0239 decision 9, to-be 47 Phase B: a delivery held past its bound is the controller's
// condition, read from mesh-delivery's `stalled`; H2 takes the table's transition through its `close`.
// ownerAnswers replaces the delivery owner with one that answers stalled with these lines and records
// what close was asked.
func ownerAnswers(t *testing.T, lines []stalledLine, closeErr error) *[]string {
t.Helper()
var closed []string
was := askDeliveryOwner
askDeliveryOwner = func(_ context.Context, _ *nats.Conn, verb string, args map[string]any) (json.RawMessage, error) {
switch verb {
case "stalled":
raw, _ := json.Marshal(lines)
return raw, nil
case "close":
closed = append(closed, args["id"].(string))
if closeErr != nil {
return nil, closeErr
}
return json.RawMessage(`"` + args["id"].(string) + `: delivering → delivered, as its walk's record says"`), nil
}
t.Fatalf("asked the owner %s", verb)
return nil, nil
}
t.Cleanup(func() { askDeliveryOwner = was })
return &closed
}
func holdTheDeliverySeat(t *testing.T, open *stores) {
t.Helper()
m := catalogue.Manifest{Module: "mesh-delivery", Version: "1",
Claims: []catalogue.Claim{{Name: catalogue.DeliverySeat, Scope: catalogue.ScopeMesh}}}
if err := open.inventory.RegisterModule(t.Context(), m, inventory.Source{Repository: "novox/mesh-catalog",
Path: "modules/mesh-delivery", BuiltFrom: "c0"}); err != nil {
t.Fatal(err)
}
if _, err := open.inventory.Assign(t.Context(), "anchor", "mesh-delivery"); err != nil {
t.Fatal(err)
}
}
var twoStalled = []stalledLine{
{ID: "novox/app@aaaaaaaaaaaa", State: "delivering", For: "3h0m0s", Bound: "2h0m0s",
H2: "close, by done or superseded or failed, when the walk's record says so", Says: "its walk runs"},
{ID: "novox/lab@bbbbbbbbbbbb", State: "held", For: "49h0m0s", Bound: "24h0m0s",
H2: "none: the state is the operator's", Says: "it waits for the operator"},
}
func TestD14SaysEveryDeliveryHeldPastItsBound(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
ownerAnswers(t, twoStalled, nil)
d := &doctor{open: open}
got, err := probeDeliveries(ctx, d)
if err != nil || len(got) != 0 {
t.Fatalf("with no holder on record D14 said %+v %v", got, err)
}
holdTheDeliverySeat(t, open)
got, err = probeDeliveries(ctx, d)
if err != nil || len(got) != 2 {
t.Fatalf("D14 said %+v %v", got, err)
}
if got[0].Key() != "delivery.novox/app_aaaaaaaaaaaa.stalled" || got[0].Severity != conditions.Warning ||
got[0].Resolver != "" || !strings.Contains(got[0].Summary, "mesh-delivery.show novox/app@aaaaaaaaaaaa") {
t.Fatalf("the first is %+v", got[0])
}
if got[1].Resolver != conditions.ResolverOperator {
t.Fatalf("a held delivery is not the operator's: %+v", got[1])
}
// The holder not running is D3's to say, not D14's.
was := askDeliveryOwner
askDeliveryOwner = func(context.Context, *nats.Conn, string, map[string]any) (json.RawMessage, error) {
return nil, link.ErrNothingServes
}
defer func() { askDeliveryOwner = was }()
if got, err := probeDeliveries(ctx, d); err != nil || len(got) != 0 {
t.Fatalf("an owner that is down was said by D14: %+v %v", got, err)
}
}
// H2 for a delivery: close asked for the one whose state the table gives H2, never for the operator's; a
// refusal is no repair.
func TestH2ClosesAStalledDeliveryThroughItsOwnerOnly(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
h, told, _ := healingOn(t, open)
closed := ownerAnswers(t, twoStalled, nil)
for _, o := range stalledObservations(twoStalled) {
o.Source = probeDeliveriesID
if _, err := conditionsFrom.Observe(ctx, o); err != nil {
t.Fatal(err)
}
}
h.tick(ctx)
if !slices.Equal(*closed, []string{"novox/app@aaaaaaaaaaaa"}) {
t.Fatalf("close was asked for %v", *closed)
}
acts := told.said()
if len(acts) != 1 || acts[0].Healer != "H2" || acts[0].Condition != "delivery.novox/app_aaaaaaaaaaaa.stalled" {
t.Fatalf("said %+v", acts)
}
// A refusal: closed nothing, said so.
c, _, _ := conditionsFrom.Get(ctx, "delivery.novox/app_aaaaaaaaaaaa.stalled")
ownerAnswers(t, twoStalled, errors.New("mesh-delivery.close refused: still delivering"))
act, said, err := repairDelivery(ctx, h, c)
if err != nil || act != "closed nothing" || !strings.Contains(said, "still delivering") {
t.Fatalf("a refused close reads %q %q %v", act, said, err)
}
}
// The controller may ask the delivery's owner what D14 and H2 ask, on its flat subjects, and nothing else
// of it; and over a real bus the answer — wrapped as the runtime wraps it — is read.
func TestTheDeliveryOwnerIsAskedOverTheBus(t *testing.T) {
granted, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindController, PasswordHash: "x"})
if err != nil {
t.Fatal(err)
}
for _, verb := range []string{"stalled", "close"} {
if !slices.Contains(granted.Publish, link.SeatToolSubject(catalogue.DeliverySeat, verb)) {
t.Errorf("the controller may not ask %s.%s", catalogue.DeliverySeat, verb)
}
}
if _, err := askDeliveryOwner(t.Context(), nil, "stop", nil); err == nil || !strings.Contains(err.Error(), "grant") {
t.Fatalf("a verb the grant does not name was asked: %v", err)
}
conn, err := nats.Connect(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
if _, err := deliveriesStalled(t.Context(), conn, true); err != nil {
t.Fatalf("an owner not running is D3's, not an error: %v", err)
}
lines, _ := json.Marshal(twoStalled)
wrapped, _ := json.Marshal(map[string]any{"content": []map[string]any{{"type": "text", "text": string(lines)}}})
sub, err := conn.Subscribe(link.SeatToolSubject(catalogue.DeliverySeat, "stalled"), func(m *nats.Msg) {
body, _ := json.Marshal(link.Answer{Result: wrapped})
_ = m.Respond(body)
})
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
defer cancel()
got, err := deliveriesStalled(ctx, conn, true)
if err != nil || len(got) != 2 || got[0].ID != "novox/app@aaaaaaaaaaaa" {
t.Fatalf("over the bus: %+v %v", got, err)
}
}
+5
View File
@@ -111,6 +111,11 @@ var probeRegistry = []probe{
kindDataUnmeasured, kindDataMissing, kindArrayDegraded, kindProtectionMissing, kindCleanupWaiting},
Phase: 2, run: probeData,
Asks: []broker.SeatVerb{{Seat: "node-backup", Verb: "backed-up"}}},
// The delivery's owner (novox/hq ADR 0239): what it holds past a bound of its own table is said here, by
// the controller, and H2 works it through the owner's `close`.
{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},
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
// The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is
+10 -2
View File
@@ -111,9 +111,11 @@ var healerRegistry = []healerRow{
{ID: "H2", Kinds: []string{"stalled"},
Condition: "a plan stalled on a wait that is superseded — a newer plan of its repository and branch " +
"exists — or already finished — every module of every tier built or failed, and every one that " +
"rolls out sent (S3)",
"rolls out sent (S3); or a delivery held past its state's bound where mesh-delivery's table lets H2 " +
"take a transition (D14, novox/hq ADR 0239)",
Repair: "close the plan with its note, as `plans close` does: superseded, naming the newer plan, or " +
"done; what it asked still builds and registers",
"done; what it asked still builds and registers — or, for a delivery, mesh-delivery's `close`, which " +
"reads its walk again and takes only the transition its table names",
Then: "resolver operator, urgent", From: "a stuck plan closed by hand (214, 254)",
Budget: 1, Window: 24 * time.Hour, Settle: 2 * time.Minute, ActsIn: actsInController, Event: link.KeyHealerActed,
applies: appliesToAStalePlan, repair: repairPlan},
@@ -652,6 +654,9 @@ func planStale(ctx context.Context, inv *inventory.Inventory, p inventory.Plan)
// appliesToAStalePlan is H2's: the plan is open, and its wait is superseded or finished.
func appliesToAStalePlan(ctx context.Context, h *healing, c conditions.Condition) (string, bool, string, error) {
if c.Subject.Scope == conditions.ScopeDelivery {
return appliesToAStalledDelivery(ctx, h, c)
}
p, err := h.open.inventory.PlanByID(ctx, c.Subject.ID)
if err != nil {
return "", false, "", err
@@ -672,6 +677,9 @@ func appliesToAStalePlan(ctx context.Context, h *healing, c conditions.Condition
// repairPlan is H2: the plan closed with its note, under the plans' hold, by compare-and-set.
func repairPlan(ctx context.Context, h *healing, c conditions.Condition) (string, string, error) {
if c.Subject.Scope == conditions.ScopeDelivery {
return repairDelivery(ctx, h, c)
}
inv := h.open.inventory
release, err := inv.HoldPlans(ctx, true)
if err != nil {
+41
View File
@@ -75,6 +75,11 @@ const (
staleRefusalsWithin = 5 * time.Minute
// leaseBound is how long the lease may go unrenewed (S12): the key's age.
leaseBound = broker.LeaseTTL
// waitBound is how long a walk may wait for its delivery's word before it is said, and waitUrgentAfter
// before it is urgent (S16, novox/hq ADR 0239): a group member may wait behind another for a while; four
// hours without a word is the delivery's owner broken, whatever it says of itself.
waitBound = 30 * time.Minute
waitUrgentAfter = 4 * time.Hour
)
// callBounds are the verbs that may run longer than callDefault, and how long (S7).
@@ -195,6 +200,15 @@ var signalsTable = []signalRow{
newest: func(f *signalFacts) time.Time {
return newestOf(f.handActs, func(a link.HandAct) time.Time { return a.At })
}},
{Row: "S16", Signal: "a walk waiting for its delivery's word is let go", Emitter: "mesh-delivery, through deliver",
Trigger: "each merge whose walk waits (novox/hq ADR 0239)",
Bound: "30 min, then urgent after 4 h — raised by the controller whatever mesh-delivery says of itself, " +
"naming `plans go` as the way on",
Kind: kindWalkWaiting, Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.plansErr }, watch: watchWaits,
newest: func(f *signalFacts) time.Time {
return newestOf(f.waits, func(w waitFacts) time.Time { return w.since })
}},
}
// watchFacts is S14: the snapshot a merge check is fed is older than its bound, or none was kept since
@@ -317,6 +331,33 @@ func watchPlans(f *signalFacts) []conditions.Observation {
return out
}
// kindWalkWaiting is S16's kind: plan.<id>.waiting.
const kindWalkWaiting = "waiting"
// watchWaits is S16: a walk that has waited for its delivery's word past its bound, said by the controller
// itself — a delivery's owner that is up and never says go (a bug, a delivery stuck in its own table) would
// otherwise leave a merge waiting for ever with nothing open (novox/hq ADR 0239).
func watchWaits(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, w := range f.waits {
in := f.now.Sub(w.since)
if in <= waitBound {
continue
}
severity := conditions.Warning
if in > waitUrgentAfter {
severity = conditions.Urgent
}
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: w.id, Kind: kindWalkWaiting,
Severity: severity,
Summary: fmt.Sprintf("the walk of %s %s has waited %s for %s's word to start: `mesh-delivery.show` for the "+
"delivery that landed as %s says why; `plans go %s --why …` starts it by hand", w.repository,
short(w.commit), ago(in), w.awaits, short(w.commit), w.id),
Said: fmt.Sprintf("waiting since %s for %s", w.since.UTC().Format(time.RFC3339), w.awaits)})
}
return out
}
func watchLoop(f *signalFacts) []conditions.Observation {
if f.loop.pending == 0 {
return nil
+26
View File
@@ -131,6 +131,17 @@ var suppressions = map[string]suppression{
inside: func(f *signalFacts) { f.facts.taken = f.now.Add(-47 * time.Hour) },
past: func(f *signalFacts) { f.facts.taken = f.now.Add(-49 * time.Hour) },
},
// A walk waiting for its delivery's word for 30 minutes (novox/hq ADR 0239).
"S16": {
inside: func(f *signalFacts) {
f.waits = []waitFacts{{id: "plan-2", repository: "novox/app", commit: "c0ffee11", awaits: "mesh-delivery",
since: f.now.Add(-29 * time.Minute)}}
},
past: func(f *signalFacts) {
f.waits = []waitFacts{{id: "plan-2", repository: "novox/app", commit: "c0ffee11", awaits: "mesh-delivery",
since: f.now.Add(-31 * time.Minute)}}
},
},
// Twice by hand within a fortnight is a healer wanted; once, or the first of two a day too old, is not.
"S15": {
inside: func(f *signalFacts) {
@@ -524,3 +535,18 @@ func TestALeaseLostOrMissingIsSaid(t *testing.T) {
t.Fatalf("serving without the lease was not said: %+v", got)
}
}
// **A walk waiting four hours for its delivery's word is urgent** (S16, novox/hq ADR 0239), and says how on.
func TestAWalkWaitingFourHoursIsUrgentAndNamesTheWayOn(t *testing.T) {
now := time.Date(2026, 10, 6, 12, 0, 0, 0, time.UTC)
f := calm(now)
f.waits = []waitFacts{{id: "plan-2", repository: "novox/app", commit: "c0ffee11", awaits: "mesh-delivery",
since: now.Add(-4*time.Hour - time.Minute)}}
got := watchWaits(f)
if len(got) != 1 || got[0].Severity != conditions.Urgent || got[0].Key() != "plan.plan-2.waiting" {
t.Fatalf("four hours of waiting raised %+v", got)
}
if !strings.Contains(got[0].Summary, "plans go plan-2") || !strings.Contains(got[0].Summary, "mesh-delivery") {
t.Fatalf("it does not say how on: %q", got[0].Summary)
}
}
+19 -8
View File
@@ -51,6 +51,8 @@ type signalFacts struct {
plans []planFacts
plansErr error
// waits are the walks waiting for their delivery's word (S16, novox/hq ADR 0239).
waits []waitFacts
loop loopFacts
loopErr error
@@ -144,6 +146,12 @@ type planFacts struct {
paused bool
}
// waitFacts is one walk waiting for its delivery's word: since its merge opened it.
type waitFacts struct {
id, repository, commit, awaits string
since time.Time
}
type loopFacts struct {
took time.Time
pending uint64
@@ -307,7 +315,7 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
}
}
f.machines, f.machinesErr = w.gatherMachines(ctx, inv, now)
f.plans, f.plansErr = gatherPlans(ctx, inv, now)
f.plans, f.waits, f.plansErr = gatherPlans(ctx, inv, now)
f.loop, f.loopErr = w.gatherLoop()
f.mergesPassed, f.merges, f.mergesErr = watchedMerges.last()
if f.mergesErr == nil && !f.mergesPassed.IsZero() && now.Sub(f.mergesPassed) > 3*mergeCatchUpEvery {
@@ -423,17 +431,17 @@ func (w *watchdogs) gatherMachines(ctx context.Context, inv *inventory.Inventory
}
// gatherPlans is every open plan, its tier's bound from what was measured, and what it waits on.
func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]planFacts, error) {
func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]planFacts, []waitFacts, error) {
plans, err := inv.OpenPlans(ctx)
if err != nil {
return nil, fmt.Errorf("the open plans cannot be read: %w", err)
return nil, nil, fmt.Errorf("the open plans cannot be read: %w", err)
}
if len(plans) == 0 {
return nil, nil
return nil, nil, nil
}
tiers, err := inv.Durations(ctx, inventory.DurationPlanTier, now.Add(-14*24*time.Hour))
if err != nil {
return nil, fmt.Errorf("the measured plan tiers cannot be read: %w", err)
return nil, nil, fmt.Errorf("the measured plan tiers cannot be read: %w", err)
}
measured := map[string][]time.Duration{}
for _, d := range tiers {
@@ -441,10 +449,13 @@ func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) (
}
pause := buildSeatPause(ctx, inv, plans)
var out []planFacts
var waits []waitFacts
for _, p := range plans {
// A walk waiting for its delivery's word is not late (novox/hq ADR 0239): its owner keeps the bound
// of that wait, and a silent owner is its own condition (D3, holder-silent).
// A walk waiting for its delivery's word is no tier late (novox/hq ADR 0239); how long it has waited is
// S16's, whatever mesh-delivery says or does not say.
if p.Waiting() {
waits = append(waits, waitFacts{id: p.ID, repository: p.Repository, commit: p.Commit,
awaits: p.Delivery.Awaits, since: p.Created})
continue
}
_, paused := pausedWaiting(p, pause, now)
@@ -452,7 +463,7 @@ func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) (
tiers: len(p.Tiers), entered: p.TierEntered, bound: max(tierAtLeast, 3*p90(measured[p.Repository])),
waiting: planLineWith(p, now, pause), paused: paused})
}
return out, nil
return out, waits, nil
}
// p90 is the ninetieth percentile of measurements; zero for none.
+10
View File
@@ -158,6 +158,12 @@ var VerbsTheSelfCheckAsks = []SeatVerb{{Seat: "node-intrusion-prevention", Verb:
// the controller's grant that acts, and only through the step a person starts.
var VerbsTheBusStepAsks = []SeatVerb{{Seat: "node-backup", Verb: "now"}}
// VerbsTheControllerAsksTheDeliveryOwner are the mesh-delivery seat's verbs the controller calls (novox/hq
// ADR 0239): its self-check reads `stalled`, and healer H2 takes the one transition the table allows
// through `close`. A mesh seat's verb is flat: no machine in the subject.
var VerbsTheControllerAsksTheDeliveryOwner = []SeatVerb{{Seat: "mesh-delivery", Verb: "stalled"},
{Seat: "mesh-delivery", Verb: "close"}}
// perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the
// holder's machine (novox/hq ADR 0219): `paused.<node>`, the build agent saying whether it takes work.
var perMachineEvents = map[string]bool{"paused.*": true}
@@ -352,6 +358,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
for _, v := range VerbsTheBusStepAsks {
pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb+".*")
}
// And the delivery's owner, a mesh seat, asked on its flat subjects (ADR 0239).
for _, v := range VerbsTheControllerAsksTheDeliveryOwner {
pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb)
}
// And asks who answers (novox/hq to-be 45 §4, D3): the self-check finds every seat's holder by
// the same discovery the console reads. The question only; the answers come to its own inbox.
pub = append(pub, "$SRV.INFO")
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$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.assignment.>", "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.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.*"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$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.assignment.>", "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-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "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.*"] }
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"] }
allow_responses: { max: 1, ttl: "1m" }
} }
+3 -1
View File
@@ -44,11 +44,13 @@ const (
ScopeCore = "core"
ScopeProbe = "probe"
ScopeMesh = "mesh"
// ScopeDelivery is a delivery mesh-delivery owns, by its id (novox/hq ADR 0239).
ScopeDelivery = "delivery"
)
// Scopes is every scope, in the order a person reads them.
var Scopes = []string{ScopeMachine, ScopePlan, ScopeCall, ScopeBuild, ScopeMerge, ScopeProvider,
ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh}
ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh, ScopeDelivery}
// Who resolves a condition.
const (
+47
View File
@@ -486,3 +486,50 @@ func PausedSaid(js *broker.JetStream, seat string, nodes []string) (map[string]H
}
return out, nil
}
// AskMeshSeatTool asks the holder of a mesh-scoped seat one of its verbs, on the seat's flat subject
// (design 33 §4: a mesh seat's verb carries no machine), and reads its answer (novox/hq ADR 0239: the
// delivery's owner's `stalled` and `close`).
func AskMeshSeatTool(ctx context.Context, conn *nats.Conn, seat, verb string, args any, timeout time.Duration) (Answer, error) {
body, err := json.Marshal(args)
if err != nil {
return Answer{}, err
}
asking, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
subject := SeatToolSubject(seat, verb)
refused, stop := refusalsOf(conn, subject)
defer stop()
type replied struct {
msg *nats.Msg
err error
}
done := make(chan replied, 1)
go func() {
msg, err := conn.RequestWithContext(asking, subject, body)
done <- replied{msg, err}
}()
var reply *nats.Msg
select {
case r := <-done:
reply, err = r.msg, r.err
case why := <-refused:
cancel()
return Answer{}, fmt.Errorf("the bus refused the controller asking %s.%s — its grants do not name %s: %v",
seat, verb, subject, why)
}
switch {
case errors.Is(err, nats.ErrNoResponders):
return Answer{}, fmt.Errorf("%w: nothing answers %s.%s — its holder is not running, or is older than the verb",
ErrNothingServes, seat, verb)
case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout):
return Answer{}, fmt.Errorf("%s did not answer %s within %s", seat, verb, timeout)
case err != nil:
return Answer{}, err
}
var answer Answer
if err := json.Unmarshal(reply.Data, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered %s with something unreadable: %w", seat, verb, err)
}
return answer, nil
}