From 8d9d33ae85865914636e98a2d7a38dbee7cfbf11 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 00:13:22 +0200 Subject: [PATCH] Report a provider that keeps failing a consumer in status (hq ADR 0224) The identity provider failed every consumer for a day and status called the mesh well (hq issue 179). The controller now follows every provider's provisioner.failing/recovered, keeps the newest failing word per provider, machine and consumer (migration 0065), and status, its JSON and node show name it until it recovers. Every module that receives contributions is granted the two events, so no manifest can forget them. --- cmd/mesh-controller/build.go | 3 + cmd/mesh-controller/modules.go | 2 +- cmd/mesh-controller/nodes.go | 13 ++ cmd/mesh-controller/push.go | 5 + cmd/mesh-controller/readable.go | 6 + cmd/mesh-controller/rollout.go | 2 +- cmd/mesh-controller/standing.go | 120 +++++++++++ cmd/mesh-controller/standing_test.go | 115 ++++++++++ cmd/mesh-controller/status.go | 12 +- internal/broker/nats.go | 5 +- internal/broker/streams.go | 17 +- internal/broker/testdata/composed.conf | 2 +- internal/catalogue/events.go | 32 +++ internal/inventory/busrecords.go | 2 +- ...r-says-which-consumer-it-keeps-failing.sql | 20 ++ internal/inventory/standing.go | 76 +++++++ internal/inventory/standing_test.go | 130 +++++++++++ internal/link/receive.go | 8 + internal/link/receive_fake_test.go | 2 + internal/link/receive_nats.go | 28 ++- internal/link/serve.go | 7 + internal/link/standing.go | 114 ++++++++++ internal/link/standing_test.go | 203 ++++++++++++++++++ 23 files changed, 913 insertions(+), 11 deletions(-) create mode 100644 cmd/mesh-controller/standing.go create mode 100644 cmd/mesh-controller/standing_test.go create mode 100644 internal/inventory/migrations/0065-a-provider-says-which-consumer-it-keeps-failing.sql create mode 100644 internal/inventory/standing.go create mode 100644 internal/inventory/standing_test.go create mode 100644 internal/link/standing.go create mode 100644 internal/link/standing_test.go diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 915fcfc..e68b09b 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -676,6 +676,9 @@ type answers struct { // refused, until the switch — and while there is any, the mesh is not all well: the order the // machines' modules are built in is the mesh's to keep, and this is where it says it is not kept. unheld []catalogue.Unheld + // failing is every consumer a provider says it keeps failing (novox/hq ADR 0224): a provider's + // journal was the only place that said so for a day (04-ISSUES/179). + failing []inventory.ProviderStanding } // heldBy is every artifact this mesh has built, for a build that may need one as its base. diff --git a/cmd/mesh-controller/modules.go b/cmd/mesh-controller/modules.go index 2042c6a..6767884 100644 --- a/cmd/mesh-controller/modules.go +++ b/cmd/mesh-controller/modules.go @@ -661,7 +661,7 @@ func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife // durable subscription nobody reads. if consumer, needed := broker.ConsumerFor(broker.Principal{ Kind: broker.KindModule, Node: node, Module: m.Module, - Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools, + Emits: m.EmitsAll(), Consumes: m.Consumes, Serves: m.Tools, }); needed { if busAddress == "" { fmt.Printf(" %s consumes; its consumer is created when the bus is reachable (`push`, then "+ diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index f025414..acfa154 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -505,6 +505,19 @@ func showNode(ctx context.Context, inv *inventory.Inventory, name string) error fmt.Printf(" public domain %s\n", domain) } + // A provider here failing a consumer, or a consumer here failed (novox/hq ADR 0224). Before the + // capabilities, because it is something not working now and they are a description. + failing, err := failingProviders(ctx, inv) + if err != nil { + return err + } + if here := failingOn(failing, name); len(here) > 0 { + fmt.Printf("\n %d consumer(s) a provider keeps failing, here or for a module here:\n", len(here)) + for _, line := range failingLines(here, time.Now()) { + fmt.Printf(" %s\n", line) + } + } + held, err := inv.Profile(ctx, name) if err != nil { return err diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 94cddfd..dc8d5bb 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -132,6 +132,11 @@ func serve(ctx context.Context) error { if err := server.Answers(following{open}); err != nil { return err } + // And what providers say about consumers they keep failing, kept for `status` (novox/hq ADR + // 0224): a provider's journal must not be the only place that says so. + if err := server.Watches(standings{inv}); err != nil { + return err + } // And the mesh's own verbs, as the seat this control plane holds (novox/hq ADR 0154). Served // from the store's row, so what the seat declares is what is answered. diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index 119f9f7..6b58a72 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -78,6 +78,11 @@ type meshStatus struct { // that machine holds, with the modules that could hold it (novox/hq ADR 0207). Absent when every // dependency is met. Reported, not refused, until the switch. Unheld []catalogue.Unheld `json:"unheld,omitempty"` + // Failing is every consumer a provider says it keeps failing, with the class of error, since + // when, and when it was last said (novox/hq ADR 0224). Absent when no provider says so. A + // document without this called the mesh well while the identity provider refused every consumer + // for a day (04-ISSUES/179). + Failing []inventory.ProviderStanding `json:"failing,omitempty"` } // machineFiltered is one rule set on a converged machine that the mesh did not write and that @@ -210,6 +215,7 @@ func statusAsJSON(asked answers) ([]byte, error) { } } out.Unheld = asked.unheld + out.Failing = asked.failing for name := range asked.refused { out.Unresolved = append(out.Unresolved, machineUnresolved{ Node: name, Problem: asked.refused[name]}) diff --git a/cmd/mesh-controller/rollout.go b/cmd/mesh-controller/rollout.go index 8903c8c..c4608ee 100644 --- a/cmd/mesh-controller/rollout.go +++ b/cmd/mesh-controller/rollout.go @@ -189,7 +189,7 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines // A third of the catalogue never does (novox/hq ADR 0120), and counting those as missing a credential // would bury the ones that matter under a list nobody can act on. func speaksOnTheBus(m catalogue.Manifest) bool { - return len(m.Emits) > 0 || len(m.Consumes) > 0 || len(m.Tools) > 0 || + return len(m.EmitsAll()) > 0 || len(m.Consumes) > 0 || len(m.Tools) > 0 || len(m.DefinesSeats) > 0 || len(m.Uses) > 0 || len(m.Claims) > 0 } diff --git a/cmd/mesh-controller/standing.go b/cmd/mesh-controller/standing.go new file mode 100644 index 0000000..24cd742 --- /dev/null +++ b/cmd/mesh-controller/standing.go @@ -0,0 +1,120 @@ +package main + +import ( + "context" + "fmt" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A provider that keeps failing a consumer is a problem the controller reports (novox/hq ADR 0224). +// +// On 2026-10-05 the identity provider's provisioner failed every consumer from shortly after midnight +// until it was fixed by hand that night — 31,000 refused logins after its database was moved and its +// admin kept an older password — and `status` called the mesh well all day (novox/hq issue 179). A +// provider now announces a consumer it has failed for minutes; the controller keeps it until the +// provider says it recovered; and `status`, its JSON and `node show` name it, breaking "all well". + +// standings keeps what providers say, in the inventory. +type standings struct{ inv *inventory.Inventory } + +func (s standings) Stood(ctx context.Context, st link.Standing) (bool, error) { + return s.inv.KeepStanding(ctx, st.Failing, inventory.ProviderStanding{ + Module: st.Module, ProviderNode: st.ProviderNode, Provision: st.Provider, + Consumer: st.Consumer, ConsumerNode: st.Node, + Class: st.Class, Error: st.Error, Since: st.Since, Attempts: st.Attempts, + }) +} + +// failingProviders is every consumer a provider still assigned where it ran says it keeps failing. +// +// **A provider no longer assigned is not asked about.** Its last word stays in the store, and is +// not a problem: nothing runs there to fail anybody. Assigned again, its first success for each +// consumer clears it. +func failingProviders(ctx context.Context, inv *inventory.Inventory) ([]inventory.ProviderStanding, error) { + all, err := inv.FailingProviders(ctx) + if err != nil { + return nil, fmt.Errorf("what providers say they keep failing cannot be read: %w", err) + } + assigned := map[string]map[string]bool{} + var out []inventory.ProviderStanding + for _, s := range all { + on, asked := assigned[s.ProviderNode] + if !asked { + modules, err := inv.Assigned(ctx, s.ProviderNode) + if err != nil { + // A provider on a machine the mesh no longer knows has nothing running to fail anybody. + modules = nil + } + on = map[string]bool{} + for _, m := range modules { + on[m] = true + } + assigned[s.ProviderNode] = on + } + if on[s.Module] { + out = append(out, s) + } + } + return out, nil +} + +// failingLines is how status says them: one consumer per entry, the error under it, and a provider +// that stopped repeating itself said so. +func failingLines(list []inventory.ProviderStanding, now time.Time) []string { + var out []string + for _, s := range list { + where := s.Module + if s.ProviderNode != "" { + where += " on " + s.ProviderNode + } + whom := s.Consumer + if s.ConsumerNode != "" { + whom += " (" + s.ConsumerNode + ")" + } + out = append(out, fmt.Sprintf(" %-24s fails %s: %s, for %s (%d attempts since %s)", + where, whom, orUnclassed(s.Class), roughly(now.Sub(s.Since)), s.Attempts, + s.Since.Local().Format("2006-01-02 15:04"))) + if e := strings.TrimSpace(s.Error); e != "" { + out = append(out, fmt.Sprintf(" %-24s %s", "", firstLine(e))) + } + if s.Quiet(now) { + out = append(out, fmt.Sprintf(" %-24s not said again for %s — the provider has stopped "+ + "saying anything, so this is its last word", "", roughly(now.Sub(s.SaidAt)))) + } + } + return out +} + +func orUnclassed(class string) string { + if class == "" { + return "failing" + } + return class +} + +// printFailing is the status section, said when there is anything to say. +func printFailing(list []inventory.ProviderStanding, now time.Time) { + if len(list) == 0 { + return + } + fmt.Printf("%d consumer(s) a provider keeps failing (ADR 0224):\n\n", len(list)) + for _, line := range failingLines(list, now) { + fmt.Println(line) + } + fmt.Printf("\n the provider's journal has every attempt; it says recovered on its next success\n\n") +} + +// failingOn is the standings that concern one machine: a provider running there, or a consumer. +func failingOn(list []inventory.ProviderStanding, node string) []inventory.ProviderStanding { + var out []inventory.ProviderStanding + for _, s := range list { + if s.ProviderNode == node || s.ConsumerNode == node { + out = append(out, s) + } + } + return out +} diff --git a/cmd/mesh-controller/standing_test.go b/cmd/mesh-controller/standing_test.go new file mode 100644 index 0000000..6a32972 --- /dev/null +++ b/cmd/mesh-controller/standing_test.go @@ -0,0 +1,115 @@ +package main + +import ( + "encoding/json" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A provider that keeps failing a consumer is a problem `status` names (novox/hq ADR 0224). On +// 2026-10-05 the identity provider refused every consumer for a day and status called the mesh well +// (04-ISSUES/179): this is that day, told to the controller the way the provider now tells it. +func TestAProviderFailingAConsumerBreaksAllWellUntilItRecovers(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + register(t, open, catalogue.Manifest{Module: "idp", Version: "1", + Receives: map[string]string{"oidc-client": "/var/lib/mesh/idp/mesh.json"}}) + if _, err := assign(ctx, open, "anchor", "idp"); err != nil { + t.Fatal(err) + } + kept := standings{open.inventory} + since := time.Now().Add(-23 * time.Hour) + failing := link.Standing{Module: "idp", Failing: true, Provider: "oidc-client", ProviderNode: "anchor", + Consumer: "mesh_laptop_dashboard", Node: "laptop", Class: "credentials-rejected", + Error: `Keycloak token request failed: 401 {"error":"invalid_grant"}`, Since: since, Attempts: 31000} + if _, err := kept.Stood(ctx, failing); err != nil { + t.Fatal(err) + } + + asked, err := theThreeQuestions(ctx, open) + if err != nil { + t.Fatal(err) + } + if asked.well() { + t.Fatal("a mesh whose identity provider fails a consumer reads as well") + } + said := printed(t, func() error { return printStatus(asked) }) + for _, want := range []string{"1 consumer(s) a provider keeps failing", "idp on anchor", + "mesh_laptop_dashboard (laptop)", "credentials-rejected", "31000 attempts", "invalid_grant"} { + if !strings.Contains(said, want) { + t.Fatalf("status does not say %q:\n%s", want, said) + } + } + if strings.Contains(said, "all doing what they were told") { + t.Fatalf("status said all well beside a failing provider:\n%s", said) + } + body, err := statusAsJSON(asked) + if err != nil { + t.Fatal(err) + } + var doc struct { + Failing []inventory.ProviderStanding `json:"failing"` + } + if err := json.Unmarshal(body, &doc); err != nil || len(doc.Failing) != 1 || doc.Failing[0].Consumer != "mesh_laptop_dashboard" { + t.Fatalf("the document does not carry it: %v\n%s", err, body) + } + // Both machines' `node show` name it: where the provider runs, and where the consumer is. + for _, node := range []string{"anchor", "laptop"} { + shown := printed(t, func() error { return showNode(ctx, open.inventory, node) }) + if !strings.Contains(shown, "a provider keeps failing") || !strings.Contains(shown, "mesh_laptop_dashboard") { + t.Fatalf("node show %s does not name it:\n%s", node, shown) + } + } + + // Recovered: gone, and the mesh may be well again as far as this is concerned. + failing.Failing = false + if cleared, err := kept.Stood(ctx, failing); err != nil || !cleared { + t.Fatalf("%v %v", cleared, err) + } + asked, err = theThreeQuestions(ctx, open) + if err != nil { + t.Fatal(err) + } + if len(asked.failing) != 0 { + t.Fatalf("a recovered consumer is still named: %+v", asked.failing) + } +} + +// A provider no longer assigned where it ran has nothing running to fail anybody: its last word is +// not a problem. +func TestAnUnassignedProvidersLastWordIsNotAProblem(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + if _, err := (standings{open.inventory}).Stood(ctx, link.Standing{Module: "gone", Failing: true, + ProviderNode: "anchor", Consumer: "x", Since: time.Now()}); err != nil { + t.Fatal(err) + } + asked, err := theThreeQuestions(ctx, open) + if err != nil { + t.Fatal(err) + } + if len(asked.failing) != 0 { + t.Fatalf("%+v", asked.failing) + } +} + +func TestAProviderThatStoppedRepeatingItselfIsSaidToHaveGoneQuiet(t *testing.T) { + now := time.Now() + lines := strings.Join(failingLines([]inventory.ProviderStanding{{ + Module: "idp", ProviderNode: "anchor", Consumer: "c", Class: "unreachable", Error: "connection refused\nmore", + Since: now.Add(-3 * time.Hour), SaidAt: now.Add(-2 * time.Hour), Attempts: 9, + }}, now), "\n") + for _, want := range []string{"unreachable, for 3h", "connection refused", "not said again for 2h"} { + if !strings.Contains(lines, want) { + t.Fatalf("%q not in:\n%s", want, lines) + } + } + if strings.Contains(lines, "more") { + t.Fatalf("more than the first line of an error:\n%s", lines) + } +} diff --git a/cmd/mesh-controller/status.go b/cmd/mesh-controller/status.go index 7cf2141..34af279 100644 --- a/cmd/mesh-controller/status.go +++ b/cmd/mesh-controller/status.go @@ -129,6 +129,10 @@ func printStatus(asked answers) error { fmt.Println() } + // A provider failing a consumer, beside machines failing what they were told: both are something + // not working now (novox/hq ADR 0224). + printFailing(asked.failing, time.Now()) + if len(quiet) > 0 { var said []string for _, n := range quiet { @@ -416,6 +420,12 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { } out.unheld = append(out.unheld, plan.Unheld...) } + // And every consumer a provider says it keeps failing (novox/hq ADR 0224). Read from what the + // providers announced: nothing else in the mesh knows whether a provision is being made. + out.failing, err = failingProviders(ctx, inv) + if err != nil { + return answers{}, err + } out.plans, err = inv.RecentPlans(ctx, 5) if err != nil { return answers{}, err @@ -528,7 +538,7 @@ func untakenModules(ctx context.Context, inv *inventory.Inventory, nodes []inven func (a answers) well() bool { return len(a.wrong) == 0 && len(a.quiet) == 0 && len(a.behind) == 0 && len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 && - len(a.filtered) == 0 && len(a.unheld) == 0 + len(a.filtered) == 0 && len(a.unheld) == 0 && len(a.failing) == 0 } // hostSplit is which machines report which host version, for every version more than one machine diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 2cca06d..8c7d9b2 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -253,10 +253,11 @@ func PermissionsFor(p Principal) (Permissions, error) { // And says so (novox/hq ADR 0197): it answers discovery for the seat it serves. sub = append(sub, announcing(ControllerSeat)...) - // The two events it reacts to, and its ack subject on the stream they arrive from + // The events it reacts to, and its ack subject on the stream they arrive from // (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the // controller a subscriber to every event in the mesh, and its permission list would stop - // saying what it is for. The ack grant below is scoped per stream because the controller's + // saying what it is for. The one wildcard is the emitter of a provider's standing (ADR + // 0224) — still two named events, from whichever module provides. The ack grant below is scoped per stream because the controller's // consumer name is the same on both and `$JS.ACK.CONTROL.controller.>` does not cover a // delivery from EVENTS — a consumer that cannot ack has every message redelivered for // ever, refused by the list it already has. diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 87a32fa..cfd7e52 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -226,8 +226,23 @@ var ControllerFollows = []string{ // registers build-agent itself comes from there. Appended, for the same reason as above; goes // with the retired seat row. seatEventSubject("mesh-build-machine", "built"), + // **Every provider's standing** (novox/hq ADR 0224): a consumer it has failed for minutes, and + // that consumer recovered. The one pattern on this list, and a narrow one — two named events, + // from whichever module provides — because the rule is about every provider, and a list of + // providers here would be a list somebody forgets to extend. On 2026-10-05 the identity provider + // failed every consumer for a day and only its journal said so (issue 179). Appended, because + // the index is a name. + moduleEventSubject("*", ProvisionerFailing), + moduleEventSubject("*", ProvisionerRecovered), } +// The provider standing events, by their local names. Written here as well as in the catalogue +// (catalogue.ProvisionerEvents), which this package cannot import; a test keeps them agreeing. +const ( + ProvisionerFailing = "provisioner.failing" + ProvisionerRecovered = "provisioner.recovered" +) + // moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so // what the controller subscribes and what the emitter is permitted to publish cannot drift apart. func moduleEventSubject(module, event string) string { @@ -271,7 +286,7 @@ func MeshConsumers() []Consumer { // client and come back to be acted on again. MaxAckPending: 1, FromNow: true, - Why: "the two events the mesh's own controller reacts to, one at a time; after " + + Why: "the events the mesh's own controller reacts to, one at a time; after " + "max-deliver it dead-letters, because an announcement it cannot act on will not " + "become actionable", }, diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index af658e6..7eaf3a6 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -25,7 +25,7 @@ accounts { users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_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.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] } - subscribe: { allow: ["$JS.API.>", "$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.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.>", "$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" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/catalogue/events.go b/internal/catalogue/events.go index 90dda63..2db9336 100644 --- a/internal/catalogue/events.go +++ b/internal/catalogue/events.go @@ -3,6 +3,7 @@ package catalogue import ( "fmt" "regexp" + "slices" "strings" ) @@ -159,3 +160,34 @@ func consumePattern(pattern string) error { } return nil } + +// The events a provider says about its consumers (novox/hq ADR 0224): a consumer it has failed +// without one success for minutes, and that consumer succeeding again or being withdrawn. The +// controller follows them from every module and `status` names a consumer failing until it recovers. +const ( + ProvisionerFailing = "provisioner.failing" + ProvisionerRecovered = "provisioner.recovered" +) + +// ProvisionerEvents are both, in the order they are said. +var ProvisionerEvents = []string{ProvisionerFailing, ProvisionerRecovered} + +// EmitsAll is every event a module may publish: what it declares and, for a module that receives +// contributions — a provider, running a provisioner over them — the provider's standing events. +// +// **Derived, not declared**, because they are the mesh's rule about every provider rather than +// anything one module chose to say: a provider whose manifest forgot them would fail its consumers +// as silently as on 2026-10-05, with its announcement refused by the bus (novox/hq issue 179). Every +// grant of a module's publishing reads this, never the declared list alone. +func (m Manifest) EmitsAll() []string { + out := append([]string(nil), m.Emits...) + if len(m.Receives) == 0 { + return out + } + for _, e := range ProvisionerEvents { + if !slices.Contains(out, e) { + out = append(out, e) + } + } + return out +} diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index b1039bd..8de4341 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -124,7 +124,7 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio d := broker.Declared{ Module: m.Module, - Emits: m.Emits, + Emits: m.EmitsAll(), Consumes: fromModules, Watches: watches, // The tools it answers, which is `tools` and not `serves`: the manifest's `serves` is the diff --git a/internal/inventory/migrations/0065-a-provider-says-which-consumer-it-keeps-failing.sql b/internal/inventory/migrations/0065-a-provider-says-which-consumer-it-keeps-failing.sql new file mode 100644 index 0000000..2673bfb --- /dev/null +++ b/internal/inventory/migrations/0065-a-provider-says-which-consumer-it-keeps-failing.sql @@ -0,0 +1,20 @@ +-- A provider says which consumer it keeps failing (novox/hq ADR 0224). +-- +-- On 2026-10-05 the identity provider's provisioner failed every consumer 31,000 times in a day and +-- only its journal said so (issue 179). A provider now announces a consumer it has failed for minutes +-- without one success, and the consumer recovering; the controller keeps the newest failing word per +-- provider module, the machine it runs on and the consumer, and removes it on recovery. `status` and +-- `node show` read this table: a row is a problem until it is gone. +create table provider_standing ( + module text not null, + provider_node text not null, + consumer text not null, + consumer_node text not null default '', + provision text not null default '', + class text not null default '', + error text not null default '', + since timestamptz not null, + attempts integer not null default 0, + said_at timestamptz not null default now(), + primary key (module, provider_node, consumer) +); diff --git a/internal/inventory/standing.go b/internal/inventory/standing.go new file mode 100644 index 0000000..8e29369 --- /dev/null +++ b/internal/inventory/standing.go @@ -0,0 +1,76 @@ +package inventory + +import ( + "context" + "time" +) + +// ProviderStanding is a consumer a provider says it keeps failing (novox/hq ADR 0224). +type ProviderStanding struct { + // Module is the provider's module, and ProviderNode the machine it runs on. + Module string `json:"module"` + ProviderNode string `json:"provider-node"` + // Provision is the interface it provides, e.g. `oidc-client`. + Provision string `json:"provision"` + // Consumer is the identity the mesh derived for the consumer, ConsumerNode its machine. + Consumer string `json:"consumer"` + ConsumerNode string `json:"consumer-node"` + // Class is what kind of failure: credentials-rejected, unreachable, secret-unreadable, refused. + Class string `json:"class"` + Error string `json:"error"` + // Since is when the unbroken run of failures began; Attempts how many it has been. + Since time.Time `json:"since"` + Attempts int `json:"attempts"` + // SaidAt is when the controller last heard it. A provider says it again every quarter of an hour + // while it lasts, so an old one is a provider that stopped saying anything. + SaidAt time.Time `json:"said-at"` +} + +// SayAgainWithin is how long a failing standing stays current without being said again: twice the +// quarter of an hour a provider repeats it at. Older, and status says the provider has gone quiet. +const SayAgainWithin = 30 * time.Minute + +// Quiet says the provider has not repeated this standing for longer than it would while it lasts. +func (s ProviderStanding) Quiet(now time.Time) bool { return now.Sub(s.SaidAt) > SayAgainWithin } + +// KeepStanding records a provider's newest word: failing keeps it, recovered removes it, and says +// whether a recovery removed anything. +func (i *Inventory) KeepStanding(ctx context.Context, failing bool, s ProviderStanding) (bool, error) { + if !failing { + tag, err := i.store.Pool().Exec(ctx, + `delete from provider_standing where module = $1 and provider_node = $2 and consumer = $3`, + s.Module, s.ProviderNode, s.Consumer) + return err == nil && tag.RowsAffected() > 0, err + } + _, err := i.store.Pool().Exec(ctx, ` + insert into provider_standing + (module, provider_node, consumer, consumer_node, provision, class, error, since, attempts, said_at) + values ($1, $2, $3, $4, $5, $6, $7, $8, $9, now()) + on conflict (module, provider_node, consumer) do update set + consumer_node = excluded.consumer_node, provision = excluded.provision, + class = excluded.class, error = excluded.error, since = excluded.since, + attempts = excluded.attempts, said_at = excluded.said_at`, + s.Module, s.ProviderNode, s.Consumer, s.ConsumerNode, s.Provision, s.Class, s.Error, s.Since, s.Attempts) + return false, err +} + +// FailingProviders is every consumer a provider last said it keeps failing, oldest run first. +func (i *Inventory) FailingProviders(ctx context.Context) ([]ProviderStanding, error) { + rows, err := i.store.Pool().Query(ctx, ` + select module, provider_node, provision, consumer, consumer_node, class, error, since, attempts, said_at + from provider_standing order by since, module, consumer`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []ProviderStanding + for rows.Next() { + var s ProviderStanding + if err := rows.Scan(&s.Module, &s.ProviderNode, &s.Provision, &s.Consumer, &s.ConsumerNode, + &s.Class, &s.Error, &s.Since, &s.Attempts, &s.SaidAt); err != nil { + return nil, err + } + out = append(out, s) + } + return out, rows.Err() +} diff --git a/internal/inventory/standing_test.go b/internal/inventory/standing_test.go new file mode 100644 index 0000000..fe5fa3f --- /dev/null +++ b/internal/inventory/standing_test.go @@ -0,0 +1,130 @@ +package inventory + +import ( + "slices" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" +) + +// A provider's standing (novox/hq ADR 0224), from the grant that lets it say so to the row status +// reads. + +// The broker spells the events itself because it cannot import the catalogue; the two agree. +func TestTheBrokerAndTheCatalogueNameTheSameStandingEvents(t *testing.T) { + if broker.ProvisionerFailing != catalogue.ProvisionerFailing || + broker.ProvisionerRecovered != catalogue.ProvisionerRecovered { + t.Fatal("the broker and the catalogue disagree about what a provider's standing is called") + } +} + +// **Every provider may say it, whatever its manifest lists**: a provider whose manifest forgot the +// events would have its announcement refused by the bus, and fail its consumers as silently as on +// 2026-10-05 (issue 179). A module that receives no contributions provides nothing and is given +// nothing. +func TestEveryProviderIsGrantedItsStandingAndNothingElseIs(t *testing.T) { + provider := catalogue.Manifest{Module: "keycloak", Version: "1", + Emits: []string{"client.created"}, Receives: map[string]string{"oidc-client": "/x/mesh.json"}} + consumer := catalogue.Manifest{Module: "grafana", Version: "1", Emits: []string{"dashboard.saved"}} + + d := declaredFor(provider, nil) + for _, e := range []string{"client.created", catalogue.ProvisionerFailing, catalogue.ProvisionerRecovered} { + if !slices.Contains(d.Emits, e) { + t.Fatalf("a provider is not granted %s: %v", e, d.Emits) + } + } + perms, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindModule, Node: "anchor", + Module: "keycloak", Emits: d.Emits}) + if err != nil { + t.Fatal(err) + } + if !slices.Contains(perms.Publish, "mesh.mod.keycloak.event.provisioner.failing") { + t.Fatalf("the bus would refuse a provider's standing: %v", perms.Publish) + } + + if got := declaredFor(consumer, nil).Emits; slices.Contains(got, catalogue.ProvisionerFailing) { + t.Fatalf("a module that provides nothing was granted a provider's standing: %v", got) + } + // Declared by hand as well: said once. + provider.Emits = append(provider.Emits, catalogue.ProvisionerFailing) + n := 0 + for _, e := range provider.EmitsAll() { + if e == catalogue.ProvisionerFailing { + n++ + } + } + if n != 1 { + t.Fatalf("%v", provider.EmitsAll()) + } +} + +// And the controller may hear it from every provider, and only those two events. +func TestTheControllerHearsEveryProvidersStanding(t *testing.T) { + perms, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindController}) + if err != nil { + t.Fatal(err) + } + for _, want := range []string{"mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered"} { + if !slices.Contains(perms.Subscribe, want) { + t.Fatalf("the controller may not hear %s: %v", want, perms.Subscribe) + } + } + if slices.Contains(perms.Subscribe, "mesh.mod.*.event.>") { + t.Fatal("the controller hears every event in the mesh") + } +} + +func TestAFailingStandingIsKeptUntilItRecovers(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + since := time.Date(2026, 10, 5, 0, 49, 0, 0, time.UTC) + s := ProviderStanding{Module: "keycloak", ProviderNode: "anchor", Provision: "oidc-client", + Consumer: "mesh_home_grafana", ConsumerNode: "home-server", Class: "credentials-rejected", + Error: "401 invalid_grant", Since: since, Attempts: 60} + if _, err := inv.KeepStanding(ctx, true, s); err != nil { + t.Fatal(err) + } + // Said again: one row, the newest word. + s.Attempts = 31000 + if _, err := inv.KeepStanding(ctx, true, s); err != nil { + t.Fatal(err) + } + got, err := inv.FailingProviders(ctx) + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].Attempts != 31000 || !got[0].Since.Equal(since) || got[0].ConsumerNode != "home-server" || + got[0].Class != "credentials-rejected" || got[0].SaidAt.IsZero() { + t.Fatalf("%+v", got) + } + if got[0].Quiet(time.Now()) { + t.Fatal("a standing just said reads as quiet") + } + if !got[0].Quiet(time.Now().Add(SayAgainWithin + time.Minute)) { + t.Fatal("a standing not said again for longer than a provider repeats it does not read as quiet") + } + + // The same consumer from another machine's provider is its own row. + other := s + other.ProviderNode = "laptop" + if _, err := inv.KeepStanding(ctx, true, other); err != nil { + t.Fatal(err) + } + cleared, err := inv.KeepStanding(ctx, false, s) + if err != nil || !cleared { + t.Fatalf("recovered cleared nothing: %v %v", cleared, err) + } + cleared, err = inv.KeepStanding(ctx, false, s) + if err != nil || cleared { + t.Fatalf("a recovery for nothing kept said it cleared something: %v %v", cleared, err) + } + got, err = inv.FailingProviders(ctx) + if err != nil { + t.Fatal(err) + } + if len(got) != 1 || got[0].ProviderNode != "laptop" { + t.Fatalf("%+v", got) + } +} diff --git a/internal/link/receive.go b/internal/link/receive.go index 962e9d4..98447b7 100644 --- a/internal/link/receive.go +++ b/internal/link/receive.go @@ -34,6 +34,9 @@ const ( // built without anybody telling the mesh (novox/hq 04-ISSUES/131). KindSourceMoved = "source-moved" KindCatchUp = "catch-up" + // KindProvisioner is a provider saying a consumer has failed for minutes, or recovered + // (novox/hq ADR 0224). + KindProvisioner = "provisioner" ) // Control is one thing a node or a module said, as the controller must act on it. @@ -54,6 +57,11 @@ type Control interface { // Body is the message itself — the payload alone, never the envelope. Body() []byte + // Subject is where it was published. For a module's event it names the emitter, which the bus + // enforces (only a module may publish into its own namespace), so who said it is read from here + // and never from the body. + Subject() string + // Redelivered says the bus has handed this message over before. An enrolment cares and // nothing else does: one already spent is not finished a second time. Redelivered() bool diff --git a/internal/link/receive_fake_test.go b/internal/link/receive_fake_test.go index ad71890..f7023f3 100644 --- a/internal/link/receive_fake_test.go +++ b/internal/link/receive_fake_test.go @@ -61,6 +61,7 @@ func (c *fakeInbound) retries(ctx context.Context, s *Server) { type fakeControl struct { kind string + subject string body []byte tag uint64 to *settled @@ -71,6 +72,7 @@ type fakeControl struct { func (m *fakeControl) Kind() string { return m.kind } func (m *fakeControl) Body() []byte { return m.body } +func (m *fakeControl) Subject() string { return m.subject } func (m *fakeControl) Redelivered() bool { return m.redelivered } func (m *fakeControl) About(string) {} func (m *fakeControl) Answer(context.Context, []byte) error { return nil } diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 5472069..d4f67fa 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -53,7 +53,7 @@ func Nats(js *broker.JetStream) Inbound { // whatever was asked for — and not at all when nothing was. func (n *natsInbound) Also(kind string) error { switch kind { - case KindModuleMoved, KindCatchUp, KindSourceMoved: + case KindModuleMoved, KindCatchUp, KindSourceMoved, KindProvisioner: n.follows[kind] = true return nil default: @@ -241,9 +241,30 @@ func kindOfSubject(subject string) (string, bool) { // runs (ADR 0190): the old builder still answers on the retired seat until it is unassigned. return KindBuilt, true } + if _, ok := ProvisionerEmitter(subject); ok { + return KindProvisioner, true + } return "", false } +// ProvisionerEmitter is the module a provider's standing event came from, read from its subject +// (`mesh.mod..event.provisioner.`); false for any other subject. The +// controller's own follow pattern, with `*` for the module, decodes too. +func ProvisionerEmitter(subject string) (string, bool) { + rest, ok := strings.CutPrefix(subject, "mesh.mod.") + if !ok { + return "", false + } + module, event, ok := strings.Cut(rest, ".event.") + if !ok || module == "" || strings.Contains(module, ".") { + return "", false + } + if event != broker.ProvisionerFailing && event != broker.ProvisionerRecovered { + return "", false + } + return module, true +} + // natsControl is one message from the bus being built, as the controller reads it. type natsControl struct { kind string @@ -256,8 +277,9 @@ type natsControl struct { delivered uint64 } -func (m *natsControl) Kind() string { return m.kind } -func (m *natsControl) Body() []byte { return m.msg.Data } +func (m *natsControl) Kind() string { return m.kind } +func (m *natsControl) Body() []byte { return m.msg.Data } +func (m *natsControl) Subject() string { return m.msg.Subject } // Redelivered is what the server counted, not what the controller remembers. Which is the answer to // a question the AMQP side could only guess at across a restart: an enrolment redelivered because diff --git a/internal/link/serve.go b/internal/link/serve.go index 1099873..920d6bf 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -79,6 +79,8 @@ type Server struct { recorder Recorder upgrader Upgrader replayer Replayer + // standings keeps what providers say about their consumers (novox/hq ADR 0224). + standings Standings log *log.Logger // giveUp is how long one message is held for the store; zero means GiveUpAfter. @@ -153,6 +155,9 @@ func (s *Server) Serve(ctx context.Context) error { if s.replayer != nil { s.log.Printf("answering %s", KindCatchUp) } + if s.standings != nil { + s.log.Printf("keeping every provider's %s", KindProvisioner) + } return s.inbound.Receive(ctx, s.act) } @@ -173,6 +178,8 @@ func (s *Server) act(ctx context.Context, m Control) { s.sourceMoved(ctx, m) case KindCatchUp: s.catchingUp(ctx, m) + case KindProvisioner: + s.provisioner(ctx, m) default: // Dropped: a message nothing understands will not be understood on the next attempt // either, and asking for it again would spin. diff --git a/internal/link/standing.go b/internal/link/standing.go new file mode 100644 index 0000000..b861847 --- /dev/null +++ b/internal/link/standing.go @@ -0,0 +1,114 @@ +package link + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/broker" +) + +// A provider's standing (novox/hq ADR 0224). +// +// **A provider that keeps failing a consumer is a problem the controller reports**, not a line in a +// journal. On 2026-10-05 the identity provider's provisioner failed every consumer 31,000 times in a +// day — its admin no longer took the mesh's secret once its database was moved — and every surface +// the mesh has called the mesh well (novox/hq issue 179). A provider now says, as an event, a +// consumer it has failed for minutes without one success, and the consumer recovering; the +// controller keeps the newest word per provider, machine and consumer, and `status` names each one +// still failing. + +// Standing is one provider's word about one consumer. +type Standing struct { + // Module is the emitter, read from the subject the bus let it publish on — never from the body. + Module string `json:"-"` + // Failing is which of the two it said: failing, or recovered. + Failing bool `json:"-"` + + Provider string `json:"provider"` + ProviderNode string `json:"provider-node"` + Consumer string `json:"consumer"` + Node string `json:"node"` + Class string `json:"class,omitempty"` + Error string `json:"error,omitempty"` + Since time.Time `json:"since"` + Attempts int `json:"attempts"` + // Why is said with a recovery that is not a success: `withdrawn`, a consumer no longer asked for. + Why string `json:"why,omitempty"` +} + +// Standings keeps what providers say about their consumers. +type Standings interface { + // Stood records a provider's newest word about a consumer: a failing one kept, a recovered one + // cleared — and says whether a recovery cleared anything, since a provider announces its first + // success for every consumer after it starts. An error the store is away for is held and asked + // again, like a report. + Stood(ctx context.Context, s Standing) (cleared bool, err error) +} + +// Watches says where providers' standings are kept, and asks for them to be delivered. +func (s *Server) Watches(st Standings) error { + if err := s.inbound.Also(KindProvisioner); err != nil { + return err + } + s.standings = st + return nil +} + +// ReadStanding is one standing event as the controller understands it, from its subject and body. +func ReadStanding(subject string, body []byte) (Standing, error) { + module, ok := ProvisionerEmitter(subject) + if !ok { + return Standing{}, fmt.Errorf("%s is not a provider's standing", subject) + } + var st Standing + if err := json.Unmarshal(body, &st); err != nil { + return Standing{}, fmt.Errorf("%s's standing could not be read: %w", module, err) + } + if st.Consumer == "" { + return Standing{}, fmt.Errorf("%s's standing named no consumer", module) + } + st.Module = module + st.Failing = strings.HasSuffix(subject, "."+broker.ProvisionerFailing) + return st, nil +} + +// provisioner acts on one standing event. +// +// **A recovery must not be lost.** A failing standing is said again every quarter of an hour while +// it lasts, so one dropped is replaced; a recovery is said once, and dropping it would leave status +// naming a consumer that is fine. So a store that is away holds the message, as a report is held. +func (s *Server) provisioner(ctx context.Context, m Control) { + if s.standings == nil { + // Delivered because the consumer's filter names it, with nothing here keeping it: taken, + // because handing it back would not give it anywhere to go. + _ = m.Took() + return + } + st, err := ReadStanding(m.Subject(), m.Body()) + if err != nil { + s.log.Printf("%v; ignored", err) + _ = m.Took() + return + } + cleared, err := s.standings.Stood(ctx, st) + what := fmt.Sprintf("%s's standing for %s", st.Module, st.Consumer) + switch s.decide(ctx, m, what, "", "", err) { + case Hold: + return + case Stale, GiveUp: + _ = m.Took() + return + } + if err != nil { + s.log.Printf("%s could not be kept: %v", what, err) + } else if st.Failing { + s.log.Printf("%s on %s is FAILING %s on %s (%s, %d attempts since %s): %s", st.Module, + st.ProviderNode, st.Consumer, st.Node, st.Class, st.Attempts, st.Since.Format(time.RFC3339), st.Error) + } else if cleared { + s.log.Printf("%s on %s recovered %s", st.Module, st.ProviderNode, st.Consumer) + } + _ = m.Took() +} diff --git a/internal/link/standing_test.go b/internal/link/standing_test.go new file mode 100644 index 0000000..6a2d545 --- /dev/null +++ b/internal/link/standing_test.go @@ -0,0 +1,203 @@ +package link + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/broker" +) + +// A provider's standing (novox/hq ADR 0224): who said it is read from the subject the bus let it +// publish on, a failing one is kept and a recovery cleared, and a recovery is never lost to a store +// that is away — said once, it would leave status naming a consumer that is fine. + +type keptStandings struct { + kept []Standing + err error +} + +func (k *keptStandings) Stood(_ context.Context, s Standing) (bool, error) { + if k.err != nil { + return false, k.err + } + k.kept = append(k.kept, s) + return !s.Failing, nil +} + +func standingSays(t *testing.T, in *fakeInbound, to *settled, subject string, body map[string]any) Control { + t.Helper() + m := in.sends(t, to, KindProvisioner, body).(*fakeControl) + m.subject = subject + return m +} + +func TestTheControllerFollowsEveryProvidersStandingAndNothingElse(t *testing.T) { + for subject, want := range map[string]string{ + "mesh.mod.keycloak.event.provisioner.failing": "keycloak", + "mesh.mod.postgres.event.provisioner.recovered": "postgres", + "mesh.mod.*.event.provisioner.failing": "*", + } { + got, ok := ProvisionerEmitter(subject) + if !ok || got != want { + t.Errorf("%s: %q %v", subject, got, ok) + } + if kind, _ := kindOfSubject(subject); kind != KindProvisioner { + t.Errorf("%s decodes to %q", subject, kind) + } + } + for _, subject := range []string{ + "mesh.mod.keycloak.event.provisioner.other", + "mesh.mod.keycloak.event.client.created", + "mesh.mod.a.b.event.provisioner.failing", + "mesh.seat.keycloak.event.provisioner.failing", + } { + if _, ok := ProvisionerEmitter(subject); ok { + t.Errorf("%s read as a provider's standing", subject) + } + } + var follows int + for _, s := range broker.ControllerFollows { + if kind, _ := kindOfSubject(s); kind == KindProvisioner { + follows++ + } + } + if follows != 2 { + t.Fatalf("the controller follows %d standing subjects, want failing and recovered", follows) + } +} + +func TestAFailingStandingIsKeptNamingTheEmitterFromTheSubject(t *testing.T) { + s, in := serving() + kept := &keptStandings{} + if err := s.Watches(kept); err != nil { + t.Fatal(err) + } + to := &settled{} + since := time.Date(2026, 10, 5, 0, 49, 0, 0, time.UTC) + s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing", map[string]any{ + "provider": "oidc-client", "provider-node": "anchor", "consumer": "mesh_home_grafana", + "node": "home-server", "class": "credentials-rejected", "error": "401 invalid_grant", + "since": since, "attempts": 31000, + // A body naming another module is not believed: the subject is the bus's word. + "module": "postgres", + })) + if !to.acked || len(kept.kept) != 1 { + t.Fatalf("settled %+v, kept %+v", to, kept.kept) + } + got := kept.kept[0] + if got.Module != "keycloak" || !got.Failing || got.Consumer != "mesh_home_grafana" || got.Node != "home-server" || + got.ProviderNode != "anchor" || got.Class != "credentials-rejected" || got.Attempts != 31000 || !got.Since.Equal(since) { + t.Fatalf("%+v", got) + } + + to = &settled{} + s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.recovered", map[string]any{ + "provider": "oidc-client", "provider-node": "anchor", "consumer": "mesh_home_grafana", + })) + if !to.acked || len(kept.kept) != 2 || kept.kept[1].Failing { + t.Fatalf("settled %+v, kept %+v", to, kept.kept) + } +} + +func TestARecoveryIsHeldWhileTheStoreIsAway(t *testing.T) { + s, in := serving() + kept := &keptStandings{err: restarting} + if err := s.Watches(kept); err != nil { + t.Fatal(err) + } + to := &settled{} + s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.recovered", + map[string]any{"consumer": "mesh_home_grafana"})) + if !to.unsettled() || len(in.held) != 1 { + t.Fatalf("a recovery was settled while the store was away: %+v", to) + } + kept.err = nil + in.retries(t.Context(), s) + if !to.acked || len(kept.kept) != 1 { + t.Fatalf("the held recovery was not kept when the store came back: %+v %+v", to, kept.kept) + } +} + +func TestAStandingThatNamesNoConsumerIsTakenAndForgotten(t *testing.T) { + s, in := serving() + kept := &keptStandings{} + if err := s.Watches(kept); err != nil { + t.Fatal(err) + } + to := &settled{} + s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing", map[string]any{})) + if !to.acked || len(kept.kept) != 0 { + t.Fatalf("%+v %+v", to, kept.kept) + } +} + +func TestAStandingWithNothingKeepingItIsTaken(t *testing.T) { + s, in := serving() + to := &settled{} + s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing", + map[string]any{"consumer": "x"})) + if !to.acked { + t.Fatal("a standing nothing keeps was left for the bus to hand over again") + } +} + +// Over a real bus: a provider's standing published under its own module's namespace reaches the +// controller through the events consumer's filter — the one wildcard filter on it — names the emitter +// from the subject, and is acknowledged. +func TestNatsAProvidersStandingReachesTheController(t *testing.T) { + js := aBus(t) + kept := &lockedStandings{} + s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()} + if err := s.Follows(&toldAbout{}); err != nil { + t.Fatal(err) + } + if err := s.Watches(kept); err != nil { + t.Fatal(err) + } + ctx, stop := context.WithCancel(context.Background()) + defer stop() + go func() { _ = s.Serve(ctx) }() + eventually(t, "the controller's event consumer being made", func() bool { + _, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName) + return err == nil + }) + + for _, event := range []string{broker.ProvisionerFailing, broker.ProvisionerRecovered} { + if _, err := js.Context().Publish("mesh.mod.keycloak.event."+event, + []byte(`{"consumer":"mesh_home_grafana","provider-node":"anchor","class":"credentials-rejected"}`)); err != nil { + t.Fatal(err) + } + } + // Somebody else's event under the same prefix is not the controller's to hear. + if _, err := js.Context().Publish("mesh.mod.keycloak.event.client.created", []byte(`{}`)); err != nil { + t.Fatal(err) + } + eventually(t, "both standings being kept, in order, naming the emitter", func() bool { + got := kept.all() + return len(got) == 2 && got[0].Module == "keycloak" && got[0].Failing && !got[1].Failing + }) + eventually(t, "both being acknowledged and nothing else delivered", func() bool { + info, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName) + return err == nil && info.NumAckPending == 0 && info.Delivered.Consumer == 2 + }) +} + +type lockedStandings struct { + mu sync.Mutex + kept []Standing +} + +func (l *lockedStandings) Stood(_ context.Context, s Standing) (bool, error) { + l.mu.Lock() + defer l.mu.Unlock() + l.kept = append(l.kept, s) + return !s.Failing, nil +} + +func (l *lockedStandings) all() []Standing { + l.mu.Lock() + defer l.mu.Unlock() + return append([]Standing(nil), l.kept...) +}