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...) +}