Grant the self-check its ban-list question, say a refusal at once, judge the engine by its delivered version (hq to-be 45 Phase 1)
Live on 2026-10-06, two of the first self-check's findings were its own:
- D8 asked every machine's node-intrusion-prevention.banned, and the
controller's grant did not name the subject: the bus refused it 24 times
and D8 timed out after thirty seconds instead of saying so. The verbs the
self-check asks are named in broker.VerbsTheSelfCheckAsks and granted
(mesh.seat.<seat>.tool.<verb>.*); each probe declares the seat verbs it
calls, askSeatTool refuses an undeclared one, and a test over the
registry fails a probe whose question the controller is not granted.
AskSeatTool now returns a refused publish at once ("the bus refused…")
instead of waiting out its timeout; D8 asks the machines in parallel.
- D10 read every machine as behind right after a push: a node-engine says
its version as the directory it is delivered into, the archive's digest
(31045596c83a, catalogue versionOf), and D10 compared that with the
build's commit (1545b00a). It now compares with the versions the
registered build is delivered as, and a hand-placed engine's commit.
This commit is contained in:
@@ -61,7 +61,11 @@ type probe struct {
|
|||||||
Phase int
|
Phase int
|
||||||
// Deferred says why it is not run yet; empty for one that is.
|
// Deferred says why it is not run yet; empty for one that is.
|
||||||
Deferred string
|
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
|
// 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",
|
{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},
|
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",
|
{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",
|
{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},
|
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 " +
|
{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)
|
wg.Add(1)
|
||||||
go func(i int, p probe) {
|
go func(i int, p probe) {
|
||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
probing, cancel := context.WithTimeout(ctx, probeWithin)
|
probing, cancel := context.WithTimeout(context.WithValue(ctx, probeAsksKey{}, p), probeWithin)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
began := time.Now()
|
began := time.Now()
|
||||||
done := make(chan result, 1)
|
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() })
|
sort.Slice(obs, func(i, j int) bool { return obs[i].Key() < obs[j].Key() })
|
||||||
return obs
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ import (
|
|||||||
"golang.org/x/net/dns/dnsmessage"
|
"golang.org/x/net/dns/dnsmessage"
|
||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/broker"
|
"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/conditions"
|
||||||
"github.com/novox/mesh-controller/internal/link"
|
"github.com/novox/mesh-controller/internal/link"
|
||||||
)
|
)
|
||||||
@@ -345,3 +346,80 @@ func TestNatsTheBusSaysAConsumerGaveUpAndOneWasDeleted(t *testing.T) {
|
|||||||
t.Fatalf("%v", kinds)
|
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)
|
||||||
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"slices"
|
"slices"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
@@ -576,9 +577,22 @@ func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error)
|
|||||||
}
|
}
|
||||||
sort.Strings(holders)
|
sort.Strings(holders)
|
||||||
var out []conditions.Observation
|
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
|
var unasked []string
|
||||||
for _, node := range holders {
|
for i, node := range holders {
|
||||||
answer, err := askSeatTool(ctx, d.js.Conn(), "node-intrusion-prevention", "banned", node)
|
answer, err := answers[i], errs[i]
|
||||||
if err != nil {
|
if err != nil {
|
||||||
unasked = append(unasked, err.Error())
|
unasked = append(unasked, err.Error())
|
||||||
continue
|
continue
|
||||||
@@ -650,6 +664,11 @@ func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
shelf, err := inv.Catalogue(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
hostVersions := deliveredVersions(shelf[hostModule])
|
||||||
rolling := map[string]bool{}
|
rolling := map[string]bool{}
|
||||||
for _, p := range plans {
|
for _, p := range plans {
|
||||||
for m := range p.Modules {
|
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
|
continue // a machine not heard from is S1's
|
||||||
}
|
}
|
||||||
var behind []string
|
var behind []string
|
||||||
host := current[hostModule].Commit
|
// **A node-engine says its version as the directory it was delivered into** — its archive's
|
||||||
if host != "" && n.HostVersion != "" && !rolling[hostModule] && !sameCommit(n.HostVersion, host) {
|
// digest, twelve characters (catalogue.versionOf, ADR 0141) — not the commit it was built
|
||||||
behind = append(behind, fmt.Sprintf("the node-engine %s, the mesh holds %s", short(n.HostVersion), short(host)))
|
// 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)
|
assigned, err := inv.Assigned(ctx, n.Name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -708,6 +732,34 @@ func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation,
|
|||||||
return out, nil
|
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.
|
// hostModule is the node-engine's module.
|
||||||
const hostModule = "mesh-host"
|
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
|
// askSeatTool asks one machine's holder of a node seat a verb and answers its result; the holder's
|
||||||
// own refusal is an error.
|
// own refusal is an error.
|
||||||
func askSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string) (json.RawMessage, 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)
|
answer, err := link.AskSeatTool(ctx, conn, seat, verb, node, map[string]any{}, 10*time.Second)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
@@ -125,6 +125,16 @@ type Principal struct {
|
|||||||
// goes with the retired seat row.
|
// goes with the retired seat row.
|
||||||
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
|
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
|
// perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the
|
||||||
// holder's machine (novox/hq ADR 0219): `paused.<node>`, the build agent saying whether it takes work.
|
// holder's machine (novox/hq ADR 0219): `paused.<node>`, the build agent saying whether it takes work.
|
||||||
var perMachineEvents = map[string]bool{"paused.*": true}
|
var perMachineEvents = map[string]bool{"paused.*": true}
|
||||||
@@ -266,6 +276,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
sub = append(sub, "mesh.seat."+ControllerSeat+".tool.>")
|
sub = append(sub, "mesh.seat."+ControllerSeat+".tool.>")
|
||||||
// And says so (novox/hq ADR 0197): it answers discovery for the seat it serves.
|
// And says so (novox/hq ADR 0197): it answers discovery for the seat it serves.
|
||||||
sub = append(sub, announcing(ControllerSeat)...)
|
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
|
// 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.
|
// the same discovery the console reads. The question only; the answers come to its own inbox.
|
||||||
pub = append(pub, "$SRV.INFO")
|
pub = append(pub, "$SRV.INFO")
|
||||||
|
|||||||
+1
-1
@@ -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.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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
|
|||||||
@@ -225,3 +225,59 @@ func connectionAdvisory(sub *nats.Subscription, err error) (Advisory, bool) {
|
|||||||
}
|
}
|
||||||
return Advisory{}, false
|
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()
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,6 +1,14 @@
|
|||||||
package link
|
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
|
// **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
|
// 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)
|
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:
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+24
-1
@@ -413,7 +413,30 @@ func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string,
|
|||||||
}
|
}
|
||||||
asking, cancel := context.WithTimeout(ctx, timeout)
|
asking, cancel := context.WithTimeout(ctx, timeout)
|
||||||
defer cancel()
|
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 {
|
switch {
|
||||||
case errors.Is(err, nats.ErrNoResponders):
|
case errors.Is(err, nats.ErrNoResponders):
|
||||||
return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+
|
return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+
|
||||||
|
|||||||
Reference in New Issue
Block a user