Compare commits

..
Author SHA1 Message Date
mesh-admin dac7e5f7f6 Merge pull request 'Say a part waiting for the operator as needs-operator and pass its gate; add secret families (issue 386, ADR 0283)' (#209) from feat/386-a-wait-for-the-operator into main 2026-10-10 19:49:29 +00:00
mesh-admin c06a2cd972 Merge pull request 'Say a walk's phases out of order unknown, never a negative time (hq issue 382)' (#208) from fix/382-phases-out-of-order into main 2026-10-10 19:36:52 +00:00
jochen d5f6b67b94 Say a part waiting for the operator as needs-operator and pass its gate; add secret families (issue 386, hq ADR 0283)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
A module waiting for the operator's secret failed its first-node gate, held
every later walk and was said as 'nothing for you to do'. The node-engine's
new waiting state is checked against the manifest and the secrets given,
read by the gate as a wait for a person, and raised as needs-operator naming
the act. A secret family gives each part its own one-line secret, which the
mesh never makes, so the desk prompt can take each password.
2026-10-10 21:19:58 +02:00
jochen 12d4025d30 Say a walk's phases out of order unknown, never a negative time (hq issue 382)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
#206's own walk kept its first send 783 ms before its build, and delivery
walks said send-first measured at -783 ms. A moment recorded before the
one before it now makes both phases unknown: the earlier one's time joins
the unknown time, the later one has none, and the walk goes on from the
later moment. Measured phases and unknown time still add up to the total.
2026-10-10 21:15:48 +02:00
mesh-admin 1ca4d8ce60 Merge pull request 'Test that times is optional on the delivery seat (hq issue 382)' (#207) from test/382-times-optional into main 2026-10-10 19:13:56 +00:00
jochen a4f93b52a8 Test that a delivery seat holder without times still holds the seat, and one serving it is not refused (hq issue 382, review of #206)
mesh/merge-gate pass: builds build-agent, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 o…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
2026-10-10 20:55:43 +02:00
mesh-admin a2a671adb7 Merge pull request 'Keep each walk's phases and say a delivery over its budget (hq ADR 0282 slice 1, issue 382)' (#206) from feat/382-walk-phases into main 2026-10-10 18:49:18 +00:00
jochen 412f11b079 Mark times optional on the delivery seat until mesh-delivery serves it (hq ADR 0282, design 33 §7)
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 pass: its merge-check.sh passed
mesh/delivery delivered
2026-10-10 20:35:24 +02:00
jochen 5d4eb974ab Promise times among the mesh-delivery seat's verbs, so its holder may serve it (hq ADR 0282)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.75s)
mesh/delivery superseded: a newer head of the same pull request
2026-10-10 20:22:34 +02:00
jochen 6805e2bcb3 Walk phases: a report past the silent bound never moves a walk's end; keep phases outside the hold on the plans; first-node gate in words (review of #206)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.71s)
2026-10-10 20:20:07 +02:00
jochen c640341c50 Keep each walk's phases and say a delivery over its budget (hq ADR 0282 slice 1, issue 382)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.62s)
The operator's budget (a leaf module running everywhere within five minutes
of its merge, a core module within ten) cannot be held to without knowing
where a walk's time goes. A walk now keeps when its batch's window closed,
when it was cut and its class (migration 0093), each gate reading, and what
each machine of the rest was sent; 'delivery walks' and plan-moved say its
phases from the merge to every machine of the rest reporting the build
applied, and ended walks keep them as walk-phase durations per class. Probe
D16 reads mesh-delivery's 'times' and raises delivery.<class>.over-budget.
A phase that cannot be measured is said unknown, never zero. Measurement
only: no walk is held, sent or judged differently.
2026-10-10 20:15:01 +02:00
mesh-admin 690b75f659 Merge pull request 'Say where a command line's quote is left open and how a quote is written (hq issue 294)' (#205) from fix/294-quote-in-a-quoted-value into main 2026-10-10 17:21:57 +00:00
jochen d52de218c8 Quote none of the line in the unclosed-quote refusal: a refusal is kept on the bus as the call's answer (hq issue 294)
mesh/merge-gate pass: builds mesh-controller → novox; no bus step; every machine composes with the change as it did without (4 of 4 compose)
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
2026-10-10 18:39:20 +02:00
jochen 6188e6b1da Say where a command line's quote is left open and how a quote is written (hq issue 294)
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate pass: builds mesh-controller → novox; no bus step; every machine composes with the change as it did without (4 of 4 compose)
mesh/repo-check pass: its merge-check.sh passed
An apostrophe inside a single-quoted value ends that quote, and the line then failed as a bare
'unclosed quote' with no position and no way out named. The splitter's rules stay the shell's; the
refusal now names the character, shows the line from there and gives the two ways to write a quote.
2026-10-10 18:36:47 +02:00
mesh-admin 798738cc12 Merge pull request 'A shrink is read against the item's own path: a moved directory starts its size history again (hq issue 368)' (#204) from fix/368-data-shrank-path-change into main 2026-10-10 15:56:12 +00:00
jochen babfef4f3a Keep one reading when two paths are measured at one moment, rather than failing the machine's record (hq issue 368 review)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
2026-10-10 17:48:49 +02:00
jochen 96b0966ad8 Read a shrink against the item's own path, so a moved directory starts its size history again (hq issue 368)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
A reading now keeps the path it was measured at, and the peak is read only from
readings at the item's path now. Before, an agent's home that moved to its own
account was compared with the operator's home it left, and data-shrank was
raised for data that was never lost. Existing readings take their item's path
now, so a shrink that is real stays raised.
2026-10-10 17:42:46 +02:00
46 changed files with 2615 additions and 639 deletions
+48 -1
View File
@@ -478,7 +478,8 @@ func spelledOf(m inventory.BatchedMerge) string {
// namedOf is one merge as a walk or batch names it: its commit, and its pull request and moves as announced. // namedOf is one merge as a walk or batch names it: its commit, and its pull request and moves as announced.
func namedOf(m inventory.BatchedMerge, repository string) inventory.PlanMerge { func namedOf(m inventory.BatchedMerge, repository string) inventory.PlanMerge {
k := announcedOf(m) k := announcedOf(m)
return inventory.PlanMerge{Repository: repository, Commit: m.Commit, Number: k.Number, Title: k.Title, Moves: k.Moves} return inventory.PlanMerge{Repository: repository, Commit: m.Commit, Number: k.Number, Title: k.Title, Moves: k.Moves,
Merged: m.Merged, Heard: m.Heard}
} }
// laterMerge says a was merged after b: by the forge's merge time, then by when each was heard. // laterMerge says a was merged after b: by the forge's merge time, then by when each was heard.
@@ -772,6 +773,8 @@ func planBatch(ctx context.Context, open *stores, batch *inventory.Plan, carry [
} }
plan.Delivery.Merges = named plan.Delivery.Merges = named
plan.Delivery.Alone = batch.Delivery != nil && batch.Delivery.Alone plan.Delivery.Alone = batch.Delivery != nil && batch.Delivery.Alone
// **Its moments and its class** (novox/hq ADR 0282 decision 6): measured, never acted on.
plan.Times = walkTimesAtCut(*batch, plan, entries, now)
if len(moved) == 0 { if len(moved) == 0 {
plan.State = inventory.PlanDone plan.State = inventory.PlanDone
plan.Tiers = [][]string{} plan.Tiers = [][]string{}
@@ -820,6 +823,50 @@ func planBatch(ctx context.Context, open *stores, batch *inventory.Plan, carry [
return nil return nil
} }
// walkTimesAtCut is a walk's own moments as it is cut (novox/hq ADR 0282 decision 6): when its batch's window
// closed — no merge for the window's length, or its maximum, whichever came first — when it was cut, and its
// class. A merge walked alone had no window.
func walkTimesAtCut(batch, walk inventory.Plan, entries []inventory.Entry, now time.Time) *inventory.PlanTimes {
cut := now
t := &inventory.PlanTimes{Cut: &cut, Class: classOf(walk, entries)}
if batch.Delivery != nil && batch.Delivery.Batch != nil {
w := batch.Delivery.Batch
closed := w.ClosesAt
if !w.AtMost.IsZero() && w.AtMost.Before(closed) {
closed = w.AtMost
}
if !closed.IsZero() {
if closed.After(now) {
closed = now
}
t.WindowClosed = &closed
}
}
return t
}
// resolverSeat is the seat of the mesh's resolver: a module claiming it is a core module (ADR 0282 decision 1).
const resolverSeat = "mesh-dns-resolver"
// classOf is a walk's class (novox/hq ADR 0282 decision 1): core when it walks a module on the controller's own
// path or one holding the mesh's resolver, leaf otherwise.
func classOf(walk inventory.Plan, entries []inventory.Entry) string {
resolvers := map[string]bool{}
for _, e := range entries {
for _, c := range e.Manifest.Claims {
if c.Name == resolverSeat {
resolvers[e.Manifest.Module] = true
}
}
}
for m := range walk.Modules {
if _, own := onTheControllersPath[m]; own || resolvers[m] {
return inventory.ClassCore
}
}
return inventory.ClassLeaf
}
// combinedMerges is a batch's merges as one merge per repository's branch: the latest, with every file the // combinedMerges is a batch's merges as one merge per repository's branch: the latest, with every file the
// merges of it changed and said to be a module's — on a linear trunk the latest contains the others. A file is // merges of it changed and said to be a module's — on a linear trunk the latest contains the others. A file is
// removed when the last merge of the batch that changed it removed it. merges are in the order they were made. // removed when the last merge of the batch that changed it removed it. merges are in the order they were made.
+2 -9
View File
@@ -151,13 +151,6 @@ func busCommand(ctx context.Context, args []string) error {
if len(args) > 0 && !strings.HasPrefix(args[0], "-") { if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
sub, args = args[0], args[1:] sub, args = args[0], args[1:]
} }
// The view's credential, a terminal line like a person's (bus_view.go).
switch sub {
case "view-credential":
return busViewCredential(ctx, args)
case "view-revoke":
return busViewRevoke(ctx, args)
}
set := flag.NewFlagSet("bus", flag.ContinueOnError) set := flag.NewFlagSet("bus", flag.ContinueOnError)
snapshot := set.String("snapshot-taken", "", "where the streams' snapshot a person took is, while the mesh takes none itself") snapshot := set.String("snapshot-taken", "", "where the streams' snapshot a person took is, while the mesh takes none itself")
reversible := set.Bool("reversible", false, "the new version can be undone by putting the old one back") reversible := set.Bool("reversible", false, "the new version can be undone by putting the old one back")
@@ -167,14 +160,14 @@ func busCommand(ctx context.Context, args []string) error {
if rest, err := parseAround(set, args); err != nil { if rest, err := parseAround(set, args); err != nil {
return err return err
} else if len(rest) > 0 { } else if len(rest) > 0 {
return errors.New(busUsage) return errors.New("bus [upgrade --why … --reversible|--irreversible [--snapshot-taken <where>]]")
} }
switch sub { switch sub {
case "": case "":
return busStatus(ctx) return busStatus(ctx)
case "upgrade": case "upgrade":
default: default:
return fmt.Errorf("bus says what a bus upgrade would do, or `bus upgrade`, `bus view-credential`, `bus view-revoke` — not %q", sub) return fmt.Errorf("bus says what a bus upgrade would do, or `bus upgrade` — not %q", sub)
} }
// Everything refused before anything is done. // Everything refused before anything is done.
if err := why.require("bus upgrade"); err != nil { if err := why.require("bus upgrade"); err != nil {
-117
View File
@@ -1,117 +0,0 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"net"
"strconv"
"strings"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/inventory"
)
// The view's credential: the one read-only user a page in a browser connects to the bus as, over the
// bus module's WebSocket listener (novox/hq research 036, gap G1; broker.KindView).
//
// **A terminal line, like a person's credential** (operator.go): printed once, never stored — the mesh
// keeps a hash — and revoked by forgetting the row, which the next composition of the user list makes
// real. There is one view; issuing it again rotates its password.
const busUsage = "bus [upgrade --why … --reversible|--irreversible [--snapshot-taken <where>] | view-credential | view-revoke]"
// busWebSocketPort is the port the bus module's WebSocket listener is published on, mirrored from the
// nats module's manifest (its `bus-websocket` opening), because the credential names where to connect
// and the controller does not read the module's configuration. Reached across the overlay only: the
// opening is from the mesh, and the mesh's filter admits nothing else.
const busWebSocketPort = 4223
func busViewCredential(ctx context.Context, args []string) error {
if len(args) != 0 {
return errors.New("bus view-credential takes nothing: there is one view, and this prints its credential once")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
// Refused here rather than at the next composition, where it would stop the whole file.
if _, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindView, PasswordHash: "x"}); err != nil {
return err
}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: broker.ViewUser, Kind: inventory.BusView})
if err != nil {
return err
}
where, err := broker.FromEnvironment()
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
return err
}
host := where.Address
if h, _, err := net.SplitHostPort(where.Address); err == nil {
host = h
}
websocket := ""
if host != "" {
websocket = "ws://" + net.JoinHostPort(host, strconv.Itoa(busWebSocketPort))
}
held, err := json.Marshal(struct {
WebSocket string `json:"websocket,omitempty"`
URL string `json:"url,omitempty"`
Fingerprint string `json:"fingerprint,omitempty"`
User string `json:"user"`
Password string `json:"password"`
InboxPrefix string `json:"inbox_prefix"`
Hears []string `json:"hears"`
Reads string `json:"reads"`
HowToRead string `json:"how_to_read"`
}{
WebSocket: websocket, URL: "nats://" + where.Address, Fingerprint: where.Fingerprint,
User: broker.ViewUser, Password: password,
// The client must make its inboxes under the view's own prefix: its subscribe grant is
// `_INBOX.view.>` and no wider (design 25 §4), and a client's default inbox is not under it.
InboxPrefix: "_INBOX." + broker.ViewUser,
Hears: broker.ViewHears, Reads: broker.ViewBucket,
HowToRead: "direct reads only, no watch (a consumer is refused): list with a request to $JS.API.DIRECT.GET.KV_" +
broker.ViewBucket + ` carrying {"multi_last":["$KV.` + broker.ViewBucket + `.>"]}, answered until a 204 status; ` +
"read one key with $JS.API.DIRECT.GET.KV_" + broker.ViewBucket + ".$KV." + broker.ViewBucket + ".<number> " +
"(nats.js: kvm.open(bucket, {allow_direct: true}), never create); re-read the key an event's number names",
})
if err != nil {
return err
}
fmt.Printf("issued the view, which hears %s and reads the bucket %s, and nothing else\n",
strings.Join(broker.ViewHears, ", "), broker.ViewBucket)
fmt.Println(" this is the only time the credential is printed; the mesh keeps a hash")
fmt.Println(" it works once the bus has been told, which is the next push to the machine holding mesh-broker —")
fmt.Println(" and while a new build of the bus module waits for that machine, the next `bus upgrade` a person starts,")
fmt.Println(" which is also what brings the WebSocket listener it connects through")
fmt.Println()
fmt.Println(string(held))
return nil
}
func busViewRevoke(ctx context.Context, args []string) error {
if len(args) != 0 {
return errors.New("bus view-revoke takes nothing: there is one view")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
if err := open.inventory.ForgetBusUser(ctx, broker.ViewUser); err != nil {
return err
}
// **Revoked at the next composition, not now** — as a person is (operator revoke): the bus's users
// are a file, and the credential stops working when the file no longer names it.
fmt.Println("the view is forgotten, and its credential stops working at the next composition — " +
"push the machine holding mesh-broker to make it so")
return nil
}
+38
View File
@@ -169,6 +169,44 @@ func TestAShrinkOfMoreThanHalfIsUrgent(t *testing.T) {
} }
} }
// THE FALSE ALARM (issue 368), replayed through the store D13 reads: an agent's home moved from the
// operator's own home (94.7 MB) to the agent account's fresh one (490 B), and `data-shrank` was raised
// for data that was never lost. A moved item is read against its new path only, so nothing is raised —
// and a genuine shrink at the new path, a week of history later, still is.
func TestAMovedPathIsNoShrinkAndAShrinkThereStillIs(t *testing.T) {
inv := inventory.ForTest(t)
ctx := t.Context()
shelf := shelfFor(t, houseManifest)
declared := []inventory.DeclaredData{{Module: "house", Item: "config", Class: "irreplaceable", Owned: true}}
start := time.Now().Add(-6 * time.Hour)
measure := func(at time.Time, path string, size int64) []conditions.Observation {
t.Helper()
if _, err := inv.RecordData(ctx, "home", declared, map[string]map[string]inventory.Measurement{"house": {
"config": {Path: path, Size: bytesOf(size), MeasuredAt: when(at), LastWrite: when(at),
LastBackup: when(at)}}}, "", at); err != nil {
t.Fatal(err)
}
records, err := inv.Data(ctx)
if err != nil {
t.Fatal(err)
}
peaks, err := inv.DataPeaks(ctx, at.Add(-shrinkWindow))
if err != nil {
t.Fatal(err)
}
return dataFindings(records, peaks, shelf, nil, nil, at)
}
measure(start, "/home/operator/.claude", 94_700_000)
if got := findingsByKind(measure(start.Add(10*time.Minute), "/home/agent/.claude", 490)); got[kindDataShrank].Kind != "" {
t.Fatalf("a moved path raised a shrink: %+v", got[kindDataShrank])
}
measure(start.Add(2*time.Hour), "/home/agent/.claude", 300<<20)
got := findingsByKind(measure(start.Add(4*time.Hour), "/home/agent/.claude", 1<<20))[kindDataShrank]
if got.Severity != conditions.Urgent || !strings.Contains(got.Summary, "shrank") {
t.Fatalf("a genuine shrink at the new path was not raised: %+v", got)
}
}
// Data said to be written all the time and not written; data with no backup or an old one — urgent when // Data said to be written all the time and not written; data with no backup or an old one — urgent when
// irreplaceable, a warning when valuable; and a new item given its bound before it is said. // irreplaceable, a warning when valuable; and a new item given its bound before it is said.
func TestQuietDataAndMissingBackupsAreSaidByClass(t *testing.T) { func TestQuietDataAndMissingBackupsAreSaidByClass(t *testing.T) {
+72
View File
@@ -596,12 +596,81 @@ func deliveryCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
// Each walk's phases, with the rest's reports (novox/hq ADR 0282 decision 6).
now := time.Now()
for i := range walks {
walks[i].Phases = walks[i].WalkPhases(now, appliedFrom(ctx, inv))
}
return answer(map[string]any{"held": deliverySeatHeld(entries), "walks": walks, return answer(map[string]any{"held": deliverySeatHeld(entries), "walks": walks,
"own-path": sortedKeysOf(ownPathWords())}) "own-path": sortedKeysOf(ownPathWords())})
} }
return fmt.Errorf("delivery %s: plan, order, check, go, stop or walks", sub) return fmt.Errorf("delivery %s: plan, order, check, go, stop or walks", sub)
} }
// appliedFrom answers a machine's first report after a send from the controller's `apply` durations; a lookup
// that fails is a report not read, which leaves the walk's end unknown rather than wrong.
func appliedFrom(ctx context.Context, inv *inventory.Inventory) inventory.AppliedLookup {
return func(node string, sent time.Time) (inventory.AppliedReport, bool) {
r, ok, err := inv.FirstAppliedAfter(ctx, node, sent)
if err != nil {
fmt.Fprintf(os.Stderr, "the report of %s after %s could not be read: %v\n", node, sent.Format(time.RFC3339), err)
return inventory.AppliedReport{}, false
}
return r, ok
}
}
// recordWalkPhases keeps, once, each phase of every walk ended lately whose end is known, as a duration of kind
// walk-phase per class (novox/hq ADR 0282 decision 6): what `durations` summarises. A walk whose end is unknown
// past ApplySilentAfter keeps its measured phases without its total.
func recordWalkPhases(ctx context.Context, inv *inventory.Inventory, now time.Time) error {
recent, err := inv.RecentPlans(ctx, 30)
if err != nil {
return err
}
for _, p := range recent {
if p.State != inventory.PlanDone || p.Release != nil || now.Sub(p.Updated) > 2*time.Hour {
continue
}
anyKept, totalKept, err := inv.WalkPhasesKept(ctx, p.ID)
if err != nil {
return err
}
if totalKept {
continue
}
ph := p.WalkPhases(now, appliedFrom(ctx, inv))
if ph != nil && ph.End == nil && anyKept {
continue // kept without its end; kept again only once its end is known
}
if ph == nil || (ph.End == nil && now.Sub(p.Updated) < inventory.ApplySilentAfter+time.Minute) {
continue
}
class := ph.Class
if class == "" {
class = "unclassed"
}
for _, x := range ph.Phases {
if x.State != inventory.PhaseMeasured || x.Start == nil {
continue
}
if err := inv.RecordDuration(ctx, inventory.Duration{Kind: inventory.DurationWalkPhase,
Subject: class + "/" + x.Name, Ref: fmt.Sprintf("%s/%s/%d", p.ID, x.Name, x.Tier), Started: *x.Start,
Took: time.Duration(x.TookMS) * time.Millisecond, Detail: p.Named()}); err != nil {
return err
}
}
if ph.End != nil && ph.From != nil {
if err := inv.RecordDuration(ctx, inventory.Duration{Kind: inventory.DurationWalkPhase,
Subject: class + "/total", Ref: p.ID + "/total", Started: *ph.From,
Took: time.Duration(ph.TotalMS) * time.Millisecond, Detail: ph.Said}); err != nil {
return err
}
}
}
return nil
}
// ownPathWords is the controller's own path as words, for an answer. // ownPathWords is the controller's own path as words, for an answer.
func ownPathWords() map[string]string { return onTheControllersPath } func ownPathWords() map[string]string { return onTheControllersPath }
@@ -728,6 +797,9 @@ func sayPlanMoved(ctx context.Context, bus link.Bus, p inventory.Plan) {
} }
func publishPlanMoved(ctx context.Context, bus link.Bus, p inventory.Plan) { func publishPlanMoved(ctx context.Context, bus link.Bus, p inventory.Plan) {
// Its phases so far (novox/hq ADR 0282 decision 6): the rest's reports come after the walk ends, and are
// read by whoever asks for the walk (`delivery walks`).
p.Phases = p.WalkPhases(time.Now(), nil)
body, err := json.Marshal(p) body, err := json.Marshal(p)
if err != nil { if err != nil {
return return
+5
View File
@@ -116,6 +116,11 @@ var probeRegistry = []probe{
{ID: probeDeliveriesID, Asserts: "no delivery is held past its state's bound unsaid: mesh-delivery's " + {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", "`stalled`, each with the transition its table lets healer H2 take", From: "ADR 0239",
Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries}, Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries},
// The delivery budgets (novox/hq ADR 0282 decision 7): the newest delivery of each class within its budget,
// read from mesh-delivery's `times`; a measurement said, never a delivery held.
{ID: probeBudgetsID, Asserts: "the newest delivery of a leaf module ran on every machine within five minutes of " +
"its merge, and of a core module within ten: mesh-delivery's `times`", From: "ADR 0282",
Kind: kindOverBudget, Phase: 3, run: probeBudgets},
// A client of the bus reconnecting in a loop (novox/hq issue 327), from the server's record of closed // A client of the bus reconnecting in a loop (novox/hq issue 327), from the server's record of closed
// connections, which the bus's own module reads. // connections, which the bus's own module reads.
{ID: probeReconnectsID, Asserts: "no user of the bus had its connection dropped more than twelve times in the " + {ID: probeReconnectsID, Asserts: "no user of the bus had its connection dropped more than twelve times in the " +
+54 -4
View File
@@ -125,6 +125,9 @@ type gateFacts struct {
// groupsAdded is, per module, whether the move judged puts an account in a group its previous build did // groupsAdded is, per module, whether the move judged puts an account in a group its previous build did
// not (issue 318 review): the only move whose wait for a new login is excused. // not (issue 318 review): the only move whose wait for a new login is excused.
groupsAdded map[string]bool groupsAdded map[string]bool
// waits is what the controller holds to check a module's wait for the operator (novox/hq ADR 0283): the
// manifest of each module's build judged, and the secrets given on each machine.
waits operatorWaitFacts
// sent is, per machine, the declaration the gate's own send carried there (novox/hq issue 352): a // sent is, per machine, the declaration the gate's own send carried there (novox/hq issue 352): a
// report is held against it, never against the send made last. sentBuilds is what each machine was // report is held against it, never against the send made last. sentBuilds is what each machine was
// last sent of every module, and judged the commit of each module this gate judges: a machine last // last sent of every module, and judged the commit of each module this gate judges: a machine last
@@ -257,10 +260,11 @@ func judgeHealth(module, component string, m catalogue.Manifest, machine string,
firstLine(f.openErr.Error()) firstLine(f.openErr.Error())
} }
for _, c := range f.open { for _, c := range f.open {
// A wait for a person's new login, or for a directory used as found to be handed over, is the module's // A wait for a person's new login, for a directory used as found to be handed over, or for the
// reading, not a fault raised since the send: the gate reads it from the statement below (ADR 0254, // operator's secret or setting, is the module's reading, not a fault raised since the send: the gate
// novox/hq issue 339). // reads it from the statement below (ADR 0254, novox/hq issue 339, ADR 0283).
if c.Source == gateProbe || c.OpenAt(since) || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound { if c.Source == gateProbe || c.OpenAt(since) || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound ||
c.Kind == kindNeedsOperator {
continue continue
} }
onIt := c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine) || onIt := c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine) ||
@@ -477,6 +481,7 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
return "", err return "", err
} }
facts.groupsAdded = movesAddingGroups(ctx, open.inventory, g, pairs, shelf) facts.groupsAdded = movesAddingGroups(ctx, open.inventory, g, pairs, shelf)
facts.waits = gateWaitFacts(ctx, open.inventory, g, pairs, shelf, facts.health)
facts.sent = g.Sent facts.sent = g.Sent
facts.commits, facts.sentBuilds = judgedCommits(g, pairs), map[string]map[string]string{} facts.commits, facts.sentBuilds = judgedCommits(g, pairs), map[string]map[string]string{}
// A module this gate put back at once (putBackBroken) was sent its earlier build by the gate itself: // A module this gate put back at once (putBackBroken) was sent its earlier build by the gate itself:
@@ -603,6 +608,7 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
pastBound := now.Sub(*g.Since) > gateBound pastBound := now.Sub(*g.Since) > gateBound
switch { switch {
case worst == healthBroken: case worst == healthBroken:
g.Read(now, false, g.BrokenWhy)
var judging []string var judging []string
for _, m := range modules { for _, m := range modules {
if reading[m] != healthBroken && !passedAlone(m) { if reading[m] != healthBroken && !passedAlone(m) {
@@ -621,8 +627,10 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
// Waiting on a provider that is unhealthy: not a pass, and not a failure at the bound either — // Waiting on a provider that is unhealthy: not a pass, and not a failure at the bound either —
// the provider's own condition says what is wrong (ADR 0240 rule 5). // the provider's own condition says what is wrong (ADR 0240 rule 5).
g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing
g.Read(now, false, why)
case worst == healthNotYet: case worst == healthNotYet:
g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing g.Passes, g.LastPass, g.Last, g.Failing = 0, nil, why, failing
g.Read(now, false, why)
if pastBound { if pastBound {
fail(fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why)) fail(fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why))
} }
@@ -630,6 +638,7 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs
// Healthy, or waiting for a person (ADR 0254): a pass, the wait carried along in the verdict. // Healthy, or waiting for a person (ADR 0254): a pass, the wait carried along in the verdict.
g.Passes++ g.Passes++
g.LastPass, g.Last, g.Failing = &now, "", nil g.LastPass, g.Last, g.Failing = &now, "", nil
g.Read(now, true, "")
if g.Passes >= gatePasses && settled { if g.Passes >= gatePasses && settled {
decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes, decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes,
now.Sub(*g.Since).Round(time.Second))+waitsSaid(g.Waits), now) now.Sub(*g.Since).Round(time.Second))+waitsSaid(g.Waits), now)
@@ -747,6 +756,47 @@ func movesAddingGroups(ctx context.Context, inv *inventory.Inventory, g *invento
return out return out
} }
// gateWaitFacts reads what checks the waits for the operator of the modules a gate judges (novox/hq ADR 0283): the
// manifest of the build judged — the one it moves to, else the catalogue's — and the secrets given on each machine
// that says a module of them waits.
func gateWaitFacts(ctx context.Context, inv *inventory.Inventory, g *inventory.PlanGate, pairs []judged,
shelf map[string]catalogue.Manifest, health map[string]inventory.NodeHealth) operatorWaitFacts {
var f operatorWaitFacts
judgedManifests := map[string]catalogue.Manifest{}
byMachine := map[string][]string{}
for _, j := range pairs {
waiting := false
for _, r := range health[j.node].Resources {
if r.Module == j.module && r.State == link.StateWaiting {
waiting = true
}
}
if !waiting {
continue
}
byMachine[j.node] = append(byMachine[j.node], j.module)
if _, done := judgedManifests[j.module]; done {
continue
}
to := g.To
for _, c := range g.Carried {
if c.Module == j.module {
to = c.To
break
}
}
if m, found, err := inv.ManifestAt(ctx, j.module, to); err == nil && found {
judgedManifests[j.module] = m
} else if m, known := shelf[j.module]; known {
judgedManifests[j.module] = m
}
}
for machine, modules := range byMachine {
readWaitFacts(ctx, inv, machine, modules, judgedManifests, &f)
}
return f
}
// decide sets a gate's verdict. // decide sets a gate's verdict.
func decide(g *inventory.PlanGate, verdict, why string, now time.Time) { func decide(g *inventory.PlanGate, verdict, why string, now time.Time) {
g.Verdict, g.Why, g.JudgedAt = verdict, why, &now g.Verdict, g.Why, g.JudgedAt = verdict, why, &now
+1 -1
View File
@@ -311,7 +311,7 @@ func usage() {
the self-check: the last verdict, a run now, the probes, the signals' ages the self-check: the last verdict, a run now, the probes, the signals' ages
healers [--days N] [--json] the healers, what they did lately, and their brake (to-be 45 §7) healers [--days N] [--json] the healers, what they did lately, and their brake (to-be 45 §7)
durations [--kind K] [--days N] [--json] durations [--kind K] [--days N] [--json]
apply, heartbeat, plan-tier and build durations, per machine or module apply, heartbeat, plan-tier, build and walk-phase durations
collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go
builder issue <name> a broker account for a build machine, scoped to build work, builder issue <name> a broker account for a build machine, scoped to build work,
delivered as the builder module's broker secret (module add it first) delivered as the builder module's broker secret (module add it first)
+59 -6
View File
@@ -74,9 +74,21 @@ func stateHealth(ctx context.Context, inv *inventory.Inventory, k *conditions.Ke
kept := inventory.ResourceHealth{Module: r.Module, Resource: r.Resource, Kind: r.Kind, Target: r.Target, kept := inventory.ResourceHealth{Module: r.Module, Resource: r.Resource, Kind: r.Kind, Target: r.Target,
State: r.State, Reason: r.Reason, Since: r.Since, Streak: r.Streak, Restarts: r.Restarts, State: r.State, Reason: r.Reason, Since: r.Since, Streak: r.Streak, Restarts: r.Restarts,
Check: r.Check, Needs: r.Needs, Account: r.Account, Root: r.Root} Check: r.Check, Needs: r.Needs, Account: r.Account, Root: r.Root}
for _, w := range r.Waits {
kept.Waits = append(kept.Waits, inventory.Wait{Part: w.Part, Secret: w.Secret, Setting: w.Setting, What: w.What})
}
resources = append(resources, kept) resources = append(resources, kept)
if r.State == link.StateUnhealthy && r.Module != "" { }
unhealthy[r.Module] = append(unhealthy[r.Module], kept) // **A wait for the operator is checked before it is excused** (novox/hq ADR 0283): a waiting resource whose
// wait does not check out is judged unhealthy, saying why; one that does is kept beside the unhealthy ones, so
// judgeModuleHealth can say it as needs-operator. What is stored is what the machine said.
var wf operatorWaitFacts
if mods := waitingModules(resources); len(mods) > 0 {
readWaitFacts(ctx, inv, node, mods, nil, &wf)
}
for _, r := range checkWaiting(node, resources, wf) {
if (r.State == link.StateUnhealthy || r.State == link.StateWaiting) && r.Module != "" {
unhealthy[r.Module] = append(unhealthy[r.Module], r)
} }
} }
streaks := map[string]int{} streaks := map[string]int{}
@@ -135,7 +147,8 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
} }
standing := map[string]conditions.Condition{} standing := map[string]conditions.Condition{}
for _, c := range open { for _, c := range open {
if (c.Kind == kindModuleUnhealthy || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound) && if (c.Kind == kindModuleUnhealthy || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound ||
c.Kind == kindNeedsOperator) &&
c.Subject.Machine == node { c.Subject.Machine == node {
standing[c.Key] = c standing[c.Key] = c
} }
@@ -159,6 +172,20 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
heldOn := map[string]string{} heldOn := map[string]string{}
providers := map[catalogue.Chosen]bool{} providers := map[catalogue.Chosen]bool{}
for _, m := range modules { for _, m := range modules {
// **A part that waits for the operator is said as that** (novox/hq ADR 0283): its waits already checked,
// the operator's, never urgent, its words naming the act.
if waits, waiting := operatorWait(m, unhealthy[m]); waiting {
o := needsOperatorObservation(m, node, waits, unhealthy[m])
seen[o.Key()] = true
became[m] = kindNeedsOperator
if _, isOpen := standing[o.Key()]; streaks[m] < moduleUnhealthyAfter && !isOpen {
continue
}
if _, err := k.Observe(ctx, o); err != nil {
problems = append(problems, err.Error())
}
continue
}
// **A directory used as found is said as that** (novox/hq issue 339): the operator's to hand over at the // **A directory used as found is said as that** (novox/hq issue 339): the operator's to hand over at the
// machine, never urgent — nothing is broken by the wait that a person was not told of — and its own kind, // machine, never urgent — nothing is broken by the wait that a person was not told of — and its own kind,
// so the gate never reads it as a fault of the build that happened to be sent beside it. // so the gate never reads it as a fault of the build that happened to be sent beside it.
@@ -228,6 +255,9 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
if c.Kind == kindUsedAsFound { if c.Kind == kindUsedAsFound {
why = fmt.Sprintf("%s says no directory of %s is used as found any more", node, module) why = fmt.Sprintf("%s says no directory of %s is used as found any more", node, module)
} }
if c.Kind == kindNeedsOperator {
why = fmt.Sprintf("%s says %s no longer waits for the operator", node, module)
}
if on, held := heldOn[key]; held { if on, held := heldOn[key]; held {
why = fmt.Sprintf("what %s finds on %s waits on %s, which is unhealthy: held under its condition", module, node, on) why = fmt.Sprintf("what %s finds on %s waits on %s, which is unhealthy: held under its condition", module, node, on)
} }
@@ -238,6 +268,12 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
case c.Kind == kindModuleUnhealthy && became[module] == kindReloginNeeded: case c.Kind == kindModuleUnhealthy && became[module] == kindReloginNeeded:
why = fmt.Sprintf("%s on %s now waits only for a new login", module, node) why = fmt.Sprintf("%s on %s now waits only for a new login", module, node)
resolved = fmt.Sprintf("%s on %s now waits only for a new login", module, node) resolved = fmt.Sprintf("%s on %s now waits only for a new login", module, node)
case c.Kind == kindModuleUnhealthy && became[module] == kindNeedsOperator:
why = fmt.Sprintf("%s on %s now only waits for the operator", module, node)
resolved = fmt.Sprintf("%s on %s now only waits for you", module, node)
case c.Kind == kindNeedsOperator && became[module] == kindModuleUnhealthy:
why = fmt.Sprintf("%s on %s no longer only waits for the operator, and is not healthy", module, node)
resolved = fmt.Sprintf("What %s on %s waited for is given, and it still does not work", module, node)
case c.Kind == kindReloginNeeded && became[module] == kindModuleUnhealthy: case c.Kind == kindReloginNeeded && became[module] == kindModuleUnhealthy:
why = fmt.Sprintf("%s on %s no longer waits for a new login, and is not healthy", module, node) why = fmt.Sprintf("%s on %s no longer waits for a new login, and is not healthy", module, node)
resolved = fmt.Sprintf("The new login on %s is done, and %s still does not work", node, module) resolved = fmt.Sprintf("The new login on %s is done, and %s still does not work", node, module)
@@ -420,6 +456,11 @@ func reasonWords(r inventory.ResourceHealth) string {
case "": case "":
return "is unhealthy" return "is unhealthy"
} }
// A wait for the operator that did not check out is said as the controller found it (ADR 0283): names of
// secrets and settings only, never what the check itself said.
if strings.HasPrefix(r.Reason, waitRefusedPrefix) {
return r.Reason
}
// What a declared check found says an endpoint, a path or an address: evidence, never the summary the // What a declared check found says an endpoint, a path or an address: evidence, never the summary the
// operator's channel carries (ADR 0234 §6). The summary names the check. // operator's channel carries (ADR 0234 §6). The summary names the check.
if r.Check != "" { if r.Check != "" {
@@ -511,7 +552,9 @@ func moduleHealthWord(module, machine string, since time.Time, f gateFacts) (hea
if h.HeardAt.Before(since) { if h.HeardAt.Before(since) {
return healthNotYet, fmt.Sprintf("%s has not said how what %s runs is since it was sent", machine, module) return healthNotYet, fmt.Sprintf("%s has not said how what %s runs is since it was sent", machine, module)
} }
wait, waits := personWait(module, machine, h.Resources) // **A wait for the operator is checked first** (novox/hq ADR 0283): one that does not check out is unhealthy.
resources := checkWaiting(machine, h.Resources, f.waits)
wait, waits := personWait(module, machine, resources)
// **Only a build whose own send put the account in a new group is excused** (issue 318 review): read from // **Only a build whose own send put the account in a new group is excused** (issue 318 review): read from
// what the controller sent, never from when the machine says the wait began — that time is the engine's // what the controller sent, never from when the machine says the wait began — that time is the engine's
// memory, reset by its restart and moved by a change of words. A build that adds no account group cannot // memory, reset by its restart and moved by a change of words. A build that adds no account group cannot
@@ -520,10 +563,17 @@ func moduleHealthWord(module, machine string, since time.Time, f gateFacts) (hea
waits = false waits = false
} }
var found []string var found []string
for _, r := range h.Resources { var forOperator []inventory.Wait
for _, r := range resources {
if r.Module != module { if r.Module != module {
continue continue
} }
// **Any build is excused while its wait for the operator checks out** (ADR 0283 decision 4): a secret not
// given is owed by every build alike, so it is no fault of this one, and the verdict carries it.
if r.State == link.StateWaiting {
forOperator = append(forOperator, r.Waits...)
continue
}
if waits && r.State == link.StateUnhealthy { if waits && r.State == link.StateUnhealthy {
continue continue
} }
@@ -550,11 +600,14 @@ func moduleHealthWord(module, machine string, since time.Time, f gateFacts) (hea
reasonAfter(r.Reason)) reasonAfter(r.Reason))
} }
} }
if waits || len(found) > 0 { if waits || len(found) > 0 || len(forOperator) > 0 {
var said []string var said []string
if waits { if waits {
said = append(said, wait) said = append(said, wait)
} }
if len(forOperator) > 0 {
said = append(said, operatorWaitSaid(module, machine, forOperator))
}
if len(found) > 0 { if len(found) > 0 {
said = append(said, fmt.Sprintf("on %s, %s uses %s as found and waits for the operator to hand it over "+ said = append(said, fmt.Sprintf("on %s, %s uses %s as found and waits for the operator to hand it over "+
"(`nox node hand-over %s <directory>` on the control-node)", machine, module, strings.Join(found, ", "), machine)) "(`nox node hand-over %s <directory>` on the control-node)", machine, module, strings.Join(found, ", "), machine))
+261
View File
@@ -0,0 +1,261 @@
package main
import (
"context"
"fmt"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A part that waits for the operator's secret or setting (novox/hq ADR 0283, issue 386).
//
// **A module's tool check may say it waits**: nothing of it is wrong but a part that cannot work until the operator
// gives one of its own secrets or one of its settings. The node-engine states such a resource `waiting`, with what
// it waits for. A module saying so is an assertion, so **the controller checks each wait before it excuses it**:
//
// - a wait for a secret names an own secret the module's manifest declares — by its name, or as a member of a
// secret family — said `"issued-by": "outside"`, and the store holds no value a person gave for it on that
// machine;
// - a wait for a setting names a setting the manifest declares (checked by name only: the controller cannot tell
// whether a free-form value covers a part, and the needs-operator condition is where a false one shows).
//
// An excused wait is read by the first-node gate as *waits for a person* (ADR 0254), a pass carried in the verdict,
// for any build of the module — a secret not given is owed by every build alike. It is said to the operator as
// `module.<module>.<machine>.needs-operator`, naming the act. A wait that fails the check is judged unhealthy, saying
// why, and raises the module's `unhealthy` condition.
// kindNeedsOperator is a module's condition while a part of it waits for the operator's secret or setting.
const kindNeedsOperator = "needs-operator"
// needsOperatorKey is a module's needs-operator condition on a machine.
func needsOperatorKey(module, node string) string {
return conditions.Key(conditions.ScopeModule, module+"."+node, kindNeedsOperator)
}
// operatorWaitFacts is what the controller holds to check a module's waits: the manifest judged per module, and per
// "<module>@<machine>" the own secrets a person gave there, with when. A module or a machine absent is not known,
// and no wait of it is excused.
type operatorWaitFacts struct {
manifests map[string]catalogue.Manifest
given map[string]map[string]time.Time
}
// checkWait is nil when a wait is excused, and otherwise why not, in words. Pure.
func checkWait(module, machine string, w inventory.Wait, f operatorWaitFacts) error {
m, known := f.manifests[module]
if !known {
return fmt.Errorf("says it waits for %s, and the mesh holds no manifest of %s to check it against", waitNames(w), module)
}
switch {
case w.Secret != "" && w.Setting != "", w.Secret == "" && w.Setting == "":
return fmt.Errorf("says it waits, naming %s, where a wait names one secret or one setting", waitNames(w))
case w.Setting != "":
if _, declared := m.Settings[w.Setting]; !declared {
return fmt.Errorf("says it waits for the setting %s, which %s does not declare", w.Setting, module)
}
return nil
}
own, _, declared := m.OwnSecrets.Lookup(w.Secret)
if !declared {
return fmt.Errorf("says it waits for the secret %s, which %s does not declare", w.Secret, module)
}
if own.IssuedBy != catalogue.IssuedOutside {
return fmt.Errorf("says it waits for the secret %s, which the mesh makes itself: only a secret issued outside "+
"the mesh waits for the operator", w.Secret)
}
given, readable := f.given[module+"@"+machine]
if !readable {
return fmt.Errorf("says it waits for the secret %s, and what was given on %s could not be read", w.Secret, machine)
}
if at, was := given[w.Secret]; was {
return fmt.Errorf("says it waits for the secret %s, which was given at %s", w.Secret,
at.UTC().Format("2006-01-02 15:04 MST"))
}
return nil
}
// waitRefusedPrefix opens every reason checkWait gives, so the words of a refused wait are told from a check's own.
const waitRefusedPrefix = "says it waits"
// waitNames is what a wait names, as "the secret x" or "the setting y".
func waitNames(w inventory.Wait) string {
switch {
case w.Secret != "" && w.Setting != "":
return "the secret " + w.Secret + " and the setting " + w.Setting
case w.Secret != "":
return "the secret " + w.Secret
case w.Setting != "":
return "the setting " + w.Setting
}
return "nothing"
}
// checkWaiting reads one machine's resources against the facts: every waiting resource whose waits all check out is
// kept as said; one with a wait that does not, or with no wait at all, is answered as unhealthy with why. Pure; the
// statement as kept is not changed.
func checkWaiting(machine string, rs []inventory.ResourceHealth, f operatorWaitFacts) []inventory.ResourceHealth {
out := make([]inventory.ResourceHealth, 0, len(rs))
for _, r := range rs {
if r.State == link.StateWaiting {
var why error
if len(r.Waits) == 0 {
why = fmt.Errorf("says it waits, and names nothing it waits for")
}
for _, w := range r.Waits {
if why == nil {
why = checkWait(r.Module, machine, w, f)
}
}
if why != nil {
r.State, r.Reason = link.StateUnhealthy, why.Error()
}
}
out = append(out, r)
}
return out
}
// operatorWait is whether everything not healthy of a module on a machine is waiting with its waits checked
// (checkWaiting already applied), and those waits. A module with anything unhealthy, starting or unknown beside it
// does not wait: it is judged as before.
func operatorWait(module string, rs []inventory.ResourceHealth) ([]inventory.Wait, bool) {
var waits []inventory.Wait
for _, r := range rs {
if r.Module != module {
continue
}
switch r.State {
case link.StateHealthy:
case link.StateWaiting:
waits = append(waits, r.Waits...)
default:
return nil, false
}
}
return waits, len(waits) > 0
}
// operatorWaitSaid is a module's wait for the operator in one sentence, for the gate's verdict and the condition's
// summary: what the operator gives and what it names, and for a secret the line that opens the desk prompt.
func operatorWaitSaid(module, machine string, waits []inventory.Wait) string {
var parts []string
for _, w := range waits {
part := fmt.Sprintf("%s (%s", w.What, waitNames(w))
if w.Secret != "" {
part += fmt.Sprintf(", given with `nox secret ask %s %s %s`", machine, module, w.Secret)
}
parts = append(parts, part+")")
}
return fmt.Sprintf("%s on %s waits for the operator: %s", module, machine, strings.Join(parts, "; "))
}
// needsOperatorObservation is a module whose only parts not healthy wait for the operator (ADR 0283): the operator's,
// a warning however long it stands, its plain words naming the act and never saying there is nothing to do.
func needsOperatorObservation(module, node string, waits []inventory.Wait, rs []inventory.ResourceHealth) conditions.Observation {
o := moduleUnhealthyObservation(module, node, rs)
o.Token, o.Kind, o.Resolver, o.Severity = kindNeedsOperator, kindNeedsOperator, conditions.ResolverOperator, conditions.Warning
o.Summary = operatorWaitSaid(module, node, waits)
w := needsOperatorWords(module, node, waits)
o.Headline, o.Explanation, o.Needs, o.Resolved, o.Actions = w.Headline, w.Explanation, w.Needs, w.Resolved, nil
return o
}
// needsOperatorWords is what the operator reads of a module waiting for them (ADR 0253, ADR 0283): the act, for a
// secret typed at the machine's desk prompt and for a setting approved when an agent proposes it. The secret's and
// the setting's names, and the line, are in the summary for whoever looks closer.
func needsOperatorWords(module, node string, waits []inventory.Wait) words {
var acts []string
seen := map[string]bool{}
secret := false
for _, w := range waits {
var act string
switch {
case w.Secret != "":
act, secret = fmt.Sprintf("type %s at %s's desk prompt", w.What, node), true
case w.Setting != "":
act = fmt.Sprintf("approve %s of %s on %s when it is proposed to you", w.Setting, module, node)
}
if act != "" && !seen[act] {
seen[act] = true
acts = append(acts, act)
}
}
needs := strings.Join(acts, "; and ") + "."
// Plain words hold one sentence of at most conditions.NeedsMax characters: several acts are named in the
// summary instead.
if len(acts) == 0 || len(needs) > conditions.NeedsMax {
needs = fmt.Sprintf("give what %s waits for on %s; the details name each secret and setting.", module, node)
}
explanation := fmt.Sprintf("Part of %s on %s cannot work until you give what it waits for.", module, node)
if secret {
explanation += " A hidden prompt opens at the desk when the secret is asked for, and what you type there " +
"is sealed to the machine."
}
explanation += " Its update is in place and nothing was undone; it carries on by itself once it is given."
return words{
Headline: fmt.Sprintf("%s waits for you on %s", module, node),
Needs: needs,
Explanation: explanation,
Resolved: fmt.Sprintf("%s on %s no longer waits for you", module, node),
}
}
// readWaitFacts reads what the controller holds to check the waits of the modules named on one machine: the
// manifests (the catalogue's, or those given) and the secrets given there. A read that fails leaves that module
// unknown, so none of its waits is excused.
func readWaitFacts(ctx context.Context, inv *inventory.Inventory, machine string, modules []string,
manifests map[string]catalogue.Manifest, f *operatorWaitFacts) {
if f.manifests == nil {
f.manifests = map[string]catalogue.Manifest{}
}
if f.given == nil {
f.given = map[string]map[string]time.Time{}
}
if inv == nil {
return
}
var shelf map[string]catalogue.Manifest
sort.Strings(modules)
for _, module := range modules {
if _, has := f.manifests[module]; !has {
if m, given := manifests[module]; given {
f.manifests[module] = m
} else {
if shelf == nil {
var err error
if shelf, err = inv.Catalogue(ctx); err != nil {
shelf = map[string]catalogue.Manifest{}
}
}
if m, known := shelf[module]; known {
f.manifests[module] = m
}
}
}
if _, read := f.given[module+"@"+machine]; read {
continue
}
if given, err := inv.GivenOwnSecrets(ctx, machine, module); err == nil {
f.given[module+"@"+machine] = given
}
}
}
// waitingModules is every module with a waiting resource in a statement.
func waitingModules(rs []inventory.ResourceHealth) []string {
seen := map[string]bool{}
var out []string
for _, r := range rs {
if r.State == link.StateWaiting && r.Module != "" && !seen[r.Module] {
seen[r.Module] = true
out = append(out, r.Module)
}
}
return out
}
+264
View File
@@ -0,0 +1,264 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A part that waits for the operator's secret or setting (novox/hq ADR 0283, issue 386).
var mountsManifest = catalogue.Manifest{Module: "mounts", Version: "1",
Settings: map[string]catalogue.SettingDeclaration{"smb-users": {}, "sources": {}},
OwnSecrets: catalogue.OwnSecrets{
"smb-password-*": {Path: "/s/smb-password-*.secret", IssuedBy: catalogue.IssuedOutside},
"broker": {Path: "/s/broker"},
"made": {Path: "/s/made", Taken: catalogue.TakenAtStart},
"licence": {Path: "/s/licence", IssuedBy: catalogue.IssuedOutside},
}}
var passwordWait = inventory.Wait{Part: "the source games", Secret: "smb-password-games",
What: "the password of the source games"}
var usernameWait = inventory.Wait{Part: "the source games", Setting: "smb-users", What: "the username of the source games"}
func waitingResource(waits ...inventory.Wait) inventory.ResourceHealth {
return inventory.ResourceHealth{Module: "mounts", Resource: "mounts.watch", Kind: "process",
Target: "mesh-mounts-watch.service", State: link.StateWaiting, Check: "tool",
Reason: "the source games waits for its password", Waits: waits}
}
func factsGiven(given map[string]time.Time) operatorWaitFacts {
return operatorWaitFacts{manifests: map[string]catalogue.Manifest{"mounts": mountsManifest},
given: map[string]map[string]time.Time{"mounts@workstation": given}}
}
// Rule 3: a wait is excused only when what it names is the module's, issued outside the mesh, and not given there.
func TestAWaitIsExcusedOnlyWhenItChecksOut(t *testing.T) {
at := time.Date(2026, 10, 10, 15, 8, 0, 0, time.UTC)
for _, c := range []struct {
name string
w inventory.Wait
f operatorWaitFacts
says string // "" when excused
}{
{"a member of an outside family, not given", passwordWait, factsGiven(nil), ""},
{"an outside secret by name, not given", inventory.Wait{Part: "p", Secret: "licence", What: "w"}, factsGiven(nil), ""},
{"a declared setting", usernameWait, factsGiven(nil), ""},
{"a member given", passwordWait, factsGiven(map[string]time.Time{"smb-password-games": at}), "was given at 2026-10-10 15:08"},
{"a secret not declared", inventory.Wait{Part: "p", Secret: "smb-credentials", What: "w"}, factsGiven(nil), "does not declare"},
{"a secret the mesh makes", inventory.Wait{Part: "p", Secret: "made", What: "w"}, factsGiven(nil), "mesh makes itself"},
{"the bus account", inventory.Wait{Part: "p", Secret: "broker", What: "w"}, factsGiven(nil), "mesh makes itself"},
{"a setting not declared", inventory.Wait{Part: "p", Setting: "logins", What: "w"}, factsGiven(nil), "does not declare"},
{"both", inventory.Wait{Part: "p", Secret: "licence", Setting: "smb-users", What: "w"}, factsGiven(nil), "one secret or one setting"},
{"neither", inventory.Wait{Part: "p", What: "w"}, factsGiven(nil), "one secret or one setting"},
{"no manifest known", passwordWait, operatorWaitFacts{}, "no manifest"},
{"what was given cannot be read", passwordWait,
operatorWaitFacts{manifests: map[string]catalogue.Manifest{"mounts": mountsManifest}}, "could not be read"},
} {
err := checkWait("mounts", "workstation", c.w, c.f)
switch {
case c.says == "" && err != nil:
t.Errorf("%s: refused: %v", c.name, err)
case c.says != "" && (err == nil || !strings.Contains(err.Error(), c.says)):
t.Errorf("%s: %v; want a refusal saying %q", c.name, err, c.says)
}
}
}
// A waiting resource whose wait does not check out, or that names nothing, is judged unhealthy, saying why.
func TestAWaitThatFailsItsCheckIsUnhealthy(t *testing.T) {
given := factsGiven(map[string]time.Time{"smb-password-games": time.Now()})
got := checkWaiting("workstation", []inventory.ResourceHealth{waitingResource(passwordWait)}, given)
if got[0].State != link.StateUnhealthy || !strings.Contains(got[0].Reason, "was given") {
t.Fatalf("a wait for a secret given: %+v", got[0])
}
got = checkWaiting("workstation", []inventory.ResourceHealth{waitingResource()}, factsGiven(nil))
if got[0].State != link.StateUnhealthy || !strings.Contains(got[0].Reason, "names nothing") {
t.Fatalf("a wait naming nothing: %+v", got[0])
}
got = checkWaiting("workstation", []inventory.ResourceHealth{waitingResource(passwordWait, usernameWait)}, factsGiven(nil))
if got[0].State != link.StateWaiting {
t.Fatalf("two waits that check out: %+v", got[0])
}
}
// Rule 4: the gate passes a module whose only parts not healthy wait for the operator, carrying the wait; anything
// else beside it is judged as before; and any build is excused, not only one that added something.
func TestTheGatePassesAWaitForTheOperatorCarriedAlong(t *testing.T) {
now := time.Now()
since := now.Add(-time.Minute)
healthy := inventory.ResourceHealth{Module: "mounts", Resource: "mounts.apply", Kind: "process",
Target: "mesh-mounts-apply.service", State: link.StateHealthy}
f := gateFacts{now: now, waits: factsGiven(nil), groupsAdded: map[string]bool{"mounts": false},
health: map[string]inventory.NodeHealth{"workstation": {Node: "workstation", HeardAt: now,
Resources: []inventory.ResourceHealth{healthy, waitingResource(passwordWait)}}}}
h, why := moduleHealthWord("mounts", "workstation", since, f)
if h != healthPerson || !strings.Contains(why, "waits for the operator: the password of the source games") ||
!strings.Contains(why, "nox secret ask workstation mounts smb-password-games") {
t.Fatalf("an excused wait reads %v %q; want a wait for a person naming the act", h, why)
}
// A second resource unhealthy beside it: not yet, as before.
down := healthy
down.State, down.Reason = link.StateUnhealthy, "down"
f.health["workstation"] = inventory.NodeHealth{Node: "workstation", HeardAt: now,
Resources: []inventory.ResourceHealth{down, waitingResource(passwordWait)}}
if h, why := moduleHealthWord("mounts", "workstation", since, f); h != healthNotYet {
t.Fatalf("a resource down beside the wait reads %v %q", h, why)
}
// The password given and the module still saying it waits: not excused.
f.waits = factsGiven(map[string]time.Time{"smb-password-games": now})
f.health["workstation"] = inventory.NodeHealth{Node: "workstation", HeardAt: now,
Resources: []inventory.ResourceHealth{healthy, waitingResource(passwordWait)}}
if h, why := moduleHealthWord("mounts", "workstation", since, f); h != healthNotYet || !strings.Contains(why, "was given") {
t.Fatalf("a wait for a secret given reads %v %q", h, why)
}
// Facts never read (no manifest): never excused.
f.waits = operatorWaitFacts{}
if h, _ := moduleHealthWord("mounts", "workstation", since, f); h != healthNotYet {
t.Fatalf("a wait nothing could check reads %v", h)
}
}
func needsOperatorOpen(t *testing.T, k *conditions.Keeper) (*conditions.Condition, []conditions.Condition) {
t.Helper()
open, err := k.Open(t.Context())
if err != nil {
t.Fatal(err)
}
for i, c := range open {
if c.Key == needsOperatorKey("mounts", "workstation") {
return &open[i], open
}
}
return nil, open
}
// Rule 5: two statements of an excused wait raise needs-operator, the operator's, a warning however long, naming the
// act; a statement without it clears it.
func TestTheNeedsOperatorConditionNamesTheAct(t *testing.T) {
k, _ := withConditionsInMemory(t)
ctx := t.Context()
rs := map[string][]inventory.ResourceHealth{"mounts": {waitingResource(passwordWait)}}
if err := judgeModuleHealth(ctx, nil, k, "workstation", rs, map[string]int{"mounts": 1}, time.Now()); err != nil {
t.Fatal(err)
}
if got, _ := needsOperatorOpen(t, k); got != nil {
t.Fatal("raised on one statement")
}
if err := judgeModuleHealth(ctx, nil, k, "workstation", rs, map[string]int{"mounts": 2}, time.Now()); err != nil {
t.Fatal(err)
}
got, open := needsOperatorOpen(t, k)
if got == nil {
t.Fatalf("not raised on two statements: %+v", open)
}
for _, c := range open {
if c.Kind == kindModuleUnhealthy {
t.Fatalf("raised as a fault too: %+v", c)
}
}
if got.Kind != kindNeedsOperator || got.Resolver != conditions.ResolverOperator || got.Severity != conditions.Warning {
t.Fatalf("the condition: %+v", got)
}
if !strings.Contains(got.Needs, "type the password of the source games at workstation's desk prompt") ||
!strings.Contains(got.Explanation, "hidden prompt opens at the desk") {
t.Fatalf("its needs do not name the act: %q", got.Needs)
}
if strings.Contains(strings.ToLower(got.Explanation), "nothing for you") || !strings.Contains(got.Explanation, "nothing was undone") {
t.Fatalf("its explanation: %q", got.Explanation)
}
if !strings.Contains(got.Summary, "smb-password-games") || !strings.Contains(got.Summary, "nox secret ask workstation mounts smb-password-games") {
t.Fatalf("its summary does not name the secret and the line: %q", got.Summary)
}
// Long open is still a warning: only the operator can end it.
if err := judgeModuleHealth(ctx, nil, k, "workstation", rs, map[string]int{"mounts": 3}, time.Now().Add(48*time.Hour)); err != nil {
t.Fatal(err)
}
if got, _ := needsOperatorOpen(t, k); got == nil || got.Severity == conditions.Urgent {
t.Fatalf("after two days: %+v", got)
}
// Given: the next statement does not say it, and it clears.
if err := judgeModuleHealth(ctx, nil, k, "workstation", map[string][]inventory.ResourceHealth{}, nil, time.Now()); err != nil {
t.Fatal(err)
}
if got, _ := needsOperatorOpen(t, k); got != nil {
t.Fatal("not cleared once given")
}
}
func TestASettingsWaitAsksForTheApproval(t *testing.T) {
k, _ := withConditionsInMemory(t)
ctx := t.Context()
rs := map[string][]inventory.ResourceHealth{"mounts": {waitingResource(usernameWait)}}
for i := 1; i <= 2; i++ {
if err := judgeModuleHealth(ctx, nil, k, "workstation", rs, map[string]int{"mounts": i}, time.Now()); err != nil {
t.Fatal(err)
}
}
got, _ := needsOperatorOpen(t, k)
if got == nil || !strings.Contains(got.Needs, "approve smb-users of mounts on workstation when it is proposed to you") {
t.Fatalf("the condition: %+v", got)
}
}
// A wait that fails its check is the module's own fault: unhealthy, with why, and no needs-operator.
func TestAWaitThatFailsItsCheckRaisesUnhealthy(t *testing.T) {
k, _ := withConditionsInMemory(t)
ctx := t.Context()
checked := checkWaiting("workstation", []inventory.ResourceHealth{waitingResource(passwordWait)},
factsGiven(map[string]time.Time{"smb-password-games": time.Now()}))
rs := map[string][]inventory.ResourceHealth{"mounts": checked}
for i := 1; i <= 2; i++ {
if err := judgeModuleHealth(ctx, nil, k, "workstation", rs, map[string]int{"mounts": i}, time.Now()); err != nil {
t.Fatal(err)
}
}
got, open := needsOperatorOpen(t, k)
if got != nil {
t.Fatalf("a wait for a secret given raised needs-operator: %+v", got)
}
unhealthy := false
for _, c := range open {
unhealthy = unhealthy || c.Key == moduleUnhealthyKey("mounts", "workstation")
}
if !unhealthy {
t.Fatalf("not raised as unhealthy: %+v", open)
}
}
// A module that waited and is then broken says so when the wait clears, and the other way round.
func TestANeedsOperatorThatBecameUnhealthySaysSo(t *testing.T) {
k, _ := withConditionsInMemory(t)
ctx := t.Context()
waiting := map[string][]inventory.ResourceHealth{"mounts": {waitingResource(passwordWait)}}
for i := 1; i <= 2; i++ {
if err := judgeModuleHealth(ctx, nil, k, "workstation", waiting, map[string]int{"mounts": i}, time.Now()); err != nil {
t.Fatal(err)
}
}
broken := waitingResource()
broken.State, broken.Reason, broken.Waits = link.StateUnhealthy, "the source games refused its login", nil
for i := 3; i <= 4; i++ {
if err := judgeModuleHealth(ctx, nil, k, "workstation", map[string][]inventory.ResourceHealth{"mounts": {broken}},
map[string]int{"mounts": i}, time.Now()); err != nil {
t.Fatal(err)
}
}
got, open := needsOperatorOpen(t, k)
if got != nil {
t.Fatal("needs-operator still open after it became a fault")
}
found := false
for _, c := range open {
found = found || c.Key == moduleUnhealthyKey("mounts", "workstation")
}
if !found {
t.Fatalf("the fault is not raised: %+v", open)
}
}
+124
View File
@@ -0,0 +1,124 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A delivery over its budget is loud (novox/hq ADR 0282 decision 7, issue 382): the self-check reads
// mesh-delivery's `times` and raises the warning `delivery.<class>.over-budget` when the newest delivery of a
// class whose delivery time is known took longer than its class's budget — five minutes for a leaf module, ten
// for a core module — naming the delivery and its longest phase. It clears when the next delivery of that class
// lands within its budget. Measurement only: the condition says, it never holds or acts on a delivery.
// probeBudgetsID is the probe that reads the delivery times.
const probeBudgetsID = "D16"
// kindOverBudget is what a delivery over its class's budget raises.
const kindOverBudget = "over-budget"
// deliveryBudgets are the budgets by class (ADR 0282 decision 1); the probe raises nothing for another class.
var deliveryBudgets = map[string]time.Duration{inventory.ClassLeaf: 5 * time.Minute, inventory.ClassCore: 10 * time.Minute}
// timesAnswer is what mesh-delivery's `times` answers, as far as the probe reads it.
type timesAnswer struct {
Classes []timesClass `json:"classes"`
}
// timesClass is one class's line of `times`.
type timesClass struct {
Class string `json:"class"`
Latest *timesLatest `json:"latest,omitempty"`
}
// timesLatest is the newest delivery of a class whose delivery time is known.
type timesLatest struct {
ID string `json:"id"`
TookMS int64 `json:"took_ms"`
Longest string `json:"longest,omitempty"`
Landed string `json:"landed,omitempty"`
}
// deliveryTimes is what the delivery's owner says of its delivery times; nothing when no holder is on record or
// none answers (D3 says that one).
func deliveryTimes(ctx context.Context, conn *nats.Conn, held bool) (*timesAnswer, error) {
if !held {
return nil, nil
}
raw, err := askDeliveryOwner(ctx, conn, "times", map[string]any{})
if errors.Is(err, link.ErrNothingServes) {
return nil, nil
}
if err != nil {
return nil, err
}
var a timesAnswer
if err := json.Unmarshal(raw, &a); err != nil {
return nil, fmt.Errorf("%s.times answered something unreadable: %w", catalogue.DeliverySeat, err)
}
return &a, nil
}
// overBudgetObservations are the conditions of the classes whose newest delivery took longer than its budget:
// strictly longer, so a delivery of exactly its budget is within it.
func overBudgetObservations(a *timesAnswer) []conditions.Observation {
if a == nil {
return nil
}
var out []conditions.Observation
for _, c := range a.Classes {
budget, ok := deliveryBudgets[c.Class]
if !ok || c.Latest == nil {
continue
}
took := time.Duration(c.Latest.TookMS) * time.Millisecond
if took <= budget {
continue
}
longest := c.Latest.Longest
if longest == "" {
longest = "not known"
}
out = append(out, conditions.Observation{Scope: conditions.ScopeDelivery, ID: c.Class, Kind: kindOverBudget,
Severity: conditions.Warning,
Summary: fmt.Sprintf("the %s delivery %s took %s from its merge to running everywhere, over its budget of %s; "+
"its longest phase: %s — `mesh-delivery.times`", c.Class, c.Latest.ID, humanDuration(took),
humanDuration(budget), longest),
Said: fmt.Sprintf("%s took %s (budget %s), longest phase %s", c.Latest.ID, took.Round(time.Second), budget, longest),
Headline: fmt.Sprintf("A %s delivery took %s, over its %s budget", c.Class, humanDuration(took),
humanDuration(budget)),
Explanation: fmt.Sprintf("The delivery %s took %s from its merge until every machine ran it; a %s module is "+
"held to %s (ADR 0282). Most of the time went to %s. Nothing was held or changed because of this: it is "+
"a measurement.", c.Latest.ID, humanDuration(took), c.Class, humanDuration(budget), longest),
Resolved: fmt.Sprintf("the next %s delivery lands within %s", c.Class, humanDuration(budget)),
})
}
return out
}
// probeBudgets is D16: the newest delivery of each class lands within its class's budget.
func probeBudgets(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
entries, err := d.open.inventory.Catalogued(ctx)
if err != nil {
return nil, err
}
var conn *nats.Conn
if d.js != nil {
conn = d.js.Conn()
}
a, err := deliveryTimes(ctx, conn, deliverySeatHeld(entries))
if err != nil {
return nil, err
}
return overBudgetObservations(a), nil
}
+111
View File
@@ -0,0 +1,111 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// A leaf delivery of 5 minutes 10 seconds raises delivery.leaf.over-budget naming its longest phase; one of
// exactly five minutes, a core one of nine and a class without a budget raise nothing (ADR 0282 decision 7).
func TestADeliveryOverItsBudgetIsLoud(t *testing.T) {
a := &timesAnswer{Classes: []timesClass{
{Class: inventory.ClassLeaf, Latest: &timesLatest{ID: "novox/mesh-catalog@abc", TookMS: (5*time.Minute + 10*time.Second).Milliseconds(),
Longest: "judgement 2m40s"}},
{Class: inventory.ClassCore, Latest: &timesLatest{ID: "novox/mesh-controller@def", TookMS: (9 * time.Minute).Milliseconds(),
Longest: "build 1m"}},
{Class: "unclassed", Latest: &timesLatest{ID: "x@y", TookMS: time.Hour.Milliseconds()}},
}}
obs := overBudgetObservations(a)
if len(obs) != 1 {
t.Fatalf("one class over its budget, got %d: %+v", len(obs), obs)
}
o := obs[0]
if o.Key() != "delivery.leaf.over-budget" || !strings.Contains(o.Summary, "judgement 2m40s") ||
!strings.Contains(o.Summary, "novox/mesh-catalog@abc") {
t.Fatalf("the condition names its delivery and longest phase: %s — %s", o.Key(), o.Summary)
}
// The next leaf delivery within its budget clears it: the probe raises nothing for the class.
a.Classes[0].Latest = &timesLatest{ID: "novox/mesh-catalog@ghi", TookMS: (5 * time.Minute).Milliseconds()}
if obs := overBudgetObservations(a); len(obs) != 0 {
t.Fatalf("a delivery of exactly its budget is within it: %+v", obs)
}
a.Classes[1].Latest.TookMS = (10*time.Minute + time.Second).Milliseconds()
if obs := overBudgetObservations(a); len(obs) != 1 || obs[0].Key() != "delivery.core.over-budget" {
t.Fatalf("a core delivery over ten minutes: %+v", obs)
}
if obs := overBudgetObservations(nil); obs != nil {
t.Fatal("no answer raises nothing")
}
// No delivery of a class with a known time yet: nothing said of it.
if obs := overBudgetObservations(&timesAnswer{Classes: []timesClass{{Class: inventory.ClassLeaf}}}); len(obs) != 0 {
t.Fatalf("a class with no delivery: %+v", obs)
}
}
// The probe's question is one the controller's grant names: a question the bus refuses checks nothing.
func TestTheControllerMayAskForTheDeliveryTimes(t *testing.T) {
found := false
for _, v := range broker.VerbsTheControllerAsksTheDeliveryOwner {
found = found || (v.Seat == catalogue.DeliverySeat && v.Verb == "times")
}
if !found {
t.Fatal("the grant does not name mesh-delivery.times")
}
}
// A walk's class is core when it walks a module of the controller's own path or one holding the mesh's resolver,
// leaf otherwise; its window closed at its batch's window, or its maximum, never after its cut.
func TestAWalksClassAndWindowAreReadAtItsCut(t *testing.T) {
entries := []inventory.Entry{
{Manifest: catalogue.Manifest{Module: "dnsmasq", Claims: []catalogue.Claim{{Name: resolverSeat}}}},
{Manifest: catalogue.Manifest{Module: "gitea"}},
}
walk := func(modules ...string) inventory.Plan {
p := inventory.Plan{Modules: map[string]*inventory.PlanModule{}}
for _, m := range modules {
p.Modules[m] = &inventory.PlanModule{}
}
return p
}
for want, w := range map[string]inventory.Plan{
inventory.ClassLeaf: walk("gitea"), inventory.ClassCore: walk("gitea", "mesh-controller"),
} {
if got := classOf(w, entries); got != want {
t.Errorf("%v: %s, want %s", w.Modules, got, want)
}
}
if got := classOf(walk("dnsmasq"), entries); got != inventory.ClassCore {
t.Errorf("the resolver's holder is core: %s", got)
}
now := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC)
batch := inventory.Plan{Delivery: &inventory.PlanDelivery{Batch: &inventory.PlanBatch{
ClosesAt: now.Add(-20 * time.Second), AtMost: now.Add(5 * time.Minute)}}}
times := walkTimesAtCut(batch, walk("gitea"), entries, now)
if times.Cut == nil || !times.Cut.Equal(now) || times.WindowClosed == nil || !times.WindowClosed.Equal(now.Add(-20*time.Second)) ||
times.Class != inventory.ClassLeaf {
t.Fatalf("times at the cut: %+v", times)
}
batch.Delivery.Batch.AtMost = now.Add(-time.Minute)
if times := walkTimesAtCut(batch, walk("gitea"), entries, now); !times.WindowClosed.Equal(now.Add(-time.Minute)) {
t.Fatalf("a window closed at its maximum: %v", times.WindowClosed)
}
if times := walkTimesAtCut(inventory.Plan{Delivery: &inventory.PlanDelivery{Alone: true}}, walk("gitea"), entries, now); times.WindowClosed != nil {
t.Fatalf("a merge walked alone had no window: %v", times.WindowClosed)
}
}
// What each machine of the rest was sent is kept only for the machines the send reached.
func TestTheRestIsKeptForTheMachinesTheSendReached(t *testing.T) {
got := restOf([]string{"ace", "g14"}, []string{"ace", "shanks"}, map[string]inventory.SentDeclaration{"ace": {Digest: "d1"}})
if len(got) != 1 || got["ace"].Digest != "d1" {
t.Fatalf("rest: %+v", got)
}
if restOf([]string{"g14"}, []string{"ace"}, nil) != nil {
t.Fatal("no machine reached: nothing kept")
}
}
+10
View File
@@ -221,6 +221,16 @@ var plainWordings = map[string]func(conditions.Observation) words{
w := usedAsFoundObservation(orModule(module), machineOr(o, "a machine"), o.Summary, nil) w := usedAsFoundObservation(orModule(module), machineOr(o, "a machine"), o.Summary, nil)
return words{Headline: w.Headline, Explanation: w.Explanation, Needs: w.Needs, Resolved: w.Resolved} return words{Headline: w.Headline, Explanation: w.Explanation, Needs: w.Needs, Resolved: w.Resolved}
}), }),
kindNeedsOperator: worded(func(o conditions.Observation) words {
// The observation carries the act itself (ADR 0283); these are its words when only the kind is known.
module := ""
if o.Scope == conditions.ScopeModule && o.Machine != "" {
module = strings.TrimSuffix(o.ID, "."+o.Machine)
}
node := machineOr(o, "a machine")
w := needsOperatorWords(orModule(module), node, nil)
return w
}),
kindProviderFailing: worded(func(o conditions.Observation) words { kindProviderFailing: worded(func(o conditions.Observation) words {
thing, consumer := conditions.ThingWords(o), idPart(o, 2) thing, consumer := conditions.ThingWords(o), idPart(o, 2)
if consumer == "" { if consumer == "" {
+27 -1
View File
@@ -664,7 +664,33 @@ func renderingFor(ctx context.Context, open *stores, node string,
needed := map[string]map[string]string{} needed := map[string]map[string]string{}
foreseen := map[string]map[string]bool{} foreseen := map[string]map[string]bool{}
for _, m := range plan.Modules { for _, m := range plan.Modules {
for name := range m.OwnSecrets { // **A secret family is never made** (novox/hq ADR 0283): each member a person gave on this machine is
// placed, and one not given is nothing — the module says it waits for it.
for _, family := range m.OwnSecrets.Families() {
members, err := inv.GivenMembers(ctx, node, m.Module, family)
if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
}
for _, g := range members {
if _, fam, ok := m.OwnSecrets.Lookup(g.Name); !ok || fam != family {
continue // a longer family's member, or a name no longer of this family
}
if !g.Current {
if choosing == Allocating {
return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf(
"%s on %s holds %q, which was given to the mesh rather than made by it, and %s has "+
"since generated a new sealing key. The mesh cannot make another; give it again",
m.Module, node, g.Name, node)
}
continue
}
if needed[m.Module] == nil {
needed[m.Module] = map[string]string{}
}
needed[m.Module][g.Name] = g.Sealed
}
}
for name := range m.OwnSecrets.Plain() {
// Minted on the send path and only read on every other. Making one is an insert, and // Minted on the send path and only read on every other. Making one is an insert, and
// a question that writes is a question that can block against the machine it is about. // a question that writes is a question that can block against the machine it is about.
var sealed string var sealed string
+28
View File
@@ -510,6 +510,14 @@ func advancePlans(ctx context.Context, open *stores) {
advanceHeld(ctx, open) advanceHeld(ctx, open)
} }
// keepWalkPhases keeps where the walks that ended lately spent their time (novox/hq ADR 0282), outside the hold
// on the plans so it never lengthens it: measured, never acted on, and an error only said.
func keepWalkPhases(ctx context.Context, open *stores) {
if err := recordWalkPhases(ctx, open.inventory, time.Now()); err != nil {
fmt.Printf("plans: the phases of the walks ended lately could not be kept: %v\n", err)
}
}
// advanceHeld is advancePlans for a caller already holding the plans. // advanceHeld is advancePlans for a caller already holding the plans.
func advanceHeld(ctx context.Context, open *stores) { func advanceHeld(ctx context.Context, open *stores) {
inv := open.inventory inv := open.inventory
@@ -864,8 +872,12 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
strings.Join(machines, ", "), p.Tier, err) strings.Join(machines, ", "), p.Tier, err)
} }
now := time.Now().UTC() now := time.Now().UTC()
// What each machine of the rest was sent, kept for when it reports it applied (novox/hq ADR 0282 decision
// 6): measured, never acted on.
sentWhat := sentNow(ctx, open.inventory, sent)
for _, m := range rest { for _, m := range rest {
p.Modules[m].SentAt = &now p.Modules[m].SentAt = &now
p.Modules[m].Rest = restOf(restTo[m], sent, sentWhat)
} }
fmt.Printf("%s: tier %d built; sent %s to %s, one send each\n", p.ID, p.Tier, strings.Join(rest, ", "), fmt.Printf("%s: tier %d built; sent %s to %s, one send each\n", p.ID, p.Tier, strings.Join(rest, ", "),
strings.Join(sent, ", ")) strings.Join(sent, ", "))
@@ -946,6 +958,20 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
return true, nil return true, nil
} }
// restOf is, for one module, every machine of its rest that the send reached and what it carried there.
func restOf(to, sent []string, what map[string]inventory.SentDeclaration) map[string]inventory.SentDeclaration {
out := map[string]inventory.SentDeclaration{}
for _, n := range to {
if slices.Contains(sent, n) {
out[n] = what[n]
}
}
if len(out) == 0 {
return nil
}
return out
}
// firstSend sends one machine, in one send, every module of the plan's tier whose first machine it is // firstSend sends one machine, in one send, every module of the plan's tier whose first machine it is
// (novox/hq issue 281), and records the send on each: what that machine ran of it before — read once, // (novox/hq issue 281), and records the send on each: what that machine ran of it before — read once,
// before the send, so a module the same send carries is never read as already moved — the machines it // before the send, so a module the same send carries is never read as already moved — the machines it
@@ -1193,6 +1219,7 @@ func sayUnsent(p *inventory.Plan, rollsOut func(string) bool) {
// planTicker advances open plans on a timer, for the steps outcomes alone cannot take. // planTicker advances open plans on a timer, for the steps outcomes alone cannot take.
func planTicker(ctx context.Context, open *stores) { func planTicker(ctx context.Context, open *stores) {
advancePlans(ctx, open) advancePlans(ctx, open)
keepWalkPhases(ctx, open)
tick := time.NewTicker(30 * time.Second) tick := time.NewTicker(30 * time.Second)
defer tick.Stop() defer tick.Stop()
for { for {
@@ -1201,6 +1228,7 @@ func planTicker(ctx context.Context, open *stores) {
return return
case <-tick.C: case <-tick.C:
advancePlans(ctx, open) advancePlans(ctx, open)
keepWalkPhases(ctx, open)
} }
} }
} }
+14 -3
View File
@@ -1366,7 +1366,7 @@ func splitCommandLine(line string) ([]string, error) {
var words []string var words []string
var cur strings.Builder var cur strings.Builder
inWord := false inWord := false
quote := rune(0) quote, opened := rune(0), 0
runes := []rune(line) runes := []rune(line)
for i := 0; i < len(runes); i++ { for i := 0; i < len(runes); i++ {
r := runes[i] r := runes[i]
@@ -1387,7 +1387,7 @@ func splitCommandLine(line string) ([]string, error) {
cur.WriteRune(r) cur.WriteRune(r)
} }
case r == '\'' || r == '"': case r == '\'' || r == '"':
quote = r quote, opened = r, i
inWord = true inWord = true
case r == '\\' && i+1 < len(runes): case r == '\\' && i+1 < len(runes):
i++ i++
@@ -1405,7 +1405,7 @@ func splitCommandLine(line string) ([]string, error) {
} }
} }
if quote != 0 { if quote != 0 {
return nil, fmt.Errorf("command has an unclosed %c quote", quote) return nil, unclosedQuote(quote, opened)
} }
if inWord { if inWord {
words = append(words, cur.String()) words = append(words, cur.String())
@@ -1413,6 +1413,17 @@ func splitCommandLine(line string) ([]string, error) {
return words, nil return words, nil
} }
// unclosedQuote is the refusal of a line whose quote is never closed. Most often an apostrophe inside
// a single-quoted value ended that quote early and a later quote was left open, so the refusal says
// where the open quote is and how a quote is written inside a quoted value — the shell's own two
// ways, which this splitter already reads (novox/hq issue 294). It quotes none of the line: a refusal
// is kept on the bus as the call's answer, and the line may carry a setting's value.
func unclosedQuote(quote rune, opened int) error {
return fmt.Errorf("command has an unclosed %c quote, opened at character %d — an apostrophe "+
"inside a single-quoted value ends it; write a ' inside single quotes as '\\'' (it'\\''s), or use "+
"double quotes and write \\\" for a \" and \\\\ for a \\ inside them", quote, opened+1)
}
// seatAnnouncement is what the controller says it serves on the bus (novox/hq ADR 0197): the // seatAnnouncement is what the controller says it serves on the bus (novox/hq ADR 0197): the
// mesh-controller seat, one endpoint per verb it answers, each with the seat's own description and // mesh-controller seat, one endpoint per verb it answers, each with the seat's own description and
// argument schema — the same facts `tools` answers from the records, as NATS's services format. // argument schema — the same facts `tools` answers from the records, as NATS's services format.
+45
View File
@@ -415,3 +415,48 @@ func TestNoVerbSetsTheOperatorsKeyOrRevealsASecret(t *testing.T) {
t.Fatalf("behind %v: %v", behind, err) t.Fatalf("behind %v: %v", behind, err)
} }
} }
// A value with a quote in it can be said on a `command` line, as in a shell, and a line whose quote
// is left open is refused naming where it opened and how a quote is written, quoting none of it (novox/hq issue 294: an
// apostrophe inside a single-quoted JSON value cut the line, and the refusal named no cause).
func TestAQuotedValueCanHoldAQuote(t *testing.T) {
for _, c := range []struct{ line, want string }{
{`settings set claude-code '{"role":"the operator'\''s laptop"}'`, `{"role":"the operator's laptop"}`},
{`settings set claude-code "{\"role\":\"the operator's laptop\"}"`, `{"role":"the operator's laptop"}`},
{"x \"a \\\\ b\nc\"", "a \\ b\nc"},
{`x 'a\b'`, `a\b`},
} {
argv, err := splitCommandLine(c.line)
if err != nil || argv[len(argv)-1] != c.want {
t.Errorf("%s: read back %q %v, want %q", c.line, argv, err, c.want)
}
}
_, err := splitCommandLine(`settings set claude-code '{"role":"the operator's laptop"}' --node g14`)
if err == nil {
t.Fatal("an apostrophe that leaves a quote open was accepted")
}
for _, want := range []string{"character 57", `'\''`, `\"`} {
if !strings.Contains(err.Error(), want) {
t.Errorf("the refusal does not say %q: %v", want, err)
}
}
if strings.Contains(err.Error(), "laptop") || strings.Contains(err.Error(), "g14") {
t.Errorf("the refusal quotes the line, which may carry a setting's value: %v", err)
}
}
// Every value reads back as it was when written the way the refusal says: single-quoted with '\”
// for each quote, or double-quoted with \ before each " and \ — newlines and backslashes included.
func TestAQuotedValueRoundTrips(t *testing.T) {
for _, v := range []string{`plain`, `it's`, `say "hi"`, `back\slash\`, "two\nlines", `'"\'\"`, `''`, ``} {
single := "x '" + strings.ReplaceAll(v, "'", `'\''`) + "'"
double := `x "` + strings.NewReplacer(`\`, `\\`, `"`, `\"`).Replace(v) + `"`
for _, line := range []string{single, double} {
argv, err := splitCommandLine(line)
if err != nil || len(argv) != 2 || argv[1] != v {
t.Errorf("%s: read back %q %v, want %q", line, argv, err, v)
}
}
}
}
+5 -68
View File
@@ -42,64 +42,8 @@ const (
// Its authority is the union of what the modules it carries would each have had for their // Its authority is the union of what the modules it carries would each have had for their
// tools — and nothing of what they consume, because tools are what it runs, not reactions. // tools — and nothing of what they consume, because tools are what it runs, not reactions.
KindNodeTools Kind = "node-tools" KindNodeTools Kind = "node-tools"
// KindView is the one read-only principal a view onto the bus connects as (novox/hq research 036,
// gap G1): a page in a browser, over the bus module's WebSocket listener, watching the issue tracker.
// Fixed, and derived from no declaration: what it hears is ViewHears, what it reads is ViewBucket,
// and it publishes nothing but direct reads of that one bucket (ViewReads), each answered in its own
// inbox. Composed like every other user, into the same
// file, once its credential is minted (`bus view-credential`); forgotten like every other user
// (`bus view-revoke`), at the next composition.
KindView Kind = "view"
) )
// ViewUser is the view's one username: there is one view, and it is nobody's machine or module.
const ViewUser = "view"
// ViewBucket is the state the view reads: the issue tracker's issues, as the bus names the bucket
// (mesh-issues's state `issues`, novox/hq ADR 0201).
var ViewBucket = BucketName("mesh-issues", "issues")
// ViewHears are the events the view subscribes, each named: the issue tracker's own, the controller's
// walks and conditions, and the delivery owner's — what a page about issues shows beside them. Subscribe
// only, and no stream or consumer of its own: a page hears what happens while it is open, and reads the
// bucket for everything before.
var ViewHears = []string{
moduleEventSubject("mesh-issues", "opened"),
moduleEventSubject("mesh-issues", "moved"),
moduleEventSubject("mesh-issues", "noted"),
moduleEventSubject("mesh-issues", "linked"),
seatEventSubject(ControllerSeat, "plan-moved"),
seatEventSubject(ControllerSeat, "condition-raised"),
seatEventSubject(ControllerSeat, "condition-changed"),
seatEventSubject(ControllerSeat, "condition-cleared"),
moduleEventSubject("mesh-delivery", "transition"),
moduleEventSubject("mesh-delivery", "group"),
}
// ViewReads are the JetStream API requests the view makes, on ViewBucket's stream and no other: binding
// (STREAM.INFO), and direct reads — one key by its subject (`DIRECT.GET.<stream>.$KV.<bucket>.<key>`), and
// the batch form on the bare subject, which answers the newest value of every key (`multi_last`) into the
// asker's inbox. The page lists the bucket with the batch, and re-reads one key when the tracker's event
// names it (every event carries the issue's `number`).
//
// **No consumer, deliberately, and so no watch.** A KV watch is a push consumer, and a push consumer's
// deliver subject is the creator's choice, delivered by the server's own client — which the server does
// not hold to the creator's permissions. Measured on 2.11.17 (2026-10-10): the view, granted
// CONSUMER.CREATE on this stream, made a consumer delivering to `mesh.mod.mesh-issues.event.opened`, and
// a module subscribed there received the bucket's entry as the tracker's event. A grant of CONSUMER.CREATE
// is a publish to any subject in the account; the view publishes nothing, so it has none (nor
// CONSUMER.DELETE, which would let it delete a module's consumer). Nothing here is a write either: no
// `$KV.<bucket>.>`, which is what a put or a delete publishes to, and no STREAM.* that defines, purges or
// deletes.
func ViewReads() []string {
stream := "KV_" + ViewBucket
return []string{
"$JS.API.STREAM.INFO." + stream,
"$JS.API.DIRECT.GET." + stream,
"$JS.API.DIRECT.GET." + stream + ".>",
}
}
// RuntimeModule is the module that IS the node's tool runtime (novox/hq ADR 0175). Where it is // RuntimeModule is the module that IS the node's tool runtime (novox/hq ADR 0175). Where it is
// assigned, the mesh composes one runtime principal for the machine in place of that module's own, // assigned, the mesh composes one runtime principal for the machine in place of that module's own,
// and the per-module containers that served tools until then stop being the way tools reach a node. // and the per-module containers that served tools until then stop being the way tools reach a node.
@@ -288,8 +232,11 @@ func MaySubscribe(perms Permissions, subject string) bool {
// //
// And, since novox/hq ADR 0259, `release` and `stop`: the controller asks the operator for them about a // And, since novox/hq ADR 0259, `release` and `stop`: the controller asks the operator for them about a
// delivery held past its bound, and calls them on the operator's warrant, with its why. // delivery held past its bound, and calls them on the operator's warrant, with its why.
//
// And, since novox/hq ADR 0282, `times`: the self-check reads the delivery times to say a delivery over its budget.
var VerbsTheControllerAsksTheDeliveryOwner = []SeatVerb{{Seat: "mesh-delivery", Verb: "stalled"}, var VerbsTheControllerAsksTheDeliveryOwner = []SeatVerb{{Seat: "mesh-delivery", Verb: "stalled"},
{Seat: "mesh-delivery", Verb: "close"}, {Seat: "mesh-delivery", Verb: "release"}, {Seat: "mesh-delivery", Verb: "stop"}} {Seat: "mesh-delivery", Verb: "close"}, {Seat: "mesh-delivery", Verb: "release"}, {Seat: "mesh-delivery", Verb: "stop"},
{Seat: "mesh-delivery", Verb: "times"}}
// VerbsTheControllerActsOnAWarrant are the other seat verbs the controller calls when the operator's warrant // VerbsTheControllerActsOnAWarrant are the other seat verbs the controller calls when the operator's warrant
// chooses them (novox/hq ADR 0259): a machine's service restarted, and a walk started or stopped through the // chooses them (novox/hq ADR 0259): a machine's service restarted, and a walk started or stopped through the
@@ -317,8 +264,6 @@ func (p Principal) Username() string {
switch p.Kind { switch p.Kind {
case KindPerson: case KindPerson:
return "person." + p.Module return "person." + p.Module
case KindView:
return ViewUser
case KindModule, KindNodeTools: case KindModule, KindNodeTools:
// The runtime is named exactly as the module it stands for would have been: the mesh // The runtime is named exactly as the module it stands for would have been: the mesh
// issues its credential through the same path a module's takes (`module issue`), and // issues its credential through the same path a module's takes (`module issue`), and
@@ -591,14 +536,6 @@ func PermissionsFor(p Principal) (Permissions, error) {
// itself, its replies to the asker's own inbox. // itself, its replies to the asker's own inbox.
pub = append(pub, discovering()...) pub = append(pub, discovering()...)
case KindView:
// Hears what it is for and reads one bucket, and nothing else (ViewHears, ViewReads): no tool,
// no event of its own, no stream, no bucket written. Its requests are answered in its own
// inbox, granted below with the person's; a reply to anything is never permitted, because
// nothing is ever asked of it.
sub = append(sub, ViewHears...)
pub = append(pub, ViewReads()...)
case KindEnrolment: case KindEnrolment:
// A leaked token is useless for anything but enrolling: it cannot read a declaration, hear // A leaked token is useless for anything but enrolling: it cannot read a declaration, hear
// an event, or subscribe any inbox but the one its own token derives (design 25 §6). // an event, or subscribe any inbox but the one its own token derives (design 25 §6).
@@ -900,7 +837,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = unique(pub) pub = unique(pub)
} }
if p.Kind == KindPerson || p.Kind == KindView { if p.Kind == KindPerson {
// An inbox to hear answers in, and nothing else. No ack subject: a person has no durable // An inbox to hear answers in, and nothing else. No ack subject: a person has no durable
// consumer, because nothing is delivered to a person — they ask and are answered. // consumer, because nothing is delivered to a person — they ask and are answered.
sub = append(sub, p.inbox()) sub = append(sub, p.inbox())
-2
View File
@@ -31,8 +31,6 @@ func TestTheComposedConfigMatchesTheGolden(t *testing.T) {
// The bus's own module: the snapshot API and its inbox, nothing else (novox/hq ADR 0235). // The bus's own module: the snapshot API and its inbox, nothing else (novox/hq ADR 0235).
{Kind: KindModule, Node: "one", Module: "nats", SnapshotsTheBus: true, {Kind: KindModule, Node: "one", Module: "nats", SnapshotsTheBus: true,
Serves: []string{"nats_streams"}, PasswordHash: "$2a$11$bbbbbbbbbbbbbbbbbbbbbb"}, Serves: []string{"nats_streams"}, PasswordHash: "$2a$11$bbbbbbbbbbbbbbbbbbbbbb"},
// The view: hears the issue tracker and reads its bucket, writes nothing (research 036).
{Kind: KindView, PasswordHash: "$2a$11$vvvvvvvvvvvvvvvvvvvvvv"},
}) })
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
+1 -5
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_asked.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-controller.tool.plans", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.mesh-delivery.tool.stop", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*", "mesh.seat.node-launcher.tool.secret.*", "mesh.seat.node-service-manager.tool.restart.*"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_asked.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-controller.tool.plans", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.release", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.mesh-delivery.tool.stop", "mesh.seat.mesh-delivery.tool.times", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*", "mesh.seat.node-launcher.tool.secret.*", "mesh.seat.node-service-manager.tool.restart.*"] }
subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built", "mesh.seat.operator-channel.event.decided.mesh-controller"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built", "mesh.seat.operator-channel.event.decided.mesh-controller"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
@@ -56,10 +56,6 @@ accounts {
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] } subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "view", password: "$2a$11$vvvvvvvvvvvvvvvvvvvvvv", permissions: {
publish: { allow: ["$JS.API.DIRECT.GET.KV_mesh-issues_issues", "$JS.API.DIRECT.GET.KV_mesh-issues_issues.>", "$JS.API.STREAM.INFO.KV_mesh-issues_issues"] }
subscribe: { allow: ["_INBOX.view.>", "mesh.mod.mesh-delivery.event.group", "mesh.mod.mesh-delivery.event.transition", "mesh.mod.mesh-issues.event.linked", "mesh.mod.mesh-issues.event.moved", "mesh.mod.mesh-issues.event.noted", "mesh.mod.mesh-issues.event.opened", "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.plan-moved"] }
} }
] ]
} }
} }
-6
View File
@@ -67,9 +67,6 @@ type Records struct {
Enrolling []string Enrolling []string
// People is each person's name against the tools they may invoke, `*` for an administrator. // People is each person's name against the tools they may invoke, `*` for an administrator.
People map[string][]string People map[string][]string
// View says the mesh minted the view's credential (`bus view-credential`), so the one read-only
// view principal is composed (KindView); forgotten, it is left out, like a person.
View bool
// Interchangeable is each module whose definition says its instances are the same anywhere // Interchangeable is each module whose definition says its instances are the same anywhere
// (ADR 0160), which decides whether the module's plain subject is issued to every instance. // (ADR 0160), which decides whether the module's plain subject is issued to every instance.
Interchangeable map[string]bool Interchangeable map[string]bool
@@ -145,9 +142,6 @@ func Users(r Records) ([]Principal, error) {
for _, person := range sortedNames(r.People) { for _, person := range sortedNames(r.People) {
out = append(out, Principal{Kind: KindPerson, Module: person, Invokes: r.People[person]}) out = append(out, Principal{Kind: KindPerson, Module: person, Invokes: r.People[person]})
} }
if r.View {
out = append(out, Principal{Kind: KindView})
}
// Refused here rather than discovered by the server. Two users with one name is a file the // Refused here rather than discovered by the server. Two users with one name is a file the
// server reads as one of them, and which one depends on the order — so a module assigned to a // server reads as one of them, and which one depends on the order — so a module assigned to a
-230
View File
@@ -1,230 +0,0 @@
package broker
import (
"errors"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/nats-io/nats-server/v2/server"
"github.com/nats-io/nats.go"
"golang.org/x/crypto/bcrypt"
)
// The view against a real server, over WebSocket (novox/hq research 036): the composed user list is
// what the server reads, the listener is the shape the bus module declares (no TLS, compression on,
// reached across the overlay only), and the view with its credential binds the issue tracker's bucket,
// lists it with one batch read, follows a tracker event to re-read the key it names — and is refused
// every write: a put, a delete, an event, and a consumer delivering onto the tracker's event subject.
//
// go test ./internal/broker/ -run TestTheView
func TestTheViewReadsTheIssuesOverWebSocketAndWritesNothing(t *testing.T) {
hash := func(password string) string {
h, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.MinCost)
if err != nil {
t.Fatal(err)
}
return string(h)
}
const node = "anchor"
tracker := Principal{Kind: KindModule, Node: node, Module: "mesh-issues", State: []string{"issues"},
Emits: []string{"opened", "moved"}, PasswordHash: hash("tracker")}
accounts, err := ComposeAccounts([]Principal{
{Kind: KindController, PasswordHash: hash("controller")},
tracker,
{Kind: KindView, PasswordHash: hash("view")},
})
if err != nil {
t.Fatal(err)
}
conf := filepath.Join(t.TempDir(), "accounts.conf")
if err := os.WriteFile(conf, []byte(accounts), 0o600); err != nil {
t.Fatal(err)
}
// The server reads the composed file as the bus does — through its own parser — and listens as the
// bus module's configuration says: a WebSocket listener without TLS and with compression, beside
// the client port. Ports chosen by the system, so this runs beside a live bus.
opts, err := server.ProcessConfigFile(conf)
if err != nil {
t.Fatalf("the server refused the composed user list: %v", err)
}
opts.Host, opts.Port = "127.0.0.1", server.RANDOM_PORT
opts.JetStream, opts.StoreDir = true, t.TempDir()
opts.NoLog, opts.NoSigs = true, true
opts.Websocket = server.WebsocketOpts{Host: "127.0.0.1", Port: server.RANDOM_PORT, NoTLS: true, Compression: true}
s, err := server.NewServer(opts)
if err != nil {
t.Fatal(err)
}
go s.Start()
if !s.ReadyForConnections(30 * time.Second) {
s.Shutdown()
t.Fatal("the server did not come up")
}
t.Cleanup(func() { s.Shutdown(); s.WaitForShutdown() })
dial := func(url, user, password string, refused chan<- string) *nats.Conn {
t.Helper()
nc, err := nats.Connect(url, nats.UserInfo(user, password), nats.CustomInboxPrefix("_INBOX."+user),
nats.Compression(true), nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) {
if refused != nil && errors.Is(err, nats.ErrPermissionViolation) {
refused <- err.Error()
}
}))
if err != nil {
t.Fatalf("%s could not connect to %s: %v", user, url, err)
}
t.Cleanup(nc.Close)
return nc
}
// The controller defines the bucket, as it does for every module's state; the tracker writes it.
controller := dial(s.ClientURL(), "controller", "controller", nil)
cjs, _ := controller.JetStream()
if _, err := cjs.CreateKeyValue(&nats.KeyValueConfig{Bucket: ViewBucket, History: 8}); err != nil {
t.Fatalf("the controller could not define %s: %v", ViewBucket, err)
}
trackerConn := dial(s.ClientURL(), tracker.Username(), "tracker", nil)
tjs, _ := trackerConn.JetStream()
tkv, err := tjs.KeyValue(ViewBucket)
if err != nil {
t.Fatal(err)
}
if _, err := tkv.Put("365", []byte(`{"number":365,"status":"open"}`)); err != nil {
t.Fatalf("the tracker could not write its own bucket: %v", err)
}
// The view, over WebSocket with its credential.
refused := make(chan string, 8)
view := dial(s.WebsocketURL(), ViewUser, "view", refused)
if !strings.HasPrefix(view.ConnectedUrl(), "ws://") {
t.Fatalf("the view is connected to %s, not over WebSocket", view.ConnectedUrl())
}
vjs, _ := view.JetStream(nats.MaxWait(3 * time.Second))
vkv, err := vjs.KeyValue(ViewBucket)
if err != nil {
t.Fatalf("the view could not bind %s: %v", ViewBucket, err)
}
if got, err := vkv.Get("365"); err != nil {
t.Fatalf("the view could not read a key: %v", err)
} else if !strings.Contains(string(got.Value()), `"number":365`) {
t.Fatalf("the view read %q", got.Value())
}
if _, err := tkv.Put("366", []byte(`{"number":366,"status":"open"}`)); err != nil {
t.Fatal(err)
}
// The list: one batch read, the newest value of every key, into the view's own inbox, ended by the
// server's end-of-batch status (204).
listed := map[string]string{}
inbox := view.NewRespInbox()
batch, err := view.SubscribeSync(inbox)
if err != nil {
t.Fatal(err)
}
if err := view.PublishRequest("$JS.API.DIRECT.GET.KV_"+ViewBucket, inbox,
[]byte(`{"multi_last":["$KV.`+ViewBucket+`.>"]}`)); err != nil {
t.Fatal(err)
}
for {
m, err := batch.NextMsg(5 * time.Second)
if err != nil {
t.Fatalf("the view's batch read ended without its end-of-batch (%d keys so far): %v", len(listed), err)
}
if m.Header.Get("Status") == "204" {
break
}
if status := m.Header.Get("Status"); status != "" {
t.Fatalf("the batch read answered %s %s", status, m.Header.Get("Description"))
}
listed[m.Header.Get("Nats-Subject")] = string(m.Data)
}
_ = batch.Unsubscribe()
for _, key := range []string{"365", "366"} {
if !strings.Contains(listed["$KV."+ViewBucket+"."+key], `"number":`+key) {
t.Errorf("the batch read did not list %s: %v", key, listed)
}
}
// A change followed: the tracker moves 365 and says so; the view hears the event and reads the key
// it names.
moved, err := view.SubscribeSync("mesh.mod.mesh-issues.event.moved")
if err != nil {
t.Fatal(err)
}
_ = view.Flush()
if _, err := tkv.Put("365", []byte(`{"number":365,"status":"located"}`)); err != nil {
t.Fatal(err)
}
if err := trackerConn.Publish("mesh.mod.mesh-issues.event.moved", []byte(`{"number":365,"to":"located"}`)); err != nil {
t.Fatal(err)
}
if _, err := moved.NextMsg(5 * time.Second); err != nil {
t.Fatalf("the view did not hear the tracker's event: %v", err)
}
if got, err := vkv.Get("365"); err != nil || !strings.Contains(string(got.Value()), "located") {
t.Fatalf("after the event the view read %v (%v)", got, err)
}
// **The hole a watch would open, shut** (ViewReads): a consumer delivering onto the tracker's event
// subject would have the server republish the bucket there, as the tracker. Refused, and nothing
// reaches a module listening on that subject.
listener, err := trackerConn.SubscribeSync("mesh.mod.mesh-issues.event.opened")
if err != nil {
t.Fatal(err)
}
_ = trackerConn.Flush()
if _, err := vjs.AddConsumer("KV_"+ViewBucket, &nats.ConsumerConfig{Name: "w", DeliverSubject: "mesh.mod.mesh-issues.event.opened",
AckPolicy: nats.AckNonePolicy, FilterSubject: "$KV." + ViewBucket + ".>"}); err == nil {
t.Error("the view made a consumer")
}
if _, err := vkv.WatchAll(); err == nil {
t.Error("the view made a watch, which is a consumer")
}
if m, err := listener.NextMsg(2 * time.Second); err == nil {
t.Errorf("a message reached the tracker's event subject from the view: %q", m.Data)
}
// The refusals of the consumer create land as permission violations too; drained before the writes.
drain := time.After(500 * time.Millisecond)
for draining := true; draining; {
select {
case <-refused:
case <-drain:
draining = false
}
}
// And every write is refused: the server says so, and the bucket is unchanged.
if _, err := vkv.Put("367", []byte(`{"number":367}`)); err == nil {
t.Error("the view put a key")
}
if err := vkv.Delete("365"); err == nil {
t.Error("the view deleted a key")
}
if err := view.Publish("mesh.mod.mesh-issues.event.opened", []byte(`{"number":367}`)); err != nil {
t.Fatal(err)
}
_ = view.Flush()
violations := map[string]bool{}
deadline := time.After(10 * time.Second)
for len(violations) < 3 {
select {
case v := <-refused:
for _, subject := range []string{"$KV." + ViewBucket + ".367", "$KV." + ViewBucket + ".365", "mesh.mod.mesh-issues.event.opened"} {
if strings.Contains(v, subject) {
violations[subject] = true
}
}
case <-deadline:
t.Fatalf("the server refused %d of the view's 3 writes as permission violations", len(violations))
}
}
if _, err := tkv.Get("367"); !errors.Is(err, nats.ErrKeyNotFound) {
t.Errorf("after the view's put, 367 is %v", err)
}
if _, err := tkv.Get("365"); err != nil {
t.Errorf("after the view's delete, 365 is gone: %v", err)
}
}
-144
View File
@@ -1,144 +0,0 @@
package broker
import (
"reflect"
"sort"
"testing"
)
// The view (novox/hq research 036): one read-only user, composed like every other, whose whole
// authority is a list here — so a grant that is not on the list fails a test, not a review.
// Exactly what it hears, exactly what it asks, and nothing it could write or answer. A mutation that
// adds a publish grant — `$KV.<bucket>.>`, an event, a tool — fails here.
func TestTheViewHearsAndReadsAndCanPublishNothingElse(t *testing.T) {
perms, err := PermissionsFor(Principal{Kind: KindView})
if err != nil {
t.Fatal(err)
}
wantSub := append(append([]string(nil), ViewHears...), "_INBOX.view.>")
sort.Strings(wantSub)
if !reflect.DeepEqual(perms.Subscribe, wantSub) {
t.Errorf("the view subscribes\n %v\nand should subscribe exactly\n %v", perms.Subscribe, wantSub)
}
wantPub := ViewReads()
sort.Strings(wantPub)
if !reflect.DeepEqual(perms.Publish, wantPub) {
t.Errorf("the view publishes\n %v\nand should publish exactly\n %v", perms.Publish, wantPub)
}
if len(perms.PublishDeny) != 0 {
t.Errorf("the view needs no deny, because nothing it may publish reaches the controller's own: %v", perms.PublishDeny)
}
if perms.AllowResponses {
t.Error("the view may answer, and nothing is ever asked of it")
}
// Every publish grant is binding the one bucket's stream or reading it directly. **The mutation this
// holds against**: a write grant of any shape, and a consumer of any shape — a consumer's deliver
// subject is the creator's choice, so creating one is publishing anywhere (ViewReads).
stream := "KV_" + ViewBucket
for _, p := range perms.Publish {
readOnly := p == "$JS.API.STREAM.INFO."+stream ||
p == "$JS.API.DIRECT.GET."+stream ||
p == "$JS.API.DIRECT.GET."+stream+".>"
if !readOnly {
t.Errorf("the view is granted a publish on %q, which is not a read of %s", p, ViewBucket)
}
}
for _, refused := range []string{
"$KV." + ViewBucket + ".365", // a put or a delete
"$KV.mesh-controller_conditions.x", // another bucket
"$JS.API.STREAM.CREATE." + stream, // defining the stream
"$JS.API.STREAM.PURGE." + stream, // emptying it
"$JS.API.STREAM.DELETE." + stream, // deleting it
"$JS.API.STREAM.MSG.DELETE." + stream, // deleting a message
"$JS.API.CONSUMER.CREATE.KV_mesh-controller_conditions.x", // reading another bucket
"$JS.API.CONSUMER.CREATE." + stream + ".w.$KV." + ViewBucket + ".>", // a watch: delivers anywhere
"$JS.API.CONSUMER.CREATE." + stream, // an unnamed consumer
"$JS.API.CONSUMER.DELETE." + stream + ".a_mesh-issues", // a module's consumer
"$JS.FC." + stream + ".x",
"$JS.API.DIRECT.GET.KV_mesh-controller_conditions", // another bucket, directly
"$JS.API.STREAM.INFO.EVENTS", // the events stream
"$JS.API.INFO", // the account
"mesh.mod.mesh-issues.event.opened", // claiming the tracker said something
"mesh.mod.mesh-issues.tool.open", // opening an issue
"mesh.seat.issue-tracker.tool.open", // through the seat
"mesh.seat.issue-tracker.tool.open.novox", // on one machine
"mesh.seat.mesh-controller.tool.status", // the controller's verbs
"mesh.seat.mesh-controller.event.plan-moved",
"$SRV.PING",
"_INBOX.controller.x",
} {
if MayPublish(perms, refused) {
t.Errorf("the view may publish %q", refused)
}
}
for _, refused := range []string{
"mesh.mod.mesh-issues.tool.open", // a tool asked of the tracker
"mesh.mod.telegram.event.received", // another module's events
"mesh.seat.mesh-controller.event.applied",
"mesh.control.novox.report",
"_INBOX.controller.x",
"_INBOX.person.jochen.x",
"_DELIVER.controller.EVENTS",
} {
if MaySubscribe(perms, refused) {
t.Errorf("the view may subscribe %q", refused)
}
}
for _, heard := range []string{
"mesh.mod.mesh-issues.event.opened",
"mesh.mod.mesh-issues.event.moved",
"mesh.mod.mesh-issues.event.noted",
"mesh.mod.mesh-issues.event.linked",
"mesh.seat.mesh-controller.event.plan-moved",
"mesh.seat.mesh-controller.event.condition-raised",
"mesh.seat.mesh-controller.event.condition-changed",
"mesh.seat.mesh-controller.event.condition-cleared",
"mesh.mod.mesh-delivery.event.transition",
"mesh.mod.mesh-delivery.event.group",
"_INBOX.view.abc",
} {
if !MaySubscribe(perms, heard) {
t.Errorf("the view cannot subscribe %q", heard)
}
}
}
// Composed once the mesh minted its credential, and not before: its row is the whole record of it.
func TestTheViewIsComposedOnlyOnceItsCredentialIsMinted(t *testing.T) {
without, err := Users(Records{Nodes: []string{"anchor"}})
if err != nil {
t.Fatal(err)
}
for _, p := range without {
if p.Kind == KindView {
t.Fatal("the view is composed before its credential was minted")
}
}
with, err := Users(Records{Nodes: []string{"anchor"}, View: true})
if err != nil {
t.Fatal(err)
}
views := 0
for _, p := range with {
if p.Kind == KindView {
views++
if p.Username() != ViewUser {
t.Errorf("the view is called %q, and its row is %q", p.Username(), ViewUser)
}
}
}
if views != 1 {
t.Fatalf("%d view users composed; there is one view", views)
}
// And without its hash it is named as missing, like any user — never written as a user anybody is.
_, missing := WithPasswords(with, map[string]string{})
found := false
for _, m := range missing {
found = found || m == ViewUser
}
if !found {
t.Error("a view with no password was not named as missing one")
}
}
+13 -2
View File
@@ -606,7 +606,18 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
"mode": "0600", "mode": "0600",
}) })
} }
for _, name := range sortedKeys(m.OwnSecrets) { // **A secret family's members, as given** (novox/hq ADR 0283): one file each at the family's path, and
// nothing for a member not given — the mesh never makes one, and the module says it waits for it.
for _, name := range sortedKeys(with.Needed[m.Module]) {
own, family, ok := m.OwnSecrets.Lookup(name)
if !ok || family == "" || with.Needed[m.Module][name] == "" {
continue
}
first = append(first, ownedBy(m.SecretsOwner, map[string]any{
"id": NeedID(name), "type": "file", "path": own.Path, "sealed": with.Needed[m.Module][name],
}))
}
for _, name := range sortedKeys(m.OwnSecrets.Plain()) {
sealed := with.Needed[m.Module][name] sealed := with.Needed[m.Module][name]
if sealed == "" && with.Foreseen[m.Module][name] { if sealed == "" && with.Foreseen[m.Module][name] {
// Not made, and the next send makes it: composed with a stand-in so that whatever // Not made, and the next send makes it: composed with a stand-in so that whatever
@@ -2385,7 +2396,7 @@ func secretsReadBy(resource map[string]any, m Manifest) []string {
mentioned = append(mentioned, stringsIn(resource[key])...) mentioned = append(mentioned, stringsIn(resource[key])...)
} }
var out []string var out []string
for _, name := range sortedKeys(m.OwnSecrets) { for _, name := range sortedKeys(m.OwnSecrets.Plain()) {
path := m.OwnSecrets[name].Path path := m.OwnSecrets[name].Path
if path == "" { if path == "" {
continue continue
+9
View File
@@ -18,6 +18,15 @@ func deliveryVerbs() []Verb {
Input: schema(map[string]string{"state": "one state, e.g. delivering or held", Input: schema(map[string]string{"state": "one state, e.g. delivering or held",
"repository": "owner/repository", "group": "a group's id (its branch name)", "repository": "owner/repository", "group": "a group's id (its branch name)",
"all": "\"true\": the final ones of the last thirty days too"}, nil, "all")}, "all": "\"true\": the final ones of the last thirty days too"}, nil, "all")},
// Delivery time against the delivery budgets (novox/hq ADR 0282): what probe D16 reads, optional until its
// holder serves it.
{Name: "times", Description: "Delivery time: from a merge to every machine of its walk running the build. " +
"Each class's delivery budget (a leaf module five minutes, a core module ten), the last days' median and " +
"worst, how many within and over, the newest delivery of each class, every delivery over its budget with its " +
"longest phase, and the delivered ones whose time is not known. Given a delivery's id, its time phase by phase.",
Input: schema(map[string]string{"days": "how many days back (default 7)", "id": "one delivery's id: its phases"}, nil),
// Optional until mesh-delivery serves it: its holder lives in the catalogue (design 33 §7).
Optional: true},
{Name: "show", Description: "One delivery or group whole: its delivery plan (what it builds, what each " + {Name: "show", Description: "One delivery or group whole: its delivery plan (what it builds, what each " +
"machine receives, what is not an ordinary send), every transition with when and why, the machine " + "machine receives, what is not an ordinary send), every transition with when and why, the machine " +
"steps of its walk, its group and its order.", "steps of its walk, its group and its order.",
+36
View File
@@ -97,3 +97,39 @@ func TestTheDeliverySeatPromisesRetireHistoryOptionally(t *testing.T) {
t.Fatalf("a holder serving retire-history: %v", err) t.Fatalf("a holder serving retire-history: %v", err)
} }
} }
// The seat says delivery times (novox/hq ADR 0282, issue 382): `times`, over some days or for one delivery. Added
// after the holder shipped, it is optional, so the holder that does not serve it yet still holds the seat, and one
// that does is not refused for a verb the seat lacks.
func TestTheDeliverySeatPromisesTimesOptionally(t *testing.T) {
seat, _ := SeatNamed(DeliverySeat)
var times *Verb
for i := range seat.Serves {
if seat.Serves[i].Name == "times" {
times = &seat.Serves[i]
}
}
if times == nil || !times.Optional {
t.Fatalf("the delivery seat promises %v, times optionally", VerbNames(seat.Serves))
}
props, _ := times.Input["properties"].(map[string]any)
for _, arg := range []string{"days", "id"} {
if _, has := props[arg]; !has {
t.Errorf("times takes no %q", arg)
}
}
var without []string
for _, v := range VerbNames(seat.Serves) {
if v != "times" {
without = append(without, v)
}
}
m := Manifest{Module: "mesh-delivery", Claims: []Claim{{Name: DeliverySeat, Scope: ScopeMesh, Serves: without}}}
if err := CanHold(m, seat); err != nil {
t.Fatalf("a holder without times yet: %v", err)
}
m.Claims[0].Serves = VerbNames(seat.Serves)
if err := CanHold(m, seat); err != nil {
t.Fatalf("a holder serving times: %v", err)
}
}
+94
View File
@@ -1942,6 +1942,37 @@ func ParseManifest(raw []byte) (Manifest, error) {
problems = append(problems, whileStoppedProblems(m, r, hasSchedule(r))...) problems = append(problems, whileStoppedProblems(m, r, hasSchedule(r))...)
} }
for name, own := range m.OwnSecrets { for name, own := range m.OwnSecrets {
// **A family is issued outside the mesh, and lands one file per member** (novox/hq ADR 0283): the mesh
// never makes a member, so a family the mesh may make would be one nothing ever fills.
if IsFamily(name) {
if !memberRest.MatchString(strings.TrimSuffix(FamilyPrefix(name), "-")) {
problems = append(problems, fmt.Sprintf(
"%s declares the secret family %q, whose prefix is not a name", m.Module, name))
}
if own.IssuedBy != IssuedOutside {
problems = append(problems, fmt.Sprintf(
"%s declares the secret family %q without \"issued-by\": %q; the mesh never makes a "+
"member of a family, so only a party outside the mesh can fill one (novox/hq ADR 0283)",
m.Module, name, IssuedOutside))
}
if strings.Count(own.Path, "*") != 1 {
problems = append(problems, fmt.Sprintf(
"%s keeps the secret family %q at %q, which does not hold exactly one *: each member lands "+
"where the * is replaced by its name (novox/hq ADR 0283)", m.Module, name, own.Path))
}
for other := range m.OwnSecrets {
if other != name && !IsFamily(other) {
if rest := strings.TrimPrefix(other, FamilyPrefix(name)); rest != other && memberRest.MatchString(rest) {
problems = append(problems, fmt.Sprintf(
"%s declares the secret %q, which is also a member of its family %q — a name names one",
m.Module, other, name))
}
}
}
} else if strings.Contains(name, "*") {
problems = append(problems, fmt.Sprintf(
"%s declares the secret %q: a * only ends a family's name, as \"<prefix>-*\"", m.Module, name))
}
if !placedOrAbsolute(own.Path) { if !placedOrAbsolute(own.Path) {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%s needs %q at %q, which is neither an absolute path nor a placed one", m.Module, name, own.Path)) "%s needs %q at %q, which is neither an absolute path nor a placed one", m.Module, name, own.Path))
@@ -2549,6 +2580,69 @@ func (o OwnSecrets) MarshalJSON() ([]byte, error) {
return json.Marshal(entries) return json.Marshal(entries)
} }
// A secret family (novox/hq ADR 0283): an own secret declared under a name ending in FamilySuffix is one
// member per part the module's settings name — `smb-password-*` holds `smb-password-games`,
// `smb-password-library` — each given by its full name, placed at the family's path with its one `*` replaced
// by what follows the prefix. The mesh never makes a member: a family is issued outside the mesh, and a member
// not given is no file.
const FamilySuffix = "-*"
// memberRest is what may follow a family's prefix: a name's characters.
var memberRest = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`)
// IsFamily says an own secret's declared name is a family's.
func IsFamily(name string) bool { return strings.HasSuffix(name, FamilySuffix) }
// FamilyPrefix is what every member of a family begins with: the declared name without its `*`.
func FamilyPrefix(family string) string { return strings.TrimSuffix(family, "*") }
// Lookup resolves an own secret by the name it is given under: declared by that name, or a member of a family
// (the longest prefix that fits), with the family's path filled for the member. family is the family's
// declared name, or "" for a secret declared by name.
func (o OwnSecrets) Lookup(name string) (s OwnSecret, family string, ok bool) {
if s, ok := o[name]; ok && !IsFamily(name) {
return s, "", true
}
for declared, f := range o {
if !IsFamily(declared) {
continue
}
prefix := FamilyPrefix(declared)
rest := strings.TrimPrefix(name, prefix)
if !strings.HasPrefix(name, prefix) || !memberRest.MatchString(rest) {
continue
}
if family == "" || len(declared) > len(family) {
s, family, ok = OwnSecret{Path: strings.Replace(f.Path, "*", rest, 1), Taken: f.Taken, IssuedBy: f.IssuedBy}, declared, true
}
}
return s, family, ok
}
// Plain is every own secret declared by its own name: what the mesh makes when not given, and places by
// name. A family is not one, and is never made.
func (o OwnSecrets) Plain() OwnSecrets {
out := make(OwnSecrets, len(o))
for name, s := range o {
if !IsFamily(name) {
out[name] = s
}
}
return out
}
// Families is every family's declared name, sorted.
func (o OwnSecrets) Families() []string {
var out []string
for name := range o {
if IsFamily(name) {
out = append(out, name)
}
}
sort.Strings(out)
return out
}
// Paths is each own secret's path by name — the shape every placement and file walk reads. // Paths is each own secret's path by name — the shape every placement and file walk reads.
func (o OwnSecrets) Paths() map[string]string { func (o OwnSecrets) Paths() map[string]string {
out := make(map[string]string, len(o)) out := make(map[string]string, len(o))
+100
View File
@@ -0,0 +1,100 @@
package catalogue
import (
"strings"
"testing"
)
// A secret family (novox/hq ADR 0283): one own secret per part the settings name, issued outside the mesh, each
// given by its full name, never made by the mesh.
func familyManifest(t *testing.T, ownSecrets string) (Manifest, error) {
t.Helper()
return ParseManifest([]byte(`{"module":"mounts","version":"1","own-secrets":{` + ownSecrets + `}}`))
}
func TestAFamilyIsDeclaredIssuedOutsideWithOneStar(t *testing.T) {
m, err := familyManifest(t, `"smb-password-*":{"path":"/var/lib/mounts/smb-password-*.secret","issued-by":"outside"}`)
if err != nil {
t.Fatalf("a well-formed family was refused: %v", err)
}
if got := m.OwnSecrets.Families(); len(got) != 1 || got[0] != "smb-password-*" {
t.Fatalf("families: %v", got)
}
if len(m.OwnSecrets.Plain()) != 0 {
t.Fatalf("a family counted as a secret declared by name: %v", m.OwnSecrets.Plain())
}
}
func TestAMalformedFamilyIsRefused(t *testing.T) {
for name, c := range map[string]struct{ own, says string }{
"made by the mesh": {`"smb-password-*":{"path":"/s/smb-password-*.secret","taken":"at-start"}`, "issued-by"},
"no star in its path": {`"smb-password-*":{"path":"/s/smb-password.secret","issued-by":"outside"}`,
"exactly one *"},
"two stars in its path": {`"smb-password-*":{"path":"/s/*/smb-password-*.secret","issued-by":"outside"}`,
"exactly one *"},
"a star inside a name": {`"smb*password":{"path":"/s/x","issued-by":"outside"}`, "only ends a family"},
"a name that is also a member": {`"smb-password-*":{"path":"/s/p-*","issued-by":"outside"},
"smb-password-games":{"path":"/s/games","issued-by":"outside"}`, "also a member"},
} {
if _, err := familyManifest(t, c.own); err == nil || !strings.Contains(err.Error(), c.says) {
t.Errorf("%s: %v; want a refusal saying %q", name, err, c.says)
}
}
}
func TestAMemberIsFoundByItsFullNameAndLandsWhereTheStarIs(t *testing.T) {
o := OwnSecrets{
"smb-password-*": {Path: "/s/smb-password-*.secret", IssuedBy: IssuedOutside},
"smb-password-big-*": {Path: "/big/*.secret", IssuedBy: IssuedOutside},
"token": {Path: "/s/token", IssuedBy: IssuedOutside},
}
s, family, ok := o.Lookup("smb-password-games")
if !ok || family != "smb-password-*" || s.Path != "/s/smb-password-games.secret" || s.IssuedBy != IssuedOutside {
t.Fatalf("a member: %+v %q %v", s, family, ok)
}
if s, family, ok := o.Lookup("smb-password-big-one"); !ok || family != "smb-password-big-*" || s.Path != "/big/one.secret" {
t.Fatalf("the longest family that fits: %+v %q %v", s, family, ok)
}
if s, family, ok := o.Lookup("token"); !ok || family != "" || s.Path != "/s/token" {
t.Fatalf("a secret declared by name: %+v %q %v", s, family, ok)
}
for _, name := range []string{"smb-password-*", "smb-password-", "smb-password-../etc", "smb-password-a/b",
"smb-password-a.b", "smb-password-Games", "smb-password--x", "other"} {
if _, _, ok := o.Lookup(name); ok {
t.Errorf("%q was found as a secret", name)
}
}
}
// A member given is one file at its path; one not given is nothing, and the machine still composes — the mesh never
// makes a member.
func TestAMemberGivenIsPlacedAndOneNotGivenIsNothing(t *testing.T) {
m, err := familyManifest(t, `"smb-password-*":{"path":"/var/lib/mounts/smb-password-*.secret","issued-by":"outside"}`)
if err != nil {
t.Fatal(err)
}
r := Resolution{Node: "workstation", Modules: []Manifest{m}}
out, err := r.Declaration(Rendering{Needed: map[string]map[string]string{"mounts": {"smb-password-games": "sealed"}}})
if err != nil {
t.Fatalf("a module with one member given did not compose: %v", err)
}
var paths []string
for _, res := range out {
if res["sealed"] != nil {
paths = append(paths, res["path"].(string))
}
}
if len(paths) != 1 || paths[0] != "/var/lib/mounts/smb-password-games.secret" {
t.Fatalf("placed: %v", paths)
}
out, err = r.Declaration(Rendering{})
if err != nil {
t.Fatalf("a module with no member given did not compose: %v", err)
}
for _, res := range out {
if res["sealed"] != nil {
t.Fatalf("a member nobody gave was placed: %v", res)
}
}
}
+1 -1
View File
@@ -56,7 +56,7 @@ func secretsUsed(content string) []string {
// name would end that, to save writing a file. // name would end that, to save writing a file.
func sealedFor(m Manifest, needs []Needed, with Rendering) (map[string]string, error) { func sealedFor(m Manifest, needs []Needed, with Rendering) (map[string]string, error) {
sealed := map[string]string{} sealed := map[string]string{}
for name := range m.OwnSecrets { for name := range m.OwnSecrets.Plain() {
if value := with.Needed[m.Module][name]; value != "" { if value := with.Needed[m.Module][name]; value != "" {
sealed[name] = value sealed[name] = value
} }
+7 -5
View File
@@ -164,8 +164,9 @@ var ControllerVerbs = []Verb{
Input: schema(map[string]string{"plan": "the walk's id", "why": "why", "by": "who stopped the delivery"}, Input: schema(map[string]string{"plan": "the walk's id", "why": "why", "by": "who stopped the delivery"},
[]string{"plan", "why"})}, []string{"plan", "why"})},
{Name: "delivery-walks", Description: "The walks the controller keeps (novox/hq ADR 0239): every open one and the " + {Name: "delivery-walks", Description: "The walks the controller keeps (novox/hq ADR 0239): every open one and the " +
"last ended ones, each whole — its tiers, each module's state, first machines and gate — and whether the " + "last ended ones, each whole — its tiers, each module's state, first machines and first-node gate with its readings, and its " +
"delivery seat has a holder on record. Given a plan, that one.", "phases from its merge to every machine running it (novox/hq ADR 0282) — and whether the delivery seat has a " +
"holder on record. Given a plan, that one.",
Input: schema(map[string]string{"plan": "one walk's id", "limit": "how many ended walks beside the open ones (default 50)"}, Input: schema(map[string]string{"plan": "one walk's id", "limit": "how many ended walks beside the open ones (default 50)"},
nil)}, nil)},
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " + {Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " +
@@ -375,10 +376,11 @@ var ControllerVerbs = []Verb{
"recorded — who, why and the cause of each, and which causes repeat: each repeat is a healer the mesh lacks.", "recorded — who, why and the cause of each, and which causes repeat: each repeat is a healer the mesh lacks.",
Input: schema(map[string]string{"days": "how many days back (default 14)"}, nil)}, Input: schema(map[string]string{"days": "how many days back (default 14)"}, nil)},
{Name: "durations", Description: "How long things take, as the controller measured them: a send to its machine's " + {Name: "durations", Description: "How long things take, as the controller measured them: a send to its machine's " +
"report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build — per machine, " + "report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build, and each phase of an " +
"repository or module, with median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).", "ended walk per class (walk-phase, novox/hq ADR 0282) — per machine, repository, module or class/phase, with " +
"median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).",
Input: schema(map[string]string{ Input: schema(map[string]string{
"kind": "one kind: apply, heartbeat-gap, plan-tier or build; every kind when absent", "kind": "one kind: apply, heartbeat-gap, plan-tier, build or walk-phase; every kind when absent",
"days": "how many days back (default 14)", "days": "how many days back (default 14)",
}, nil)}, }, nil)},
// What is wrong, and the self-check (novox/hq to-be 45 §2, §4). // What is wrong, and the self-check (novox/hq to-be 45 §2, §4).
-9
View File
@@ -108,15 +108,6 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
for _, p := range people { for _, p := range people {
out.People[p.Name] = p.Invokes out.People[p.Name] = p.Invokes
} }
// The view is composed once its credential is minted and until it is forgotten: its row is the
// record of it, nothing else being declared about it (broker.KindView).
kept, err := i.BusUsers(ctx)
if err != nil {
return broker.Records{}, err
}
if u, minted := kept[broker.ViewUser]; minted && u.Kind == BusView {
out.View = true
}
return out, nil return out, nil
} }
-3
View File
@@ -45,9 +45,6 @@ const (
// BusNodeTools is a machine's tool runtime (novox/hq ADR 0175): named like the module it // BusNodeTools is a machine's tool runtime (novox/hq ADR 0175): named like the module it
// stands for, recorded as what it is. // stands for, recorded as what it is.
BusNodeTools = "node-tools" BusNodeTools = "node-tools"
// BusView is the one read-only view onto the bus (broker.KindView): its row is the whole record of
// it, minted by `bus view-credential` and forgotten by `bus view-revoke`.
BusView = "view"
) )
// MintBusPassword makes a bus password and records its hash under a username, replacing whatever was // MintBusPassword makes a bus password and records its hash under a username, replacing whatever was
+18 -8
View File
@@ -94,7 +94,9 @@ type DataChange struct {
} }
// readingEvery is how often a measurement is kept as a reading: the shrink is read over days, and a // readingEvery is how often a measurement is kept as a reading: the shrink is read over days, and a
// row every five minutes would say the same thing sixty times an hour. // row every five minutes would say the same thing sixty times an hour. A reading keeps the item's path
// as it stands after the measurement — the holder's, or the last one known when it named none — and an
// item at a path with no recent reading is read at once (novox/hq issue 368).
const readingEvery = 55 * time.Minute const readingEvery = 55 * time.Minute
// readingsKept is how long readings are kept. // readingsKept is how long readings are kept.
@@ -187,10 +189,14 @@ func (i *Inventory) RecordData(ctx context.Context, machine string, declared []D
} }
if m.Size != nil && m.MeasuredAt != nil && Comparable(m.Precision) { if m.Size != nil && m.MeasuredAt != nil && Comparable(m.Precision) {
if _, err := tx.Exec(ctx, ` if _, err := tx.Exec(ctx, `
insert into data_reading (machine, module, item, at, size_bytes, last_write) insert into data_reading (machine, module, item, at, size_bytes, last_write, path)
select $1, $2, $3, $4, $5, $6 select $1, $2, $3, $4, $5, $6, d.path
where not exists (select 1 from data_reading from data_item d
where machine = $1 and module = $2 and item = $3 and at > $4::timestamptz - $7::interval)`, where d.machine = $1 and d.module = $2 and d.item = $3
and not exists (select 1 from data_reading
where machine = $1 and module = $2 and item = $3 and path = d.path
and at > $4::timestamptz - $7::interval)
on conflict (machine, module, item, at) do nothing`,
machine, d.Module, d.Item, *m.MeasuredAt, *m.Size, m.LastWrite, machine, d.Module, d.Item, *m.MeasuredAt, *m.Size, m.LastWrite,
fmt.Sprintf("%d seconds", int(readingEvery.Seconds()))); err != nil { fmt.Sprintf("%d seconds", int(readingEvery.Seconds()))); err != nil {
return change, err return change, err
@@ -281,10 +287,14 @@ func (i *Inventory) DataOf(ctx context.Context, machine, module, item string) (D
return r, err return r, err
} }
// DataPeaks is each item's largest reading since a moment, keyed by DataRecord.Key. // DataPeaks is each item's largest reading since a moment, keyed by DataRecord.Key — read only from
// readings at the item's path now (novox/hq issue 368): a directory is never compared with another one
// that once held the same item, so a moved item starts its size history again at its new path.
func (i *Inventory) DataPeaks(ctx context.Context, since time.Time) (map[string]int64, error) { func (i *Inventory) DataPeaks(ctx context.Context, since time.Time) (map[string]int64, error) {
rows, err := i.store.Pool().Query(ctx, `select machine, module, item, max(size_bytes) from data_reading rows, err := i.store.Pool().Query(ctx, `select r.machine, r.module, r.item, max(r.size_bytes)
where at >= $1 group by machine, module, item`, since) from data_reading r
join data_item d on d.machine = r.machine and d.module = r.module and d.item = r.item and d.path = r.path
where r.at >= $1 group by r.machine, r.module, r.item`, since)
if err != nil { if err != nil {
return nil, err return nil, err
} }
+63
View File
@@ -128,3 +128,66 @@ func TestAPartialMeasurementIsNeverAReading(t *testing.T) {
t.Fatalf("%+v, %v", r, err) t.Fatalf("%+v, %v", r, err)
} }
} }
// A shrink is read against what the same directory held (novox/hq issue 368): an item whose path moved
// — an agent's home moved to its own account — starts its size history again at the new path, and a
// genuine shrink at the new path is still read against what that path held.
func TestThePeakIsReadAtTheItemsPathOnly(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
now := time.Date(2026, 10, 10, 0, 0, 0, 0, time.UTC)
declared := []DeclaredData{{Module: "claude-code", Item: "agent-home", Class: "valuable", Owned: true}}
measure := func(when time.Time, path string, s int64) {
t.Helper()
if _, err := inv.RecordData(ctx, "novox", declared, map[string]map[string]Measurement{"claude-code": {
"agent-home": {Path: path, Size: size(s), MeasuredAt: at(when)}}}, "", when); err != nil {
t.Fatal(err)
}
}
measure(now, "/home/operator/.claude", 94_700_000)
// The path moves ten minutes later: the new directory is measured at once, not an hour on.
measure(now.Add(10*time.Minute), "/home/agent/.claude", 490)
peaks, err := inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 490 {
t.Fatalf("after the path moved the peak is %v (%v), want 490: the old directory's size is no shrink "+
"of the new one", peaks, err)
}
// The new directory grows, then genuinely loses most of it: that is read against the new path's peak.
measure(now.Add(2*time.Hour), "/home/agent/.claude", 80_000_000)
measure(now.Add(4*time.Hour), "/home/agent/.claude", 1_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 80_000_000 {
t.Fatalf("a shrink at the new path is read against %v (%v), want 80000000", peaks, err)
}
// A measurement that names no path is the item's last known path's, not a new history.
measure(now.Add(6*time.Hour), "", 2_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 80_000_000 {
t.Fatalf("a measurement with no path started a new history: peak %v (%v)", peaks, err)
}
// Moving back to a directory measured before reads it against what it held then.
measure(now.Add(8*time.Hour), "/home/operator/.claude", 94_000_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 94_700_000 {
t.Fatalf("back at the first path the peak is %v (%v), want 94700000", peaks, err)
}
}
// Two measurements at the same moment at two paths keep one reading and fail nothing: the second
// would otherwise break the reading's key and lose the machine's whole record (issue 368 review).
func TestTwoPathsAtOneMomentFailNothing(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
now := time.Date(2026, 10, 10, 0, 0, 0, 0, time.UTC)
declared := []DeclaredData{{Module: "claude-code", Item: "agent-home", Class: "valuable", Owned: true}}
for _, path := range []string{"/home/operator/.claude", "/home/agent/.claude"} {
if _, err := inv.RecordData(ctx, "novox", declared, map[string]map[string]Measurement{"claude-code": {
"agent-home": {Path: path, Size: size(100), MeasuredAt: at(now)}}}, "", now); err != nil {
t.Fatalf("a measurement at %s failed: %v", path, err)
}
}
r, err := inv.DataOf(ctx, "novox", "claude-code", "agent-home")
if err != nil || r.Path != "/home/agent/.claude" {
t.Fatalf("%+v, %v", r, err)
}
}
+12 -1
View File
@@ -17,10 +17,13 @@ const (
DurationHeartbeatGap = "heartbeat-gap" DurationHeartbeatGap = "heartbeat-gap"
DurationPlanTier = "plan-tier" DurationPlanTier = "plan-tier"
DurationBuild = "build" DurationBuild = "build"
// DurationWalkPhase is one phase of an ended walk, per class (novox/hq ADR 0282 decision 6): subject
// "<class>/<phase>", and "<class>/total" for its merge to every machine running it.
DurationWalkPhase = "walk-phase"
) )
// DurationKinds are every kind, in the order `durations` shows them. // DurationKinds are every kind, in the order `durations` shows them.
var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild} var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild, DurationWalkPhase}
// DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it // DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it
// in Phase 1's first live week. // in Phase 1's first live week.
@@ -100,6 +103,14 @@ func (i *Inventory) Durations(ctx context.Context, kind string, since time.Time)
return out, rows.Err() return out, rows.Err()
} }
// WalkPhasesKept says whether any phase of a walk is kept as a walk-phase duration, and whether its total is.
func (i *Inventory) WalkPhasesKept(ctx context.Context, plan string) (anyKept, totalKept bool, err error) {
err = i.store.Pool().QueryRow(ctx,
`select count(*) > 0, count(*) filter (where ref = $2) > 0 from duration where kind = $1 and ref like $3`,
DurationWalkPhase, plan+"/total", plan+"/%").Scan(&anyKept, &totalKept)
return anyKept, totalKept, err
}
// ForgetOldDurations removes what is older than DurationsKeptFor, and says how many. // ForgetOldDurations removes what is older than DurationsKeptFor, and says how many.
func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) { func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) {
tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`, tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`,
+1 -1
View File
@@ -155,7 +155,7 @@ func (i *Inventory) ReplaceGivenAfterStart(ctx context.Context, node, declared s
} }
// The definition is asked again now, not only when the value was given: one that has since // The definition is asked again now, not only when the value was given: one that has since
// said the value is an outside party's, or applied, keeps it as given, and the mark stays gone. // said the value is an outside party's, or applied, keeps it as given, and the mark stays gone.
if own, ok := m.OwnSecrets[d.name]; !ok || !own.MeshMayMake() { if own, _, ok := m.OwnSecrets.Lookup(d.name); !ok || !own.MeshMayMake() {
continue continue
} }
if err := i.remakeOwn(ctx, d.nodeID, key, m, d.module, d.name); err != nil { if err := i.remakeOwn(ctx, d.nodeID, key, m, d.module, d.name); err != nil {
+11
View File
@@ -34,6 +34,17 @@ type ResourceHealth struct {
// Root is "never" on an account verdict that judged whether the account can become root without a // Root is "never" on an account verdict that judged whether the account can become root without a
// person (novox/hq ADR 0266). // person (novox/hq ADR 0266).
Root string `json:"root,omitempty"` Root string `json:"root,omitempty"`
// Waits is what a waiting resource waits for the operator to give (novox/hq ADR 0283).
Waits []Wait `json:"waits,omitempty"`
}
// Wait is one thing a waiting resource waits for the operator to give: exactly one of Secret and Setting
// (novox/hq ADR 0283), as link.Wait carries it.
type Wait struct {
Part string `json:"part"`
Secret string `json:"secret,omitempty"`
Setting string `json:"setting,omitempty"`
What string `json:"what"`
} }
// NodeHealth is a machine's newest statement, as kept. // NodeHealth is a machine's newest statement, as kept.
@@ -0,0 +1,21 @@
-- A reading says the path it was measured at (novox/hq issue 368).
--
-- A shrink is read against the largest reading of the last seven days. Readings were kept by machine,
-- module and item only, so when an item's path moved — on 2026-10-10 an agent's home moved from the
-- operator's own home to the agent account's (ADR 0266) — the new, fresh directory was read against
-- the old one's size, and `data-shrank` was raised for data that was never lost. A reading now keeps
-- its path, and the peak is read only from readings at the item's path now: a moved item starts its
-- size history again, and moving back to a path reads it against what that path held.
--
-- **Every reading kept so far is taken to be at its item's path now.** Where it was really measured
-- is not on record; taking the path now keeps every item's history, so a shrink that is real stays
-- raised. An item whose path moved before this migration stays compared with its old directory until
-- those readings leave the seven-day window. Readings of an item no longer kept keep no path, and are
-- read for nothing.
--
-- Numbered 0092, past 0091, the highest on main or any open branch when this was written.
alter table data_reading add column path text;
update data_reading r set path = d.path
from data_item d
where d.machine = r.machine and d.module = r.module and d.item = r.item;
@@ -0,0 +1,10 @@
-- A walk keeps its phases (novox/hq ADR 0282 decision 6, issue 382).
--
-- A small fix took 45 minutes to reach a node on 2026-10-10, and where the time went was read back from the
-- walks by hand: the walk kept when its modules were asked, built, sent first, judged and sent to the rest,
-- but not when its batch's window closed or when it was cut, so the window and the wait behind another walk
-- could not be told apart. A walk now keeps those moments and its class (core or leaf, read at its cut).
-- Its readings, each module's send to the rest and the times of the merges it answers are kept in the plan's
-- own records (modules, delivery). Null for a walk kept before this: its phases before the build are said
-- unknown.
alter table release_plan add column times jsonb;
+73 -8
View File
@@ -54,6 +54,25 @@ type Plan struct {
// walk can carry several. Repository and Commit above keep one of them, for a reader that knows one. Empty // walk can carry several. Repository and Commit above keep one of them, for a reader that knows one. Empty
// for a walk kept before it was: Repository and Commit are then the whole of it. // for a walk kept before it was: Repository and Commit are then the whole of it.
Commits []PlanCommit `json:"commits,omitempty"` Commits []PlanCommit `json:"commits,omitempty"`
// Times are the moments of a walk no other field keeps (novox/hq ADR 0282 decision 6): when its batch's
// window closed, when it was cut, and its class. Nil for a walk kept before they were, whose phases before
// its build are said unknown.
Times *PlanTimes `json:"times,omitempty"`
// Phases is the walk's time from its merge, phase by phase (ADR 0282 decision 6): never kept, worked out
// from the walk's own moments by whoever says the walk (`delivery walks`, `plan-moved`).
Phases *WalkPhases `json:"phases,omitempty"`
}
// PlanTimes are a walk's own moments beside its modules' (novox/hq ADR 0282 decision 6).
type PlanTimes struct {
// WindowClosed is when its batch's merge window closed: no merge for the window's length, or its maximum.
// Nil for a walk no window assembled (one merge walked alone).
WindowClosed *time.Time `json:"window_closed,omitempty"`
// Cut is when the batch became the walk.
Cut *time.Time `json:"cut,omitempty"`
// Class is the walk's module class, read at its cut: core when it moves a module on the controller's own
// path or the mesh's resolver, leaf otherwise (ADR 0282 decision 1).
Class string `json:"class,omitempty"`
} }
// PlanCommit is one repository's commit a walk carries: the latest merge of its branch in the batch. // PlanCommit is one repository's commit a walk carries: the latest merge of its branch in the batch.
@@ -76,6 +95,10 @@ type PlanMerge struct {
Number int `json:"number,omitempty"` Number int `json:"number,omitempty"`
Title string `json:"title,omitempty"` Title string `json:"title,omitempty"`
Moves []string `json:"moves,omitempty"` Moves []string `json:"moves,omitempty"`
// Merged is when the forge made the merge and Heard when the controller heard it (novox/hq ADR 0282): where
// its delivery time starts. Zero for a merge named before they were kept.
Merged time.Time `json:"merged,omitzero"`
Heard time.Time `json:"heard,omitzero"`
} }
// PlanBatch is a batch's window while it is one (novox/hq ADR 0276): when it closes unless another merge // PlanBatch is a batch's window while it is one (novox/hq ADR 0276): when it closes unless another merge
@@ -220,6 +243,9 @@ type PlanModule struct {
// went there in one send (novox/hq issue 281), and one gate judges what one send moved. Empty for // went there in one send (novox/hq issue 281), and one gate judges what one send moved. Empty for
// the module the gate is kept on, and for a plan from before tiers were sent whole. // the module the gate is kept on, and for a plan from before tiers were sent whole.
GatedBy string `json:"gated_by,omitempty"` GatedBy string `json:"gated_by,omitempty"`
// Rest is, per machine of the rest, the declaration the send after the first-node gate carried there (novox/hq ADR
// 0282 decision 6): the machine's first report of it, applied, is when the build runs there.
Rest map[string]SentDeclaration `json:"rest,omitempty"`
} }
// PlanGate is one module's rollout record at its gate (to-be 45 §8): the component, the first machine, // PlanGate is one module's rollout record at its gate (to-be 45 §8): the component, the first machine,
@@ -276,6 +302,35 @@ type PlanGate struct {
BrokenWhy string `json:"broken_why,omitempty"` BrokenWhy string `json:"broken_why,omitempty"`
// Returned names the broken modules already put back, at once, while the rest of the send is judged. // Returned names the broken modules already put back, at once, while the rest of the send is judged.
Returned []string `json:"returned,omitempty"` Returned []string `json:"returned,omitempty"`
// Readings are the judging's readings with their times (novox/hq ADR 0282 decision 6): every reading that
// counted a pass, and the first that did not after one that did. At most maxReadings, the newest kept.
Readings []GateReading `json:"readings,omitempty"`
}
// GateReading is one reading of a first-node gate.
type GateReading struct {
At time.Time `json:"at"`
Healthy bool `json:"healthy"`
// Said is what a reading that did not pass found wanting.
Said string `json:"said,omitempty"`
}
// maxReadings bounds a first-node gate's readings: a judging that never passes reads every few seconds for ten minutes.
const maxReadings = 24
// Read keeps one reading: a pass always, and a reading that did not pass only when the one before passed or
// there is none, so a judging waiting for a module to start keeps one line of it, not hundreds.
func (g *PlanGate) Read(at time.Time, healthy bool, said string) {
if !healthy && len(g.Readings) > 0 && !g.Readings[len(g.Readings)-1].Healthy {
return
}
if r := []rune(said); len(r) > 200 {
said = string(r[:200])
}
g.Readings = append(g.Readings, GateReading{At: at, Healthy: healthy, Said: said})
if len(g.Readings) > maxReadings {
g.Readings = g.Readings[len(g.Readings)-maxReadings:]
}
} }
// CarriedMove is one module's build moving on a machine with a gated send. // CarriedMove is one module's build moving on a machine with a gated send.
@@ -342,7 +397,12 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error {
if err != nil { if err != nil {
return err return err
} }
var release, delivery, commits []byte var release, delivery, commits, times []byte
if p.Times != nil {
if times, err = json.Marshal(p.Times); err != nil {
return err
}
}
if len(p.Commits) > 0 { if len(p.Commits) > 0 {
if commits, err = json.Marshal(p.Commits); err != nil { if commits, err = json.Marshal(p.Commits); err != nil {
return err return err
@@ -373,18 +433,18 @@ func (i *Inventory) SavePlan(ctx context.Context, p *Plan) error {
var revision int64 var revision int64
err = tx.QueryRow(ctx, err = tx.QueryRow(ctx,
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note,
branch, tier_entered, revision, epoch, release, delivery, merged_at, commits) branch, tier_entered, revision, epoch, release, delivery, merged_at, commits, times)
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13, $14, $15, $16, $17) values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11, 1, $13, $14, $15, $16, $17, $18)
on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier,
tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch, tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch,
tier_entered = excluded.tier_entered, revision = release_plan.revision + 1, epoch = excluded.epoch, tier_entered = excluded.tier_entered, revision = release_plan.revision + 1, epoch = excluded.epoch,
release = excluded.release, delivery = excluded.delivery, repository = excluded.repository, release = excluded.release, delivery = excluded.delivery, repository = excluded.repository,
commit_hash = excluded.commit_hash, merged_at = excluded.merged_at, commits = excluded.commits, commit_hash = excluded.commit_hash, merged_at = excluded.merged_at, commits = excluded.commits,
created = excluded.created created = excluded.created, times = excluded.times
where release_plan.revision = $12 where release_plan.revision = $12
returning revision`, returning revision`,
p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered, p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered,
p.Revision, epoch, release, delivery, mergedAt(p.Merged), commits).Scan(&revision) p.Revision, epoch, release, delivery, mergedAt(p.Merged), commits, times).Scan(&revision)
if errors.Is(err, pgx.ErrNoRows) { if errors.Is(err, pgx.ErrNoRows) {
// The row is there and at another revision — moved since this was read, or there already // The row is there and at another revision — moved since this was read, or there already
// when this one is new: either way not this writer's to overwrite. (A plan saved before plans // when this one is new: either way not this writer's to overwrite. (A plan saved before plans
@@ -467,7 +527,7 @@ func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) {
func (i *Inventory) plans(ctx context.Context, tail string, args ...any) ([]Plan, error) { func (i *Inventory) plans(ctx context.Context, tail string, args ...any) ([]Plan, error) {
rows, err := i.store.Pool().Query(ctx, rows, err := i.store.Pool().Query(ctx,
`select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch, `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch,
coalesce(tier_entered, created), revision, coalesce(epoch, 0), release, delivery, merged_at, commits coalesce(tier_entered, created), revision, coalesce(epoch, 0), release, delivery, merged_at, commits, times
from release_plan `+tail, args...) from release_plan `+tail, args...)
if err != nil { if err != nil {
return nil, err return nil, err
@@ -476,14 +536,19 @@ func (i *Inventory) plans(ctx context.Context, tail string, args ...any) ([]Plan
var out []Plan var out []Plan
for rows.Next() { for rows.Next() {
var p Plan var p Plan
var tiers, modules, release, delivery, commits []byte var tiers, modules, release, delivery, commits, times []byte
var epoch int64 var epoch int64
var merged *time.Time var merged *time.Time
if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State, if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State,
&p.Tier, &tiers, &modules, &p.Note, &p.Branch, &p.TierEntered, &p.Revision, &epoch, &release, &p.Tier, &tiers, &modules, &p.Note, &p.Branch, &p.TierEntered, &p.Revision, &epoch, &release,
&delivery, &merged, &commits); err != nil { &delivery, &merged, &commits, &times); err != nil {
return nil, err return nil, err
} }
if len(times) > 0 {
if err := json.Unmarshal(times, &p.Times); err != nil {
return nil, err
}
}
if len(commits) > 0 { if len(commits) > 0 {
if err := json.Unmarshal(commits, &p.Commits); err != nil { if err := json.Unmarshal(commits, &p.Commits); err != nil {
return nil, err return nil, err
+76
View File
@@ -0,0 +1,76 @@
package inventory
import (
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// A member of a secret family is given by its full name, at the desk too, and a name outside the family is refused
// (novox/hq ADR 0283).
func TestAMemberIsGivableAtTheDeskAndANameOutsideTheFamilyIsNot(t *testing.T) {
m := catalogue.Manifest{Module: "mounts", OwnSecrets: catalogue.OwnSecrets{
"smb-password-*": {Path: "/s/smb-password-*.secret", IssuedBy: catalogue.IssuedOutside},
}}
if err := GivableAtDesk(m, "smb-password-games"); err != nil {
t.Fatalf("a member was refused: %v", err)
}
for _, name := range []string{"smb-password-*", "smb-password-", "smb-password-a/b", "smb-password-A", "smb-credentials"} {
err := GivableAtDesk(m, name)
if err == nil {
t.Errorf("%q was givable", name)
} else if !strings.Contains(err.Error(), "smb-password-<name>") {
t.Errorf("%q: the refusal does not say the family's form: %v", name, err)
}
}
}
// A member given is read back as given, sealed to the machine's key; the mesh makes none, so nothing else is.
func TestAMemberGivenIsReadAsGivenAndNothingElseIs(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
node, err := inv.AddNode(ctx, "workstation")
if err != nil {
t.Fatal(err)
}
key, _ := aSealingKey(t)
if err := inv.RecordSealingKey(ctx, node.ID, key); err != nil {
t.Fatal(err)
}
m := catalogue.Manifest{Module: "mounts", Version: "1", OwnSecrets: catalogue.OwnSecrets{
"smb-password-*": {Path: "/s/smb-password-*.secret", IssuedBy: catalogue.IssuedOutside},
"token": {Path: "/s/token", IssuedBy: catalogue.IssuedOutside},
}}
if err := inv.RegisterModule(ctx, m, Source{}); err != nil {
t.Fatal(err)
}
if err := inv.AcceptSecretForModule(ctx, "workstation", "mounts", "smb-password-games", "hunter2"); err != nil {
t.Fatalf("a member was refused: %v", err)
}
if err := inv.AcceptSecretForModule(ctx, "workstation", "mounts", "smb-pass", "x"); err == nil {
t.Fatal("a name of no family was accepted")
}
members, err := inv.GivenMembers(ctx, "workstation", "mounts", "smb-password-*")
if err != nil {
t.Fatal(err)
}
if len(members) != 1 || members[0].Name != "smb-password-games" || !members[0].Current || members[0].Sealed == "" ||
strings.Contains(members[0].Sealed, "hunter2") {
t.Fatalf("members: %+v", members)
}
given, err := inv.GivenOwnSecrets(ctx, "workstation", "mounts")
if err != nil {
t.Fatal(err)
}
if _, was := given["smb-password-games"]; !was || len(given) != 1 {
t.Fatalf("given: %v", given)
}
// A value the mesh made is not one a person gave.
if _, err := inv.SecretForModule(ctx, "workstation", "mounts", "token"); err != nil {
t.Fatal(err)
}
if given, _ := inv.GivenOwnSecrets(ctx, "workstation", "mounts"); len(given) != 1 {
t.Fatalf("a value the mesh made counted as given: %v", given)
}
}
+77 -4
View File
@@ -7,6 +7,7 @@ import (
"slices" "slices"
"sort" "sort"
"strings" "strings"
"time"
"github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5"
@@ -446,7 +447,7 @@ func (i *Inventory) acceptOwn(ctx context.Context, node, module, name, value str
if err != nil { if err != nil {
return false, err return false, err
} }
own, declared := m.OwnSecrets[name] own, _, declared := m.OwnSecrets.Lookup(name)
if !declared { if !declared {
return false, fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m)) return false, fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m))
} }
@@ -583,7 +584,7 @@ const BrokerSecret = "broker"
// in place of a value given (catalogue.OwnSecret.MeshMayMake). The desk takes only what a person holds and // in place of a value given (catalogue.OwnSecret.MeshMayMake). The desk takes only what a person holds and
// the mesh cannot make — a bot's token — so nobody is asked to type the mesh's own credential into a prompt. // the mesh cannot make — a bot's token — so nobody is asked to type the mesh's own credential into a prompt.
func GivableAtDesk(m catalogue.Manifest, name string) error { func GivableAtDesk(m catalogue.Manifest, name string) error {
own, ok := m.OwnSecrets[name] own, _, ok := m.OwnSecrets.Lookup(name)
if !ok { if !ok {
return fmt.Errorf("%s does not declare %q as an own secret; %s", m.Module, name, declaresOwn(m)) return fmt.Errorf("%s does not declare %q as an own secret; %s", m.Module, name, declaresOwn(m))
} }
@@ -602,7 +603,79 @@ func declaresOwn(m catalogue.Manifest) string {
if len(m.OwnSecrets) == 0 { if len(m.OwnSecrets) == 0 {
return "it declares no own secrets" return "it declares no own secrets"
} }
return "it declares: " + strings.Join(sortedNames(m.OwnSecrets.Paths()), ", ") names := sortedNames(m.OwnSecrets.Plain().Paths())
// A family is said as its members are given (novox/hq ADR 0283): one per part, by its full name.
for _, f := range m.OwnSecrets.Families() {
names = append(names, catalogue.FamilyPrefix(f)+"<name> (one per part, issued outside the mesh)")
}
return "it declares: " + strings.Join(names, ", ")
}
// GivenMember is one member of a secret family given on a machine (novox/hq ADR 0283): its full name, its value
// sealed to the machine, and whether it is sealed to the key the machine holds now.
type GivenMember struct {
Name string
Sealed string
Current bool
}
// GivenMembers is every member of a module's secret family given on a machine, by name. The mesh never makes a
// member, so only a value a person gave is one; a machine with no sealing key holds none.
func (i *Inventory) GivenMembers(ctx context.Context, node, module, family string) ([]GivenMember, error) {
key, err := i.SealingKeyOf(ctx, node)
if err != nil || key == "" {
return nil, err
}
record, err := i.NodeByName(ctx, node)
if err != nil {
return nil, err
}
prefix := catalogue.FamilyPrefix(family)
rows, err := i.store.Pool().Query(ctx,
`select name, sealed, node_key from module_secret
where node = $1 and module = $2 and origin = 'accepted' and left(name, length($3)) = $3
order by name`, record.ID, module, prefix)
if err != nil {
return nil, err
}
defer rows.Close()
var out []GivenMember
for rows.Next() {
var g GivenMember
var against string
if err := rows.Scan(&g.Name, &g.Sealed, &against); err != nil {
return nil, err
}
g.Current = against == key
out = append(out, g)
}
return out, rows.Err()
}
// GivenOwnSecrets is every own secret of a module a person gave on a machine, by name, with when (novox/hq ADR
// 0283): what the controller checks a module's wait for a secret against. A value the mesh made is not one.
func (i *Inventory) GivenOwnSecrets(ctx context.Context, node, module string) (map[string]time.Time, error) {
record, err := i.NodeByName(ctx, node)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select name, coalesce(made_at, now()) from module_secret where node = $1 and module = $2 and origin = 'accepted'`,
record.ID, module)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]time.Time{}
for rows.Next() {
var name string
var at time.Time
if err := rows.Scan(&name, &at); err != nil {
return nil, err
}
out[name] = at
}
return out, rows.Err()
} }
func sortedNames(of map[string]string) []string { func sortedNames(of map[string]string) []string {
@@ -645,7 +718,7 @@ func (i *Inventory) RotateModuleSecret(ctx context.Context, node, module, name s
if err != nil { if err != nil {
return err return err
} }
own, declared := m.OwnSecrets[name] own, _, declared := m.OwnSecrets.Lookup(name)
if !declared { if !declared {
return fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m)) return fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m))
} }
+478
View File
@@ -0,0 +1,478 @@
package inventory
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5"
)
// Where a walk's time went (novox/hq ADR 0282 decision 6, issue 382): from its merge to every machine of its
// rest running its builds, phase by phase, worked out from the walk's own moments. Measurement only: nothing
// here decides anything about the walk.
//
// **A phase that cannot be measured is said unknown, never zero.** Where a moment is missing (a walk kept
// before it was recorded, a machine not yet reported), the phases around it are unknown, and the span between
// the moments on either side is counted as unknown time: the measured phases and the unknown time always add
// up to the total. A phase that did not happen (no window for a merge walked alone, no first machine for a
// module nobody runs) is said none, with no time.
// The phases, in the order a walk passes them (ADR 0282's table).
const (
PhaseWindow = "window" // the merge to its batch's window closing
PhaseQueued = "queued" // the window closed to the walk being cut: waiting behind another walk
PhaseWord = "word" // the cut to the delivery's word, for a walk that waits for one
PhaseBetween = "between-tiers" // one tier's end to the next tier's ask
PhaseBuild = "build" // a tier asked to its last module built
PhaseSend = "send-first" // built to sent to the first machines
PhaseJudge = "judgement" // sent first to the first-node gate's last verdict: its readings
PhaseRest = "send-rest" // judged to sent to the rest
PhaseApply = "apply" // the last send to every machine of the rest reporting the build applied
PhaseOpen = "open" // a walk still running: since its last moment
PhaseMeasured = "measured"
PhaseUnknown = "unknown"
PhaseNone = "none"
)
// The classes of a walk (ADR 0282 decision 1) and their delivery budgets.
const (
ClassCore = "core"
ClassLeaf = "leaf"
)
// ApplySilentAfter is how long after a send to the rest a machine that has said nothing is left out of the
// walk's end, as a machine not heard from (ADR 0282 decision 1: a sleeping laptop is listed, not waited for).
const ApplySilentAfter = 15 * time.Minute
// WalkPhases is a walk's time, phase by phase.
type WalkPhases struct {
Class string `json:"class,omitempty"`
// From is the walk's earliest merge; End when its last phase ended: every machine of its rest reported the
// build applied. Nil End while that is not known.
From *time.Time `json:"from,omitempty"`
End *time.Time `json:"end,omitempty"`
// TotalMS is End − From; zero while either is unknown, Total its words.
TotalMS int64 `json:"total_ms,omitempty"`
Total string `json:"total,omitempty"`
// UnknownMS is the time within the walk no phase could be measured over.
UnknownMS int64 `json:"unknown_ms,omitempty"`
Phases []WalkPhase `json:"phases"`
Silent []string `json:"silent,omitempty"`
Said string `json:"said"`
Merges []MergeStart `json:"merges,omitempty"`
}
// MergeStart is one merge the walk answers and when it was made: where that merge's delivery time starts.
type MergeStart struct {
Repository string `json:"repository"`
Commit string `json:"commit"`
Merged time.Time `json:"merged"`
}
// WalkPhase is one phase of a walk.
type WalkPhase struct {
Name string `json:"name"`
// Tier is the tier a tier's phase belongs to; -1 for a phase of the walk.
Tier int `json:"tier"`
State string `json:"state"`
Start *time.Time `json:"start,omitempty"`
End *time.Time `json:"end,omitempty"`
TookMS int64 `json:"took_ms,omitempty"`
Took string `json:"took,omitempty"`
Said string `json:"said,omitempty"`
}
// AppliedReport is a machine's first report, after a send, that it applied what it was sent.
type AppliedReport struct {
At time.Time
Outcome string
}
// AppliedLookup answers a machine's first report after a send to it; false when it has not reported.
type AppliedLookup func(node string, sent time.Time) (AppliedReport, bool)
// point is one moment of the walk, ending the phase named.
type point struct {
phase string
tier int
at *time.Time
none bool // the phase did not happen
said string // why it is none, or unknown
}
// Phases is the walk's time phase by phase, at now; applied answers the rest's reports (nil: none read).
func (p Plan) WalkPhases(now time.Time, applied AppliedLookup) *WalkPhases {
if p.Release != nil || p.Batch() {
return nil
}
w := &WalkPhases{}
if p.Times != nil {
w.Class = p.Times.Class
}
from := p.firstMerge(w)
if from != nil {
f := from.Truncate(time.Millisecond)
from = &f
}
if from == nil {
w.Said = "when its merge was made is not kept: its delivery time is unknown"
} else {
w.From = from
}
points := p.points(now, applied, w)
// Walk the moments: each known moment after a known one is a measured phase; a missing moment makes the
// phases up to the next known one unknown, and their span unknown time.
last := from
var pending []int
for _, pt := range points {
if pt.none {
w.Phases = append(w.Phases, WalkPhase{Name: pt.phase, Tier: pt.tier, State: PhaseNone, Said: pt.said})
continue
}
ph := WalkPhase{Name: pt.phase, Tier: pt.tier, Said: pt.said}
if pt.at == nil {
ph.State = PhaseUnknown
w.Phases = append(w.Phases, ph)
pending = append(pending, len(w.Phases)-1)
continue
}
at := pt.at.Truncate(time.Millisecond)
ph.End = &at
if last != nil && at.Before(*last) {
// **Out of order** (issue 382): this moment was recorded before the one before it — a module's first
// send kept before its build was, two clocks, a late record. Neither phase can be measured honestly:
// the one before it is said unknown and its time joins the unknown time, and this one is said unknown
// with no time. The walk goes on from the later moment, so nothing is negative and the sum holds.
for i := len(w.Phases) - 1; i >= 0; i-- {
prev := &w.Phases[i]
if prev.End == nil {
continue // a phase that did not happen, or is unknown with no moment: the one before carries last
}
if prev.State == PhaseMeasured {
w.UnknownMS += prev.TookMS
}
prev.State, prev.Start, prev.TookMS, prev.Took = PhaseUnknown, nil, 0, ""
prev.Said = "out of order: the phase after it ended first"
break
}
ph.State, ph.Said = PhaseUnknown, "out of order: it ended before the phase before it"
w.Phases = append(w.Phases, ph)
pending = nil
continue
}
if last != nil && len(pending) == 0 {
start := *last
ph.State, ph.Start = PhaseMeasured, &start
ph.TookMS = at.Sub(start).Milliseconds()
ph.Took = words(at.Sub(start))
} else {
ph.State = PhaseUnknown
if last != nil {
w.UnknownMS += at.Sub(*last).Milliseconds()
}
}
w.Phases = append(w.Phases, ph)
pending = nil
last = &at
}
if len(pending) > 0 {
// The walk's last moments are unknown: it has no end.
if w.Said == "" {
w.Said = "its end is unknown: " + w.Phases[pending[0]].Name + " " + orNot(w.Phases[pending[0]].Said)
}
return w
}
if last == nil || from == nil {
return w
}
if p.Open() {
since := *last
w.Phases = append(w.Phases, WalkPhase{Name: PhaseOpen, Tier: -1, State: PhaseOpen, Start: &since,
TookMS: now.Sub(since).Milliseconds(), Took: words(now.Sub(since)), Said: "the walk is " + p.State})
w.Said = "open: " + p.State + ", " + words(now.Sub(*from)) + " since its merge"
return w
}
if p.State != PlanDone {
w.Said = "ended " + p.State + ": no delivery time"
return w
}
end := *last
w.End = &end
w.TotalMS = end.Sub(*from).Milliseconds()
w.Total = words(end.Sub(*from))
if w.Said == "" {
if w.UnknownMS > 0 {
w.Said = fmt.Sprintf("%s from its merge to running everywhere, %s of it unknown", w.Total,
words(time.Duration(w.UnknownMS)*time.Millisecond))
} else {
w.Said = w.Total + " from its merge to running everywhere"
}
}
return w
}
// firstMerge is when the walk's earliest merge was made, and every merge's start kept on w.
func (p Plan) firstMerge(w *WalkPhases) *time.Time {
var first *time.Time
if p.Delivery != nil {
for _, m := range p.Delivery.Merges {
if m.Merged.IsZero() {
continue
}
w.Merges = append(w.Merges, MergeStart{Repository: m.Repository, Commit: m.Commit, Merged: m.Merged})
if first == nil || m.Merged.Before(*first) {
at := m.Merged
first = &at
}
}
}
if first == nil {
for _, c := range p.Carried() {
if c.Merged.IsZero() {
continue
}
w.Merges = append(w.Merges, MergeStart{Repository: c.Repository, Commit: c.Commit, Merged: c.Merged})
if first == nil || c.Merged.Before(*first) {
at := c.Merged
first = &at
}
}
}
return first
}
// points are the walk's moments in order.
func (p Plan) points(now time.Time, applied AppliedLookup, w *WalkPhases) []point {
var out []point
// The window and the wait behind another walk.
cut := p.Created
if p.Times != nil && p.Times.Cut != nil {
cut = *p.Times.Cut
}
switch {
case p.Times == nil:
out = append(out, point{phase: PhaseWindow, tier: -1, said: "not kept for a walk made before ADR 0282"},
point{phase: PhaseQueued, tier: -1, at: &cut, said: "the window and the wait behind another walk together"})
case p.Times.WindowClosed == nil:
out = append(out, point{phase: PhaseWindow, tier: -1, none: true, said: "no window: walked on its own"},
point{phase: PhaseQueued, tier: -1, at: &cut})
default:
closed := *p.Times.WindowClosed
out = append(out, point{phase: PhaseWindow, tier: -1, at: &closed}, point{phase: PhaseQueued, tier: -1, at: &cut})
}
if p.Delivery != nil && p.Delivery.Awaits != "" {
out = append(out, point{phase: PhaseWord, tier: -1, at: p.Delivery.Go, said: "waiting for " + p.Delivery.Awaits + "'s word"})
} else {
out = append(out, point{phase: PhaseWord, tier: -1, none: true, said: "waits for no word"})
}
var lastSends []restSend
for t, tier := range p.Tiers {
if t > p.Tier || (t == p.Tier && p.Tier < len(p.Tiers) && !askedAnyOf(p, tier)) {
break
}
var asked, built, first, judged, rest *time.Time
allBuilt, anyFirst, allJudged, allRest := true, false, true, true
for _, m := range tier {
s := p.Modules[m]
if s == nil {
allBuilt, allRest = false, false
continue
}
asked = earliest(asked, s.AskedAt)
if s.BuiltAt == nil {
if s.State != "deleted" {
allBuilt = false
}
} else {
built = latest(built, s.BuiltAt)
}
if s.FirstAt != nil {
anyFirst = true
first = earliest(first, s.FirstAt)
if s.Gate != nil {
if s.Gate.JudgedAt == nil {
allJudged = false
} else {
judged = latest(judged, s.Gate.JudgedAt)
}
}
}
if s.SentAt == nil {
if s.State != "deleted" {
allRest = false
}
} else {
rest = latest(rest, s.SentAt)
for node := range s.Rest {
lastSends = append(lastSends, restSend{module: m, node: node, at: *s.SentAt})
}
}
}
if t > 0 {
out = append(out, point{phase: PhaseBetween, tier: t, at: asked})
}
if !allBuilt {
built = nil
}
out = append(out, point{phase: PhaseBuild, tier: t, at: built})
if anyFirst {
if !allJudged {
judged = nil
}
out = append(out, point{phase: PhaseSend, tier: t, at: first}, point{phase: PhaseJudge, tier: t, at: judged})
} else {
out = append(out, point{phase: PhaseSend, tier: t, none: true, said: "no first machine: nothing to judge"},
point{phase: PhaseJudge, tier: t, none: true, said: "no first machine: nothing to judge"})
}
if !allRest {
rest = nil
}
out = append(out, point{phase: PhaseRest, tier: t, at: rest})
}
if p.State != PlanDone {
return out
}
// Every machine of the rest running the build: its first report after the send, applied.
if p.Times == nil {
return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's sends are not kept for a walk made before ADR 0282"})
}
if len(lastSends) == 0 {
return append(out, point{phase: PhaseApply, tier: -1, none: true,
said: "no machine beyond the first: the first-node gate's readings were its run"})
}
if applied == nil {
return append(out, point{phase: PhaseApply, tier: -1, said: "the rest's reports were not read"})
}
var end *time.Time
var waiting, failed []string
for _, s := range lastSends {
r, ok := applied(s.node, s.at)
switch {
case !ok && now.Sub(s.at) > ApplySilentAfter, ok && r.At.Sub(s.at) > ApplySilentAfter:
w.Silent = appendOnce(w.Silent, s.node)
case !ok:
waiting = appendOnce(waiting, s.node)
case r.Outcome != OutcomeApplied:
failed = appendOnce(failed, s.node+" ("+r.Outcome+")")
default:
at := r.At
end = latest(end, &at)
}
}
switch {
case len(failed) > 0:
return append(out, point{phase: PhaseApply, tier: -1, said: "not applied on " + strings.Join(failed, ", ")})
case len(waiting) > 0:
return append(out, point{phase: PhaseApply, tier: -1, said: "waiting for " + strings.Join(waiting, ", ") +
" to report it applied"})
case end == nil:
return append(out, point{phase: PhaseApply, tier: -1, said: "no machine of the rest heard from: " +
strings.Join(w.Silent, ", ")})
}
// A machine may have reported before the last tier ended: the walk ends at whichever is later.
for i := len(out) - 1; i >= 0; i-- {
if out[i].at != nil {
if out[i].at.After(*end) {
e := *out[i].at
end = &e
}
break
}
}
said := ""
if len(w.Silent) > 0 {
sort.Strings(w.Silent)
said = "not waited for, not heard from: " + strings.Join(w.Silent, ", ")
}
return append(out, point{phase: PhaseApply, tier: -1, at: end, said: said})
}
type restSend struct {
module, node string
at time.Time
}
func askedAnyOf(p Plan, tier []string) bool {
for _, m := range tier {
if s := p.Modules[m]; s != nil && s.AskedAt != nil {
return true
}
}
return false
}
func earliest(a, b *time.Time) *time.Time {
if b == nil {
return a
}
if a == nil || b.Before(*a) {
t := *b
return &t
}
return a
}
func latest(a, b *time.Time) *time.Time {
if b == nil {
return a
}
if a == nil || b.After(*a) {
t := *b
return &t
}
return a
}
func appendOnce(to []string, s string) []string {
for _, x := range to {
if x == s {
return to
}
}
return append(to, s)
}
func orNot(s string) string {
if s == "" {
return "is not known"
}
return "(" + s + ")"
}
// words is a duration as a person reads it.
func words(d time.Duration) string {
if d < 0 {
return "-" + words(-d)
}
return d.Round(100 * time.Millisecond).String()
}
// FirstAppliedAfter is a machine's first report of a send made at or after a walk's send to it, within
// ApplySilentAfter of it — the controller's `apply` durations, which measure every send to its first report (to-be
// 45 Phase 0). A declaration sent later carries the build too, so its report counts; one sent past the bound is a
// machine that slept, listed as silent and never moving the walk's end (ADR 0282 decision 1).
func (i *Inventory) FirstAppliedAfter(ctx context.Context, node string, sent time.Time) (AppliedReport, bool, error) {
var started time.Time
var ms int64
var outcome string
err := i.store.Pool().QueryRow(ctx,
`select started, took_ms, detail from duration
where kind = $1 and subject = $2 and started >= $3 and started < $4
order by started limit 1`,
DurationApply, node, sent.Add(-appliedSlack), sent.Add(ApplySilentAfter)).Scan(&started, &ms, &outcome)
if errors.Is(err, pgx.ErrNoRows) {
return AppliedReport{}, false, nil
}
if err != nil {
return AppliedReport{}, false, err
}
return AppliedReport{At: started.Add(time.Duration(ms) * time.Millisecond).UTC(), Outcome: outcome}, true, nil
}
// appliedSlack is how much earlier than a walk's record of its send the machine's own record of it may be:
// the send is recorded on the machine first, then on the walk.
const appliedSlack = 2 * time.Second
+330
View File
@@ -0,0 +1,330 @@
package inventory
import (
"encoding/json"
"strings"
"testing"
"time"
)
// A walk recorded on the live mesh on 2026-10-10 (mesh-delivery, one tier, one machine), as `delivery walks`
// gave it, with the moments ADR 0282 adds: its window closed ninety seconds after its merge was heard, and it was
// cut eleven seconds later.
const recordedWalk = `{
"id": "plan-1791654663505629616", "repository": "novox/mesh-catalog", "branch": "main",
"commit": "a14fa306f262d88bdcca83432a70df353d2eceb0", "merged": "2026-10-10T17:50:53Z",
"created": "2026-10-10T19:52:34.037755+02:00", "updated": "2026-10-10T19:55:37.974069+02:00",
"state": "done", "tier": 1, "tiers": [["mesh-delivery"]],
"modules": {"mesh-delivery": {"state": "built",
"asked_at": "2026-10-10T17:52:34.209423145Z", "built_at": "2026-10-10T17:52:50.818011314Z",
"sent_at": "2026-10-10T17:55:35.894985053Z", "first": ["novox"], "first_at": "2026-10-10T17:53:07.287136827Z",
"gate": {"machines": ["novox"], "since": "2026-10-10T17:53:07.287136827Z", "passes": 3,
"verdict": "passed", "judged_at": "2026-10-10T17:55:35.894985053Z"}}},
"delivery": {"awaits": "", "merges": [{"repository": "novox/mesh-catalog",
"commit": "a14fa306f262d88bdcca83432a70df353d2eceb0", "number": 187, "merged": "2026-10-10T17:50:53Z"}]},
"times": {"window_closed": "2026-10-10T17:52:23Z", "cut": "2026-10-10T17:52:34.037755Z", "class": "core"}
}`
func walkOf(t *testing.T, raw string) Plan {
t.Helper()
var p Plan
if err := json.Unmarshal([]byte(raw), &p); err != nil {
t.Fatal(err)
}
return p
}
// sumsUp checks the phases' measured time and the unknown time add up to the total, to the millisecond.
func sumsUp(t *testing.T, w *WalkPhases) {
t.Helper()
if w.End == nil || w.From == nil {
t.Fatalf("the walk has no end: %+v", w)
}
var sum int64
for _, ph := range w.Phases {
if ph.State == PhaseMeasured {
if ph.Start == nil || ph.End == nil || ph.End.Sub(*ph.Start).Milliseconds() != ph.TookMS {
t.Errorf("phase %s %d says %dms between %v and %v", ph.Name, ph.Tier, ph.TookMS, ph.Start, ph.End)
}
sum += ph.TookMS
}
if ph.State == PhaseUnknown && ph.TookMS != 0 {
t.Errorf("an unknown phase %s says a time: %dms", ph.Name, ph.TookMS)
}
}
if sum+w.UnknownMS != w.TotalMS || w.End.Sub(*w.From).Milliseconds() != w.TotalMS {
t.Fatalf("the phases add up to %dms and %dms unknown, the total is %dms (%s): %+v", sum, w.UnknownMS,
w.TotalMS, w.End.Sub(*w.From), w.Phases)
}
}
func phase(w *WalkPhases, name string, tier int) WalkPhase {
for _, ph := range w.Phases {
if ph.Name == name && ph.Tier == tier {
return ph
}
}
return WalkPhase{}
}
// The phases of a recorded walk add up to its measured total: its merge at 17:50:53 to its verdict on its only
// machine at 17:55:35.894, 4m42.894s.
func TestARecordedWalksPhasesAddUpToItsTotal(t *testing.T) {
w := walkOf(t, recordedWalk).WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil)
sumsUp(t, w)
if w.TotalMS != 282894 || w.Class != ClassCore || w.UnknownMS != 0 {
t.Fatalf("total %dms (want 282894), class %q, unknown %d", w.TotalMS, w.Class, w.UnknownMS)
}
for name, want := range map[string]int64{PhaseWindow: 90000, PhaseQueued: 11037, PhaseBuild: 16781,
PhaseSend: 16469, PhaseJudge: 148607, PhaseRest: 0} {
tier := 0
if name == PhaseWindow || name == PhaseQueued {
tier = -1
}
if got := phase(w, name, tier); got.State != PhaseMeasured || got.TookMS != want {
t.Errorf("%s: %s %dms, want measured %dms", name, got.State, got.TookMS, want)
}
}
// The cut to the first ask (171ms) is a time no phase of ADR 0282's table names: it is counted in the build.
if got := phase(w, PhaseWord, -1); got.State != PhaseNone {
t.Errorf("a walk that waits for no word says its word phase %s", got.State)
}
if got := phase(w, PhaseApply, -1); got.State != PhaseNone || !strings.Contains(got.Said, "no machine beyond the first") {
t.Errorf("one machine only: apply %s (%s)", got.State, got.Said)
}
}
// Two tiers, a machine of the rest each, both reports read: every phase measured, the walk ends at the last
// report applied, and the phases add up.
func TestAWalkOfTwoTiersEndsAtItsRestsLastReport(t *testing.T) {
at := func(s string) *time.Time {
v, err := time.Parse(time.RFC3339Nano, "2026-10-10T18:"+s+"Z")
if err != nil {
t.Fatal(err)
}
return &v
}
p := Plan{ID: "plan-2", State: PlanDone, Tier: 2, Tiers: [][]string{{"a"}, {"b"}}, Created: *at("01:40"),
Delivery: &PlanDelivery{Merges: []PlanMerge{{Repository: "novox/x", Commit: "c1", Merged: *at("00:00")},
{Repository: "novox/x", Commit: "c0", Merged: *at("00:30")}}},
Times: &PlanTimes{WindowClosed: at("01:30"), Cut: at("01:40"), Class: ClassLeaf},
Modules: map[string]*PlanModule{
"a": {State: "built", AskedAt: at("01:41"), BuiltAt: at("02:05"), FirstAt: at("02:20"), SentAt: at("05:00"),
Gate: &PlanGate{JudgedAt: at("04:50")}, Rest: map[string]SentDeclaration{"ace": {Digest: "d1"}}},
"b": {State: "built", AskedAt: at("05:10"), BuiltAt: at("05:40"), FirstAt: at("05:50"), SentAt: at("08:00.5"),
Gate: &PlanGate{JudgedAt: at("07:59")}, Rest: map[string]SentDeclaration{"shanks": {Digest: "d2"}}},
}}
reports := map[string]AppliedReport{"ace": {At: *at("05:12"), Outcome: OutcomeApplied},
"shanks": {At: *at("08:15.25"), Outcome: OutcomeApplied}}
w := p.WalkPhases(*at("30:00"), func(node string, _ time.Time) (AppliedReport, bool) {
r, ok := reports[node]
return r, ok
})
sumsUp(t, w)
if w.TotalMS != (8*time.Minute + 15250*time.Millisecond).Milliseconds() {
t.Fatalf("the walk's delivery time runs from its earliest merge to the last report: %s", w.Total)
}
if got := phase(w, PhaseBetween, 1); got.TookMS != 10000 {
t.Errorf("between the tiers: %dms", got.TookMS)
}
if got := phase(w, PhaseApply, -1); got.State != PhaseMeasured || got.TookMS != 14750 {
t.Errorf("apply: %s %dms", got.State, got.TookMS)
}
if len(w.Merges) != 2 {
t.Errorf("each merge's start is said, for its own delivery time: %+v", w.Merges)
}
// A report not yet in: the walk has no end, and its apply is unknown — never zero.
delete(reports, "shanks")
w = p.WalkPhases(*at("10:00"), func(node string, _ time.Time) (AppliedReport, bool) {
r, ok := reports[node]
return r, ok
})
if w.End != nil || w.TotalMS != 0 {
t.Fatalf("a walk whose rest has not reported has no end: %+v", w)
}
if got := phase(w, PhaseApply, -1); got.State != PhaseUnknown || got.TookMS != 0 || !strings.Contains(got.Said, "shanks") {
t.Errorf("apply while shanks has not reported: %+v", got)
}
// Silent past the bound: listed, not waited for.
w = p.WalkPhases(at("08:00.5").Add(ApplySilentAfter+time.Second), func(node string, _ time.Time) (AppliedReport, bool) {
r, ok := reports[node]
return r, ok
})
sumsUp(t, w)
if len(w.Silent) != 1 || w.Silent[0] != "shanks" {
t.Errorf("a machine silent past the bound is said, not waited for: %+v", w.Silent)
}
// A failed report: the walk has no end, said.
reports["shanks"] = AppliedReport{At: *at("08:10"), Outcome: OutcomeFailed}
w = p.WalkPhases(*at("30:00"), func(node string, _ time.Time) (AppliedReport, bool) {
r, ok := reports[node]
return r, ok
})
if w.End != nil || !strings.Contains(phase(w, PhaseApply, -1).Said, "shanks (failed)") {
t.Errorf("a failed apply ends nothing: %+v", phase(w, PhaseApply, -1))
}
}
// A moment that is missing makes the phases around it unknown, and their span unknown time: never a zero.
func TestAMissingMomentIsUnknownNeverZero(t *testing.T) {
p := walkOf(t, recordedWalk)
p.Modules["mesh-delivery"].BuiltAt = nil
w := p.WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil)
sumsUp(t, w)
if got := phase(w, PhaseBuild, 0); got.State != PhaseUnknown || got.TookMS != 0 {
t.Errorf("a build without its moment: %+v", got)
}
if got := phase(w, PhaseSend, 0); got.State != PhaseUnknown {
t.Errorf("the send after a build without its moment: %+v", got)
}
if w.UnknownMS != 16781+16469 {
t.Errorf("the unknown time is the span between the moments around it: %dms", w.UnknownMS)
}
// A walk kept before ADR 0282: its window and its wait together, and its apply unknown.
p = walkOf(t, recordedWalk)
p.Times = nil
w = p.WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil)
if got := phase(w, PhaseWindow, -1); got.State != PhaseUnknown {
t.Errorf("a window not kept: %+v", got)
}
if got := phase(w, PhaseApply, -1); got.State != PhaseUnknown || w.End != nil {
t.Errorf("a rest not kept: %+v, end %v", got, w.End)
}
}
// An open walk is said open since its last moment, with no end.
func TestAnOpenWalkHasNoEnd(t *testing.T) {
p := walkOf(t, recordedWalk)
p.State, p.Tier = PlanRolling, 0
p.Modules["mesh-delivery"].SentAt = nil
p.Modules["mesh-delivery"].Gate.JudgedAt = nil
w := p.WalkPhases(time.Date(2026, 10, 10, 17, 54, 0, 0, time.UTC), nil)
if w.End != nil || !strings.HasPrefix(w.Said, "its end is unknown") && !strings.HasPrefix(w.Said, "open") {
t.Fatalf("an open walk: %+v", w)
}
if got := phase(w, PhaseBuild, 0); got.State != PhaseMeasured {
t.Errorf("an open walk's phases so far are measured: %+v", got)
}
}
// A gate keeps every pass and the first reading after one that did not pass, at most maxReadings.
func TestAGateKeepsItsReadings(t *testing.T) {
var g PlanGate
t0 := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC)
g.Read(t0, false, "starting")
g.Read(t0.Add(time.Second), false, "starting")
g.Read(t0.Add(40*time.Second), true, "")
g.Read(t0.Add(80*time.Second), true, "")
g.Read(t0.Add(81*time.Second), false, "gone")
g.Read(t0.Add(82*time.Second), false, "gone")
if len(g.Readings) != 4 || g.Readings[0].Said != "starting" || !g.Readings[2].Healthy || g.Readings[3].Said != "gone" {
t.Fatalf("readings: %+v", g.Readings)
}
for i := 0; i < 50; i++ {
g.Read(t0.Add(time.Duration(100+i)*time.Second), true, "")
}
if len(g.Readings) != maxReadings {
t.Fatalf("readings are bounded: %d", len(g.Readings))
}
}
// A walk's moments are kept and read back with it; a machine's first report after a send is read from the
// controller's apply durations.
func TestAWalksMomentsAreKeptAndItsRestsReportRead(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
p := walkOf(t, recordedWalk)
p.Revision, p.Epoch = 0, 0
if err := inv.SavePlan(ctx, &p); err != nil {
t.Fatal(err)
}
back, err := inv.PlanByID(ctx, p.ID)
if err != nil || back.Times == nil || back.Times.Class != ClassCore || back.Times.WindowClosed == nil ||
!back.Times.WindowClosed.Equal(*p.Times.WindowClosed) {
t.Fatalf("the walk's moments were not kept: %v %+v", err, back.Times)
}
if back.Delivery.Merges[0].Merged.IsZero() {
t.Fatal("a merge's time was not kept")
}
sent := time.Now().UTC().Add(-time.Minute).Truncate(time.Millisecond)
if _, ok, err := inv.FirstAppliedAfter(ctx, "ace", sent); err != nil || ok {
t.Fatalf("no report yet: %v %v", ok, err)
}
for _, d := range []Duration{
{Kind: DurationApply, Subject: "ace", Node: "ace", Ref: "old@1", Started: sent.Add(-time.Hour), Took: time.Second, Detail: OutcomeApplied},
{Kind: DurationApply, Subject: "ace", Node: "ace", Ref: "d1@2", Started: sent.Add(-time.Second), Took: 12 * time.Second, Detail: OutcomeApplied},
} {
if err := inv.RecordDuration(ctx, d); err != nil {
t.Fatal(err)
}
}
r, ok, err := inv.FirstAppliedAfter(ctx, "ace", sent)
if err != nil || !ok || r.Outcome != OutcomeApplied || !r.At.Equal(sent.Add(11*time.Second)) {
t.Fatalf("the report after the send: %+v %v %v", r, ok, err)
}
if anyKept, total, err := inv.WalkPhasesKept(ctx, p.ID); err != nil || anyKept || total {
t.Fatalf("nothing kept yet: %v %v %v", anyKept, total, err)
}
if err := inv.RecordDuration(ctx, Duration{Kind: DurationWalkPhase, Subject: "core/total", Ref: p.ID + "/total",
Started: sent, Took: time.Minute}); err != nil {
t.Fatal(err)
}
if anyKept, total, err := inv.WalkPhasesKept(ctx, p.ID); err != nil || !anyKept || !total {
t.Fatalf("the total kept: %v %v %v", anyKept, total, err)
}
}
// A machine that slept past the bound and reported later is listed silent: its late report never moves the
// walk's end; and a report sent past the bound is not read as the walk's.
func TestALateReportIsSilentNotTheEnd(t *testing.T) {
sent := time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC)
p := walkOf(t, recordedWalk)
p.Modules["mesh-delivery"].SentAt = &sent
p.Modules["mesh-delivery"].Rest = map[string]SentDeclaration{"laptop": {Digest: "d"}, "server": {Digest: "e"}}
w := p.WalkPhases(sent.Add(3*time.Hour), func(node string, _ time.Time) (AppliedReport, bool) {
if node == "server" {
return AppliedReport{At: sent.Add(20 * time.Second), Outcome: OutcomeApplied}, true
}
return AppliedReport{At: sent.Add(2 * time.Hour), Outcome: OutcomeApplied}, true
})
if len(w.Silent) != 1 || w.Silent[0] != "laptop" || w.End == nil || !w.End.Equal(sent.Add(20*time.Second)) {
t.Fatalf("a late report: silent %v, end %v", w.Silent, w.End)
}
inv := ForTest(t)
ctx := t.Context()
if err := inv.RecordDuration(ctx, Duration{Kind: DurationApply, Subject: "laptop", Node: "laptop", Ref: "late@1",
Started: sent.Add(time.Hour), Took: time.Second, Detail: OutcomeApplied}); err != nil {
t.Fatal(err)
}
if _, ok, err := inv.FirstAppliedAfter(ctx, "laptop", sent); err != nil || ok {
t.Fatalf("a send past the bound read as the walk's: %v %v", ok, err)
}
}
// Moments out of order — #206's own walk on 2026-10-10 kept its first send 783 ms before its build — are said
// unknown, never a negative phase: the phase before is unknown and its time joins the unknown time, the phase
// out of order has no time, and the phases still add up to the total.
func TestMomentsOutOfOrderAreUnknownNeverNegative(t *testing.T) {
p := walkOf(t, recordedWalk)
built := time.Date(2026, 10, 10, 17, 53, 8, 70000000, time.UTC) // 783 ms after its first send
p.Modules["mesh-delivery"].BuiltAt = &built
w := p.WalkPhases(time.Date(2026, 10, 10, 18, 0, 0, 0, time.UTC), nil)
sumsUp(t, w)
for _, ph := range w.Phases {
if ph.TookMS < 0 {
t.Fatalf("a negative phase: %+v", ph)
}
}
build, send := phase(w, PhaseBuild, 0), phase(w, PhaseSend, 0)
if build.State != PhaseUnknown || send.State != PhaseUnknown || send.TookMS != 0 ||
!strings.Contains(send.Said, "out of order") || !strings.Contains(build.Said, "out of order") {
t.Fatalf("build %+v, send-first %+v", build, send)
}
// The build's span (cut to its later moment) is unknown time; the judgement is measured from the later moment.
if w.UnknownMS != built.Sub(time.Date(2026, 10, 10, 17, 52, 34, 37000000, time.UTC)).Milliseconds() {
t.Fatalf("unknown %dms", w.UnknownMS)
}
if j := phase(w, PhaseJudge, 0); j.State != PhaseMeasured || j.Start == nil || !j.Start.Equal(built) {
t.Fatalf("the judgement is measured from the later moment: %+v", j)
}
}
+16
View File
@@ -556,8 +556,21 @@ const (
StateStarting = "starting" StateStarting = "starting"
StateHeld = "held" StateHeld = "held"
StateUnknown = "unknown" StateUnknown = "unknown"
// StateWaiting is a resource whose module's tool check says it waits for the operator: a named own secret
// or setting not given yet (novox/hq ADR 0283). Not a fault; the controller checks the wait before it
// excuses it.
StateWaiting = "waiting"
) )
// Wait is one thing a waiting resource waits for the operator to give: exactly one of Secret and Setting, the
// part that waits and what the operator gives, in words (novox/hq ADR 0283; mesh-sdk go/health's shape).
type Wait struct {
Part string `json:"part"`
Secret string `json:"secret,omitempty"`
Setting string `json:"setting,omitempty"`
What string `json:"what"`
}
// ResourceHealth is one long-running resource's state. // ResourceHealth is one long-running resource's state.
type ResourceHealth struct { type ResourceHealth struct {
Module string `json:"module"` Module string `json:"module"`
@@ -581,6 +594,9 @@ type ResourceHealth struct {
// without a person (novox/hq ADR 0266): the engine judged that too, and a healthy verdict says it cannot. // without a person (novox/hq ADR 0266): the engine judged that too, and a healthy verdict says it cannot.
// Empty from an engine older than RootContract, and on every other verdict. // Empty from an engine older than RootContract, and on every other verdict.
Root string `json:"root,omitempty"` Root string `json:"root,omitempty"`
// Waits is what a resource in StateWaiting waits for (novox/hq ADR 0283). Empty otherwise, and from an
// engine older than that.
Waits []Wait `json:"waits,omitempty"`
} }
// HealthSaid is the health event's body: the machine and its statement. The machine is read from the // HealthSaid is the health event's body: the machine and its statement. The machine is read from the