Merge pull request 'Report a provider that keeps failing a consumer in status (hq ADR 0224)' (#70) from feat/a-provider-failing-a-consumer-is-reported into main

This commit was merged in pull request #70.
This commit is contained in:
2026-10-05 22:20:45 +00:00
23 changed files with 913 additions and 11 deletions
+3
View File
@@ -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.
+1 -1
View File
@@ -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 "+
+13
View File
@@ -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
+5
View File
@@ -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.
+6
View File
@@ -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]})
+1 -1
View File
@@ -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
}
+120
View File
@@ -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
}
+115
View File
@@ -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)
}
}
+11 -1
View File
@@ -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
+3 -2
View File
@@ -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.
+16 -1
View File
@@ -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",
},
+1 -1
View File
@@ -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: {
+32
View File
@@ -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
}
+1 -1
View File
@@ -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
@@ -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)
);
+76
View File
@@ -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()
}
+130
View File
@@ -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)
}
}
+8
View File
@@ -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
+2
View File
@@ -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 }
+25 -3
View File
@@ -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.<module>.event.provisioner.<failing|recovered>`); 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
+7
View File
@@ -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.
+114
View File
@@ -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()
}
+203
View File
@@ -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...)
}