From 6f1e2f5a0da507bf29926985e8fbc0a69733b6f0 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 00:12:35 +0200 Subject: [PATCH] postgres: announce a consumer failed for minutes, and its recovery MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A provider failed every consumer for a day and said so only in its journal (hq issue 179). The provisioner loop now emits provisioner.failing after five minutes without a success — create, check or secret — and repeats it every fifteen; provisioner.recovered on the next success, on withdrawal, and on the first success after a restart, so the controller can name it in status (hq ADR 0224). --- .../postgres/cmd/postgres-provider/harness.go | 189 +++++++++++++++++- .../postgres-provider/harness_same_test.go | 31 +++ .../cmd/postgres-provider/harness_test.go | 16 +- .../postgres/cmd/postgres-provider/main.go | 8 + .../cmd/postgres-provider/standing_test.go | 187 +++++++++++++++++ 5 files changed, 425 insertions(+), 6 deletions(-) create mode 100644 modules/postgres/cmd/postgres-provider/harness_same_test.go create mode 100644 modules/postgres/cmd/postgres-provider/standing_test.go diff --git a/modules/postgres/cmd/postgres-provider/harness.go b/modules/postgres/cmd/postgres-provider/harness.go index d7cceed..e7a6e45 100644 --- a/modules/postgres/cmd/postgres-provider/harness.go +++ b/modules/postgres/cmd/postgres-provider/harness.go @@ -8,6 +8,17 @@ package main // Read the contributions the mesh delivered; bring each consumer's resource into being through the // adapter, under the login and password the mesh minted; withdraw what the mesh no longer asks for. // **A provider creates the credential the mesh minted, and seals nothing (novox/hq ADR 0048).** +// +// **A provider that keeps failing a consumer says so on the bus (novox/hq ADR 0224).** A consumer +// whose create, check or secret has failed without one success in between for FailingAfter is +// announced as `provisioner.failing` — naming the consumer, its machine and the class of error — and +// again every SayAgainEvery while it lasts; the first success after that is `provisioner.recovered`. +// The controller keeps the newest per provider and consumer and `status` names it. On 2026-10-05 the +// identity provider failed every consumer 31,000 times in a day and said so only in its journal +// (novox/hq issue 179). +// +// Carried, identical, by every Go provider until the Go SDK has the loop: postgres and keycloak. +// Each module's `harness_same_test.go` fails when its copy and the other's differ. import ( "context" @@ -54,15 +65,57 @@ type Harness struct { HoldsTimeout time.Duration // 30s Log func(format string, args ...any) Now func() time.Time + // Announce publishes one of the provider's standing events; nil announces nothing. Node is the + // machine this provider runs on, said in each. + Announce func(event string, body map[string]any) + Node string + // FailingAfter is how long a consumer fails without a success before it is announced (5m); + // SayAgainEvery is how often it is announced again while it lasts (15m), so a controller that + // missed the first hears the next, and a standing nobody repeats can be told from one that holds. + FailingAfter time.Duration + SayAgainEvery time.Duration verifiedAt time.Time applied map[string]appliedEntry lost map[string]brake waiting map[string]int failing map[string]failure + trouble map[string]*standing + cleared map[string]bool lastWarning string } +// standing is one consumer's unbroken run of failures: since when, how often, and the last error. +type standing struct { + node string + since time.Time + attempts int + class string + text string + saidAt time.Time +} + +// The events a provider's standing is announced as (novox/hq ADR 0224). The controller derives the +// permission to emit them for every module that receives contributions; no manifest lists them. +const ( + EventFailing = "provisioner.failing" + EventRecovered = "provisioner.recovered" +) + +// The classes of error a standing is announced with: what a person reading `status` needs to know +// before reading the journal. An adapter may say better (Classifier). +const ( + ClassCredentials = "credentials-rejected" + ClassUnreachable = "unreachable" + ClassSecret = "secret-unreadable" + ClassRefused = "refused" +) + +// Classifier is an adapter that can say what class an error of its own is. +type Classifier interface { + Class(err error) string +} + type appliedEntry struct { hash string derived map[string]any @@ -111,6 +164,12 @@ func (h *Harness) init() { if h.Now == nil { h.Now = time.Now } + if h.FailingAfter == 0 { + h.FailingAfter = 5 * time.Minute + } + if h.SayAgainEvery == 0 { + h.SayAgainEvery = 15 * time.Minute + } if h.Log == nil { h.Log = func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) } } @@ -119,6 +178,8 @@ func (h *Harness) init() { h.lost = map[string]brake{} h.waiting = map[string]int{} h.failing = map[string]failure{} + h.trouble = map[string]*standing{} + h.cleared = map[string]bool{} } } @@ -230,6 +291,7 @@ func (h *Harness) Reconcile(ctx context.Context) { "nothing has been provisioned for this consumer and nothing will be until somebody looks. "+ "Check who owns the file and who this process runs as (novox/hq issue 225)", g.As, n, g.Secret, err) } + h.failed(g.As, g.Node, ClassSecret, fmt.Sprintf("secret not readable (%s): %v", g.Secret, err)) continue } delete(h.waiting, g.As) @@ -253,7 +315,9 @@ func (h *Harness) Reconcile(ctx context.Context) { if err != nil { // Unable to ask is not evidence of loss. A backend that timed out will time out for // the next consumer too, so the rest of this pass is not asked. - h.say("%s: could not check the backend, will ask again: %s", g.As, scrub(err, password)) + text := scrub(err, password) + h.say("%s: could not check the backend, will ask again: %s", g.As, text) + h.failed(g.As, g.Node, h.classOf(err, text), text) if timedOut { verifying = false } @@ -261,6 +325,7 @@ func (h *Harness) Reconcile(ctx context.Context) { } if held { delete(h.lost, g.As) + h.succeeded(g.As) continue } reapplying = b.times + 1 @@ -283,6 +348,7 @@ func (h *Harness) Reconcile(ctx context.Context) { if f.times == 1 || f.times%loudlyEvery == 0 { h.say("%s: create failed, will retry: %s", g.As, text) } + h.failed(g.As, g.Node, h.classOf(err, text), text) if reapplying > 0 { h.lost[g.As] = brake{times: reapplying - 1} } @@ -292,6 +358,7 @@ func (h *Harness) Reconcile(ctx context.Context) { h.say("%s: created, after %d failed attempt(s)", g.As, f.times) delete(h.failing, g.As) } + h.succeeded(g.As) h.applied[g.As] = appliedEntry{hash: hash, derived: p.Derived} if reapplying == 0 { delete(h.lost, g.As) @@ -325,6 +392,126 @@ func (h *Harness) Reconcile(ctx context.Context) { delete(h.failing, as) } } + // A consumer the mesh stopped asking for is no longer failed by anyone: said, so a standing + // the controller keeps for it is cleared rather than left naming a consumer that is gone. + for as := range h.trouble { + if !want[as] { + h.recovered(as, "withdrawn") + } + } +} + +// failed counts one more failure in a consumer's unbroken run, and announces the run once it has +// lasted FailingAfter — then again every SayAgainEvery while it lasts. +func (h *Harness) failed(as, node, class, text string) { + now := h.Now() + s := h.trouble[as] + if s == nil { + s = &standing{since: now} + h.trouble[as] = s + } + s.node, s.class, s.text = node, class, text + s.attempts++ + if now.Sub(s.since) < h.FailingAfter { + return + } + if !s.saidAt.IsZero() && now.Sub(s.saidAt) < h.SayAgainEvery { + return + } + first := s.saidAt.IsZero() + s.saidAt = now + if first { + h.say("%s: FAILING for %s (%d attempts, %s): %s. Announced as %s; `status` names it until it "+ + "succeeds (novox/hq ADR 0224)", as, now.Sub(s.since).Round(time.Second), s.attempts, class, text, EventFailing) + } + h.announce(EventFailing, map[string]any{ + "provider": h.Resource, "provider-node": h.Node, + "consumer": as, "node": node, + "class": class, "error": clip(text), + "since": s.since.UTC().Format(time.RFC3339), "attempts": s.attempts, + }) +} + +// succeeded ends a consumer's run of failures; one that was announced is announced recovered. +// +// **And the first success for a consumer since this process started is announced too**, failing or +// not: a provider that announced a failure and was restarted has forgotten it, and without this the +// controller would name the consumer failing for ever after it recovered unheard. +func (h *Harness) succeeded(as string) { + if h.trouble[as] == nil && !h.cleared[as] { + h.cleared[as] = true + h.announce(EventRecovered, map[string]any{ + "provider": h.Resource, "provider-node": h.Node, "consumer": as, "why": "first-success", + }) + return + } + h.cleared[as] = true + h.recovered(as, "") +} + +func (h *Harness) recovered(as, why string) { + s := h.trouble[as] + if s == nil { + return + } + delete(h.trouble, as) + if s.saidAt.IsZero() { + return // never announced, so there is nothing to take back + } + if why == "" { + h.say("%s: recovered after %s and %d failed attempt(s)", as, h.Now().Sub(s.since).Round(time.Second), s.attempts) + } + body := map[string]any{ + "provider": h.Resource, "provider-node": h.Node, "consumer": as, "node": s.node, + "since": s.since.UTC().Format(time.RFC3339), "attempts": s.attempts, + } + if why != "" { + body["why"] = why + } + h.announce(EventRecovered, body) +} + +func (h *Harness) announce(event string, body map[string]any) { + if h.Announce != nil { + h.Announce(event, body) + } +} + +// classOf is an error's class: the adapter's word when it has one, else read from the text. +func (h *Harness) classOf(err error, text string) string { + if c, ok := h.Adapter.(Classifier); ok { + if class := c.Class(err); class != "" { + return class + } + } + return ClassOf(text) +} + +// ClassOf reads an error's class from its text — the words the backends the mesh runs use. +func ClassOf(text string) string { + t := strings.ToLower(text) + for _, w := range []string{"invalid_grant", "invalid user credentials", "password authentication failed", + "authentication failed", "unauthorized", " 401"} { + if strings.Contains(t, w) { + return ClassCredentials + } + } + for _, w := range []string{"connection refused", "no such host", "i/o timeout", "deadline exceeded", + "connection reset", "network is unreachable", "no route to host", "eof"} { + if strings.Contains(t, w) { + return ClassUnreachable + } + } + return ClassRefused +} + +// clip keeps an announced error to what belongs in a status line. +func clip(text string) string { + const most = 300 + if len(text) <= most { + return text + } + return text[:most] + "…" } // scrub is an error's text with the consumer's password removed, raw and URL-encoded. diff --git a/modules/postgres/cmd/postgres-provider/harness_same_test.go b/modules/postgres/cmd/postgres-provider/harness_same_test.go new file mode 100644 index 0000000..f156732 --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/harness_same_test.go @@ -0,0 +1,31 @@ +package main + +// The provisioner loop is carried, identical, by every Go provider until the Go SDK has it +// (harness.go). Two copies drift the moment one is fixed and the other is not — and the one left +// behind is the provider that fails a consumer without saying so (novox/hq ADR 0224). This holds +// them to one text. Skipped where keycloak is not beside this module, as in a build of this one alone. + +import ( + "bytes" + "errors" + "io/fs" + "os" + "testing" +) + +func TestTheHarnessIsTheSameAsKeycloaks(t *testing.T) { + theirs, err := os.ReadFile("../../../keycloak/cmd/keycloak-provider/harness.go") + if errors.Is(err, fs.ErrNotExist) { + t.Skip("keycloak is not beside this module") + } + if err != nil { + t.Fatal(err) + } + ours, err := os.ReadFile("harness.go") + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(ours, theirs) { + t.Fatal("harness.go differs from keycloak/cmd/keycloak-provider/harness.go: change both, identically") + } +} diff --git a/modules/postgres/cmd/postgres-provider/harness_test.go b/modules/postgres/cmd/postgres-provider/harness_test.go index f79aab1..525a216 100644 --- a/modules/postgres/cmd/postgres-provider/harness_test.go +++ b/modules/postgres/cmd/postgres-provider/harness_test.go @@ -11,10 +11,11 @@ import ( ) type recorder struct { - created []Provision - removed []string - held bool - failing error + created []Provision + removed []string + held bool + failing error + holdsErr error } func (r *recorder) Create(_ context.Context, p Provision) error { @@ -30,7 +31,12 @@ func (r *recorder) Remove(_ context.Context, as string, _ map[string]any) error return nil } -func (r *recorder) Holds(context.Context, Provision) (bool, error) { return r.held, nil } +func (r *recorder) Holds(context.Context, Provision) (bool, error) { + if r.holdsErr != nil { + return false, r.holdsErr + } + return r.held, nil +} type world struct { t *testing.T diff --git a/modules/postgres/cmd/postgres-provider/main.go b/modules/postgres/cmd/postgres-provider/main.go index f93b7cc..57b0352 100644 --- a/modules/postgres/cmd/postgres-provider/main.go +++ b/modules/postgres/cmd/postgres-provider/main.go @@ -40,6 +40,14 @@ func main() { Receives: receives, Adapter: provisioner{pg: pg, announce: announce}, Log: func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) }, + // A consumer failed for minutes is said on the bus, where the controller hears it and + // `status` names it (novox/hq ADR 0224). + Announce: func(event string, body map[string]any) { + if err := stdio.Emit(event, body); err != nil { + say("emit %s failed: %v", event, err) + } + }, + Node: os.Getenv("MESH_NODE"), } go h.Run(context.Background()) } diff --git a/modules/postgres/cmd/postgres-provider/standing_test.go b/modules/postgres/cmd/postgres-provider/standing_test.go new file mode 100644 index 0000000..7f1b9ee --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/standing_test.go @@ -0,0 +1,187 @@ +package main + +// A provider that keeps failing a consumer says so on the bus (novox/hq ADR 0224): not on the first +// failure, which may be a restart; after FailingAfter of failures with no success between; again +// every SayAgainEvery while it lasts; and recovered on the first success, or when the consumer goes. + +import ( + "errors" + "os" + "strings" + "testing" + "time" +) + +type announced struct { + event string + body map[string]any +} + +func standingWorld(t *testing.T) (*world, *[]announced) { + w := newWorld(t) + var said []announced + w.h.Announce = func(e string, b map[string]any) { + // The first success since start is its own test's; every other test reads past it. + if b["why"] != "first-success" { + said = append(said, announced{e, b}) + } + } + w.h.Node = "anchor" + return w, &said +} + +// passes reconciles every five seconds for d, as Run would. +func (w *world) passes(d time.Duration) { + for end := w.now.Add(d); w.now.Before(end); w.now = w.now.Add(5 * time.Second) { + w.h.Reconcile(ctx) + } +} + +func events(said []announced) string { + var out []string + for _, a := range said { + out = append(out, a.event) + } + return strings.Join(out, ",") +} + +func TestAConsumerFailedForMinutesIsAnnouncedNamingItAndTheClass(t *testing.T) { + w, said := standingWorld(t) + w.a.failing = errors.New(`token request failed: 401 {"error":"invalid_grant","error_description":"Invalid user credentials"}`) + w.give(map[string]any{"as": "mesh_home_grafana", "node": "home-server"}) + + w.passes(4 * time.Minute) + if len(*said) != 0 { + t.Fatalf("announced before FailingAfter: %v", events(*said)) + } + w.passes(2 * time.Minute) + if events(*said) != EventFailing { + t.Fatalf("want one %s, got %q", EventFailing, events(*said)) + } + b := (*said)[0].body + if b["consumer"] != "mesh_home_grafana" || b["node"] != "home-server" || b["class"] != ClassCredentials || + b["provider"] != "postgres-database" || b["provider-node"] != "anchor" || b["attempts"].(int) < 60 { + t.Fatalf("%v", b) + } + + // Said again while it lasts, not every pass. + w.passes(14 * time.Minute) + if events(*said) != EventFailing { + t.Fatalf("repeated too soon: %q", events(*said)) + } + w.passes(2 * time.Minute) + if events(*said) != EventFailing+","+EventFailing { + t.Fatalf("not repeated: %q", events(*said)) + } + + // The first success takes it back. + w.a.failing = nil + w.passes(5 * time.Second) + if events(*said) != EventFailing+","+EventFailing+","+EventRecovered { + t.Fatalf("no recovery: %q", events(*said)) + } + if (*said)[2].body["consumer"] != "mesh_home_grafana" { + t.Fatal((*said)[2].body) + } +} + +func TestOneSuccessBetweenFailuresStartsTheRunAgain(t *testing.T) { + w, said := standingWorld(t) + w.a.failing = errors.New("connection refused") + w.give(map[string]any{"as": "a"}) + w.passes(4 * time.Minute) + w.a.failing = nil + w.passes(5 * time.Second) + w.give(map[string]any{"as": "a", "values": map[string]any{"name": "changed"}}) + w.a.failing = errors.New("connection refused") + w.passes(4 * time.Minute) + if len(*said) != 0 { + t.Fatalf("two runs of four minutes are not one of eight: %q", events(*said)) + } +} + +// The check that failed for a day on 2026-10-05: clients already made, every minute's check refused +// at the token. A check that cannot be asked is a failure too. +func TestACheckThatKeepsFailingIsAFailureToo(t *testing.T) { + w, said := standingWorld(t) + w.give(map[string]any{"as": "a"}) + w.h.Reconcile(ctx) + w.a.holdsErr = errors.New("401 invalid_grant") + w.passes(7 * time.Minute) + if events(*said) != EventFailing || (*said)[0].body["class"] != ClassCredentials { + t.Fatalf("%q %v", events(*said), *said) + } + w.a.holdsErr = nil + w.passes(time.Minute + 5*time.Second) + if events(*said) != EventFailing+","+EventRecovered { + t.Fatalf("%q", events(*said)) + } +} + +func TestAnUnreadableSecretIsAnnouncedAsSuch(t *testing.T) { + w, said := standingWorld(t) + w.give(map[string]any{"as": "a"}) + os.Remove(w.dir + "/a.secret") + w.passes(6 * time.Minute) + if events(*said) != EventFailing || (*said)[0].body["class"] != ClassSecret { + t.Fatalf("%q %v", events(*said), *said) + } +} + +func TestAWithdrawnConsumerIsNoLongerFailing(t *testing.T) { + w, said := standingWorld(t) + w.a.failing = errors.New("boom") + w.give(map[string]any{"as": "a"}, map[string]any{"as": "b"}) + w.passes(6 * time.Minute) + if events(*said) != EventFailing+","+EventFailing { + t.Fatalf("%q", events(*said)) + } + w.give(map[string]any{"as": "a"}) + w.passes(5 * time.Second) + last := (*said)[len(*said)-1] + if last.event != EventRecovered || last.body["consumer"] != "b" || last.body["why"] != "withdrawn" { + t.Fatalf("%v", *said) + } +} + +func TestAnErrorIsClassedByItsWords(t *testing.T) { + for text, want := range map[string]string{ + `Keycloak token request failed: 401 {"error":"invalid_grant"}`: ClassCredentials, + `FATAL: password authentication failed for user "postgres"`: ClassCredentials, + `dial tcp 127.0.0.1:5432: connect: connection refused`: ClassUnreachable, + `context deadline exceeded`: ClassUnreachable, + `extension "nope" is not available`: ClassRefused, + } { + if got := ClassOf(text); got != want { + t.Errorf("%s: %s, want %s", text, got, want) + } + } +} + +func TestAnAdapterThatClassesItsOwnErrorsIsBelieved(t *testing.T) { + w, said := standingWorld(t) + w.h.Adapter = classing{w.a} + w.a.failing = errors.New("anything") + w.give(map[string]any{"as": "a"}) + w.passes(6 * time.Minute) + if (*said)[0].body["class"] != "its-own" { + t.Fatal((*said)[0].body) + } +} + +type classing struct{ *recorder } + +func (classing) Class(error) string { return "its-own" } + +// A provider restarted after announcing a failure has forgotten it; its first success for each +// consumer is announced, so the controller clears what it kept rather than naming it for ever. +func TestTheFirstSuccessSinceStartIsAnnouncedOnce(t *testing.T) { + w := newWorld(t) + var said []announced + w.h.Announce = func(e string, b map[string]any) { said = append(said, announced{e, b}) } + w.give(map[string]any{"as": "a"}) + w.passes(3 * time.Minute) + if events(said) != EventRecovered || said[0].body["why"] != "first-success" || said[0].body["consumer"] != "a" { + t.Fatalf("%v", said) + } +}