Merge pull request 'Say a walk that waits too long (S16), and a delivery's own stalls (D14, H2) — hq ADR 0239' (#103) from feat/mesh-delivery-waits-said into main
This commit was merged in pull request #103.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
@@ -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" }
|
||||
} }
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user