diff --git a/cmd/mesh-controller/busobjects.go b/cmd/mesh-controller/busobjects.go index ffa46f2..71c9943 100644 --- a/cmd/mesh-controller/busobjects.go +++ b/cmd/mesh-controller/busobjects.go @@ -2,9 +2,11 @@ package main import ( "context" + "errors" "fmt" "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/conditions" "github.com/novox/mesh-controller/internal/inventory" ) @@ -51,14 +53,70 @@ func assertBusObjects(ctx context.Context, inv *inventory.Inventory, r broker.Ra if err != nil { return nil, err } + // Every one tried, and every failure named: one module's consumer the bus refuses is no reason + // the modules after it in the list hear nothing (novox/hq issue 208, where this runs on each send). + var failed []error for _, c := range consumers { if err := r.EnsureConsumer(c.Consumer); err != nil { - return nil, fmt.Errorf("how %s on %s hears what it consumes: %w", c.Module, c.Node, err) + failed = append(failed, fmt.Errorf("how %s on %s hears what it consumes: %w", c.Module, c.Node, err)) } } + if len(failed) > 0 { + return nil, errors.Join(failed...) + } return names, nil } +// What raises the condition a send says when the objects it implies could not be asserted, and its kind. +const ( + sourceBusObjects = "bus-objects" + kindBusObjectsUnasserted = "bus-objects-unasserted" +) + +// assertOnSend asserts, on the bus a send is about to use, every object assertBusObjects derives — +// **whenever a declaration is sent, not only when the controller starts** (novox/hq issue 208). +// +// A module assigned after the controller started was sent its declaration and found no consumer to +// bind (`consumer not found`, messenger on 2026-10-06), and a seat holder assigned after it found no +// worker: the objects a declaration implies were asserted at start and nowhere else, so they existed +// only for what was assigned before the last restart. The same derivation, not a second list of what +// a send needs: what start asserts, the self-check expects and a send asserts are one answer. Every +// part is idempotent, so asserting the whole of it again is the no-op a restart already relies on. +// +// **A failure is said and raised, and the send goes on.** The objects are the mesh's, not the +// machines' being sent: holding every machine back for one consumer that none of them may use would +// turn one fault into all of them, and the declarations are not what is wrong. It is never silent — +// said in the send's own output and raised as a condition, which the next send that asserts them +// clears — and the start-time raise still refuses to serve without them. +func assertOnSend(ctx context.Context, inv *inventory.Inventory, r broker.Raiser, indent string) error { + _, err := assertBusObjects(ctx, inv, r) + var observed []conditions.Observation + if err != nil { + fmt.Printf("%sTHE BUS DOES NOT HOLD WHAT THIS SEND IMPLIES: %v\n", indent, err) + fmt.Printf("%s a module may find no consumer to bind, or a holder no worker; sent anyway, raised as "+ + "condition %s, and asserted again by the next send\n", indent, unassertedObservation(err).Key()) + observed = append(observed, unassertedObservation(err)) + } + // Observed when it failed, cleared when it did not: a send that asserted everything is the + // observation that the bus holds what it should. + if kerr := withKeeper(ctx, func(k *conditions.Keeper) error { + return k.Reconcile(ctx, sourceBusObjects, observed) + }); kerr != nil { + fmt.Printf("%sand whether the bus holds what this send implies could not be kept as a condition: %v\n", + indent, kerr) + } + return err +} + +// unassertedObservation is a send's failure to assert the bus's objects, as a condition. +func unassertedObservation(err error) conditions.Observation { + return conditions.Observation{Scope: conditions.ScopeBus, ID: "objects", Token: "unasserted", + Kind: kindBusObjectsUnasserted, Severity: conditions.Warning, Source: sourceBusObjects, + Summary: "the bus's streams and consumers could not be asserted when a declaration was sent: " + + "a module may find no consumer to bind, or a seat's holder no worker", + Said: err.Error()} +} + // moduleConsumers is every module's durable consumer, from the records the user list is composed from. func moduleConsumers(ctx context.Context, inv *inventory.Inventory) ([]broker.ModuleConsumer, error) { records, err := inv.BusRecords(ctx) diff --git a/cmd/mesh-controller/busobjects_test.go b/cmd/mesh-controller/busobjects_test.go new file mode 100644 index 0000000..f3c7199 --- /dev/null +++ b/cmd/mesh-controller/busobjects_test.go @@ -0,0 +1,199 @@ +package main + +import ( + "errors" + "os" + "strings" + "testing" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// novox/hq issue 208: the bus's objects a declaration implies are asserted whenever one is sent, not +// only when the controller starts. + +// aCarriedConsumer is the runtime on laptop and, carried by it, a module that consumes an event: the +// shape of messenger on 2026-10-06, assigned after the controller started. +func aCarriedConsumer(t *testing.T, open *stores) broker.Consumer { + t.Helper() + register(t, open, catalogue.Manifest{Module: catalogue.RuntimeModule, Version: "1", + OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/var/lib/mesh/node-tools/broker"}}}) + register(t, open, catalogue.Manifest{Module: "messenger", Version: "1", + Consumes: []string{"billing.order.placed"}}) + for _, m := range []string{catalogue.RuntimeModule, "messenger"} { + if _, err := open.inventory.Assign(t.Context(), "laptop", m); err != nil { + t.Fatal(err) + } + } + consumers, err := moduleConsumers(t.Context(), open.inventory) + if err != nil { + t.Fatal(err) + } + for _, c := range consumers { + if c.Module == "messenger" && c.Node == "laptop" { + return c.Consumer + } + } + t.Fatalf("messenger on laptop is derived no consumer: %+v", consumers) + return broker.Consumer{} +} + +// aLateHolder is a module holding the build agent's seat on anchor, assigned after the controller +// started, and the worker its seat's queue should have for it. +func aLateHolder(t *testing.T, open *stores) broker.Consumer { + t.Helper() + register(t, open, catalogue.Manifest{Module: "late-builder", Version: "1", + Claims: []catalogue.Claim{{Name: "node-build-agent", Scope: catalogue.ScopeNode}}}) + if _, err := open.inventory.Assign(t.Context(), "anchor", "late-builder"); err != nil { + t.Fatal(err) + } + for _, s := range inventory.MeshSeats() { + if s.Name == "node-build-agent" { + c, needed := broker.HolderConsumerFor("anchor", "late-builder", s) + if !needed { + t.Fatal("the build agent's seat needs no worker") + } + return c + } + } + t.Fatal("the mesh declares no build agent's seat") + return broker.Consumer{} +} + +// asserted says whether a recording holds a consumer by stream and name. +func asserted(rec *recordingRaiser, want broker.Consumer) bool { + for _, c := range rec.consumers { + if c.Stream == want.Stream && c.Name == want.Name { + return true + } + } + return false +} + +// What a send asserts includes what was assigned after the start's assertion: a carried module's +// consumer and a late holder's worker. The derivation is the start's own, so this holds without a bus. +func TestASendAssertsWhatWasAssignedAfterStart(t *testing.T) { + open := aMesh(t) + var atStart recordingRaiser + if _, err := assertBusObjects(t.Context(), open.inventory, &atStart); err != nil { + t.Fatal(err) + } + consumer := aCarriedConsumer(t, open) + worker := aLateHolder(t, open) + if asserted(&atStart, consumer) || asserted(&atStart, worker) { + t.Fatal("the start asserted what was not yet assigned; the test proves nothing") + } + var onSend recordingRaiser + if err := assertOnSend(t.Context(), open.inventory, &onSend, ""); err != nil { + t.Fatal(err) + } + if !asserted(&onSend, consumer) { + t.Fatalf("messenger's consumer %s on %s is not asserted by the send: %+v", consumer.Name, consumer.Stream, onSend.consumers) + } + if !asserted(&onSend, worker) { + t.Fatalf("the late holder's worker %s on %s is not asserted by the send: %+v", worker.Name, worker.Stream, onSend.consumers) + } +} + +// failingRaiser refuses one consumer by name and records everything it was asked. +type failingRaiser struct { + recordingRaiser + refuse string +} + +func (f *failingRaiser) EnsureConsumer(c broker.Consumer) error { + f.consumers = append(f.consumers, c) + if c.Name == f.refuse { + return errors.New("nats: API error: code=503 description=insufficient resources") + } + return nil +} + +// A send whose objects cannot be asserted says so in its output and raises a condition — and the next +// send that asserts them clears it. The consumers after the refused one are still asked for. +func TestASendThatCannotAssertTheBusSaysSoAndRaisesACondition(t *testing.T) { + open := aMesh(t) + consumer := aCarriedConsumer(t, open) + worker := aLateHolder(t, open) + failing := &failingRaiser{refuse: consumer.Name} + + var sendErr error + out := stdoutOf(t, func() error { + sendErr = assertOnSend(t.Context(), open.inventory, failing, " ") + return nil + }) + if sendErr == nil || !strings.Contains(sendErr.Error(), "how messenger on laptop hears what it consumes") { + t.Fatalf("the failure is not answered: %v", sendErr) + } + if !strings.Contains(out, "THE BUS DOES NOT HOLD WHAT THIS SEND IMPLIES") || + !strings.Contains(out, "insufficient resources") || !strings.Contains(out, "bus.objects.unasserted") { + t.Fatalf("the failure is not said in the send's output:\n%s", out) + } + if !asserted(&failing.recordingRaiser, worker) { + t.Fatal("the seat's worker was not asked for") + } + c, open1, err := conditionsFrom.Get(t.Context(), "bus.objects.unasserted") + if err != nil || !open1 { + t.Fatalf("no condition raised: %v", err) + } + if c.Kind != kindBusObjectsUnasserted || c.Source != sourceBusObjects || + len(c.Evidence) == 0 || !strings.Contains(c.Evidence[0].Said, "insufficient resources") { + t.Fatalf("the condition does not say what failed: %+v", c) + } + + // The next send asserts them, and that observation clears it. + if err := assertOnSend(t.Context(), open.inventory, &recordingRaiser{}, ""); err != nil { + t.Fatal(err) + } + if _, stillOpen, err := conditionsFrom.Get(t.Context(), "bus.objects.unasserted"); err != nil || stillOpen { + t.Fatalf("a send that asserted everything left the condition open: %v", err) + } +} + +// Against a real bus, through the send's own grant: a module assigned after start has its consumer +// once a declaration is sent, and a holder assigned after start its seat's worker. +func TestNatsAModuleAssignedAfterStartGetsItsConsumerAtItsFirstSend(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + open := aMesh(t) + js, err := broker.Dial(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(js.Close) + for _, s := range []string{"CONTROL", "NODES", "ASSIGNMENTS", "EVENTS"} { + _ = js.Context().DeleteStream(s) + } + // The controller starting, before either was assigned. + if _, err := assertBusObjects(t.Context(), open.inventory, js); err != nil { + t.Fatal(err) + } + consumer := aCarriedConsumer(t, open) + worker := aLateHolder(t, open) + _ = js.Context().DeleteConsumer(worker.Stream, worker.Name) + for _, c := range []broker.Consumer{consumer, worker} { + if _, err := js.Context().ConsumerInfo(c.Stream, c.Name); !errors.Is(err, nats.ErrConsumerNotFound) { + t.Fatalf("%s on %s exists before any send; the test proves nothing: %v", c.Name, c.Stream, err) + } + } + + server := link.ConnectNats(js, nil, nil) + if err := (overTheBus{open: open, server: server}).grant(t.Context(), nil); err != nil { + t.Fatal(err) + } + for _, c := range []broker.Consumer{consumer, worker} { + if _, err := js.Context().ConsumerInfo(c.Stream, c.Name); err != nil { + t.Fatalf("%s on %s is not on the bus after a send: %v", c.Name, c.Stream, err) + } + } + if _, raised, _ := conditionsFrom.Get(t.Context(), "bus.objects.unasserted"); raised { + t.Fatal("a send that asserted everything raised a condition") + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 764bf93..d034370 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -776,6 +776,13 @@ type overTheBus struct { } func (b overTheBus) grant(ctx context.Context, sending []readyNode) error { + // The bus's objects first, which every declaration implies (novox/hq issue 208): said and raised + // when they cannot be, and never what holds the send back (assertOnSend). + if bus, ok := b.server.Bus().(link.OverNATS); ok { + js := broker.OnConn(bus.Conn) + js.Note = func(format string, args ...any) { fmt.Printf(b.indent+" "+format+"\n", args...) } + _ = assertOnSend(ctx, b.open.inventory, js, b.indent) + } return issueMemberships(ctx, b.open, b.server, sending) }