postgres: announce a consumer failed for minutes, and its recovery

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).
This commit is contained in:
jochen
2026-10-06 00:13:42 +02:00
parent ed6384feb0
commit 6f1e2f5a0d
5 changed files with 425 additions and 6 deletions
@@ -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.
@@ -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")
}
}
@@ -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
@@ -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())
}
@@ -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)
}
}