Files
mesh-catalog/modules/postgres/cmd/postgres-provider/harness.go
T

628 lines
21 KiB
Go

package main
// The reconcile loop every provider shares, as the TypeScript SDK's runProvisioner runs it
// (@novox/mesh-sdk/provisioner, 0.1.11). 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; 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).
//
// **A reconcile that would withdraw more than its bound stops, and says so (novox/hq to-be 45 Phase 2,
// ADR 0227 rule 4).** Withdrawing more than WithdrawAtOnce consumers in one pass — or more than
// WithdrawFraction of those this process holds — is the shape of issue 241, where one misread file
// withdrew seven at once. Such a pass withdraws nothing: each consumer it would have withdrawn is kept,
// announced `provisioner.failing` with the class `withdrawal-braked` (the controller raises it as a
// condition), and said. While the same consumers stay unasked for, one is released every ReleaseEvery,
// said and announced as it goes, so an intended unassignment of many completes without a hand and a
// mistaken one costs at most one consumer an hour while the operator is told. Withdrawal never destroys
// data (issue 241's second half), so a release is a login locked, not a database dropped.
//
// 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"
"encoding/json"
"errors"
"fmt"
"net/url"
"os"
"sort"
"strings"
"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(ctx context.Context, p Provision) error
// Remove withdraws what Create made; derived is what the mesh last derived, remembered here.
Remove(ctx context.Context, as string, derived map[string]any) error
// Holds says whether the backend still holds the consumer exactly as p says. Read-only.
Holds(ctx context.Context, p Provision) (bool, 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
// WithdrawAtOnce (1) and WithdrawFraction (0.5) bound what one pass may withdraw: more consumers
// than WithdrawAtOnce, or a larger share of those held than WithdrawFraction, brakes the pass.
// ReleaseEvery (1h) is how often a braked withdrawal lets one consumer go.
WithdrawAtOnce int
WithdrawFraction float64
ReleaseEvery time.Duration
verifiedAt time.Time
// braked is every consumer a braked pass kept, by when it was first kept; releasedAt is when the
// brake last let one go, and brakeSaid the set it last said, so a pass repeats nothing.
braked map[string]time.Time
releasedAt time.Time
brakeSaid string
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"
// ClassWithdrawalBraked is a consumer the mesh no longer asks for, kept because the pass that would
// withdraw it would withdraw more than its bound (ADR 0227 rule 4).
ClassWithdrawalBraked = "withdrawal-braked"
)
// 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.WithdrawAtOnce == 0 {
h.WithdrawAtOnce = 1
}
if h.WithdrawFraction == 0 {
h.WithdrawFraction = 0.5
}
if h.ReleaseEvery == 0 {
h.ReleaseEvery = time.Hour
}
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.braked = map[string]time.Time{}
}
}
// 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 withdraw.
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.init()
given := h.readContributions()
if given == nil {
return
}
want := map[string]bool{}
for _, g := range given {
want[g.As] = true
}
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}
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))
}
}
}
// Withdraw every login this process made that the mesh no longer asks for — within the bound.
var withdrawing []string
for as := range h.applied {
if !want[as] {
withdrawing = append(withdrawing, as)
}
}
sort.Strings(withdrawing)
for as := range h.braked {
if want[as] {
// Asked for again: the brake held what the mesh still wanted. Its standing was ended by this
// pass's success above, as any consumer's is.
delete(h.braked, as)
}
}
if h.overTheBound(len(withdrawing), len(h.applied)) {
withdrawing = h.brakeWithdrawal(withdrawing)
} else if len(h.braked) > 0 {
h.say("the withdrawal is within its bound again: %s withdrawn as asked", strings.Join(withdrawing, ", "))
h.braked, h.brakeSaid = map[string]time.Time{}, ""
}
for _, as := range withdrawing {
was := h.applied[as]
h.say("%s: no longer in %s; withdrawing it from the backend", as, h.Receives)
if err := h.Adapter.Remove(ctx, as, was.derived); err != nil {
h.say("%s: remove failed, will retry: %v", as, err)
continue
}
delete(h.applied, as)
delete(h.lost, as)
delete(h.braked, as)
}
for as := range h.failing {
if !want[as] {
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. One the
// brake holds is still kept, and its standing stays.
for as := range h.trouble {
if _, held := h.braked[as]; !want[as] && !held {
h.recovered(as, "withdrawn")
}
}
}
// overTheBound says a pass withdrawing n of the held consumers would withdraw more than it may.
func (h *Harness) overTheBound(n, held int) bool {
if n == 0 {
return false
}
return n > h.WithdrawAtOnce || (held > 1 && float64(n) > h.WithdrawFraction*float64(held))
}
// brakeWithdrawal keeps every consumer a pass over its bound would withdraw, announces each as failing
// with the class withdrawal-braked, and answers the one it releases now, if one is due.
func (h *Harness) brakeWithdrawal(withdrawing []string) []string {
now := h.Now()
set := strings.Join(withdrawing, ", ")
if set != h.brakeSaid {
h.say("WITHDRAWAL BRAKED: this pass would withdraw %d of the %d consumer(s) this provider holds (%s), "+
"more than %d at once or %.0f%% of them. Nothing is withdrawn; each is announced as %s (%s), and one "+
"is let go every %s while the mesh goes on not asking for them (novox/hq ADR 0227 rule 4)",
len(withdrawing), len(h.applied), set, h.WithdrawAtOnce, h.WithdrawFraction*100, EventFailing,
ClassWithdrawalBraked, h.ReleaseEvery)
h.brakeSaid = set
}
if len(h.braked) == 0 {
// The release clock starts with the brake, not at the last release of an earlier one.
h.releasedAt = now
}
for _, as := range withdrawing {
if _, kept := h.braked[as]; !kept {
h.braked[as] = now
}
text := fmt.Sprintf("the mesh no longer asks for it, and the pass that would withdraw it would withdraw %d "+
"consumers at once: kept until released (%s)", len(withdrawing), set)
h.failed(as, h.applied[as].node, ClassWithdrawalBraked, text)
}
if now.Sub(h.releasedAt) < h.ReleaseEvery {
return nil
}
h.releasedAt = now
release := withdrawing[0]
h.say("%s: released by the withdrawal brake after %s; %d more kept", release,
now.Sub(h.braked[release]).Round(time.Second), len(withdrawing)-1)
h.recovered(release, "withdrawn")
return []string{release}
}
// 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
}