Files
mesh-catalog/modules/postgres/cmd/postgres-provider/harness.go
T
jochen e68ef88333 Retire a consumer the mesh stops asking for, and delete only on a person's word (hq ADR 0230)
The hourly release of ADR 0229's brake still ended in the mesh acting alone on
a mistake. A consumer now stays active until the same unasked set holds for
five passes, waits for a person past three or half of those held, is disabled
and marked rather than withdrawn, comes back as it was when asked again, and is
deleted only through the provider's delete tool. The backend keeps the mark, so
a restart forgets nothing and finds what was withdrawn before.
2026-10-06 13:54:10 +02:00

564 lines
19 KiB
Go

package main
// The reconcile loop every provider shares, as the TypeScript SDK's runProvisioner runs it
// (@novox/mesh-sdk/provisioner, 0.1.12). The Go SDK has no provisioner yet, so this module carries
// the loop itself, line for line in behaviour; when the Go SDK grows one, this file is what moves
// there (novox/hq ADR 0039: the loop is the SDK's, the adapter is the module's).
//
// Read the contributions the mesh delivered; bring each consumer's resource into being through the
// adapter, under the login and password the mesh minted; retire 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).
//
// **A consumer the mesh stops asking for is retired, not withdrawn, and deleted only by a person
// (novox/hq ADR 0230, the operator's model of 2026-10-06).** Retiring disables its access — reversibly —
// and marks its login and data "to delete" with when and why; nothing is deleted. Asked for again, the
// ordinary create re-enables it as it was. It is retired only once the same set has gone unasked in
// StablePasses consecutive passes, and a set larger than the bound waits for a person. retirement.go.
//
// 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 (this file and
// retirement.go).
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/url"
"os"
"strings"
"sync"
"time"
)
// Provision is one consumer's resource to bring into being — everything the mesh derived and delivered.
type Provision struct {
// As is the login the mesh derived and gave the consumer to present.
As string
// Password is the one the mesh minted, read from the file the host unsealed.
Password string
// Values are what the consumer contributed (e.g. {"name": "letta", "extensions": ["vector"]}).
Values map[string]any
// Derived is what this provider's own definition derives for the consumer (novox/hq ADR 0201).
Derived map[string]any
// At is where the consumer is; Consumer is its node.
At string
Consumer string
}
// Adapter is the per-service half.
type Adapter interface {
// Create brings the consumer into being under the mesh's login and password — and, for one that
// was retired, re-enables it as it was and clears its mark to delete.
Create(ctx context.Context, p Provision) error
// Retire disables the consumer's access, reversibly, and marks its login and data to delete with
// when and why. It deletes nothing (novox/hq ADR 0230). derived is what the mesh last derived,
// remembered here; nil for a consumer this process did not make.
Retire(ctx context.Context, as string, derived map[string]any, why string, at time.Time) error
// Holds says whether the backend still holds the consumer exactly as p says. Read-only.
Holds(ctx context.Context, p Provision) (bool, error)
// Inventory is what the backend holds that the mesh made: consumers active and retired, read from
// the backend itself so a restart forgets nothing. Read-only.
Inventory(ctx context.Context) (Inventory, error)
// Delete removes one retired consumer — its login and its data — for good. Only ever asked by a
// person, through `cleanup delete`; the harness has checked it is retired and not asked for.
Delete(ctx context.Context, r Retired) (freedBytes int64, err error)
}
// Harness is the loop's settings and memory.
type Harness struct {
Resource string
Receives string
Adapter Adapter
Every time.Duration // 5s
VerifyEvery time.Duration // 60s
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
// StablePasses (5) is how many consecutive passes must see the same set no longer asked for before
// it is retired. RetireAtOnce (3) and RetireFraction (0.5) bound what is retired without a person:
// more consumers than RetireAtOnce, or — where more than one is held — a larger share of those held
// than RetireFraction, waits for `retire approve` (novox/hq ADR 0230).
StablePasses int
RetireAtOnce int
RetireFraction float64
// mu is held by a pass and by every tool a person asks, so the two never interleave.
mu sync.Mutex
verifiedAt time.Time
r retirement
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
node string
}
type brake struct {
times int
nextAt time.Time
}
type failure struct {
text string
times int
}
// The longest a consumer whose create keeps failing to satisfy holds waits between checks.
const maxBackoff = time.Hour
// How many passes a secret may be unreadable before it stops being called a race (issue 225), and
// once said loudly, how often it is repeated. The same cadence quiets a create that keeps failing
// the same way.
const (
patiently = 12
loudlyEvery = 240
)
type contribution struct {
As string `json:"as"`
Secret string `json:"secret"`
Node string `json:"node"`
At string `json:"at"`
Values map[string]any `json:"values"`
Derived map[string]any `json:"derived"`
}
func (h *Harness) init() {
if h.Every == 0 {
h.Every = 5 * time.Second
}
if h.VerifyEvery == 0 {
h.VerifyEvery = time.Minute
}
if h.HoldsTimeout == 0 {
h.HoldsTimeout = 30 * time.Second
}
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.StablePasses == 0 {
h.StablePasses = 5
}
if h.RetireAtOnce == 0 {
h.RetireAtOnce = 3
}
if h.RetireFraction == 0 {
h.RetireFraction = 0.5
}
if h.Log == nil {
h.Log = func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) }
}
if h.applied == nil {
h.applied = map[string]appliedEntry{}
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{}
h.r.retired = map[string]retiredEntry{}
}
}
// Run reconciles until ctx ends. One consumer's failure never stops the others'.
func (h *Harness) Run(ctx context.Context) {
h.init()
for {
h.Reconcile(ctx)
select {
case <-ctx.Done():
return
case <-time.After(h.Every):
}
}
}
func (h *Harness) say(format string, args ...any) {
h.Log("[provisioner:"+h.Resource+"] "+format, args...)
}
func (h *Harness) warn(why string) {
if why == h.lastWarning {
return
}
if why != "" {
h.say("%s; nothing applied or removed until it can be read", why)
} else {
h.say("contributions file readable again")
}
h.lastWarning = why
}
// readContributions answers the consumers asked for, or nil when the file says nothing usable.
// **Nothing read is not nobody asking** (novox/hq issue 241): only a file that was read can retire anything.
func (h *Harness) readContributions() []contribution {
raw, err := os.ReadFile(h.Receives)
if err != nil {
h.warn(fmt.Sprintf("contributions file unreadable (%s): %v", h.Receives, err))
return nil
}
var doc struct {
Requirement string `json:"requirement"`
Given json.RawMessage `json:"given"`
}
if err := json.Unmarshal(raw, &doc); err != nil {
h.warn(fmt.Sprintf("contributions file is not JSON (%s): %v", h.Receives, err))
return nil
}
if doc.Requirement != "" && doc.Requirement != h.Resource {
h.warn(fmt.Sprintf("%s is for %s, not %s", h.Receives, doc.Requirement, h.Resource))
return nil
}
var given []contribution
if len(doc.Given) == 0 || string(doc.Given) == "null" || json.Unmarshal(doc.Given, &given) != nil {
h.warn(fmt.Sprintf("%s has no given list", h.Receives))
return nil
}
h.warn("")
out := []contribution{}
for _, g := range given {
// No `as` is not a credential grant: nothing to create for it.
if g.As != "" && g.Secret != "" {
out = append(out, g)
}
}
return out
}
func hashOf(as, password string, values, derived map[string]any) string {
// Derived is in the hash: a provider that renames what it derives gave a different resource.
b, _ := json.Marshal([]any{as, password, orEmpty(values), orEmpty(derived)})
return string(b)
}
func orEmpty(m map[string]any) map[string]any {
if m == nil {
return map[string]any{}
}
return m
}
// Reconcile is one pass.
func (h *Harness) Reconcile(ctx context.Context) {
h.mu.Lock()
defer h.mu.Unlock()
h.init()
given := h.readContributions()
if given == nil {
// Not a result: the passes that must agree before anything is retired start again.
h.r.key, h.r.count = "", 0
return
}
want := map[string]bool{}
for _, g := range given {
want[g.As] = true
}
h.seed(ctx, want)
verifying := h.Now().Sub(h.verifiedAt) >= h.VerifyEvery
if verifying {
h.verifiedAt = h.Now()
}
for _, g := range given {
raw, err := os.ReadFile(g.Secret)
if err != nil {
// A secret the host has not written yet is a race on the first pass; past a minute it is
// a person's to look at, and said so (novox/hq issue 225).
n := h.waiting[g.As] + 1
h.waiting[g.As] = n
if n <= patiently {
h.say("%s: secret not readable yet (%s): %v", g.As, g.Secret, err)
} else if n == patiently+1 || n%loudlyEvery == 0 {
h.say("%s: CANNOT READ the secret after %d attempts (%s): %v. This is not a race any more — "+
"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)
password := strings.TrimSuffix(string(raw), "\n")
p := Provision{As: g.As, Password: password, Values: orEmpty(g.Values), Derived: orEmpty(g.Derived), At: g.At, Consumer: g.Node}
hash := hashOf(g.As, password, g.Values, g.Derived)
reapplying := 0
if was, ok := h.applied[g.As]; ok && was.hash == hash {
if !verifying {
continue
}
b, braked := h.lost[g.As]
if braked && h.Now().Before(b.nextAt) {
continue
}
hctx, cancel := context.WithTimeout(ctx, h.HoldsTimeout)
held, err := h.Adapter.Holds(hctx, p)
timedOut := errors.Is(hctx.Err(), context.DeadlineExceeded)
cancel()
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.
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
}
continue
}
if held {
delete(h.lost, g.As)
h.succeeded(g.As)
continue
}
reapplying = b.times + 1
if reapplying == 1 {
h.say("%s: the backend no longer holds it; applying again", g.As)
} else {
h.say("%s: still not held after being applied again (%d times in a row) — create does not "+
"produce what holds checks; applying again", g.As, reapplying)
}
}
if err := h.Adapter.Create(ctx, p); err != nil {
text := scrub(err, password)
f := h.failing[g.As]
if f.text != text {
f = failure{text: text}
}
f.times++
h.failing[g.As] = f
// Said each time it changes, and while it stays the same, as rarely as a lost secret.
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}
}
continue
}
if f, was := h.failing[g.As]; was {
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, node: g.Node}
h.reenabled(g.As, g.Node)
if reapplying == 0 {
delete(h.lost, g.As)
} else {
wait := h.VerifyEvery << (reapplying - 1)
if wait > maxBackoff || wait <= 0 {
wait = maxBackoff
}
h.lost[g.As] = brake{times: reapplying, nextAt: h.Now().Add(wait)}
if reapplying > 1 {
h.say("%s: next check in %s", g.As, wait.Round(time.Second))
}
}
}
// Retire what the mesh has stably stopped asking for — within the bound, or with a person.
h.retireUnasked(ctx, want)
h.r.lastWant = want
for as := range h.failing {
if !want[as] {
delete(h.failing, as)
}
}
// A consumer the mesh stopped asking for that this provider holds nothing 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. One still held is said recovered when it is retired.
for as := range h.trouble {
if _, held := h.applied[as]; !want[as] && !held {
h.recovered(as, "no longer asked for")
}
}
}
// 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.
func scrub(err error, password string) string {
text := err.Error()
if password == "" {
return text
}
for _, form := range []string{password, url.QueryEscape(password), url.PathEscape(password)} {
text = strings.ReplaceAll(text, form, "***")
}
return text
}