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).
This commit is contained in:
jochen
2026-10-06 14:41:15 +02:00
parent 298daa6fbe
commit 68009b16fe
23 changed files with 1678 additions and 27 deletions
+2
View File
@@ -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},
}
+12
View File
@@ -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 <name> roll-out [--together] ...send it to the machines running it
upgrade <name> 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 <node> <module> --why <text> retire what it waits with: access off, data kept
retire reject <node> <module> --why <text> keep them active; a warning stays open
cleanup [list] [--json] every retired consumer per provider: age, size, why
cleanup delete <node> <module> <consumer> --why <text> the provider deletes that one retired consumer
cleanup delete --older-than <days> --why <text> [--confirm] list those older; delete only with --confirm
seats [--json] every seat this mesh defines, what it delivers, and who holds it
seat rename <from> <to> rename a seat; its former name still resolves (ADR 0122)
seat <name> --to <node>/<module> hand a seat to that assignment as one act; never empty in between (ADR 0131)
+6
View File
@@ -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.
+434
View File
@@ -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 <node> <module> --why <text> | retire reject <node> <module> --why <text>"
const cleanupUsage = "cleanup [list] [--json] | cleanup delete <node> <module> <consumer> --why <text> | " +
"cleanup delete --older-than <days> --why <text> [--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
}
+343
View File
@@ -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 <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
}
+460
View File
@@ -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")
}
}
+49 -1
View File
@@ -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 ""
}
@@ -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"},
+30 -8
View File
@@ -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.<module>.<node>.<consumer>.failing`, and a retirement's `provider.<module>.<node>.<token>`,
// 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
+6 -1
View File
@@ -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)
+7 -2
View File
@@ -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
+1 -1
View File
@@ -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: {
+4 -3
View File
@@ -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",
+1 -1
View File
@@ -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{
+10 -4
View File
@@ -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.
@@ -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")
}
}
+27
View File
@@ -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{
+6 -2
View File
@@ -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"}},
}
+3 -2
View File
@@ -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.<module>.event.provisioner.<failing|recovered>`); false for any other subject. The
// (`mesh.mod.<module>.event.provisioner.<failing|recovered|retirement>`); 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
+241
View File
@@ -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"`
}
+2
View File
@@ -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.
+4
View File
@@ -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.
+6 -2
View File
@@ -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")
}
}