From 68009b16fed0bf4a98d8f9b017795bf5e2c00482 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 14:24:58 +0200 Subject: [PATCH] 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). --- cmd/mesh-controller/doctor.go | 2 + cmd/mesh-controller/main.go | 12 + cmd/mesh-controller/push.go | 6 + cmd/mesh-controller/retire_verbs.go | 434 +++++++++++++++++ cmd/mesh-controller/retirement.go | 343 ++++++++++++++ cmd/mesh-controller/retirement_test.go | 460 +++++++++++++++++++ cmd/mesh-controller/seatverbs.go | 50 +- cmd/mesh-controller/seatverbs_schema_test.go | 2 + cmd/mesh-controller/standing.go | 38 +- cmd/mesh-controller/watchdogs.go | 7 +- internal/broker/streams.go | 9 +- internal/broker/testdata/composed.conf | 2 +- internal/broker/writers.go | 7 +- internal/broker/writers_test.go | 2 +- internal/catalogue/events.go | 14 +- internal/catalogue/retirement_event_test.go | 22 + internal/catalogue/verbs.go | 27 ++ internal/link/contracts.go | 8 +- internal/link/receive_nats.go | 5 +- internal/link/retirement.go | 241 ++++++++++ internal/link/serve.go | 2 + internal/link/standing.go | 4 + internal/link/standing_test.go | 8 +- 23 files changed, 1678 insertions(+), 27 deletions(-) create mode 100644 cmd/mesh-controller/retire_verbs.go create mode 100644 cmd/mesh-controller/retirement.go create mode 100644 cmd/mesh-controller/retirement_test.go create mode 100644 internal/catalogue/retirement_event_test.go create mode 100644 internal/link/retirement.go diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index b7e0d3a..24149ea 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -97,6 +97,8 @@ var probeRegistry = []probe{ From: "issue 265", Kind: "status-slow", Phase: 1, run: probeStatus}, {ID: "D10", Asserts: "every machine runs the node-engine and node tools builds the mesh holds, or is inside " + "a plan's window", From: "the version split", Kind: "core-behind", Phase: 1, run: probeCoreBuilds}, + {ID: "D11", Asserts: "no provider holds a consumer retired more than thirty days without a person deciding " + + "its cleanup", From: "ADR 0230", Kind: kindCleanupWaiting, Phase: 2, run: probeRetired}, {ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals", From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs}, } diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index cce4f30..c6b7147 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -164,6 +164,12 @@ func run() error { // What the healers did, and their brake (novox/hq to-be 45 §7). case "healers": return healersCommand(ctx, args[1:]) + // A consumer the mesh stopped asking for: retired, waiting for a person, deleted only by one + // (novox/hq ADR 0230). + case "retire": + return retireCommand(ctx, args[1:]) + case "cleanup": + return cleanupCommand(ctx, args[1:]) case "version": fmt.Println(version) return nil @@ -210,6 +216,12 @@ func usage() { upgrade roll-out [--together] ...send it to the machines running it upgrade record ...record that they are behind, and send nothing status [--json] what is wrong, what is quiet, and what is out of date + retire [--json] every provider waiting for a person to approve a retirement (ADR 0230) + retire approve --why retire what it waits with: access off, data kept + retire reject --why keep them active; a warning stays open + cleanup [list] [--json] every retired consumer per provider: age, size, why + cleanup delete --why the provider deletes that one retired consumer + cleanup delete --older-than --why [--confirm] list those older; delete only with --confirm seats [--json] every seat this mesh defines, what it delivers, and who holds it seat rename rename a seat; its former name still resolves (ADR 0122) seat --to / hand a seat to that assignment as one act; never empty in between (ADR 0131) diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 45a5431..72d5fc0 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -181,6 +181,12 @@ func serve(ctx context.Context) (err error) { if err := server.Watches(standings{keeper: func() *conditions.Keeper { return conditionsFrom }}); err != nil { return err } + // And what they say about consumers the mesh stopped asking for: one waiting for a person is an + // urgent condition, and an act asked of a provider some other way is recorded by hand (ADR 0230). + if err := server.KeepsRetirements(retirements{keeper: func() *conditions.Keeper { return conditionsFrom }, + record: recordHandActOnTheBus}); err != nil { + return err + } // And the mesh's own verbs, as the seat this control plane holds (novox/hq ADR 0154). Served // from the store's row, so what the seat declares is what is answered. diff --git a/cmd/mesh-controller/retire_verbs.go b/cmd/mesh-controller/retire_verbs.go new file mode 100644 index 0000000..681b3a7 --- /dev/null +++ b/cmd/mesh-controller/retire_verbs.go @@ -0,0 +1,434 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "sort" + "strconv" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// The verbs a person answers a provider's retirement with, and cleans up with (novox/hq ADR 0230). +// +// **The controller asks; the provider acts.** Approving, rejecting and deleting are each a question to +// the provider on the machine it runs on — its `provisioner_*` tools — because the provider owns its +// backend and the controller touches none. Every one says why, is written in the hand-act log before +// it is asked, and the provider says what it did as its `provisioner.retirement` event. + +const retireUsage = "retire [--json] | retire approve --why | retire reject --why " + +const cleanupUsage = "cleanup [list] [--json] | cleanup delete --why | " + + "cleanup delete --older-than --why [--confirm]" + +func isNothingServes(err error) bool { return errors.Is(err, link.ErrNothingServes) } + +func unmarshalAnswer(a link.Answer, v any) error { + if len(a.Result) == 0 { + return errors.New("an empty answer") + } + return json.Unmarshal(a.Result, v) +} + +// retireCommand is `retire`, `retire approve` and `retire reject`. +func retireCommand(ctx context.Context, args []string) error { + sub := "list" + if len(args) > 0 && !strings.HasPrefix(args[0], "-") { + sub, args = args[0], args[1:] + } + switch sub { + case "list": + set := flag.NewFlagSet("retire", flag.ContinueOnError) + asJSON := set.Bool("json", false, "as data") + if rest, err := parseAround(set, args); err != nil { + return err + } else if len(rest) > 0 { + return errors.New(retireUsage) + } + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + return onTheBus(func(conn *nats.Conn) error { return listRetiring(ctx, open.inventory, conn, *asJSON) }) + case "approve", "reject": + set := flag.NewFlagSet("retire "+sub, flag.ContinueOnError) + f := addHandActFlags(set) + rest, err := parseAround(set, args) + if err != nil { + return err + } + if len(rest) != 2 { + return errors.New(retireUsage) + } + if err := f.require("retire " + sub); err != nil { + return err + } + return onTheBus(func(conn *nats.Conn) error { + return answerRetirement(ctx, conn, providerInstance{Node: rest[0], Module: rest[1]}, sub == "approve", f) + }) + } + return errors.New(retireUsage) +} + +// retiringRow is one provider's waiting or rejected set, as `retire` lists it. +type retiringRow struct { + Node string `json:"node"` + Module string `json:"module"` + Waiting []link.RetiredConsumer `json:"waiting,omitempty"` + Since string `json:"since,omitempty"` + Held int `json:"held,omitempty"` + Bound string `json:"bound,omitempty"` + Rejected []link.RetiredConsumer `json:"rejected,omitempty"` + RejectWhy string `json:"rejected-why,omitempty"` + Unasked string `json:"unasked,omitempty"` +} + +func listRetiring(ctx context.Context, inv *inventory.Inventory, conn *nats.Conn, asJSON bool) error { + instances, err := providerInstances(ctx, inv) + if err != nil { + return err + } + var rows []retiringRow + for _, p := range instances { + row := retiringRow{Node: p.Node, Module: p.Module} + state, err := askRetirement(ctx, conn, p) + switch { + case isNothingServes(err): + row.Unasked = "answers no retirement question: it predates ADR 0230" + case err != nil: + row.Unasked = err.Error() + default: + if state.Waiting != nil { + row.Waiting, row.Since, row.Held = state.Waiting.Consumers, state.Waiting.Since, state.Waiting.Held + } + if state.Rejected != nil { + row.Rejected, row.RejectWhy = state.Rejected.Consumers, orWhy(state.Rejected.By, state.Rejected.Why) + } + row.Bound = state.Bound + if row.Waiting == nil && row.Rejected == nil { + continue + } + } + rows = append(rows, row) + } + if asJSON { + if rows == nil { + rows = []retiringRow{} + } + return printJSON(map[string]any{"providers": rows}) + } + said := false + for _, r := range rows { + switch { + case r.Unasked != "": + fmt.Printf("%s on %s: not asked — %s\n", r.Module, r.Node, r.Unasked) + default: + if r.Waiting != nil { + said = true + fmt.Printf("%s on %s WAITS since %s to retire %d of the %d it holds (%s): %s\n"+ + " retire approve %s %s --why … | retire reject %s %s --why …\n", + r.Module, r.Node, r.Since, len(r.Waiting), r.Held, orBound(r.Bound), consumerList(r.Waiting), + r.Node, r.Module, r.Node, r.Module) + } + if r.Rejected != nil { + said = true + fmt.Printf("%s on %s keeps active, by a rejection (%s): %s\n", r.Module, r.Node, r.RejectWhy, + consumerList(r.Rejected)) + } + } + } + if !said { + fmt.Println("no provider waits for a person to approve a retirement") + } + return nil +} + +// answerRetirement approves or rejects what one provider waits with. **The set sent is the set the +// provider says it waits with, read now**, and the provider refuses any other: a person approves what +// they were shown, never a set that moved since. +func answerRetirement(ctx context.Context, conn *nats.Conn, p providerInstance, approve bool, f handActFlags) error { + state, err := askRetirement(ctx, conn, p) + if err != nil { + return err + } + verb, tool := "retire reject", link.ToolRetireReject + var set []link.RetiredConsumer + if state.Waiting != nil { + set = state.Waiting.Consumers + } + if approve { + verb, tool = "retire approve", link.ToolRetireApprove + if set == nil && state.Rejected != nil { + // A rejection can be taken back: the consumers it kept are retired after all. + set = state.Rejected.Consumers + } + } + if len(set) == 0 { + return fmt.Errorf("%s on %s waits for nobody to approve or reject a retirement. Nothing was done", p.Module, p.Node) + } + names := make([]string, 0, len(set)) + for _, c := range set { + names = append(names, c.Consumer) + } + if strings.TrimSpace(*f.cause) == "" { + *f.cause = kindRetireWaiting + } + if strings.TrimSpace(*f.condition) == "" { + *f.condition = retireWaitingKey(p.Module, p.Node) + } + f.record(ctx, verb, append([]string{p.Node, p.Module}, names...)) + answer, err := link.AskModuleToolOn(ctx, conn, p.Module, tool, p.Node, map[string]any{ + "consumers": names, "why": strings.TrimSpace(*f.why), "by": link.Caller(), "via": link.ViaController, + }, retirementAsk) + if err != nil { + return err + } + if answer.Error != "" { + return fmt.Errorf("%s on %s refused: %s", p.Module, p.Node, answer.Error) + } + if approve { + fmt.Printf("%s on %s retired %s: access disabled, data kept; `cleanup list` shows them\n", p.Module, p.Node, + strings.Join(names, ", ")) + } else { + fmt.Printf("%s on %s keeps %s active; the warning stays open until the mesh asks for them again or the "+ + "retirement is approved\n", p.Module, p.Node, strings.Join(names, ", ")) + } + return nil +} + +// cleanupCommand is `cleanup list` and `cleanup delete`. +func cleanupCommand(ctx context.Context, args []string) error { + sub := "list" + if len(args) > 0 && !strings.HasPrefix(args[0], "-") { + sub, args = args[0], args[1:] + } + switch sub { + case "list": + set := flag.NewFlagSet("cleanup", flag.ContinueOnError) + asJSON := set.Bool("json", false, "as data") + if rest, err := parseAround(set, args); err != nil { + return err + } else if len(rest) > 0 { + return errors.New(cleanupUsage) + } + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + return onTheBus(func(conn *nats.Conn) error { + listing, err := gatherRetired(ctx, open.inventory, conn, time.Now()) + if err != nil { + return err + } + return printRetired(listing, *asJSON) + }) + case "delete": + set := flag.NewFlagSet("cleanup delete", flag.ContinueOnError) + f := addHandActFlags(set) + olderThan := set.Int("older-than", 0, "every retired consumer older than this many days") + confirm := set.Bool("confirm", false, "with --older-than: delete what is listed, rather than only list it") + rest, err := parseAround(set, args) + if err != nil { + return err + } + if err := f.require("cleanup delete"); err != nil { + return err + } + if strings.TrimSpace(*f.cause) == "" { + *f.cause = kindCleanupWaiting + } + switch { + case *olderThan > 0 && len(rest) == 0: + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + return onTheBus(func(conn *nats.Conn) error { + return deleteOlderThan(ctx, open.inventory, conn, *olderThan, *confirm, f, time.Now()) + }) + case *olderThan == 0 && len(rest) == 3 && !*confirm: + return onTheBus(func(conn *nats.Conn) error { + return deleteRetired(ctx, conn, providerInstance{Node: rest[0], Module: rest[1]}, rest[2], f) + }) + } + return errors.New(cleanupUsage) + } + return errors.New(cleanupUsage) +} + +// retiredRow is one retired consumer, as `cleanup list` shows it. +type retiredRow struct { + Node string `json:"node"` + Module string `json:"module"` + Consumer string `json:"consumer"` + // ConsumerNode is where the consumer was, when the provider knows. + ConsumerNode string `json:"consumer-node,omitempty"` + Kind string `json:"kind,omitempty"` + RetiredAt string `json:"retired-at,omitempty"` + // AgeDays is whole days since it was retired; -1 when the provider could not say when. + AgeDays int `json:"age-days"` + SizeBytes *int64 `json:"size-bytes,omitempty"` + Why string `json:"why,omitempty"` +} + +// retiredListing is every provider's retired consumers, and the providers that could not say. +type retiredListing struct { + Retired []retiredRow `json:"retired"` + // Unasked names each provider not asked, and why: one older than the question, or one that failed. + Unasked []string `json:"unasked,omitempty"` +} + +func gatherRetired(ctx context.Context, inv *inventory.Inventory, conn *nats.Conn, now time.Time) (retiredListing, error) { + instances, err := providerInstances(ctx, inv) + if err != nil { + return retiredListing{}, err + } + return retiredOf(ctx, conn, instances, now), nil +} + +func retiredOf(ctx context.Context, conn *nats.Conn, instances []providerInstance, now time.Time) retiredListing { + out := retiredListing{Retired: []retiredRow{}} + for _, p := range instances { + state, err := askRetirement(ctx, conn, p) + if err != nil { + if isNothingServes(err) { + out.Unasked = append(out.Unasked, fmt.Sprintf("%s on %s answers no retirement question: it predates ADR 0230", p.Module, p.Node)) + } else { + out.Unasked = append(out.Unasked, err.Error()) + } + continue + } + for _, c := range state.Retired { + row := retiredRow{Node: p.Node, Module: p.Module, Consumer: c.Consumer, ConsumerNode: c.Node, Kind: c.Kind, + RetiredAt: c.RetiredAt, AgeDays: -1, SizeBytes: c.SizeBytes, Why: c.Why} + if at := retiredAt(c); !at.IsZero() { + row.AgeDays = int(now.Sub(at).Hours() / 24) + } + out.Retired = append(out.Retired, row) + } + } + sort.SliceStable(out.Retired, func(i, j int) bool { return out.Retired[i].AgeDays > out.Retired[j].AgeDays }) + return out +} + +func printRetired(l retiredListing, asJSON bool) error { + if asJSON { + return printJSON(l) + } + if len(l.Retired) == 0 { + fmt.Println("no provider holds a retired consumer") + } + for _, r := range l.Retired { + age := "age unknown" + if r.AgeDays >= 0 { + age = strconv.Itoa(r.AgeDays) + " day(s)" + } + kind := "" + if r.Kind != "" && r.Kind != "consumer" { + kind = " [" + r.Kind + "]" + } + fmt.Printf("%s on %s: %s%s — retired %s, %s, %s\n %s\n", r.Module, r.Node, r.Consumer, kind, age, + sizeWords(r.SizeBytes), r.RetiredAt, orWhy("", r.Why)) + } + for _, u := range l.Unasked { + fmt.Printf("not asked: %s\n", u) + } + return nil +} + +// deleteRetired asks one provider to delete one consumer it holds retired — never an active one: the +// provider refuses that, and this refuses it first, from what the provider says it holds. +func deleteRetired(ctx context.Context, conn *nats.Conn, p providerInstance, consumer string, f handActFlags) error { + state, err := askRetirement(ctx, conn, p) + if err != nil { + return err + } + found := false + for _, c := range state.Retired { + found = found || c.Consumer == consumer + } + if !found { + return fmt.Errorf("%s on %s holds no retired consumer %s — only a retired consumer is deleted. Nothing was done", + p.Module, p.Node, consumer) + } + f.record(ctx, "cleanup delete", []string{p.Node, p.Module, consumer}) + answer, err := link.AskModuleToolOn(ctx, conn, p.Module, link.ToolRetiredDelete, p.Node, map[string]any{ + "consumer": consumer, "confirm": consumer, "why": strings.TrimSpace(*f.why), "by": link.Caller(), + "via": link.ViaController, + }, retirementAsk) + if err != nil { + return err + } + if answer.Error != "" { + return fmt.Errorf("%s on %s refused to delete %s: %s", p.Module, p.Node, consumer, answer.Error) + } + var done struct { + FreedBytes *int64 `json:"freed_bytes"` + } + _ = unmarshalAnswer(answer, &done) + fmt.Printf("%s on %s deleted %s (%s freed)\n", p.Module, p.Node, consumer, sizeWords(done.FreedBytes)) + return nil +} + +// deleteOlderThan lists every consumer retired more than days ago, and deletes them only with confirm. +// One whose age the provider cannot say is never in it. +func deleteOlderThan(ctx context.Context, inv *inventory.Inventory, conn *nats.Conn, days int, confirm bool, + f handActFlags, now time.Time) error { + listing, err := gatherRetired(ctx, inv, conn, now) + if err != nil { + return err + } + return deleteFrom(ctx, conn, listing, days, confirm, f) +} + +func deleteFrom(ctx context.Context, conn *nats.Conn, listing retiredListing, days int, confirm bool, f handActFlags) error { + var due []retiredRow + unknown := 0 + for _, r := range listing.Retired { + switch { + case r.AgeDays < 0: + unknown++ + case r.AgeDays > days: + due = append(due, r) + } + } + for _, u := range listing.Unasked { + fmt.Printf("not asked: %s\n", u) + } + if unknown > 0 { + fmt.Printf("%d retired consumer(s) whose age their provider cannot say are left out\n", unknown) + } + if len(due) == 0 { + fmt.Printf("nothing has been retired more than %d day(s)\n", days) + return nil + } + fmt.Printf("retired more than %d day(s):\n", days) + for _, r := range due { + fmt.Printf(" %s on %s: %s — %d day(s), %s\n", r.Module, r.Node, r.Consumer, r.AgeDays, sizeWords(r.SizeBytes)) + } + if !confirm { + fmt.Printf("nothing was deleted: add --confirm to delete these %d\n", len(due)) + return nil + } + var failed []string + for _, r := range due { + if err := deleteRetired(ctx, conn, providerInstance{Node: r.Node, Module: r.Module}, r.Consumer, f); err != nil { + failed = append(failed, err.Error()) + } + } + if len(failed) > 0 { + return fmt.Errorf("%d of %d not deleted: %s", len(failed), len(due), strings.Join(failed, "; ")) + } + return nil +} diff --git a/cmd/mesh-controller/retirement.go b/cmd/mesh-controller/retirement.go new file mode 100644 index 0000000..c61207a --- /dev/null +++ b/cmd/mesh-controller/retirement.go @@ -0,0 +1,343 @@ +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 --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 +} diff --git a/cmd/mesh-controller/retirement_test.go b/cmd/mesh-controller/retirement_test.go new file mode 100644 index 0000000..8b3a089 --- /dev/null +++ b/cmd/mesh-controller/retirement_test.go @@ -0,0 +1,460 @@ +package main + +import ( + "context" + "encoding/json" + "flag" + "os" + "slices" + "sort" + "strings" + "sync" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/conditions" + "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): what a provider says becomes a condition a person answers, and the verbs that +// answer it ask the provider — never anything else — for exactly what it said. + +func waitingWord(module, node string, names ...string) link.Retirement { + r := link.Retirement{Module: module, Provider: "postgres-database", ProviderNode: node, Change: link.RetireWaiting, + Held: 7, Bound: "more than 3, or more than half of the 7 held", At: time.Now(), Since: time.Now()} + for _, n := range names { + r.Consumers = append(r.Consumers, link.RetiredConsumer{Consumer: n, Node: "laptop"}) + } + return r +} + +func openKeys(t *testing.T, k *conditions.Keeper) map[string]conditions.Condition { + t.Helper() + all, err := k.Open(t.Context()) + if err != nil { + t.Fatal(err) + } + out := map[string]conditions.Condition{} + for _, c := range all { + out[c.Key] = c + } + return out +} + +func TestARetirementWaitingIsUrgentARejectionAWarningAndARetirementClears(t *testing.T) { + k, _ := withConditionsInMemory(t) + var recorded []link.HandAct + r := retirements{keeper: func() *conditions.Keeper { return k }, + record: func(_ context.Context, a link.HandAct) error { recorded = append(recorded, a); return nil }} + ctx := t.Context() + waiting := waitingWord("postgres", "anchor", "a", "b", "c", "d") + if err := r.Retired(ctx, waiting); err != nil { + t.Fatal(err) + } + open := openKeys(t, k) + c, ok := open["provider.postgres.anchor.retire"] + if !ok || c.Kind != kindRetireWaiting || c.Severity != conditions.Urgent { + t.Fatalf("waiting is not an urgent retire-waiting condition: %+v", open) + } + for _, want := range []string{"a (laptop), b (laptop), c (laptop), d (laptop)", "retire approve anchor postgres", + "retire reject anchor postgres", "more than 3"} { + if !strings.Contains(c.Summary, want) { + t.Errorf("the condition does not say %q: %s", want, c.Summary) + } + } + + // Rejected, through the controller: the urgent one becomes a warning, and nothing is recorded here — + // the verb recorded it before it asked. + rejected := waiting + rejected.Change, rejected.By, rejected.Why, rejected.Via = link.RetireRejected, "operator", "moving them", link.ViaController + if err := r.Retired(ctx, rejected); err != nil { + t.Fatal(err) + } + open = openKeys(t, k) + if _, still := open["provider.postgres.anchor.retire"]; still { + t.Fatal("a rejection left the provider waiting") + } + if c, ok := open["provider.postgres.anchor.retire-rejected"]; !ok || c.Severity != conditions.Warning || + c.Kind != kindRetireRejected || !strings.Contains(c.Summary, "moving them") { + t.Fatalf("a rejection is not a warning saying why: %+v", open) + } + if len(recorded) != 0 { + t.Fatalf("an act through the controller was recorded twice: %+v", recorded) + } + + // Approved afterwards, asked of the provider directly: both cleared, and recorded by hand here. + approved := waiting + approved.Change, approved.By, approved.Why, approved.Via = link.RetireApproved, "someone", "done moving", "" + if err := r.Retired(ctx, approved); err != nil { + t.Fatal(err) + } + if open := openKeys(t, k); len(open) != 0 { + t.Fatalf("an approval left conditions open: %+v", open) + } + if len(recorded) != 1 || recorded[0].Verb != "retire approve" || recorded[0].By != "someone" || + !strings.Contains(recorded[0].Why, "directly") || !slices.Contains(recorded[0].Args, "d") { + t.Fatalf("an approval outside the controller was not recorded: %+v", recorded) + } + + // Waiting again, then the set settles (asked for again): cleared. And retired clears too. + for _, change := range []string{link.RetireSettled, link.RetireRetired} { + if err := r.Retired(ctx, waiting); err != nil { + t.Fatal(err) + } + done := waiting + done.Change = change + if err := r.Retired(ctx, done); err != nil { + t.Fatal(err) + } + if open := openKeys(t, k); len(open) != 0 { + t.Fatalf("%s left conditions open: %+v", change, open) + } + } + + // A deletion asked of the provider directly is recorded too. + deleted := link.Retirement{Module: "postgres", ProviderNode: "anchor", Change: link.RetireDeleted, By: "x", + Why: "gone for good", At: time.Now(), Consumers: []link.RetiredConsumer{{Consumer: "a"}}} + if err := r.Retired(ctx, deleted); err != nil { + t.Fatal(err) + } + if len(recorded) != 2 || recorded[1].Verb != "cleanup delete" || recorded[1].Cause != kindCleanupWaiting { + t.Fatalf("a deletion outside the controller was not recorded: %+v", recorded) + } +} + +func TestARetirementWordIsKeptAsItsConditionNamingTheEmitterFromTheSubject(t *testing.T) { + body, _ := json.Marshal(map[string]any{"provider": "oidc-client", "provider-node": "anchor", "change": "waiting", + "held": 2, "consumers": []map[string]any{{"consumer": "x"}, {"consumer": "y"}}, "module": "liar"}) + r, err := link.ReadRetirement("mesh.mod.idp.event.provisioner.retirement", body) + if err != nil || r.Module != "idp" || len(r.Consumers) != 2 { + t.Fatalf("%+v %v", r, err) + } + if _, err := link.ReadRetirement("mesh.mod.idp.event.provisioner.retirement", + []byte(`{"provider-node":"anchor","change":"vanished"}`)); err == nil { + t.Fatal("a change the mesh has no name for was read") + } + if _, err := link.ReadRetirement("mesh.mod.idp.event.provisioner.failing", body); err == nil { + t.Fatal("a failing word was read as a retirement") + } + k, _ := withConditionsInMemory(t) + if err := (retirements{keeper: func() *conditions.Keeper { return k }}).Retired(t.Context(), r); err != nil { + t.Fatal(err) + } + if _, ok := openKeys(t, k)["provider.idp.anchor.retire"]; !ok { + t.Fatal("not kept under the emitter the subject names") + } +} + +func TestARetirementConditionOfAnUnassignedProviderClears(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + register(t, open, catalogue.Manifest{Module: "pg", Version: "1", + Receives: map[string]string{"postgres-database": "/var/lib/mesh/pg/mesh.json"}}) + if _, err := assign(ctx, open, "anchor", "pg"); err != nil { + t.Fatal(err) + } + r := retirements{keeper: func() *conditions.Keeper { return conditionsFrom }} + for _, w := range []link.Retirement{waitingWord("pg", "anchor", "a", "b", "c", "d"), waitingWord("gone", "anchor", "a", "b", "c", "d")} { + if err := r.Retired(ctx, w); err != nil { + t.Fatal(err) + } + } + if err := conditionsFrom.Reconcile(ctx, "D11", []conditions.Observation{ + cleanupObservation(providerInstance{Node: "anchor", Module: "gone"}, + []link.RetiredConsumer{{Consumer: "a", RetiredAt: time.Now().Add(-40 * 24 * time.Hour).Format(time.RFC3339)}}, time.Now()), + }); err != nil { + t.Fatal(err) + } + all, _ := conditionsFrom.Open(ctx) + if len(all) != 3 { + t.Fatalf("%+v", all) + } + if err := unassignedProviders(ctx, open.inventory, conditionsFrom, providerConditions(all)); err != nil { + t.Fatal(err) + } + left := openKeys(t, conditionsFrom) + if len(left) != 1 || left["provider.pg.anchor.retire"].Key == "" { + t.Fatalf("only the assigned provider's waiting should stay: %+v", left) + } +} + +// fakeProvider answers the retirement tools on one machine, over a real bus, keeping what it was asked. +type fakeProvider struct { + mu sync.Mutex + state link.RetirementState + asked []map[string]any + deleted []string + approved []string + rejected []string +} + +func (f *fakeProvider) serve(t *testing.T, conn *nats.Conn, module, node string) { + t.Helper() + answer := func(m *nats.Msg, result any, refusal string) { + body, _ := json.Marshal(map[string]any{"result": result, "error": refusal, "node": node}) + _ = m.Respond(body) + } + names := func(args map[string]any) []string { + var out []string + for _, v := range args["consumers"].([]any) { + out = append(out, v.(string)) + } + sort.Strings(out) + return out + } + set := func(cs []link.RetiredConsumer) []string { + var out []string + for _, c := range cs { + out = append(out, c.Consumer) + } + sort.Strings(out) + return out + } + for _, tool := range []string{link.ToolRetirement, link.ToolRetireApprove, link.ToolRetireReject, link.ToolRetiredDelete} { + tool := tool + sub, err := conn.Subscribe(link.ModuleToolOn(module, tool, node), func(m *nats.Msg) { + f.mu.Lock() + defer f.mu.Unlock() + var args map[string]any + _ = json.Unmarshal(m.Data, &args) + f.asked = append(f.asked, map[string]any{"tool": tool, "args": args}) + switch tool { + case link.ToolRetirement: + answer(m, f.state, "") + case link.ToolRetireApprove: + if f.state.Waiting == nil || !slices.Equal(names(args), set(f.state.Waiting.Consumers)) { + answer(m, nil, "not the set I wait with") + return + } + f.approved = names(args) + answer(m, map[string]any{"retired": f.approved}, "") + case link.ToolRetireReject: + if f.state.Waiting == nil || !slices.Equal(names(args), set(f.state.Waiting.Consumers)) { + answer(m, nil, "not the set I wait with") + return + } + f.rejected = names(args) + answer(m, map[string]any{"kept": f.rejected}, "") + case link.ToolRetiredDelete: + if args["confirm"] != args["consumer"] || args["why"] == "" || args["via"] != link.ViaController { + answer(m, nil, "confirm, why and via") + return + } + f.deleted = append(f.deleted, args["consumer"].(string)) + answer(m, map[string]any{"deleted": args["consumer"], "freed_bytes": 1024}, "") + } + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sub.Unsubscribe() }) + } + if err := conn.Flush(); err != nil { + t.Fatal(err) + } +} + +func onATestBus(t *testing.T) *nats.Conn { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + js, err := broker.Dial(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(js.Close) + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + before := handActConn + handActConn = js.Conn() + t.Cleanup(func() { handActConn = before }) + return js.Conn() +} + +func whyFlags(t *testing.T, why string) handActFlags { + t.Helper() + set := flag.NewFlagSet("t", flag.ContinueOnError) + f := addHandActFlags(set) + if err := set.Parse([]string{"--why", why}); err != nil { + t.Fatal(err) + } + return f +} + +func handActsBy(t *testing.T, conn *nats.Conn, verb string) []link.HandAct { + t.Helper() + acts, err := link.HandActs(t.Context(), conn, time.Now().Add(-time.Minute)) + if err != nil { + t.Fatal(err) + } + var out []link.HandAct + for _, a := range acts { + if a.Verb == verb { + out = append(out, a) + } + } + return out +} + +func retiredDaysAgo(name string, days int) link.RetiredConsumer { + size := int64(4096) + return link.RetiredConsumer{Consumer: name, Kind: "consumer", Why: "the mesh stopped asking for it", + RetiredAt: time.Now().Add(-time.Duration(days) * 24 * time.Hour).UTC().Format(time.RFC3339), SizeBytes: &size} +} + +func TestNatsApproveSendsTheExactSetTheProviderWaitsWith(t *testing.T) { + conn := onATestBus(t) + fake := &fakeProvider{} + fake.state.Waiting = &link.RetirementWaiting{Consumers: []link.RetiredConsumer{{Consumer: "d"}, {Consumer: "b"}, {Consumer: "a"}, {Consumer: "c"}}, Held: 6} + fake.serve(t, conn, "pg-approve", "anchor") + p := providerInstance{Node: "anchor", Module: "pg-approve"} + said := printed(t, func() error { return answerRetirement(t.Context(), conn, p, true, whyFlags(t, "moved them")) }) + if !slices.Equal(fake.approved, []string{"a", "b", "c", "d"}) || !strings.Contains(said, "retired d, b, a, c") { + t.Fatalf("approved %v; said %s", fake.approved, said) + } + acts := handActsBy(t, conn, "retire approve") + if len(acts) == 0 || acts[len(acts)-1].Cause != kindRetireWaiting || + acts[len(acts)-1].Condition != "provider.pg-approve.anchor.retire" { + t.Fatalf("the approval is not in the hand-act log: %+v", acts) + } + + // Reject, on another provider: the same set, and nothing approved. + other := &fakeProvider{state: fake.state} + other.serve(t, conn, "pg-reject", "anchor") + printed(t, func() error { + return answerRetirement(t.Context(), conn, providerInstance{Node: "anchor", Module: "pg-reject"}, false, + whyFlags(t, "still moving")) + }) + if !slices.Equal(other.rejected, []string{"a", "b", "c", "d"}) || other.approved != nil { + t.Fatalf("rejected %v approved %v", other.rejected, other.approved) + } + + // Nothing waiting: refused, nothing asked but the question. + idle := &fakeProvider{} + idle.serve(t, conn, "pg-idle", "anchor") + if err := answerRetirement(t.Context(), conn, providerInstance{Node: "anchor", Module: "pg-idle"}, true, + whyFlags(t, "x")); err == nil || !strings.Contains(err.Error(), "waits for nobody") { + t.Fatalf("%v", err) + } + if len(idle.asked) != 1 { + t.Fatalf("an idle provider was asked more than its state: %+v", idle.asked) + } +} + +func TestNatsCleanupDeletesOnlyTheNamedRetiredConsumer(t *testing.T) { + conn := onATestBus(t) + fake := &fakeProvider{} + fake.state.Held = []string{"active"} + fake.state.Retired = []link.RetiredConsumer{retiredDaysAgo("old", 40), retiredDaysAgo("young", 2)} + fake.serve(t, conn, "pg-clean", "anchor") + p := providerInstance{Node: "anchor", Module: "pg-clean"} + said := printed(t, func() error { return deleteRetired(t.Context(), conn, p, "old", whyFlags(t, "not needed")) }) + if !slices.Equal(fake.deleted, []string{"old"}) || !strings.Contains(said, "deleted old (1.0 KB freed)") { + t.Fatalf("deleted %v; said %s", fake.deleted, said) + } + // An active consumer, or one it does not hold, is refused before the provider is asked to delete. + for _, name := range []string{"active", "nobody"} { + if err := deleteRetired(t.Context(), conn, p, name, whyFlags(t, "x")); err == nil || + !strings.Contains(err.Error(), "only a retired consumer is deleted") { + t.Fatalf("%s: %v", name, err) + } + } + if !slices.Equal(fake.deleted, []string{"old"}) { + t.Fatalf("more was deleted: %v", fake.deleted) + } + if acts := handActsBy(t, conn, "cleanup delete"); len(acts) == 0 || acts[len(acts)-1].Args[2] != "old" { + t.Fatalf("the deletion is not in the hand-act log: %+v", acts) + } + + // Older than: listed, and nothing deleted without confirm; with it, only the old one. + listing := retiredOf(t.Context(), conn, []providerInstance{p}, time.Now()) + if len(listing.Retired) != 2 || listing.Retired[0].Consumer != "old" || listing.Retired[0].AgeDays != 40 { + t.Fatalf("%+v", listing) + } + fake.deleted = nil + said = printed(t, func() error { return deleteFrom(t.Context(), conn, listing, 30, false, whyFlags(t, "tidy")) }) + if fake.deleted != nil || !strings.Contains(said, "nothing was deleted: add --confirm") || !strings.Contains(said, "old") || + strings.Contains(said, "young") { + t.Fatalf("deleted %v; said %s", fake.deleted, said) + } + printed(t, func() error { return deleteFrom(t.Context(), conn, listing, 30, true, whyFlags(t, "tidy")) }) + if !slices.Equal(fake.deleted, []string{"old"}) { + t.Fatalf("confirmed, deleted %v", fake.deleted) + } +} + +func TestNatsD11SaysCleanupWaitsAfterThirtyDays(t *testing.T) { + conn := onATestBus(t) + old := &fakeProvider{} + old.state.Retired = []link.RetiredConsumer{retiredDaysAgo("mesh_a_letta", 31), retiredDaysAgo("fresh", 1)} + old.serve(t, conn, "pg-d11-old", "anchor") + young := &fakeProvider{} + young.state.Retired = []link.RetiredConsumer{retiredDaysAgo("x", 29)} + young.serve(t, conn, "pg-d11-young", "anchor") + instances := []providerInstance{{Node: "anchor", Module: "pg-d11-old"}, {Node: "anchor", Module: "pg-d11-young"}, + // Nothing serves this one: a provider older than the question is passed over, not a failure. + {Node: "anchor", Module: "pg-d11-predates"}} + found, err := retiredTooLong(t.Context(), conn, instances, time.Now()) + if err != nil { + t.Fatal(err) + } + if len(found) != 1 || found[0].Key() != "provider.pg-d11-old.anchor.cleanup" || found[0].Kind != kindCleanupWaiting || + found[0].Severity != conditions.Warning || !strings.Contains(found[0].Summary, "mesh_a_letta (31 days") || + strings.Contains(found[0].Summary, "fresh") { + t.Fatalf("%+v", found) + } +} + +func TestTheRetireAndCleanupVerbsComposeTheirCommandLines(t *testing.T) { + for _, c := range []struct { + verb string + args map[string]any + want string + }{ + {"retire", map[string]any{}, "retire --json"}, + {"retire", map[string]any{"answer": "approve", "node": "anchor", "module": "postgres", "why": "moved"}, + "retire approve anchor postgres --why moved"}, + {"retire", map[string]any{"answer": "reject", "node": "anchor", "module": "postgres", "why": "no"}, + "retire reject anchor postgres --why no"}, + {"cleanup", map[string]any{}, "cleanup list --json"}, + {"cleanup", map[string]any{"node": "anchor", "module": "postgres", "consumer": "x", "why": "gone"}, + "cleanup delete anchor postgres x --why gone"}, + {"cleanup", map[string]any{"older-than": "30", "why": "tidy"}, "cleanup delete --older-than 30 --why tidy"}, + {"cleanup", map[string]any{"older-than": "30", "why": "tidy", "confirm": "true"}, + "cleanup delete --older-than 30 --why tidy --confirm"}, + } { + argv, err := argvFor(c.verb, c.args) + if err != nil || strings.Join(argv, " ") != c.want { + t.Errorf("%s %v: %q %v, want %q", c.verb, c.args, argv, err, c.want) + } + } + for _, args := range []map[string]any{ + {"answer": "approve", "node": "anchor", "module": "postgres"}, // no why + {"answer": "maybe", "node": "anchor", "module": "postgres", "why": "x"}, + } { + if _, err := argvFor("retire", args); err == nil { + t.Errorf("retire %v was composed", args) + } + } + for _, args := range []map[string]any{ + {"node": "anchor", "module": "postgres", "consumer": "x"}, // no why + {"older-than": "30"}, + {"consumer": "x", "older-than": "30", "node": "a", "module": "b", "why": "y"}, + } { + if _, err := argvFor("cleanup", args); err == nil { + t.Errorf("cleanup %v was composed", args) + } + } + if repairingCommand([]string{"cleanup", "delete", "a", "b", "c"}) == "" || + repairingCommand([]string{"retire", "approve", "a", "b"}) == "" || repairingCommand([]string{"retire"}) != "" { + t.Error("the generic verb would let a retirement or a deletion through without a why") + } +} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 5b5e2a4..9813068 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -455,6 +455,50 @@ func (a *verbArguments) commandLine() ([]string, error) { } } return argv, nil + case "retire": + answer := str("answer") + if answer == "" { + return []string{"retire", "--json"}, nil + } + if answer != "approve" && answer != "reject" { + return nil, fmt.Errorf("retire answers approve or reject, not %q", answer) + } + if err := need("node", "module", "why"); err != nil { + return nil, err + } + argv := []string{"retire", answer, str("node"), str("module"), "--why", str("why")} + if c := str("cause"); c != "" { + argv = append(argv, "--cause", c) + } + return argv, nil + case "cleanup": + consumer, older := str("consumer"), str("older-than") + switch { + case consumer != "" && older != "": + return nil, errors.New("cleanup deletes one consumer or those older than some days, not both") + case consumer != "": + if err := need("node", "module", "why"); err != nil { + return nil, err + } + argv := []string{"cleanup", "delete", str("node"), str("module"), consumer, "--why", str("why")} + if c := str("cause"); c != "" { + argv = append(argv, "--cause", c) + } + return argv, nil + case older != "": + if err := need("why"); err != nil { + return nil, err + } + argv := []string{"cleanup", "delete", "--older-than", older, "--why", str("why")} + if on("confirm") { + argv = append(argv, "--confirm") + } + if c := str("cause"); c != "" { + argv = append(argv, "--cause", c) + } + return argv, nil + } + return []string{"cleanup", "list", "--json"}, nil case "doctor": which := 0 argv := []string{"doctor"} @@ -554,7 +598,7 @@ func (a *verbArguments) commandLine() ([]string, error) { // jsonVerbs are the verbs whose command speaks JSON, so the answer carries it as data as well. var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "collection": true, - "hand-acts": true, "durations": true, "conditions": true, "doctor": true} + "hand-acts": true, "durations": true, "conditions": true, "doctor": true, "retire": true, "cleanup": true} // repairingCommand names a command line that repairs by hand, and so says why: a push, a plan stopped // or closed, a consumer re-made (novox/hq to-be 45 §7). Empty for any other. @@ -570,6 +614,10 @@ func repairingCommand(argv []string) string { return "hand-act record" case argv[0] == "conditions" && len(argv) > 1 && argv[1] == "silence": return "conditions silence" + case argv[0] == "retire" && len(argv) > 1 && (argv[1] == "approve" || argv[1] == "reject"): + return "retire " + argv[1] + case argv[0] == "cleanup" && len(argv) > 1 && argv[1] == "delete": + return "cleanup delete" } return "" } diff --git a/cmd/mesh-controller/seatverbs_schema_test.go b/cmd/mesh-controller/seatverbs_schema_test.go index df4c877..7363ebb 100644 --- a/cmd/mesh-controller/seatverbs_schema_test.go +++ b/cmd/mesh-controller/seatverbs_schema_test.go @@ -274,6 +274,8 @@ var accountedFlags = map[string]map[string]string{ }, "hand-acts": {"json": "set by the verb: the answer is data"}, "conditions": {"json": "set by the verb: the answer is data"}, + "retire": {"json": "set by the verb: the answer is data"}, + "cleanup": {"json": "set by the verb: the answer is data"}, "conditions history": {"json": "set by the verb: the answer is data"}, "conditions show": {"json": "set by the verb: the answer is data"}, "healers": {"json": "set by the verb: the answer is data"}, diff --git a/cmd/mesh-controller/standing.go b/cmd/mesh-controller/standing.go index 58df850..07ce8ac 100644 --- a/cmd/mesh-controller/standing.go +++ b/cmd/mesh-controller/standing.go @@ -103,16 +103,38 @@ func providerStandings(open []conditions.Condition) []conditions.Condition { return out } -// providerOf reads a standing's provider module and machine back from its key. -func providerOf(c conditions.Condition) (module, node, consumer string, ok bool) { - parts := strings.Split(c.Key, ".") - if len(parts) != 5 || parts[0] != conditions.ScopeProvider { - return "", "", "", false +// providerConditions is every open condition a provider's word raised: a consumer failing (ADR 0224), +// and a retirement waiting for a person, kept by a rejection, or cleanup waiting (ADR 0230). +func providerConditions(open []conditions.Condition) []conditions.Condition { + var out []conditions.Condition + for _, c := range open { + switch c.Kind { + case kindProviderFailing, kindRetireWaiting, kindRetireRejected, kindCleanupWaiting: + out = append(out, c) + } } - return parts[1], parts[2], parts[3], true + return out } -// unassignedProviders clears the standing of every provider no longer assigned where it ran. +// providerOf reads a provider condition's module and machine back from its key: a standing's +// `provider....failing`, and a retirement's `provider...`, +// which names no consumer. +func providerOf(c conditions.Condition) (module, node, consumer string, ok bool) { + parts := strings.Split(c.Key, ".") + if parts[0] != conditions.ScopeProvider { + return "", "", "", false + } + switch len(parts) { + case 5: + return parts[1], parts[2], parts[3], true + case 4: + return parts[1], parts[2], "", true + } + return "", "", "", false +} + +// unassignedProviders clears the standing of every provider no longer assigned where it ran — and its +// retirement and cleanup conditions with it (ADR 0230): nothing runs there to retire or delete anything. // // **A provider no longer assigned is not asked about** (ADR 0224 §4): nothing runs there to fail // anybody, and nothing there will ever say it recovered. The observation that resolves it is the @@ -120,7 +142,7 @@ func providerOf(c conditions.Condition) (module, node, consumer string, ok bool) func unassignedProviders(ctx context.Context, inv *inventory.Inventory, k *conditions.Keeper, open []conditions.Condition) error { assigned := map[string]map[string]bool{} - for _, c := range providerStandings(open) { + for _, c := range providerConditions(open) { module, node, _, ok := providerOf(c) if !ok { continue diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index b121bd6..656169c 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -66,6 +66,10 @@ type signalFacts struct { standings []conditions.Condition standingsErr error + // providerWords are every open condition a provider's word raised — failing, waiting for a person + // to approve a retirement, kept by a rejection, cleanup waiting (ADR 0224, 0230) — so a provider no + // longer assigned has them all cleared. + providerWords []conditions.Condition advisories []link.Advisory lostConsumers map[string]bool @@ -228,7 +232,7 @@ func (w *watchdogs) see(running context.Context, f *signalFacts) { } // A provider no longer assigned where it ran: its standing is resolved by the assignment. if f.standingsErr == nil && w.open != nil { - if err := unassignedProviders(running, w.open.inventory, w.keeper, f.standings); err != nil { + if err := unassignedProviders(running, w.open.inventory, w.keeper, f.providerWords); err != nil { problems = append(problems, "S8: "+err.Error()) } } @@ -297,6 +301,7 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts { f.standingsErr = err } else { f.standings = providerStandings(open) + f.providerWords = providerConditions(open) } f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet)) f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories) diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 5f92855..8e34fa5 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -255,13 +255,18 @@ var ControllerFollows = []string{ // the index is a name. moduleEventSubject("*", ProvisionerFailing), moduleEventSubject("*", ProvisionerRecovered), + // **What becomes of a consumer the mesh stopped asking for** (novox/hq ADR 0230): retired, waiting + // for a person, re-enabled, deleted — a provider's third word, from whichever module provides. + // Appended, because the index is a name. + moduleEventSubject("*", ProvisionerRetirement), } // The provider standing events, by their local names. Written here as well as in the catalogue // (catalogue.ProvisionerEvents), which this package cannot import; a test keeps them agreeing. const ( - ProvisionerFailing = "provisioner.failing" - ProvisionerRecovered = "provisioner.recovered" + ProvisionerFailing = "provisioner.failing" + ProvisionerRecovered = "provisioner.recovered" + ProvisionerRetirement = "provisioner.retirement" ) // moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 782634a..850937b 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -25,7 +25,7 @@ accounts { users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } - subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } + subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/broker/writers.go b/internal/broker/writers.go index 63d0862..21b5c37 100644 --- a/internal/broker/writers.go +++ b/internal/broker/writers.go @@ -117,9 +117,10 @@ var WritersTable = []WriterRow{ {State: "a merge announced", Writer: "one announcer per forge (the hook, or the poll when the hook is absent — never both)", KeptIn: "the bus", Others: "—", Subjects: []string{"mesh.mod.*.event.pull.merged"}, Writes: ownModule}, {State: "a provider's standing", Writer: "the provider", KeptIn: "the provider's events", - Others: "the controller keeps the newest word as a condition", - Subjects: []string{"mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered"}, - Writes: ownModule}, + Others: "the controller keeps the newest word as a condition", + Subjects: []string{"mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", + "mesh.mod.*.event.provisioner.retirement"}, + Writes: ownModule}, {State: "the operator-channel's open messages", Writer: "the seat's holder", KeptIn: "its own key-value state", Others: "—"}, {State: "the facts snapshot", Writer: "controller", KeptIn: "the artifact store, facts/latest", diff --git a/internal/broker/writers_test.go b/internal/broker/writers_test.go index b4f17d9..8b3bb8e 100644 --- a/internal/broker/writers_test.go +++ b/internal/broker/writers_test.go @@ -55,7 +55,7 @@ func TestAWholeMeshComposesWithOneWriterPerState(t *testing.T) { {Kind: KindEnrolment, Node: "two", PasswordHash: "x"}, {Kind: KindPerson, Module: "jochen", Invokes: []string{"*"}, PasswordHash: "x"}, {Kind: KindModule, Node: "one", Module: "gitea", Emits: []string{"pull.merged"}, PasswordHash: "x"}, - {Kind: KindModule, Node: "one", Module: "postgres", Emits: []string{"provisioner.failing", "provisioner.recovered"}, + {Kind: KindModule, Node: "one", Module: "postgres", Emits: []string{"provisioner.failing", "provisioner.recovered", "provisioner.retirement"}, PasswordHash: "x"}, {Kind: KindModule, Node: "one", Module: "build-agent", Holds: []Seat{builder}, PasswordHash: "x"}, {Kind: KindNodeTools, Node: "one", Module: RuntimeModule, PasswordHash: "x", Carries: []Declared{ diff --git a/internal/catalogue/events.go b/internal/catalogue/events.go index 2db9336..ce8b9dd 100644 --- a/internal/catalogue/events.go +++ b/internal/catalogue/events.go @@ -164,13 +164,19 @@ func consumePattern(pattern string) error { // The events a provider says about its consumers (novox/hq ADR 0224): a consumer it has failed // without one success for minutes, and that consumer succeeding again or being withdrawn. The // controller follows them from every module and `status` names a consumer failing until it recovers. +// +// **And what becomes of a consumer the mesh stopped asking for** (novox/hq ADR 0230): retired — its +// access disabled and its data kept — after the same answer in five passes, waiting for a person when +// more would go than the bound allows, re-enabled when asked for again, and deleted only by a person's +// `cleanup delete`. One event, its `change` saying which. const ( - ProvisionerFailing = "provisioner.failing" - ProvisionerRecovered = "provisioner.recovered" + ProvisionerFailing = "provisioner.failing" + ProvisionerRecovered = "provisioner.recovered" + ProvisionerRetirement = "provisioner.retirement" ) -// ProvisionerEvents are both, in the order they are said. -var ProvisionerEvents = []string{ProvisionerFailing, ProvisionerRecovered} +// ProvisionerEvents are all three, in the order they are said. +var ProvisionerEvents = []string{ProvisionerFailing, ProvisionerRecovered, ProvisionerRetirement} // EmitsAll is every event a module may publish: what it declares and, for a module that receives // contributions — a provider, running a provisioner over them — the provider's standing events. diff --git a/internal/catalogue/retirement_event_test.go b/internal/catalogue/retirement_event_test.go new file mode 100644 index 0000000..ec02978 --- /dev/null +++ b/internal/catalogue/retirement_event_test.go @@ -0,0 +1,22 @@ +package catalogue + +import ( + "slices" + "testing" +) + +// Every provider may say what becomes of a consumer the mesh stopped asking for (novox/hq ADR 0230), +// whatever its manifest lists — as it may say a consumer it keeps failing (ADR 0224) — and a module +// that receives nothing may not. +func TestEveryProviderMaySayWhatItRetires(t *testing.T) { + provider := Manifest{Module: "pg", Receives: map[string]string{"postgres-database": "/x/mesh.json"}, + Emits: []string{"database.provisioned"}} + for _, e := range []string{ProvisionerFailing, ProvisionerRecovered, ProvisionerRetirement, "database.provisioned"} { + if !slices.Contains(provider.EmitsAll(), e) { + t.Errorf("a provider may not emit %s: %v", e, provider.EmitsAll()) + } + } + if slices.Contains(Manifest{Module: "app", Emits: []string{"x"}}.EmitsAll(), ProvisionerRetirement) { + t.Error("a module that provides nothing may say what it retires") + } +} diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 025bb7a..c54d7c4 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -267,6 +267,33 @@ var ControllerVerbs = []Verb{ "probes": "\"true\": the registry — what each probe asserts, and the condition it raises", "signals": "\"true\": the signals table, each row with the age of its newest signal", }, nil, "run", "probes", "signals")}, + // A consumer the mesh stopped asking for: retired, waiting for a person, deleted only by one + // (novox/hq ADR 0230). + {Name: "retire", Description: "A consumer the mesh stops asking for is retired by its provider — access " + + "disabled, data kept — after the same answer in five passes; more than three at once, or more than half " + + "of those held, waits for a person. With no answer, every provider that waits and what it would retire. " + + "With answer approve or reject, a node and a module: retire exactly what that provider waits with, or " + + "keep it active — a hand act, which says why (novox/hq ADR 0230).", + Input: schema(map[string]string{ + "answer": "approve or reject: answer what the provider waits with (needs node, module and why)", + "node": "with answer: the machine the provider runs on", + "module": "with answer: the provider module", + "why": "with answer: why — required, and recorded in the hand-act log", + "cause": "with answer: the cause in a word (retire-waiting when absent)", + }, nil)}, + {Name: "cleanup", Description: "Every consumer a provider holds retired — its age, its size where the " + + "backend knows, and why it was retired. With consumer (and node, module): the provider deletes that one " + + "retired consumer — never an active one. With older-than: every retired consumer older than that many " + + "days, listed; deleted only with confirm. Deleting is a hand act, which says why (novox/hq ADR 0230).", + Input: schema(map[string]string{ + "node": "with consumer: the machine the provider runs on", + "module": "with consumer: the provider module", + "consumer": "delete this retired consumer (needs node, module and why)", + "older-than": "delete every consumer retired more than this many days (needs why; lists only without confirm)", + "confirm": "\"true\": with older-than, delete what is listed", + "why": "with consumer or older-than: why — required, and recorded in the hand-act log", + "cause": "with consumer or older-than: the cause in a word (cleanup-waiting when absent)", + }, nil, "confirm")}, {Name: "build", Description: "Have the build machine build a repository. Answers at once with the build's id: " + "`builds` with that id follows it line by line, and the module is registered when the outcome comes.", Input: schema(map[string]string{ diff --git a/internal/link/contracts.go b/internal/link/contracts.go index bcf44cd..a2a54b2 100644 --- a/internal/link/contracts.go +++ b/internal/link/contracts.go @@ -42,6 +42,10 @@ var Contracts = map[string]Contract{ "catalogue's record, which is the order, so a late one reads the same record"}, KindCatchUp: {Unordered: "a catalogue asking what it missed: answered from the record, whenever asked"}, KindProvisioner: {Unordered: "a provider's newest word about a consumer, said again every fifteen minutes " + - "while it holds (ADR 0224): the condition keeps the last observed, and S8 says when the words stop", - Tests: []string{"TestAProviderFailingAConsumerBreaksAllWellUntilItRecovers"}}, + "while it holds (ADR 0224): the condition keeps the last observed, and S8 says when the words stop. " + + "Its retirement word (ADR 0230) is unordered too: a waiting set is said again every fifteen minutes, " + + "and every other change is checked against the provider itself — `retire` and `cleanup` ask it, " + + "and the self-check's D11 asks it every run — so an older word read late is corrected by the next", + Tests: []string{"TestAProviderFailingAConsumerBreaksAllWellUntilItRecovers", + "TestARetirementWordIsKeptAsItsConditionNamingTheEmitterFromTheSubject"}}, } diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index cace2dc..6c330ed 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -262,7 +262,7 @@ func kindOfSubject(subject string) (string, bool) { } // ProvisionerEmitter is the module a provider's standing event came from, read from its subject -// (`mesh.mod..event.provisioner.`); false for any other subject. The +// (`mesh.mod..event.provisioner.`); false for any other subject. The // controller's own follow pattern, with `*` for the module, decodes too. func ProvisionerEmitter(subject string) (string, bool) { rest, ok := strings.CutPrefix(subject, "mesh.mod.") @@ -273,7 +273,8 @@ func ProvisionerEmitter(subject string) (string, bool) { if !ok || module == "" || strings.Contains(module, ".") { return "", false } - if event != broker.ProvisionerFailing && event != broker.ProvisionerRecovered { + if event != broker.ProvisionerFailing && event != broker.ProvisionerRecovered && + event != broker.ProvisionerRetirement { return "", false } return module, true diff --git a/internal/link/retirement.go b/internal/link/retirement.go new file mode 100644 index 0000000..9a1f70b --- /dev/null +++ b/internal/link/retirement.go @@ -0,0 +1,241 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// What becomes of a consumer the mesh stopped asking for (novox/hq ADR 0230). +// +// **A consumer is retired, not withdrawn, and deleted only by a person.** A provider that has seen the +// same consumers go unasked for in five passes disables their access and keeps their data, marked in +// its backend with when and why; asked for again, it enables them as they were. When more would go +// than its bound allows — more than three, or more than half of those it holds — it retires nothing and +// waits for `retire approve` or `retire reject`. A retired consumer is deleted only by `cleanup delete`, +// which the provider executes. Each of these is the provider's `provisioner.retirement` event, its +// `change` saying which; the controller keeps the ones that need a person as conditions. + +// The changes a provider says. +const ( + RetireWaiting = "waiting" + RetireSettled = "settled" + RetireApproved = "approved" + RetireRejected = "rejected" + RetireRetired = "retired" + RetireReenabled = "reenabled" + RetireDeleted = "deleted" + RetireAdopted = "adopted" +) + +// ViaController is what the controller's verbs say they came through, so an act a provider was asked +// for some other way is told apart and recorded by hand. +const ViaController = "mesh-controller" + +// RetiredConsumer is one consumer in a retirement word, or in a provider's answer. +type RetiredConsumer struct { + Consumer string `json:"consumer"` + Node string `json:"node,omitempty"` + RetiredAt string `json:"retired_at,omitempty"` + Why string `json:"why,omitempty"` + // SizeBytes is what it keeps on the backend; -1 or absent where the backend cannot say. + SizeBytes *int64 `json:"size_bytes,omitempty"` + // Kind is "consumer", or a backend's own word for something set aside ("set-aside-database"). + Kind string `json:"kind,omitempty"` +} + +// Retirement is one provider's word about consumers the mesh stopped asking for. +type Retirement struct { + // Module is the emitter, read from the subject the bus let it publish on — never from the body. + Module string `json:"-"` + + Provider string `json:"provider"` + ProviderNode string `json:"provider-node"` + Change string `json:"change"` + Consumers []RetiredConsumer `json:"consumers"` + Held int `json:"held"` + Bound string `json:"bound,omitempty"` + Why string `json:"why,omitempty"` + By string `json:"by,omitempty"` + Via string `json:"via,omitempty"` + At time.Time `json:"at"` + Since time.Time `json:"since,omitempty"` +} + +// Names are the consumers it names, in its order. +func (r Retirement) Names() []string { + out := make([]string, 0, len(r.Consumers)) + for _, c := range r.Consumers { + out = append(out, c.Consumer) + } + return out +} + +// Retirements keeps what providers say about consumers they retire. +type Retirements interface { + // Retired records one word. An error the store is away for is held and asked again, like a report. + Retired(ctx context.Context, r Retirement) error +} + +// KeepsRetirements says where retirement words are kept, and asks for them to be delivered. +func (s *Server) KeepsRetirements(r Retirements) error { + if err := s.inbound.Also(KindProvisioner); err != nil { + return err + } + s.retirements = r + return nil +} + +// IsRetirement says a subject is a provider's retirement word. +func IsRetirement(subject string) bool { + _, ok := ProvisionerEmitter(subject) + return ok && strings.HasSuffix(subject, ".event."+broker.ProvisionerRetirement) +} + +// ReadRetirement is one retirement word as the controller understands it. +func ReadRetirement(subject string, body []byte) (Retirement, error) { + module, ok := ProvisionerEmitter(subject) + if !ok || !IsRetirement(subject) { + return Retirement{}, fmt.Errorf("%s is not a provider's retirement word", subject) + } + var r Retirement + if err := json.Unmarshal(body, &r); err != nil { + return Retirement{}, fmt.Errorf("%s's retirement word could not be read: %w", module, err) + } + switch r.Change { + case RetireWaiting, RetireSettled, RetireApproved, RetireRejected, RetireRetired, RetireReenabled, + RetireDeleted, RetireAdopted: + default: + return Retirement{}, fmt.Errorf("%s said a retirement change the mesh has no name for: %q", module, r.Change) + } + if r.ProviderNode == "" { + return Retirement{}, fmt.Errorf("%s's retirement word named no machine", module) + } + r.Module = module + return r, nil +} + +// retirement acts on one retirement word. Like a recovery, a word that clears a condition is said +// once, so a store that is away holds the message rather than dropping it. +func (s *Server) retirement(ctx context.Context, m Control) { + if s.retirements == nil { + _ = m.Took() + return + } + r, err := ReadRetirement(m.Subject(), m.Body()) + if err != nil { + s.log.Printf("%v; ignored", err) + _ = m.Took() + return + } + err = s.retirements.Retired(ctx, r) + what := fmt.Sprintf("%s's retirement word (%s)", r.Module, r.Change) + switch s.decide(ctx, m, what, "", "", err) { + case Hold: + return + case Stale, GiveUp: + _ = m.Took() + return + } + if err != nil { + s.log.Printf("%s could not be kept: %v", what, err) + } else { + s.log.Printf("%s on %s %s %s", r.Module, r.ProviderNode, r.Change, strings.Join(r.Names(), ", ")) + } + _ = m.Took() +} + +// The tools every provider serves about its retired consumers (ADR 0230), asked of one machine. +const ( + ToolRetirement = "provisioner_retirement" + ToolRetireApprove = "provisioner_retire_approve" + ToolRetireReject = "provisioner_retire_reject" + ToolRetiredDelete = "provisioner_delete" +) + +// ErrNothingServes is a tool nothing on that machine serves: the module is not running there, or +// is older than the tool. +var ErrNothingServes = errors.New("nothing serves it") + +// ModuleToolOn is where one machine's instance of a module answers a tool: the module's tool subject +// with the machine as its last token — the subject the node tools bind beside the shared one. +func ModuleToolOn(module, tool, node string) string { return ToolSubject(module, tool) + "." + node } + +// AskModuleToolOn asks one machine's instance of a module one of its tools and reads its answer. +func AskModuleToolOn(ctx context.Context, conn *nats.Conn, module, tool, node string, args any, + timeout time.Duration) (Answer, error) { + body, err := json.Marshal(args) + if err != nil { + return Answer{}, err + } + asking, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + subject := ModuleToolOn(module, tool, node) + refused, stop := refusalsOf(conn, subject) + defer stop() + type replied struct { + msg *nats.Msg + err error + } + done := make(chan replied, 1) + go func() { + msg, err := conn.RequestWithContext(asking, subject, body) + done <- replied{msg, err} + }() + var reply *nats.Msg + select { + case r := <-done: + reply, err = r.msg, r.err + case why := <-refused: + cancel() + return Answer{}, fmt.Errorf("the bus refused the controller asking %s.%s on %s: %v", module, tool, node, why) + } + switch { + case errors.Is(err, nats.ErrNoResponders): + return Answer{}, fmt.Errorf("%w: %s on %s does not answer %s — it is not running there, or is "+ + "older than the tool", ErrNothingServes, module, node, tool) + case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout): + return Answer{}, fmt.Errorf("%s on %s did not answer %s within %s", module, node, tool, timeout) + case err != nil: + return Answer{}, err + } + var answer Answer + if err := json.Unmarshal(reply.Data, &answer); err != nil { + return Answer{}, fmt.Errorf("%s on %s answered %s with something unreadable: %w", module, node, tool, err) + } + return answer, nil +} + +// RetirementWaiting is the set a provider waits with for a person. +type RetirementWaiting struct { + Consumers []RetiredConsumer `json:"consumers"` + Since string `json:"since"` + Held int `json:"held"` +} + +// RetirementRejected is the set a person's rejection keeps active. +type RetirementRejected struct { + Consumers []RetiredConsumer `json:"consumers"` + By string `json:"by"` + Why string `json:"why"` + At string `json:"at"` +} + +// RetirementState is a provider's answer to provisioner_retirement. +type RetirementState struct { + Resource string `json:"resource"` + Node string `json:"node"` + Held []string `json:"held"` + StablePasses int `json:"stable_passes"` + Bound string `json:"bound"` + Waiting *RetirementWaiting `json:"waiting"` + Rejected *RetirementRejected `json:"rejected"` + Retired []RetiredConsumer `json:"retired"` +} diff --git a/internal/link/serve.go b/internal/link/serve.go index 6094329..d8da197 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -81,6 +81,8 @@ type Server struct { replayer Replayer // standings keeps what providers say about their consumers (novox/hq ADR 0224). standings Standings + // retirements keeps what providers say about consumers the mesh stopped asking for (ADR 0230). + retirements Retirements log *log.Logger // giveUp is how long one message is held for the store; zero means GiveUpAfter. diff --git a/internal/link/standing.go b/internal/link/standing.go index b861847..3add14e 100644 --- a/internal/link/standing.go +++ b/internal/link/standing.go @@ -81,6 +81,10 @@ func ReadStanding(subject string, body []byte) (Standing, error) { // it lasts, so one dropped is replaced; a recovery is said once, and dropping it would leave status // naming a consumer that is fine. So a store that is away holds the message, as a report is held. func (s *Server) provisioner(ctx context.Context, m Control) { + if IsRetirement(m.Subject()) { + s.retirement(ctx, m) + return + } if s.standings == nil { // Delivered because the consumer's filter names it, with nothing here keeping it: taken, // because handing it back would not give it anywhere to go. diff --git a/internal/link/standing_test.go b/internal/link/standing_test.go index 6a2d545..4d7b3f3 100644 --- a/internal/link/standing_test.go +++ b/internal/link/standing_test.go @@ -37,6 +37,7 @@ func TestTheControllerFollowsEveryProvidersStandingAndNothingElse(t *testing.T) for subject, want := range map[string]string{ "mesh.mod.keycloak.event.provisioner.failing": "keycloak", "mesh.mod.postgres.event.provisioner.recovered": "postgres", + "mesh.mod.minio.event.provisioner.retirement": "minio", "mesh.mod.*.event.provisioner.failing": "*", } { got, ok := ProvisionerEmitter(subject) @@ -63,8 +64,11 @@ func TestTheControllerFollowsEveryProvidersStandingAndNothingElse(t *testing.T) follows++ } } - if follows != 2 { - t.Fatalf("the controller follows %d standing subjects, want failing and recovered", follows) + if follows != 3 { + t.Fatalf("the controller follows %d provider subjects, want failing, recovered and retirement (ADR 0230)", follows) + } + if !IsRetirement("mesh.mod.postgres.event.provisioner.retirement") || IsRetirement("mesh.mod.postgres.event.provisioner.failing") { + t.Fatal("a retirement word is not told from a standing") } }