diff --git a/cmd/mesh-controller/delivery_conditions.go b/cmd/mesh-controller/delivery_conditions.go new file mode 100644 index 0000000..439caf8 --- /dev/null +++ b/cmd/mesh-controller/delivery_conditions.go @@ -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..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 +} diff --git a/cmd/mesh-controller/delivery_conditions_test.go b/cmd/mesh-controller/delivery_conditions_test.go new file mode 100644 index 0000000..27be921 --- /dev/null +++ b/cmd/mesh-controller/delivery_conditions_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index fdb9cd3..6615843 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -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 diff --git a/cmd/mesh-controller/healers.go b/cmd/mesh-controller/healers.go index a7471a7..db9797b 100644 --- a/cmd/mesh-controller/healers.go +++ b/cmd/mesh-controller/healers.go @@ -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 { diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index 834cabf..1ac39ca 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -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..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 diff --git a/cmd/mesh-controller/signals_test.go b/cmd/mesh-controller/signals_test.go index d66ecff..068438d 100644 --- a/cmd/mesh-controller/signals_test.go +++ b/cmd/mesh-controller/signals_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index b305a1c..dc196b3 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -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. diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 5b88417..44f7343 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -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.`, 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") diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 46c75b6..2891cc5 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.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" } } } diff --git a/internal/conditions/condition.go b/internal/conditions/condition.go index ced3126..11b6659 100644 --- a/internal/conditions/condition.go +++ b/internal/conditions/condition.go @@ -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 ( diff --git a/internal/link/queue.go b/internal/link/queue.go index 2322bb1..55190e1 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -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 +}