Files
mesh-controller/cmd/mesh-controller/retirement.go
T
jochen 68009b16fe Hear what providers retire, and let a person approve, reject and delete (hq ADR 0230)
A provider now waits for a person before retiring more than three consumers
or half of what it holds, and deletes only when asked. The controller is that
person's way in: it keeps waiting and rejected sets as conditions, answers them
with retire approve|reject, lists and deletes retired consumers through the
provider's own tools on its machine, records each act in the hand-act log, and
probes for anything retired longer than thirty days (D11).
2026-10-06 14:41:15 +02:00

344 lines
13 KiB
Go

package main
import (
"context"
"fmt"
"sort"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A consumer the mesh stops asking for is retired, not withdrawn, and deleted only by a person
// (novox/hq ADR 0230, the operator's decision of 2026-10-06; it replaces ADR 0229's withdrawal brake).
//
// A provider retires a consumer — disables its access, keeps its data, marks it with when and why —
// once it has seen the same consumers go unasked for in five passes. More than three at once, or more
// than half of those it holds where it holds more than one, is a person actively working on the mesh,
// so the provider retires nothing and waits: the controller keeps that as an urgent condition naming
// what would go and the two verbs that answer it. A retired consumer is deleted only by `cleanup
// delete`, which the provider executes on its own backend — the controller never touches one — and
// one retired more than thirty days is a warning that cleanup is waiting (D11).
// The kinds a provider's retirement word raises.
const (
kindRetireWaiting = "retire-waiting"
kindRetireRejected = "retire-rejected"
kindCleanupWaiting = "cleanup-waiting"
// sourceRetirement is what raised a retirement condition: the provider's own event.
sourceRetirement = "provisioner.retirement"
)
// cleanupAfter is how long a consumer may stay retired before cleanup is said to be waiting (D11).
const cleanupAfter = 30 * 24 * time.Hour
// retirementAsk is how long one provider is given to answer one question about its retired consumers.
var retirementAsk = 20 * time.Second
func retireWaitingKey(module, node string) string {
return conditions.Key(conditions.ScopeProvider, module+"."+node, "retire")
}
func retireRejectedKey(module, node string) string {
return conditions.Key(conditions.ScopeProvider, module+"."+node, "retire-rejected")
}
// retirements keeps what providers say about consumers they retire, as conditions.
type retirements struct {
keeper func() *conditions.Keeper
// record writes an act done by hand that reached a provider some other way than this controller's
// verbs; nil records nothing (a test).
record func(ctx context.Context, act link.HandAct) error
}
func consumerList(cs []link.RetiredConsumer) string {
parts := make([]string, 0, len(cs))
for _, c := range cs {
p := c.Consumer
if c.Node != "" {
p += " (" + c.Node + ")"
}
parts = append(parts, p)
}
return strings.Join(parts, ", ")
}
// waitingObservation is a provider waiting for a person to approve or reject a retirement.
func waitingObservation(r link.Retirement) conditions.Observation {
since := r.Since
if since.IsZero() {
since = r.At
}
summary := fmt.Sprintf("%s on %s would retire %d consumer(s) the mesh no longer asks for — %s — more than its "+
"bound (%s), so it retires nothing until a person answers: `retire approve %s %s --why …` or `retire reject "+
"%s %s --why …`", r.Module, r.ProviderNode, len(r.Consumers), consumerList(r.Consumers), orBound(r.Bound),
r.ProviderNode, r.Module, r.ProviderNode, r.Module)
said := fmt.Sprintf("waiting since %s, %d held: %s", since.UTC().Format("2006-01-02 15:04 MST"), r.Held,
consumerList(r.Consumers))
return conditions.Observation{Scope: conditions.ScopeProvider, ID: r.Module + "." + r.ProviderNode,
Token: "retire", Kind: kindRetireWaiting, Machine: r.ProviderNode, Also: consumerNodes(r),
Severity: conditions.Urgent, Summary: summary, Said: said, Source: sourceRetirement,
Resolver: conditions.ResolverOperator}
}
// rejectedObservation is consumers kept active by a person's rejection though the mesh asks for them no more.
func rejectedObservation(r link.Retirement) conditions.Observation {
summary := fmt.Sprintf("%s on %s keeps %s active although the mesh no longer asks for them: a person rejected "+
"their retirement (%s). Assign them again, or `retire approve %s %s --why …`", r.Module, r.ProviderNode,
consumerList(r.Consumers), orWhy(r.By, r.Why), r.ProviderNode, r.Module)
return conditions.Observation{Scope: conditions.ScopeProvider, ID: r.Module + "." + r.ProviderNode,
Token: "retire-rejected", Kind: kindRetireRejected, Machine: r.ProviderNode, Also: consumerNodes(r),
Severity: conditions.Warning, Summary: summary, Said: "rejected: " + orWhy(r.By, r.Why),
Source: sourceRetirement, Resolver: conditions.ResolverOperator}
}
func orBound(b string) string {
if b == "" {
return "more than 3, or more than half of those held"
}
return b
}
func orWhy(by, why string) string {
switch {
case by != "" && why != "":
return by + ": " + why
case why != "":
return why
case by != "":
return "by " + by
}
return "no reason given"
}
// consumerNodes are the other machines a retirement concerns: where its consumers are.
func consumerNodes(r link.Retirement) []string {
seen := map[string]bool{r.ProviderNode: true}
var out []string
for _, c := range r.Consumers {
if c.Node != "" && !seen[c.Node] {
seen[c.Node] = true
out = append(out, c.Node)
}
}
sort.Strings(out)
return out
}
// Retired keeps one word. Waiting raises the urgent condition; a rejection turns it into a warning; a
// retirement, an approval or the set settling clears both. An approval, a rejection or a deletion that
// did not come through this controller's verbs is recorded in the hand-act log here, so every one is.
func (s retirements) Retired(ctx context.Context, r link.Retirement) error {
k := s.keeper()
if k == nil {
return fmt.Errorf("the condition store is not open in this controller: %w", link.ErrTryAgain)
}
waiting, rejected := retireWaitingKey(r.Module, r.ProviderNode), retireRejectedKey(r.Module, r.ProviderNode)
clear := func(keys ...string) error {
for _, key := range keys {
if _, err := k.Clear(ctx, key, fmt.Sprintf("%s on %s says %s", r.Module, r.ProviderNode, r.Change)); err != nil {
return storeAway(err)
}
}
return nil
}
var err error
switch r.Change {
case link.RetireWaiting:
if _, err = k.Observe(ctx, waitingObservation(r)); err != nil {
return storeAway(err)
}
case link.RetireRejected:
if err = clear(waiting); err != nil {
return err
}
if _, err = k.Observe(ctx, rejectedObservation(r)); err != nil {
return storeAway(err)
}
case link.RetireApproved, link.RetireRetired, link.RetireSettled:
if err = clear(waiting, rejected); err != nil {
return err
}
}
switch r.Change {
case link.RetireApproved, link.RetireRejected, link.RetireDeleted:
if r.Via != link.ViaController && s.record != nil {
verb := map[string]string{link.RetireApproved: "retire approve", link.RetireRejected: "retire reject",
link.RetireDeleted: "cleanup delete"}[r.Change]
cause := kindRetireWaiting
if r.Change == link.RetireDeleted {
cause = kindCleanupWaiting
}
why := r.Why
if strings.TrimSpace(why) == "" {
why = "no reason given to the provider"
}
act := link.HandAct{Verb: verb, Args: append([]string{r.ProviderNode, r.Module}, r.Names()...),
Why: why + " (asked of the provider directly, not through the controller)", Cause: cause,
By: r.By, At: r.At}
if act.By == "" {
act.By = "unknown, asked of " + r.Module + " on " + r.ProviderNode + " directly"
}
if err := s.record(ctx, act); err != nil {
return fmt.Errorf("%s on %s %s by hand, and it could not be recorded: %v: %w", r.Module,
r.ProviderNode, r.Change, err, link.ErrTryAgain)
}
}
}
return nil
}
// recordHandActOnTheBus writes an act through the serving controller's connection.
func recordHandActOnTheBus(ctx context.Context, act link.HandAct) error {
return onTheBus(func(conn *nats.Conn) error {
_, err := link.RecordHandAct(ctx, conn, act)
return err
})
}
// providerInstance is one provider module on one machine.
type providerInstance struct {
Node, Module string
}
// providerInstances are every module assigned on every machine whose manifest receives contributions:
// every provider, which serves the retirement tools (ADR 0230) — or is older than them.
func providerInstances(ctx context.Context, inv *inventory.Inventory) ([]providerInstance, error) {
shelf, err := inv.Catalogue(ctx)
if err != nil {
return nil, err
}
nodes, err := inv.Nodes(ctx)
if err != nil {
return nil, err
}
var out []providerInstance
for _, n := range nodes {
modules, err := inv.Assigned(ctx, n.Name)
if err != nil {
return nil, fmt.Errorf("what %s is assigned cannot be read: %w", n.Name, err)
}
for _, m := range modules {
if man, ok := shelf[m]; ok && len(man.Receives) > 0 {
out = append(out, providerInstance{Node: n.Name, Module: m})
}
}
}
sort.Slice(out, func(i, j int) bool {
if out[i].Node != out[j].Node {
return out[i].Node < out[j].Node
}
return out[i].Module < out[j].Module
})
return out, nil
}
// askRetirement asks one provider for what it holds retired and what waits.
func askRetirement(ctx context.Context, conn *nats.Conn, p providerInstance) (link.RetirementState, error) {
var state link.RetirementState
answer, err := link.AskModuleToolOn(ctx, conn, p.Module, link.ToolRetirement, p.Node, map[string]any{}, retirementAsk)
if err != nil {
return state, err
}
if answer.Error != "" {
return state, fmt.Errorf("%s on %s answered %s with an error: %s", p.Module, p.Node, link.ToolRetirement, answer.Error)
}
if err := unmarshalAnswer(answer, &state); err != nil {
return state, fmt.Errorf("%s on %s answered %s with something unreadable: %w", p.Module, p.Node, link.ToolRetirement, err)
}
return state, nil
}
// retiredAt reads a retired consumer's moment; zero when the provider could not say.
func retiredAt(c link.RetiredConsumer) time.Time {
t, err := time.Parse(time.RFC3339, c.RetiredAt)
if err != nil {
return time.Time{}
}
return t
}
// cleanupObservation is a provider holding consumers retired longer than cleanupAfter.
func cleanupObservation(p providerInstance, old []link.RetiredConsumer, now time.Time) conditions.Observation {
parts := make([]string, 0, len(old))
for _, c := range old {
parts = append(parts, fmt.Sprintf("%s (%d days, %s)", c.Consumer, int(now.Sub(retiredAt(c)).Hours()/24),
sizeWords(c.SizeBytes)))
}
return conditions.Observation{Scope: conditions.ScopeProvider, ID: p.Module + "." + p.Node, Token: "cleanup",
Kind: kindCleanupWaiting, Machine: p.Node, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s on %s holds %d consumer(s) retired more than %d days, waiting for a person to "+
"delete or bring them back: %s — `cleanup list`, then `cleanup delete %s %s <consumer> --why …`",
p.Module, p.Node, len(old), int(cleanupAfter.Hours()/24), strings.Join(parts, ", "), p.Node, p.Module),
Resolver: conditions.ResolverOperator}
}
func sizeWords(size *int64) string {
if size == nil || *size < 0 {
return "size unknown"
}
b := float64(*size)
for _, unit := range []string{"B", "KB", "MB", "GB"} {
if b < 1024 || unit == "GB" {
if unit == "B" {
return fmt.Sprintf("%d B", *size)
}
return fmt.Sprintf("%.1f %s", b, unit)
}
b /= 1024
}
return ""
}
// probeRetired is D11: no provider holds a consumer retired more than thirty days. Every provider
// assigned is asked what it holds retired; one that does not serve the question — an older build, a
// TypeScript provider not yet on the SDK that retires — has nothing it can say and is passed over,
// which is not a failure of the probe. A provider that cannot be asked otherwise is.
func probeRetired(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
if d.js == nil {
return nil, fmt.Errorf("no bus to ask the providers over")
}
instances, err := providerInstances(ctx, d.open.inventory)
if err != nil {
return nil, err
}
return retiredTooLong(ctx, d.js.Conn(), instances, time.Now())
}
// retiredTooLong asks each provider and answers one observation per provider holding anything retired
// longer than cleanupAfter.
func retiredTooLong(ctx context.Context, conn *nats.Conn, instances []providerInstance,
now time.Time) ([]conditions.Observation, error) {
var out []conditions.Observation
var problems []string
for _, p := range instances {
state, err := askRetirement(ctx, conn, p)
if err != nil {
if isNothingServes(err) {
continue
}
problems = append(problems, err.Error())
continue
}
var old []link.RetiredConsumer
for _, c := range state.Retired {
if at := retiredAt(c); !at.IsZero() && now.Sub(at) > cleanupAfter {
old = append(old, c)
}
}
if len(old) > 0 {
out = append(out, cleanupObservation(p, old, now))
}
}
if len(problems) > 0 {
// Not "nothing retired": a provider that could not be asked may hold the oldest of all.
return nil, fmt.Errorf("%s", strings.Join(problems, "; "))
}
return out, nil
}