diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index f4d02ea..3baf24c 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -61,7 +61,11 @@ type probe struct { Phase int // Deferred says why it is not run yet; empty for one that is. Deferred string - run func(ctx context.Context, d *doctor) ([]conditions.Observation, error) + // Asks are the seat verbs it calls. A probe may call no other (askSeatTool refuses), and the + // controller's grant names every one (a test over this registry): a probe whose question the bus + // refuses checks nothing (D8, 2026-10-06). + Asks []broker.SeatVerb + run func(ctx context.Context, d *doctor) ([]conditions.Observation, error) } // probeRegistry is the registry, in to-be 45's order. **The registry is the design's live form**: a @@ -83,7 +87,8 @@ var probeRegistry = []probe{ {ID: "D7", Asserts: "every stream the controller defines exists with its definition, and its own buckets", From: "issue 208", Kind: "stream-wrong", Phase: 1, run: probeStreams}, {ID: "D8", Asserts: "no address the mesh owns — a machine's private address or its endpoint — is in a ban list", - From: "issue 238", Kind: "own-address-banned", Phase: 1, run: probeBans}, + From: "issue 238", Kind: "own-address-banned", Phase: 1, run: probeBans, + Asks: []broker.SeatVerb{{Seat: "node-intrusion-prevention", Verb: "banned"}}}, {ID: "D9", Asserts: "status answers in full within ten seconds, from a summary composed lately", From: "issue 265", Kind: "status-slow", Phase: 1, run: probeStatus}, {ID: "D10", Asserts: "every machine runs the node-engine and node tools builds the mesh holds, or is inside " + @@ -209,7 +214,7 @@ func (d *doctor) runOnce(ctx context.Context, why string) doctorRun { wg.Add(1) go func(i int, p probe) { defer wg.Done() - probing, cancel := context.WithTimeout(ctx, probeWithin) + probing, cancel := context.WithTimeout(context.WithValue(ctx, probeAsksKey{}, p), probeWithin) defer cancel() began := time.Now() done := make(chan result, 1) @@ -512,3 +517,20 @@ func sortedFound(obs []conditions.Observation) []conditions.Observation { sort.Slice(obs, func(i, j int) bool { return obs[i].Key() < obs[j].Key() }) return obs } + +// probeAsksKey carries the running probe, so a seat verb it calls is checked against what it declares. +type probeAsksKey struct{} + +// declaredBy says whether the probe running in ctx declared a seat verb; outside a probe, false. +func declaredBy(ctx context.Context, seat, verb string) (string, bool) { + p, ok := ctx.Value(probeAsksKey{}).(probe) + if !ok { + return "a caller outside the self-check", false + } + for _, v := range p.Asks { + if v.Seat == seat && v.Verb == verb { + return p.ID, true + } + } + return p.ID, false +} diff --git a/cmd/mesh-controller/doctor_test.go b/cmd/mesh-controller/doctor_test.go index ca4405d..7a238a6 100644 --- a/cmd/mesh-controller/doctor_test.go +++ b/cmd/mesh-controller/doctor_test.go @@ -17,6 +17,7 @@ import ( "golang.org/x/net/dns/dnsmessage" "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/conditions" "github.com/novox/mesh-controller/internal/link" ) @@ -345,3 +346,80 @@ func TestNatsTheBusSaysAConsumerGaveUpAndOneWasDeleted(t *testing.T) { t.Fatalf("%v", kinds) } } + +// **D10 compares a node-engine with what it is delivered as, not with its commit** (2026-10-06: every +// machine read as behind right after a push sent it the current build — it says the digest-named +// directory it runs from, and the mesh holds a commit). +func TestANodeEngineIsJudgedByTheVersionItIsDeliveredAs(t *testing.T) { + m := catalogue.Manifest{Module: "mesh-host", Resources: []map[string]any{ + {"id": "launcher", "type": "file", "path": "/usr/lib/nox-mesh-host/launch"}, + {"id": "host", "type": "archive", "path": "/usr/lib/nox-mesh-host/versions/31045596c83a"}, + {"id": "unfilled", "type": "archive", "path": "/usr/lib/x/versions/${version}"}, + }} + delivered := deliveredVersions(m) + if !slices.Equal(delivered, []string{"31045596c83a"}) { + t.Fatalf("%v", delivered) + } + commit := "1545b00a9f0c" + for _, c := range []struct { + reported string + behind bool + }{ + {"31045596c83a", false}, // the live case: current, and was called behind + {"0123456789ab", true}, // another delivery + {"1545b00a", false}, // placed by hand, stamped with the commit + {"", false}, // not said + } { + if got := engineBehind(c.reported, delivered, commit); got != c.behind { + t.Errorf("%q behind = %v, want %v", c.reported, got, c.behind) + } + } + if engineBehind("31045596c83a", nil, commit) { + t.Error("behind a mesh that holds no delivered build") + } +} + +// **Every seat verb a probe calls is one it declares, and one the controller is granted** — derived +// from the registry, so a probe added with a question the bus would refuse fails here, not live. +func TestEverySeatVerbAProbeAsksIsGranted(t *testing.T) { + granted, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindController, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + asked := 0 + for _, p := range probeRegistry { + for _, v := range p.Asks { + asked++ + subject := link.NodeSeatToolSubject(v.Seat, v.Verb, "anchor") + if !slices.ContainsFunc(granted.Publish, func(pattern string) bool { return subjectMatches(pattern, subject) }) { + t.Errorf("%s asks %s.%s and the controller may not publish %s", p.ID, v.Seat, v.Verb, subject) + } + if !slices.Contains(broker.VerbsTheSelfCheckAsks, v) { + t.Errorf("%s asks %s.%s, which broker.VerbsTheSelfCheckAsks does not name", p.ID, v.Seat, v.Verb) + } + } + } + if asked == 0 { + t.Fatal("no probe asks a seat verb: D8 lost its declaration") + } + // And a probe asking what it did not declare is refused before anything is sent. + ctx := context.WithValue(t.Context(), probeAsksKey{}, probe{ID: "DX"}) + if _, err := askSeatTool(ctx, nil, "node-intrusion-prevention", "banned", "anchor"); err == nil || + !strings.Contains(err.Error(), "does not declare") { + t.Fatalf("an undeclared question was asked: %v", err) + } +} + +// subjectMatches is the bus's matching of a permission pattern against a subject. +func subjectMatches(pattern, subject string) bool { + p, s := strings.Split(pattern, "."), strings.Split(subject, ".") + for i, tok := range p { + if tok == ">" { + return len(s) > i + } + if i >= len(s) || (tok != "*" && tok != s[i]) { + return false + } + } + return len(p) == len(s) +} diff --git a/cmd/mesh-controller/probes.go b/cmd/mesh-controller/probes.go index b4cbb80..9146a51 100644 --- a/cmd/mesh-controller/probes.go +++ b/cmd/mesh-controller/probes.go @@ -10,6 +10,7 @@ import ( "slices" "sort" "strings" + "sync" "time" "github.com/nats-io/nats.go" @@ -576,9 +577,22 @@ func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error) } sort.Strings(holders) var out []conditions.Observation + // Every machine asked at once: one after another, four holders that each wait their bound + // outlast the probe's thirty seconds, and the probe says nothing about any of them. + answers := make([]json.RawMessage, len(holders)) + errs := make([]error, len(holders)) + var wg sync.WaitGroup + for i, node := range holders { + wg.Add(1) + go func(i int, node string) { + defer wg.Done() + answers[i], errs[i] = askSeatTool(ctx, d.js.Conn(), "node-intrusion-prevention", "banned", node) + }(i, node) + } + wg.Wait() var unasked []string - for _, node := range holders { - answer, err := askSeatTool(ctx, d.js.Conn(), "node-intrusion-prevention", "banned", node) + for i, node := range holders { + answer, err := answers[i], errs[i] if err != nil { unasked = append(unasked, err.Error()) continue @@ -650,6 +664,11 @@ func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation, if err != nil { return nil, err } + shelf, err := inv.Catalogue(ctx) + if err != nil { + return nil, err + } + hostVersions := deliveredVersions(shelf[hostModule]) rolling := map[string]bool{} for _, p := range plans { for m := range p.Modules { @@ -675,9 +694,14 @@ func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation, continue // a machine not heard from is S1's } var behind []string - host := current[hostModule].Commit - if host != "" && n.HostVersion != "" && !rolling[hostModule] && !sameCommit(n.HostVersion, host) { - behind = append(behind, fmt.Sprintf("the node-engine %s, the mesh holds %s", short(n.HostVersion), short(host))) + // **A node-engine says its version as the directory it was delivered into** — its archive's + // digest, twelve characters (catalogue.versionOf, ADR 0141) — not the commit it was built + // from. Compared as a commit, every machine read as behind right after a push sent it the + // current one (2026-10-06). An engine placed by hand reports its link-time stamp instead, + // which a commit can match. + if !rolling[hostModule] && engineBehind(n.HostVersion, hostVersions, current[hostModule].Commit) { + behind = append(behind, fmt.Sprintf("the node-engine %s, the mesh holds %s (built from %s)", + n.HostVersion, strings.Join(hostVersions, " or "), short(current[hostModule].Commit))) } assigned, err := inv.Assigned(ctx, n.Name) if err != nil { @@ -708,6 +732,34 @@ func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation, return out, nil } +// deliveredVersions are the versions a module's registered build is delivered as: the last element +// of every resource path under a `versions/` directory, which registration filled from the artifact's +// digest (catalogue `${version}`). The node-engine names itself by that directory. +func deliveredVersions(m catalogue.Manifest) []string { + var out []string + for _, r := range m.Resources { + path, _ := r["path"].(string) + before, version, found := strings.Cut(path, "/versions/") + if !found || before == "" || version == "" || strings.Contains(version, "/") || strings.Contains(version, "$") { + continue + } + if !slices.Contains(out, version) { + out = append(out, version) + } + } + return out +} + +// engineBehind says a node-engine's reported version is not the build the mesh holds: neither the +// directory that build is delivered as, nor (for an engine placed by hand) its commit. A machine that +// has not said, or a mesh that holds no delivered build, is not behind anything. +func engineBehind(reported string, delivered []string, commit string) bool { + if reported == "" || len(delivered) == 0 { + return false + } + return !slices.Contains(delivered, reported) && !sameCommit(reported, commit) +} + // hostModule is the node-engine's module. const hostModule = "mesh-host" @@ -750,6 +802,10 @@ func probeWatchdogs(_ context.Context, d *doctor) ([]conditions.Observation, err // askSeatTool asks one machine's holder of a node seat a verb and answers its result; the holder's // own refusal is an error. func askSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string) (json.RawMessage, error) { + if who, declared := declaredBy(ctx, seat, verb); !declared { + return nil, fmt.Errorf("%s asks %s.%s, which it does not declare in the probe registry — and so the "+ + "controller is not granted it", who, seat, verb) + } answer, err := link.AskSeatTool(ctx, conn, seat, verb, node, map[string]any{}, 10*time.Second) if err != nil { return nil, err diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 3fb374a..00e9703 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -125,6 +125,16 @@ type Principal struct { // goes with the retired seat row. var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"} +// SeatVerb is one verb of one seat, on every machine holding it. +type SeatVerb struct{ Seat, Verb string } + +// VerbsTheSelfCheckAsks are the seat verbs the controller's self-check and watchdogs call (novox/hq +// to-be 45 §4): D8 reads every machine's ban list. **Named one by one, and the test that holds them +// to the probe registry is the reason they cannot drift** — a probe that calls a verb its grant does +// not name is refused by the bus on every run (found live on 2026-10-06: D8 timed out on each +// machine, refused). Asked of any machine (`.*`), read-only verbs, nothing else of the seat. +var VerbsTheSelfCheckAsks = []SeatVerb{{Seat: "node-intrusion-prevention", Verb: "banned"}} + // perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the // holder's machine (novox/hq ADR 0219): `paused.`, the build agent saying whether it takes work. var perMachineEvents = map[string]bool{"paused.*": true} @@ -266,6 +276,10 @@ func PermissionsFor(p Principal) (Permissions, error) { sub = append(sub, "mesh.seat."+ControllerSeat+".tool.>") // And says so (novox/hq ADR 0197): it answers discovery for the seat it serves. sub = append(sub, announcing(ControllerSeat)...) + // And the verbs the self-check reads with, of any machine's holder (novox/hq to-be 45 §4). + for _, v := range VerbsTheSelfCheckAsks { + pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb+".*") + } // And asks who answers (novox/hq to-be 45 §4, D3): the self-check finds every seat's holder by // the same discovery the console reads. The question only; the answers come to its own inbox. pub = append(pub, "$SRV.INFO") diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 344f363..f90ceab 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "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.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.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "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.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.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/link/advisories.go b/internal/link/advisories.go index f14f3ad..36fc6bb 100644 --- a/internal/link/advisories.go +++ b/internal/link/advisories.go @@ -225,3 +225,59 @@ func connectionAdvisory(sub *nats.Subscription, err error) (Advisory, bool) { } return Advisory{}, false } + +// Refusals of a publish, by subject, for a caller waiting on an answer to it. +var refusalWaiters = struct { + sync.Mutex + hooked map[*nats.Conn]bool + by map[string][]chan error +}{hooked: map[*nats.Conn]bool{}, by: map[string][]chan error{}} + +// refusalsOf is told when the bus refuses this connection a publish to subject, until stop is called. +// The connection's error handler is chained once, before whatever it had, which still runs. +func refusalsOf(conn *nats.Conn, subject string) (<-chan error, func()) { + ch := make(chan error, 1) + refusalWaiters.Lock() + defer refusalWaiters.Unlock() + if !refusalWaiters.hooked[conn] { + refusalWaiters.hooked[conn] = true + before := conn.ErrorHandler() + conn.SetErrorHandler(func(c *nats.Conn, sub *nats.Subscription, err error) { + if m := refusedPublish.FindStringSubmatch(errString(err)); m != nil { + refusalWaiters.Lock() + for _, w := range refusalWaiters.by[m[1]] { + select { + case w <- err: + default: + } + } + refusalWaiters.Unlock() + } + if before != nil { + before(c, sub, err) + } + }) + } + refusalWaiters.by[subject] = append(refusalWaiters.by[subject], ch) + return ch, func() { + refusalWaiters.Lock() + defer refusalWaiters.Unlock() + waiting := refusalWaiters.by[subject] + for i, w := range waiting { + if w == ch { + refusalWaiters.by[subject] = append(waiting[:i], waiting[i+1:]...) + break + } + } + if len(refusalWaiters.by[subject]) == 0 { + delete(refusalWaiters.by, subject) + } + } +} + +func errString(err error) string { + if err == nil { + return "" + } + return err.Error() +} diff --git a/internal/link/advisories_test.go b/internal/link/advisories_test.go index e3a5359..381d7ba 100644 --- a/internal/link/advisories_test.go +++ b/internal/link/advisories_test.go @@ -1,6 +1,14 @@ package link -import "testing" +import ( + "fmt" + "os" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" +) // **A reader's own consumer gone is not a consumer lost**: every bucket watch and stream read-back // makes and deletes one, many a minute, and the bus says so each time (found running the controller @@ -31,3 +39,54 @@ func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) { t.Fatalf("%+v", a) } } + +// **A question the bus refuses is answered as refused at once**, not after its timeout: against a +// server whose only user may not publish where it asks. +func TestNatsARefusedQuestionIsSaidAtOnce(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS_REFUSING") + if url == "" { + t.Skip("MESH_TEST_NATS_REFUSING unset: a server whose user may not publish mesh.seat.>") + } + conn, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + started := time.Now() + _, err = AskSeatTool(t.Context(), conn, "node-intrusion-prevention", "banned", "anchor", map[string]any{}, 10*time.Second) + if err == nil || !strings.Contains(err.Error(), "the bus refused") { + t.Fatalf("answered %v", err) + } + if took := time.Since(started); took > 3*time.Second { + t.Fatalf("a refusal took %s to be known", took) + } +} + +// The refusal reaches whoever waits on that subject, and only them. +func TestNatsARefusalReachesItsWaiter(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + conn, err := nats.Connect(url) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + mine, stop := refusalsOf(conn, "mesh.seat.s.tool.v.anchor") + defer stop() + other, stopOther := refusalsOf(conn, "mesh.seat.s.tool.v.laptop") + defer stopOther() + conn.ErrorHandler()(conn, nil, fmt.Errorf("%w: Permissions Violation for Publish to %q", + nats.ErrPermissionViolation, "mesh.seat.s.tool.v.anchor")) + select { + case <-mine: + case <-time.After(time.Second): + t.Fatal("the waiter was not told") + } + select { + case <-other: + t.Fatal("another subject's waiter was told") + default: + } +} diff --git a/internal/link/queue.go b/internal/link/queue.go index 254dc52..2322bb1 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -413,7 +413,30 @@ func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string, } asking, cancel := context.WithTimeout(ctx, timeout) defer cancel() - reply, err := conn.RequestWithContext(asking, NodeSeatToolSubject(seat, verb, node), body) + subject := NodeSeatToolSubject(seat, verb, node) + // **A refused question is known at once** (novox/hq to-be 45, D8 live on 2026-10-06): the bus + // says it to the connection and the request would otherwise wait out its whole timeout, reading + // as a machine that did not answer. + refused, stop := refusalsOf(conn, subject) + defer stop() + type replied struct { + msg *nats.Msg + err error + } + done := make(chan replied, 1) + go func() { + msg, err := conn.RequestWithContext(asking, subject, body) + done <- replied{msg, err} + }() + var reply *nats.Msg + select { + case r := <-done: + reply, err = r.msg, r.err + case why := <-refused: + cancel() + return Answer{}, fmt.Errorf("the bus refused the controller asking %s.%s of %s — its grants do not "+ + "name %s: %v", seat, verb, node, subject, why) + } switch { case errors.Is(err, nats.ErrNoResponders): return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+