letta crash-loops on 'type "vector" does not exist': pgvector is not a trusted extension, so only the provider's superuser can create it, and the provisioner never did. A contribution may now name extensions; the provider creates each (IF NOT EXISTS, available ones only) in the consumer's database on every pass. Go per the standing rule for a TypeScript module that changes. letta asks for vector.
341 lines
10 KiB
Go
341 lines
10 KiB
Go
package main
|
|
|
|
// The reconcile loop every provider shares, as the TypeScript SDK's runProvisioner runs it
|
|
// (@novox/mesh-sdk/provisioner, 0.1.10). 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).**
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/url"
|
|
"os"
|
|
"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
|
|
|
|
verifiedAt time.Time
|
|
applied map[string]appliedEntry
|
|
lost map[string]brake
|
|
waiting map[string]int
|
|
failing map[string]failure
|
|
lastWarning string
|
|
}
|
|
|
|
type appliedEntry struct {
|
|
hash string
|
|
derived map[string]any
|
|
}
|
|
|
|
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.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{}
|
|
}
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
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.
|
|
h.say("%s: could not check the backend, will ask again: %s", g.As, scrub(err, password))
|
|
if timedOut {
|
|
verifying = false
|
|
}
|
|
continue
|
|
}
|
|
if held {
|
|
delete(h.lost, 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)
|
|
}
|
|
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.applied[g.As] = appliedEntry{hash: hash, derived: p.Derived}
|
|
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.
|
|
for as, was := range h.applied {
|
|
if want[as] {
|
|
continue
|
|
}
|
|
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)
|
|
}
|
|
for as := range h.failing {
|
|
if !want[as] {
|
|
delete(h.failing, as)
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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
|
|
}
|