diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 10460a6..a4b0a3f 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -683,6 +683,10 @@ type answers struct { // (novox/hq ADR 0225): its provider leaves it out of the grants and composes everything else, so // this is the one place it is said across the mesh. Not well while there is any. overflowing []catalogue.Overflow + // handActs is how many acts were done by hand in the last seven days (novox/hq to-be 45 §7), nil + // where the log is not on hand; handActsUnread why it could not be read when it could not. + handActs *int + handActsUnread string } // heldBy is every artifact this mesh has built, for a build that may need one as its base. diff --git a/cmd/mesh-controller/calls_verb_test.go b/cmd/mesh-controller/calls_verb_test.go index 56c3a19..1aaea74 100644 --- a/cmd/mesh-controller/calls_verb_test.go +++ b/cmd/mesh-controller/calls_verb_test.go @@ -35,10 +35,10 @@ func TestAPushAnswersBeforeItSends(t *testing.T) { args map[string]any want bool }{ - {"push", map[string]any{"node": "anchor"}, true}, - {"push", map[string]any{}, true}, - {"command", map[string]any{"command": "push anchor"}, true}, - {"command", map[string]any{"command": "push --behind"}, true}, + {"push", map[string]any{"node": "anchor", "why": "w"}, true}, + {"push", map[string]any{"why": "w"}, true}, + {"command", map[string]any{"command": "push anchor --why w"}, true}, + {"command", map[string]any{"command": "push --behind --why=w"}, true}, {"command", map[string]any{"command": "builds"}, false}, {"status", map[string]any{}, false}, {"assign", map[string]any{"node": "anchor", "module": "m"}, false}, diff --git a/cmd/mesh-controller/durations.go b/cmd/mesh-controller/durations.go new file mode 100644 index 0000000..bb7e545 --- /dev/null +++ b/cmd/mesh-controller/durations.go @@ -0,0 +1,170 @@ +package main + +import ( + "context" + "encoding/json" + "flag" + "fmt" + "os" + "slices" + "sort" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// What the core's bounds are set from (novox/hq to-be 45 Phase 0). +// +// **A bound is set from what was measured, not from what seemed reasonable.** Phase 1 puts a watchdog +// on each row of the signals table, and each has a bound: S1 three heartbeat intervals, S2 three times +// a machine's last apply, S3 a tier's build and apply time, S6 a build's timeout. Marked provisional +// in the design until a fortnight of these says what the mesh actually takes. Recorded by the serving +// controller as it hears each — a send's first report, a machine's next word, a plan leaving a tier, a +// build's outcome — and summarised here per machine, repository or module. + +// recordBuildDuration measures one build from its ask to its outcome heard. +func recordBuildDuration(ctx context.Context, inv *inventory.Inventory, result link.BuildResult, asked time.Time) { + if asked.IsZero() || result.ID == "" { + return + } + subject := result.Module + if subject == "" { + subject = result.Repository + } + detail := "built" + if result.Failed != "" { + detail = "failed: " + firstLine(result.Failed) + } + if err := inv.RecordDuration(ctx, inventory.Duration{Kind: inventory.DurationBuild, Subject: subject, + Node: result.On, Ref: result.ID, Started: asked, Took: time.Since(asked), Detail: detail}); err != nil { + fmt.Fprintf(os.Stderr, "%s: how long it took could not be recorded: %v\n", result.ID, err) + } +} + +// durationSummary is one subject's measurements of one kind. +type durationSummary struct { + Kind string `json:"kind"` + Subject string `json:"subject"` + Count int `json:"count"` + Median string `json:"median"` + P90 string `json:"p90"` + Max string `json:"max"` + // Bound is what to-be 45's rule would make of these, where the rule is a multiple of a measured + // time: three times the slowest apply (S2), three times the median word interval (S1). + Suggests string `json:"suggests,omitempty"` +} + +func summarise(ds []inventory.Duration) []durationSummary { + type key struct{ kind, subject string } + by := map[key][]time.Duration{} + for _, d := range ds { + k := key{d.Kind, d.Subject} + by[k] = append(by[k], d.Took) + } + var out []durationSummary + for k, took := range by { + slices.Sort(took) + at := func(q float64) time.Duration { return took[int(q*float64(len(took)-1))] } + s := durationSummary{Kind: k.kind, Subject: k.subject, Count: len(took), + Median: round(at(0.5)), P90: round(at(0.9)), Max: round(took[len(took)-1])} + switch k.kind { + case inventory.DurationApply: + s.Suggests = "S2 bound max(2m, 3×last apply) ≈ " + round(max(2*time.Minute, 3*at(0.9))) + " at the p90" + case inventory.DurationHeartbeatGap: + s.Suggests = "S1 bound 3×interval ≈ " + round(3*at(0.5)) + } + out = append(out, s) + } + sort.Slice(out, func(i, j int) bool { + ki, kj := slices.Index(inventory.DurationKinds, out[i].Kind), slices.Index(inventory.DurationKinds, out[j].Kind) + if ki != kj { + return ki < kj + } + return out[i].Subject < out[j].Subject + }) + return out +} + +func round(d time.Duration) string { + switch { + case d < time.Second: + return d.Round(time.Millisecond).String() + case d < time.Minute: + return d.Round(100 * time.Millisecond).String() + default: + return d.Round(time.Second).String() + } +} + +// durationsCommand is `durations`: the summary per kind and subject, or every measurement as data. +func durationsCommand(ctx context.Context, args []string) error { + set := flag.NewFlagSet("durations", flag.ContinueOnError) + kind := set.String("kind", "", "one kind: "+strings.Join(inventory.DurationKinds, ", ")) + days := set.Int("days", 14, "how many days back") + asJSON := set.Bool("json", false, "the summary as data") + all := set.Bool("all", false, "every measurement rather than the summary") + if _, err := parseAround(set, args); err != nil { + return err + } + if *kind != "" && !slices.Contains(inventory.DurationKinds, *kind) { + return fmt.Errorf("%q is not a kind of duration: %s", *kind, strings.Join(inventory.DurationKinds, ", ")) + } + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + ds, err := open.inventory.Durations(ctx, *kind, time.Now().Add(-time.Duration(*days)*24*time.Hour)) + if err != nil { + return err + } + if *all { + body, err := json.MarshalIndent(ds, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + return nil + } + summary := summarise(ds) + if *asJSON { + body, err := json.MarshalIndent(map[string]any{"days": *days, "durations": summary}, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + return nil + } + if len(summary) == 0 { + fmt.Printf("nothing measured in the last %d day(s): the serving controller records apply, heartbeat-gap, "+ + "plan-tier and build durations as it hears them\n", *days) + return nil + } + fmt.Printf("durations over the last %d day(s) — what the core's bounds are set from (to-be 45 Phase 0)\n\n", *days) + fmt.Printf(" %-14s %-28s %6s %10s %10s %10s\n", "kind", "of", "count", "median", "p90", "max") + for _, s := range summary { + fmt.Printf(" %-14s %-28s %6d %10s %10s %10s\n", s.Kind, s.Subject, s.Count, s.Median, s.P90, s.Max) + if s.Suggests != "" { + fmt.Printf(" %-14s %-28s %s\n", "", "", s.Suggests) + } + } + return nil +} + +// forgettingOldDurations removes what is older than a month, at start and daily after. +func forgettingOldDurations(ctx context.Context, inv *inventory.Inventory) { + for { + if n, err := inv.ForgetOldDurations(ctx); err != nil { + fmt.Fprintf(os.Stderr, "durations older than %s could not be removed: %v\n", inventory.DurationsKeptFor, err) + } else if n > 0 { + fmt.Printf("removed %d duration(s) older than %s\n", n, inventory.DurationsKeptFor) + } + select { + case <-ctx.Done(): + return + case <-time.After(24 * time.Hour): + } + } +} diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go new file mode 100644 index 0000000..cf87327 --- /dev/null +++ b/cmd/mesh-controller/handacts.go @@ -0,0 +1,197 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "os" + "sort" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/link" +) + +// Acts done by hand, and why (novox/hq to-be 45 §7). +// +// **Which verbs ask why.** `plans close` and `plans stop`, `broker consumer-reset` and `hand-act +// record` refuse without it everywhere: nothing automated runs them, so a call without a reason is a +// person who has not given one. A named `push` asks for it through the mesh-controller seat, which is +// how a person or an agent acts by hand on the mesh; at a shell `--why` is recorded when given and not +// required, because the installer and the lab push by command line as a step of what they do, and a +// step of a procedure is not a repair. `conditions silence` joins them when the condition store does +// (Phase 1). + +// handActFlags are the flags every repairing verb takes. +type handActFlags struct { + why, cause, condition *string +} + +func addHandActFlags(set *flag.FlagSet) handActFlags { + return handActFlags{ + why: set.String("why", "", "why this is done by hand — recorded in the hand-act log (novox/hq to-be 45 §7)"), + cause: set.String("cause", "", "the cause, in a word or a condition's kind; the verb's own name when not given"), + condition: set.String("condition", "", "the key of the condition this act addresses, if any"), + } +} + +// given is whether a reason was given. +func (f handActFlags) given() bool { return strings.TrimSpace(*f.why) != "" } + +// require refuses an act without a reason, before anything is done. +func (f handActFlags) require(verb string) error { + if f.given() { + return nil + } + return fmt.Errorf("%s is a repair done by hand, and says why: --why (recorded in the hand-act "+ + "log, novox/hq to-be 45 §7). Nothing was done", verb) +} + +// handActConn is the serving controller's connection, for what it reads of the log itself; a +// command dials its own. +var handActConn *nats.Conn + +// onTheBus runs f with a connection to the bus: the serving controller's, or one of its own. +func onTheBus(f func(*nats.Conn) error) error { + if handActConn != nil { + return f(handActConn) + } + address, err := broker.BusAddress() + if err != nil { + return err + } + js, err := broker.Dial(address) + if err != nil { + return fmt.Errorf("cannot reach the bus: %w", err) + } + defer js.Close() + return f(js.Conn()) +} + +// record writes the entry for an act about to be done. **Before the act, and never instead of it**: +// a log that cannot be written is said loudly, and the repair it was about still happens — a mesh +// whose bus is down is exactly the mesh somebody is repairing by hand. +func (f handActFlags) record(ctx context.Context, verb string, args []string) { + if !f.given() { + return + } + act := link.HandAct{Verb: verb, Args: args, Why: strings.TrimSpace(*f.why), + Cause: strings.TrimSpace(*f.cause), Condition: strings.TrimSpace(*f.condition)} + err := onTheBus(func(conn *nats.Conn) error { + written, err := link.RecordHandAct(ctx, conn, act) + act = written + return err + }) + if err != nil { + fmt.Fprintf(os.Stderr, "this act by hand could NOT be recorded in the hand-act log, and is done anyway: %v\n", err) + return + } + fmt.Printf("recorded as %s in the hand-act log: %s, because %q (cause: %s)\n", act.ID, act.By, act.Why, act.Cause) +} + +// handActCommand is `hand-act record` and `hand-acts`. +func handActCommand(ctx context.Context, args []string) error { + if len(args) > 0 && args[0] == "record" { + set := flag.NewFlagSet("hand-act record", flag.ContinueOnError) + f := addHandActFlags(set) + positionals, err := parseAround(set, args[1:]) + if err != nil { + return err + } + what := strings.TrimSpace(strings.Join(positionals, " ")) + if what == "" { + return errors.New("hand-act record --why [--cause ] [--condition ]") + } + if err := f.require("hand-act record"); err != nil { + return err + } + if strings.TrimSpace(*f.cause) == "" { + return errors.New("hand-act record says the cause too: --cause , the word a second " + + "act for the same reason will use — it is how a repair done twice is found") + } + act := link.HandAct{Verb: "hand-act record", Args: []string{what}, Why: strings.TrimSpace(*f.why), + Cause: strings.TrimSpace(*f.cause), Condition: strings.TrimSpace(*f.condition)} + return onTheBus(func(conn *nats.Conn) error { + written, err := link.RecordHandAct(ctx, conn, act) + if err != nil { + return fmt.Errorf("the act could not be recorded: %w", err) + } + fmt.Printf("recorded as %s: %s did %q, because %q (cause: %s)\n", written.ID, written.By, what, + written.Why, written.Cause) + return nil + }) + } + if len(args) > 0 && args[0] != "list" && !strings.HasPrefix(args[0], "-") { + return errors.New("hand-act record --why --cause | hand-acts [--days N] [--json]") + } + if len(args) > 0 && args[0] == "list" { + args = args[1:] + } + set := flag.NewFlagSet("hand-acts", flag.ContinueOnError) + days := set.Int("days", 14, "how many days back") + asJSON := set.Bool("json", false, "as data") + if _, err := parseAround(set, args); err != nil { + return err + } + return onTheBus(func(conn *nats.Conn) error { + now := time.Now() + acts, err := link.HandActs(ctx, conn, now.Add(-time.Duration(*days)*24*time.Hour)) + if err != nil { + return err + } + repeated := link.RepeatedCauses(acts, now) + if *asJSON { + body, err := json.MarshalIndent(map[string]any{"acts": acts, "repeated": repeated}, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + return nil + } + if len(acts) == 0 { + fmt.Printf("nothing was done by hand in the last %d day(s)\n", *days) + return nil + } + for i := len(acts) - 1; i >= 0; i-- { + a := acts[i] + fmt.Printf("%s %s %s %s\n by %s — %s (cause: %s", a.At.Local().Format("2006-01-02 15:04"), a.ID, + a.Verb, strings.Join(a.Args, " "), a.By, a.Why, a.Cause) + if a.Condition != "" { + fmt.Printf(", condition %s", a.Condition) + } + fmt.Println(")") + } + if len(repeated) > 0 { + causes := make([]string, 0, len(repeated)) + for c, n := range repeated { + causes = append(causes, fmt.Sprintf("%s ×%d", c, n)) + } + sort.Strings(causes) + fmt.Printf("\ndone by hand more than once in a fortnight — a healer is wanted (to-be 45 S15): %s\n", + strings.Join(causes, ", ")) + } + return nil + }) +} + +// handActsThisWeek is how many acts were done by hand in the last seven days, for `status`; -1 when +// the log could not be read, which status says rather than reading as none. +func handActsThisWeek(ctx context.Context) (int, string) { + n := -1 + err := onTheBus(func(conn *nats.Conn) error { + reading, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + acts, err := link.HandActs(reading, conn, time.Now().Add(-7*24*time.Hour)) + n = len(acts) + return err + }) + if err != nil { + return -1, err.Error() + } + return n, "" +} diff --git a/cmd/mesh-controller/handacts_test.go b/cmd/mesh-controller/handacts_test.go new file mode 100644 index 0000000..37756d7 --- /dev/null +++ b/cmd/mesh-controller/handacts_test.go @@ -0,0 +1,94 @@ +package main + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// **Every verb that repairs by hand takes a required why** (novox/hq to-be 45 §7): refused before +// anything is done, through the seat, through `command`, and at a shell where nothing automated runs +// the verb. +func TestARepairByHandWithoutAReasonIsRefused(t *testing.T) { + for _, c := range []struct { + verb string + args map[string]any + }{ + {"push", map[string]any{"node": "anchor"}}, + {"push", map[string]any{}}, + {"plans", map[string]any{"close": "plan-1"}}, + {"plans", map[string]any{"stop": "plan-1"}}, + {"hand-act", map[string]any{"what": "restarted the proxy", "cause": "proxy-stuck"}}, + {"command", map[string]any{"command": "push anchor"}}, + {"command", map[string]any{"command": "plans close plan-1"}}, + {"command", map[string]any{"command": "broker consumer-reset EVENTS controller"}}, + {"command", map[string]any{"command": "hand-act record restarted --cause x"}}, + } { + argv, err := argvFor(c.verb, c.args) + if c.verb == "plans" && err == nil { + // The seat composes the command line; the command refuses it, before opening anything. + err = plansCommand(context.Background(), argv[1:]) + } + if err == nil || !strings.Contains(err.Error(), "why") { + t.Errorf("%s %v was not refused for want of why: %v %v", c.verb, c.args, argv, err) + } + } + for _, args := range [][]string{{"EVENTS", "controller"}} { + if err := consumerReset(context.Background(), args); err == nil || !strings.Contains(err.Error(), "--why") { + t.Errorf("consumer-reset without why: %v", err) + } + } + if err := handActCommand(context.Background(), []string{"record", "restarted the proxy", "--cause", "x"}); err == nil || + !strings.Contains(err.Error(), "--why") { + t.Errorf("hand-act record without why: %v", err) + } + if err := handActCommand(context.Background(), []string{"record", "restarted the proxy", "--why", "it hung"}); err == nil || + !strings.Contains(err.Error(), "--cause") { + t.Errorf("hand-act record without a cause: %v", err) + } +} + +// With a reason, the seat passes it to the command, and a verb that only reads is not held to one. +func TestARepairByHandCarriesItsReason(t *testing.T) { + for _, c := range []struct { + verb string + args map[string]any + want string + }{ + {"push", map[string]any{"node": "anchor", "why": "stuck", "cause": "sent-not-reported"}, + "push anchor --wait 0 --why stuck --cause sent-not-reported"}, + {"plans", map[string]any{"close": "plan-1", "why": "the report will not come"}, + "plans close plan-1 --why the report will not come"}, + {"plans", map[string]any{"retry": "plan-1"}, "plans retry plan-1"}, + {"hand-act", map[string]any{"what": "restarted", "why": "hung", "cause": "proxy", "condition": "machine.a.silent"}, + "hand-act record restarted --why hung --cause proxy --condition machine.a.silent"}, + {"command", map[string]any{"command": "push anchor --why stuck"}, "push anchor --why stuck"}, + {"command", map[string]any{"command": "plans plan-1"}, "plans plan-1"}, + } { + argv, err := argvFor(c.verb, c.args) + if err != nil || strings.Join(argv, " ") != c.want { + t.Errorf("%s %v: %v %v, want %q", c.verb, c.args, argv, err, c.want) + } + } +} + +// The summary of durations says, per kind and subject, what a bound would be set from. +func TestDurationsAreSummarisedPerSubject(t *testing.T) { + var ds []inventory.Duration + for i := 1; i <= 10; i++ { + ds = append(ds, inventory.Duration{Kind: inventory.DurationApply, Subject: "anchor", + Took: time.Duration(i) * time.Second}) + } + ds = append(ds, inventory.Duration{Kind: inventory.DurationHeartbeatGap, Subject: "anchor", Took: time.Minute}) + got := summarise(ds) + if len(got) != 2 || got[0].Kind != inventory.DurationApply || got[0].Count != 10 || + got[0].Max != "10s" || got[0].Median != "5s" || got[0].P90 != "9s" { + t.Fatalf("%+v", got) + } + if !strings.Contains(got[1].Suggests, "3m0s") { + t.Fatalf("a minute between words suggests %q", got[1].Suggests) + } +} diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 16a49f1..7880440 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -145,6 +145,14 @@ func run() error { return seatCommand(ctx, args[1:]) case "status": return statusCommand(ctx, args[1:]) + // Acts done by hand, and why (novox/hq to-be 45 §7). + case "hand-act": + return handActCommand(ctx, args[1:]) + case "hand-acts": + return handActCommand(ctx, append([]string{"list"}, args[1:]...)) + // What the mesh's bounds will be set from (novox/hq to-be 45 Phase 0). + case "durations": + return durationsCommand(ctx, args[1:]) case "version": fmt.Println(version) return nil @@ -226,6 +234,12 @@ func usage() { kill end a build where it runs; recorded failed, killed by hand pause [] / resume [] the build seat's holder there, or every holder, takes nothing new / again plans retry ask a failed plan's failed builds again, and carry the plan on + plans stop|close --why end a plan by hand; recorded in the hand-act log + hand-act record --why --cause [--condition ] + record an act done by hand outside the mesh (to-be 45 §7) + hand-acts [--days N] [--json] what was done by hand lately, why, and which causes repeat + durations [--kind K] [--days N] [--json] + apply, heartbeat, plan-tier and build durations, per machine or module collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go builder issue a broker account for a build machine, scoped to build work, delivered as the builder module's broker secret (module add it first) @@ -238,7 +252,8 @@ func usage() { which provider this one gets a provision from: the module, and its node unpin put that question back plan [--files|--json] what that node would run, and why - push [] [--behind] send a node everything it should be, or only those that need it + push [] [--behind] [--why ] send a node everything it should be, or only those + that need it; --why records it in the hand-act log version what this binary is Each context reaches its own store through its own credential (novox/hq ADR 0008), named @@ -291,6 +306,9 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { // When it was asked, so a plan takes as its outcome only a build asked for it or after it // (novox/hq 04-ISSUES/219). Zero when the id does not say. asked, _ := link.BuildAskedAt(result.ID) + // And how long it took, asked to heard, which a build's bound will be set from (novox/hq to-be 45 + // Phase 0). Said if lost; never a reason not to take the build in. + recordBuildDuration(ctx, b.inv, result, asked) switch { case err != nil && result.Failed != "": fmt.Printf("%s: %v\n", result.ID, err) @@ -317,5 +335,7 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit)) saysWhenThePolicyActs(ctx, b.inv, manifest.Module) planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked, result.ID) + // A module registered may be one a machine is now behind: `status` is composed again. + statusFrom.nudge() return nil } diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index acfa154..eb81031 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -387,7 +387,7 @@ func brokerCommand(ctx context.Context, args []string) error { return busAccounts(ctx, args[1:]) } if len(args) > 0 && args[0] == "consumer-reset" { - return consumerReset(args[1:]) + return consumerReset(ctx, args[1:]) } if len(args) == 0 || args[0] != "show" { return errors.New("broker show | broker certificate [--check] --into | broker accounts --into | " + @@ -414,10 +414,21 @@ func brokerCommand(ctx context.Context, args []string) error { // consumerReset re-makes one consumer on a stream that keeps history to start from now (novox/hq issue // 248): the way out of a consumer replaying a week of announcements, said rather than done by hand. A // person's act — what was pending is dropped — so it is a command, and nothing calls it on its own. -func consumerReset(args []string) error { - if len(args) != 2 { - return errors.New("broker consumer-reset , e.g. broker consumer-reset EVENTS controller") +func consumerReset(ctx context.Context, args []string) error { + set := flag.NewFlagSet("broker consumer-reset", flag.ContinueOnError) + // A repair by hand, which says why (novox/hq to-be 45 §7). + why := addHandActFlags(set) + args, err := parseAround(set, args) + if err != nil { + return err } + if len(args) != 2 { + return errors.New("broker consumer-reset --why , e.g. broker consumer-reset EVENTS controller --why ...") + } + if err := why.require("broker consumer-reset"); err != nil { + return err + } + why.record(ctx, "broker consumer-reset", args) address, err := broker.BusAddress() if err != nil { return err diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index b7ef31f..eb5e56b 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -105,7 +105,10 @@ func serve(ctx context.Context) error { work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known, OnNATS: true} - server, err := connectLink(ctx, inv, work, work) + // `status` from a summary kept current here (novox/hq to-be 45 Phase 0): a machine saying + // something new is one thing that moves it, so the listener nudges it. + statusFrom = newStatusSummary(composeStatus(open)) + server, err := connectLink(ctx, inv, work, nudgingListener{work, statusFrom}) if err != nil { return err } @@ -121,6 +124,8 @@ func serve(ctx context.Context) error { // Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for // machines to report moves when they have, and a plan left by a replaced controller resumes. go planTicker(ctx, open) + // The durations the core's bounds are set from are kept a month (novox/hq to-be 45 Phase 0). + go forgettingOldDurations(ctx, inv) // And what the catalogue decided a build meant. The builder's own result is already handled // above; this is the other half — the control plane is the only one of the three that knows // which machines run the thing, so it is the one that acts (novox/hq ADR 0072). @@ -158,8 +163,21 @@ func serve(ctx context.Context) error { if !isNATS { return errors.New("the mesh's verbs are served over the bus, and this control plane is not on it") } + // The hand-act log is counted for `status` on this connection rather than a new one a minute. + handActConn = bus.Conn + // Composed now and kept current, before the verb that answers from it is served. + go statusFrom.keep(ctx) // A call that outlasts its caller's patience is followed by `calls` (novox/hq issue 265). link.Calls.Follow = catalogue.ControllerSeatName + ".calls" + // And every call is kept on the bus, so a restart of this process keeps what came of each + // (novox/hq to-be 45 §6). A bus without the bucket is said and served from memory, as before: + // answering no calls at all would be worse than answering them without the record. + said := log.New(os.Stdout, "", log.LstdFlags) + if keeper, err := link.CallsOnTheBus(ctx, bus.Conn); err != nil { + fmt.Printf("calls are kept in memory only, and lost when this controller stops: %v\n", err) + } else if err := link.Calls.Durably(ctx, keeper, controllerProcess(), said); err != nil { + fmt.Printf("calls are kept on the bus from now on; the ones kept before could not be read: %v\n", err) + } stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags)) if err != nil { return err @@ -260,6 +278,9 @@ func pushCommand(ctx context.Context, args []string) error { // (novox/hq ADR 0010). 0 waits for nothing, which is the old fire-and-forget. wait := set.Duration("wait", 0, "for a named node, how long to wait for it to report applying what it was sent (0: do not wait)") + // A push by hand is a repair, and says why (novox/hq to-be 45 §7): required through the seat, + // recorded when given at a shell — see handacts.go for why a shell is not refused. + why := addHandActFlags(set) positionals, err := parseAround(set, args) if err != nil { return err @@ -274,6 +295,11 @@ func pushCommand(ctx context.Context, args []string) error { return errors.New("push or push --behind, not both: one names a machine and the " + "other asks which machines need one") } + recorded := append([]string(nil), args...) + if *behind { + recorded = append(recorded, "--behind") + } + why.record(ctx, "push", recorded) open, err := openStores(ctx) if err != nil { return err @@ -1039,6 +1065,20 @@ func digestOf(body []byte) string { // be worked out" is a different problem with a different remedy, and `plan` is where it is said. func wouldSend(ctx context.Context, open *stores, nodes []inventory.Node) (map[string]string, error) { + return wouldSendFrom(ctx, open, nodes, nil) +} + +// planned is one machine's plan as planFor answered it, for a caller that already asked. +type planned struct { + plan catalogue.Resolution + settings catalogue.SettingsBy +} + +// wouldSendFrom is wouldSend reusing the plans a caller worked out a moment before: resolving a +// machine is most of what `status` costs, and it used to resolve every machine twice (novox/hq +// to-be 45 Phase 0). A machine absent from plans is worked out here. +func wouldSendFrom(ctx context.Context, open *stores, + nodes []inventory.Node, plans map[string]planned) (map[string]string, error) { gens, err := generators(ctx, open) if err != nil { @@ -1046,9 +1086,12 @@ func wouldSend(ctx context.Context, open *stores, } out := map[string]string{} for _, n := range nodes { - plan, settings, err := planFor(ctx, open, n.Name) - if err != nil { - continue + known, have := plans[n.Name] + plan, settings := known.plan, known.settings + if !have { + if plan, settings, err = planFor(ctx, open, n.Name); err != nil { + continue + } } declared, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading) if err != nil { @@ -1136,6 +1179,11 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) if err := broker.RaiseCancelledSets(js, inventory.MeshSeats()); err != nil { return err } + // And the controller's own buckets (novox/hq to-be 45 §1): the calls it serves and the acts done + // by hand, kept where a restart of this process does not take them. + if err := js.EnsureControllerBuckets(); err != nil { + return err + } // Every module's state (novox/hq ADR 0201), from the catalogue: a bucket exists from // registration, so a module reading one may watch it before its owner runs anywhere. One that // nothing declares any more is said and kept — what it holds is data. @@ -1271,3 +1319,11 @@ func reportUnheldPushed(w io.Writer, named bool, asked []string, unheld map[stri fmt.Fprintf(w, "%s: %d unmet seat dependenc(ies) — see `status`\n", node, len(lines)) } } + +// controllerProcess names this serving process among controllers: the machine, the process and when +// it started — what a call kept on the bus carries, so the next controller can tell a call this one +// left running from one it is running itself (novox/hq to-be 45 §6). +func controllerProcess() string { + host, _ := os.Hostname() + return fmt.Sprintf("controller@%s pid %d since %s", host, os.Getpid(), time.Now().UTC().Format(time.RFC3339)) +} diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index 6db9d54..0115c6c 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -78,6 +78,11 @@ type meshStatus struct { // that machine holds, with the modules that could hold it (novox/hq ADR 0207). Absent when every // dependency is met. Reported, not refused, until the switch. Unheld []catalogue.Unheld `json:"unheld,omitempty"` + // HandActsThisWeek is how many acts were done by hand in the last seven days (novox/hq to-be 45 + // §7): every one is a repair a healer could have made. Absent where the log is not on hand; + // HandActsUnread says why when it could not be read, rather than reading as none. + HandActsThisWeek *int `json:"handActsThisWeek,omitempty"` + HandActsUnread string `json:"handActsUnread,omitempty"` // Failing is every consumer a provider says it keeps failing, with the class of error, since // when, and when it was last said (novox/hq ADR 0224). Absent when no provider says so. A // document without this called the mesh well while the identity provider refused every consumer @@ -218,6 +223,7 @@ func statusAsJSON(asked answers) ([]byte, error) { } } out.Unheld = asked.unheld + out.HandActsThisWeek, out.HandActsUnread = asked.handActs, asked.handActsUnread out.Failing = asked.failing out.Overflowing = asked.overflowing for name := range asked.refused { diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 0f527fa..bc79cf9 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -986,10 +986,18 @@ func plansCommand(ctx context.Context, args []string) error { whatIf := set.String("what-if", "", "owner/repository: the plan a merge there would produce, saving nothing — with --paths or --modules") paths := set.String("paths", "", "the files the merge would change, comma-separated, from the repository's root") modules := set.String("modules", "", "or the modules it would change, comma-separated") + // Ending a plan by hand is a repair, and says why (novox/hq to-be 45 §7). + why := addHandActFlags(set) positionals, err := parseAround(set, args) if err != nil { return err } + if len(positionals) == 2 && (positionals[0] == "stop" || positionals[0] == "close") { + // Refused before anything is opened: a repair by hand says why. + if err := why.require("plans " + positionals[0]); err != nil { + return err + } + } open, err := openStores(ctx) if err != nil { return err @@ -1061,8 +1069,9 @@ func plansCommand(ctx context.Context, args []string) error { if !p.Open() { return fmt.Errorf("%s is already %s", p.ID, p.State) } + why.record(ctx, "plans "+positionals[0], positionals[1:]) p.State = inventory.PlanFailed - p.Note = how + " by hand at tier " + fmt.Sprint(p.Tier) + p.Note = how + " by hand at tier " + fmt.Sprint(p.Tier) + ": " + strings.TrimSpace(*why.why) sayUnsent(&p, func(m string) bool { u, err := inv.UpgradeOf(ctx, m) return err == nil && u.RollOut diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 5c706b3..0e0922b 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -9,9 +9,11 @@ import ( "github.com/nats-io/nats.go/micro" "os" "os/exec" + "slices" "sort" "strings" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/link" ) @@ -239,6 +241,12 @@ func (a *verbArguments) commandLine() ([]string, error) { if len(argv) == 0 { return nil, errors.New("command names no command") } + // The generic verb is no way round the hand-act log (novox/hq to-be 45 §7): a repair through + // it says why, as it would through its own verb. + if repair := repairingCommand(argv); repair != "" && !slices.ContainsFunc(argv, isWhyFlag) { + return nil, fmt.Errorf("%s is a repair done by hand, and says why: add --why to the command "+ + "line (recorded in the hand-act log). Nothing was done", repair) + } return argv, nil case "tools": return nil, errors.New("tools is answered from the records, not by a command") @@ -280,7 +288,18 @@ func (a *verbArguments) commandLine() ([]string, error) { } for _, act := range []string{"stop", "close", "retry"} { if id := str(act); id != "" { - return []string{"plans", act, id}, nil + argv := []string{"plans", act, id} + if act == "retry" { + return argv, nil + } + // Ending a plan by hand says why (novox/hq to-be 45 §7); the command refuses it without. + if w := str("why"); w != "" { + argv = append(argv, "--why", w) + } + if c := str("cause"); c != "" { + argv = append(argv, "--cause", c) + } + return argv, nil } } if id := str("id"); id != "" { @@ -355,16 +374,48 @@ func (a *verbArguments) commandLine() ([]string, error) { // Sent and not waited for: the asker reads `status` for what the machine did, which is // what a person at a shell does too. A tool call that blocked for a push's whole apply would // time out on every machine that takes a minute, and say nothing about the ones that did not. + // A push through the seat is a push by hand, and says why (novox/hq to-be 45 §7). + if err := need("why"); err != nil { + return nil, fmt.Errorf("%w: a push by hand is a repair, recorded in the hand-act log with why", err) + } + why := []string{"--why", str("why")} + if c := str("cause"); c != "" { + why = append(why, "--cause", c) + } if n := str("node"); n != "" { // behind is not read here: given with a machine, it is refused as passed over — naming // a machine and asking for every machine behind are two requests, and guessing one // would push a machine nobody named, or not push one somebody did. - return []string{"push", n, "--wait", "0"}, nil + return append([]string{"push", n, "--wait", "0"}, why...), nil } // No machine: the whole mesh, whether or not behind said so. The command's answer says it // first, so a caller who meant one machine reads that it was not one. on("behind") - return []string{"push", "--behind", "--wait", "0"}, nil + return append([]string{"push", "--behind", "--wait", "0"}, why...), nil + case "hand-act": + if err := need("what", "why", "cause"); err != nil { + return nil, err + } + argv := []string{"hand-act", "record", str("what"), "--why", str("why"), "--cause", str("cause")} + if c := str("condition"); c != "" { + argv = append(argv, "--condition", c) + } + return argv, nil + case "hand-acts": + argv := []string{"hand-acts", "--json"} + if d := str("days"); d != "" { + argv = append(argv, "--days", d) + } + return argv, nil + case "durations": + argv := []string{"durations", "--json"} + if k := str("kind"); k != "" { + argv = append(argv, "--kind", k) + } + if d := str("days"); d != "" { + argv = append(argv, "--days", d) + } + return argv, nil case "rotate": if p := str("provision"); p != "" { argv := []string{"rotate", p} @@ -442,7 +493,28 @@ 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} +var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "collection": true, + "hand-acts": true, "durations": 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. +func repairingCommand(argv []string) string { + switch { + case argv[0] == "push": + return "push" + case argv[0] == "plans" && len(argv) > 1 && (argv[1] == "stop" || argv[1] == "close"): + return "plans " + argv[1] + case argv[0] == "broker" && len(argv) > 1 && argv[1] == "consumer-reset": + return "broker consumer-reset" + case argv[0] == "hand-act": + return "hand-act record" + } + return "" +} + +func isWhyFlag(word string) bool { + return word == "--why" || word == "-why" || strings.HasPrefix(word, "--why=") || strings.HasPrefix(word, "-why=") +} // runVerb runs this binary with the given command line and gathers what it said. func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { @@ -454,6 +526,12 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { // The same environment: the stores' credentials, the bus, the broker — everything a command run // from a shell in this container would have, because it is that. cmd.Env = os.Environ() + // And who asked, so an act it does by hand is recorded as theirs (novox/hq to-be 45 §7). + caller := link.CallerIn(ctx) + if caller == "" { + caller = "a seat call whose caller the bus did not name" + } + cmd.Env = append(cmd.Env, link.CallerVar+"="+caller+", through the "+catalogue.ControllerSeatName+" seat") // Two buffers, one answer. What the command *says* is both streams, in the order a person at // a shell would read them; what it *answers as data* is standard output alone — `status --json` // prints its warnings beside the document, and a JSON parsed from the two together parsed @@ -544,6 +622,14 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { if err != nil { return nil, err } + if verb == "status" && statusFrom != nil { + // At once, from the summary the serving controller keeps (novox/hq to-be 45 Phase 0). + return statusFrom.answer(ctx) + } + if !readingVerbs[verb] && !(verb == "plans" && !actsOnAPlan(args)) { + // Whatever it did, `status` is composed again once it has. + defer statusFrom.nudge() + } if answersFirst(argv) { // Before anything is sent: a push sends the bus's own machine first, and a broker // reloading its user list forgets the answer it was about to permit (novox/hq issue 265). @@ -555,6 +641,16 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { return handlers, behind, nil } +// actsOnAPlan is `plans` asked to stop, close or retry one rather than to show them. +func actsOnAPlan(args map[string]any) bool { + for _, act := range []string{"stop", "close", "retry"} { + if v, _ := args[act].(string); strings.TrimSpace(v) != "" { + return true + } + } + return false +} + // inProcess are the verbs answered by this process rather than by a command it runs: `tools` from // the records, `calls` from what this process served. var inProcess = map[string]bool{"tools": true, "calls": true} @@ -568,22 +664,41 @@ func answersFirst(argv []string) bool { } // callsAnswer is what `calls` answers: the kept calls, newest first, without their answers — or -// one call whole. +// one call whole. Kept on the bus, so a call a controller before this one served is answered too +// (novox/hq to-be 45 §6); where the bus cannot be read, what this process served is answered and +// the reason said beside it. func callsAnswer(log *link.CallLog, id string) (any, error) { if id != "" { - c, ok := log.Get(id) + c, ok, err := log.Get(id) + if err != nil { + return nil, fmt.Errorf("call %s is not in this controller's memory, and the calls kept on the "+ + "bus could not be read: %w", id, err) + } if !ok { + if log.IsDurable() { + return nil, fmt.Errorf("no call %s is kept: the bus keeps the last %d calls, or %s, and this "+ + "is not among them — `calls` lists them", id, broker.KeptCallsDurably, broker.CallsKeptFor) + } return nil, fmt.Errorf("no call %s is kept here: calls are kept by the controller that "+ "answered them, the last %d, and not across a restart — `calls` lists them", id, link.KeptCalls) } return c, nil } - recent := log.Recent() + recent, err := log.Recent() for i := range recent { recent[i].Answer = nil } - return map[string]any{"calls": recent, "kept": link.KeptCalls, - "note": "newest first; `calls` with a call's id gives its whole answer"}, nil + answer := map[string]any{"calls": recent, "note": "newest first; `calls` with a call's id gives its whole answer"} + if log.IsDurable() { + answer["kept"] = fmt.Sprintf("the last %d calls, or %s, on the bus — across a restart of the controller", + broker.KeptCallsDurably, broker.CallsKeptFor) + } else { + answer["kept"] = fmt.Sprintf("the last %d calls this controller served, in its memory only", link.KeptCalls) + } + if err != nil { + answer["unread"] = err.Error() + } + return answer, nil } // seatTools is what `tools` answers: every seat with a protocol, and the tools each serves, from the diff --git a/cmd/mesh-controller/seatverbs_schema_test.go b/cmd/mesh-controller/seatverbs_schema_test.go index d9a5d21..ea666a3 100644 --- a/cmd/mesh-controller/seatverbs_schema_test.go +++ b/cmd/mesh-controller/seatverbs_schema_test.go @@ -145,17 +145,17 @@ func TestAnArgumentAVerbDoesNotDeclareIsRefused(t *testing.T) { // The push that was the cause: a machine named is that machine; none named is the whole mesh, and // naming one beside behind is refused rather than one of the two guessed. func TestAPushIsOneMachineOrSaysItIsTheWholeMesh(t *testing.T) { - argv, err := argvFor("push", map[string]any{"node": "g1"}) - if err != nil || strings.Join(argv, " ") != "push g1 --wait 0" { + argv, err := argvFor("push", map[string]any{"node": "g1", "why": "w"}) + if err != nil || strings.Join(argv, " ") != "push g1 --wait 0 --why w" { t.Fatalf("a named push: %v %v", argv, err) } - for _, args := range []map[string]any{{}, {"behind": "true"}} { + for _, args := range []map[string]any{{"why": "w"}, {"behind": "true", "why": "w"}} { argv, err := argvFor("push", args) - if err != nil || strings.Join(argv, " ") != "push --behind --wait 0" { + if err != nil || strings.Join(argv, " ") != "push --behind --wait 0 --why w" { t.Fatalf("a push of the whole mesh %v: %v %v", args, argv, err) } } - if _, err := argvFor("push", map[string]any{"node": "g1", "behind": "true"}); err == nil || + if _, err := argvFor("push", map[string]any{"node": "g1", "behind": "true", "why": "w"}); err == nil || !strings.Contains(err.Error(), `"behind"`) { t.Fatalf("a named push with behind was taken: %v", err) } @@ -268,6 +268,11 @@ var accountedFlags = map[string]map[string]string{ }, "builds": {"n": "=limit"}, "plans": {"n": "=limit", "what-if": "=repository"}, + "durations": { + "json": "set by the verb: the answer is data", + "all": "withheld: every measurement of a fortnight is more than a call should carry; `command` reaches it", + }, + "hand-acts": {"json": "set by the verb: the answer is data"}, } // **Every flag of the command a verb runs is in the verb's schema, or accounted for here.** Derived diff --git a/cmd/mesh-controller/seatverbs_test.go b/cmd/mesh-controller/seatverbs_test.go index fa7f899..4709f5b 100644 --- a/cmd/mesh-controller/seatverbs_test.go +++ b/cmd/mesh-controller/seatverbs_test.go @@ -107,8 +107,8 @@ func TestAVerbMissingWhatItNeedsIsRefused(t *testing.T) { // A push and a build are sent, not waited for: the asker reads status, or the build's log by its // id, for what happened. A repository given as a forge path is said to be one (issue 176). func TestActsDoNotBlockTheCall(t *testing.T) { - argv, _ := argvFor("push", map[string]any{"node": "one"}) - if strings.Join(argv, " ") != "push one --wait 0" { + argv, _ := argvFor("push", map[string]any{"node": "one", "why": "w"}) + if strings.Join(argv, " ") != "push one --wait 0 --why w" { t.Fatalf("push waits: %v", argv) } argv, _ = argvFor("build", map[string]any{"repository": "novox/x", "path": "modules/x"}) diff --git a/cmd/mesh-controller/status.go b/cmd/mesh-controller/status.go index 7291043..4d0afe5 100644 --- a/cmd/mesh-controller/status.go +++ b/cmd/mesh-controller/status.go @@ -8,6 +8,7 @@ import ( "strings" "time" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/overlay" ) @@ -315,6 +316,16 @@ func printStatus(asked answers) error { fmt.Printf("\n a shorter `slug` in the module's definition fits it; `module check` refuses one before merge\n\n") } + // Repairs done by hand this week (novox/hq to-be 45 §7). Not a fault, so it does not break "all + // well"; each is a healer the mesh does not have yet, and the count is how that is watched. + switch { + case asked.handActsUnread != "": + fmt.Printf("the hand-act log could not be read, so how much was done by hand this week is not known: %s\n\n", + asked.handActsUnread) + case asked.handActs != nil && *asked.handActs > 0: + fmt.Printf("%d act(s) done by hand in the last seven days — `hand-acts` lists them, and why\n\n", *asked.handActs) + } + if adopted := adoptedNodes(nodes); len(adopted) > 0 { // Said, because nothing forces the flip: a node left adopted is visible here rather than // read as converged (novox/hq ADR 0100). Not a fault, so it does not break "all well". @@ -423,14 +434,16 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { // ADR 0207). Each machine resolved again rather than threaded through whoResolves, whose answer // the private network is built from and should say nothing else; a machine that does not // resolve is already in refused, and is passed over here. + plans := map[string]planned{} for _, n := range out.nodes { - plan, _, err := planFor(ctx, open, n.Name) + plan, settings, err := planFor(ctx, open, n.Name) if err != nil { if unresolvable(err) { continue } return answers{}, err } + plans[n.Name] = planned{plan, settings} out.unheld = append(out.unheld, plan.Unheld...) // And which of its modules a provider leaves out of its grants, for an identity too long // for what the provision keeps (novox/hq ADR 0225) — judged from the consumer's own @@ -449,6 +462,16 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { } // Whether the build seat takes work, for a plan waiting on it (novox/hq ADR 0219). out.paused = buildSeatPause(ctx, inv, out.plans) + // And how many repairs were done by hand this week (novox/hq to-be 45 §7) — where there is a bus + // to read the log from; a process with none has no log to count. + if _, onBus := broker.BusAddress(); onBus == nil { + n, unread := handActsThisWeek(ctx) + if unread != "" { + out.handActsUnread = unread + } else { + out.handActs = &n + } + } // And which machines are not running what the mesh would send them. The same question as a // module being behind its source, one level down: that one says the catalogue is out of date, @@ -460,7 +483,7 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { // said nothing at all and `status --json` emitted prose to stderr and no JSON anywhere. The // reason is kept and reported as data; every question that does not depend on it is still // answered. - would, err := wouldSend(ctx, open, out.nodes) + would, err := wouldSendFrom(ctx, open, out.nodes, plans) if err != nil { out.network = err.Error() would = map[string]string{} diff --git a/cmd/mesh-controller/status_summary.go b/cmd/mesh-controller/status_summary.go new file mode 100644 index 0000000..b814cf5 --- /dev/null +++ b/cmd/mesh-controller/status_summary.go @@ -0,0 +1,188 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "sync" + "time" + + "github.com/novox/mesh-controller/internal/link" +) + +// `status` answered from a summary the serving controller keeps current (novox/hq to-be 45 Phase 0, +// §4's D9 and §8's health of the controller). +// +// **Asked, it was composed: every machine resolved, twice, while its caller waited.** On 2026-10-06 +// the verb took eighteen seconds on a mesh of four machines, so its caller read "still running" and +// had to ask `calls` for the answer to "is the mesh alright" — the one question that must answer at +// once, and the one a self-check and a rollout gate will ask every few minutes. So the serving +// controller composes it in the background — at its start, after anything that changes what it says +// (a machine's report, a build, a verb that acts), and every minute regardless — and the verb answers +// the last composition at once, saying when it was composed and how long that took. A caller who +// needs it newer than that reads the time and asks again; nothing is answered as current that is not. + +// statusEvery is how often the summary is composed with nothing having nudged it; statusSettle how +// long a nudge waits for the next, so a push answered by four machines is composed once. +var ( + statusEvery = time.Minute + statusSettle = 2 * time.Second + // statusComposeWithin bounds one composition, so a store that hangs cannot stop the summary for + // good; the attempt is said as failed, and the last summary stands with its age. + statusComposeWithin = 2 * time.Minute +) + +// statusSummary is the last composed `status --json`, and when. +type statusSummary struct { + compose func(context.Context) ([]byte, error) + // every, settle and within are the clocks above, read once when it is made. + every, settle, within time.Duration + + mu sync.Mutex + body []byte + composedAt time.Time + took time.Duration + failed string + failedAt time.Time + started time.Time + first chan struct{} // closed when the first attempt ends, either way + + nudged chan struct{} +} + +func newStatusSummary(compose func(context.Context) ([]byte, error)) *statusSummary { + return &statusSummary{compose: compose, started: time.Now(), first: make(chan struct{}), + nudged: make(chan struct{}, 1), every: statusEvery, settle: statusSettle, within: statusComposeWithin} +} + +// statusFrom is the serving controller's summary; nil in any other process, where `status` is +// composed when asked, as at a shell. +var statusFrom *statusSummary + +// nudge asks for a composition soon. Never blocks: one pending is as good as many. +func (s *statusSummary) nudge() { + if s == nil { + return + } + select { + case s.nudged <- struct{}{}: + default: + } +} + +// keep composes until ctx ends: now, on a nudge once things settle, and every statusEvery. +func (s *statusSummary) keep(ctx context.Context) { + once := sync.Once{} + for { + s.composeOnce(ctx) + once.Do(func() { close(s.first) }) + timer := time.NewTimer(s.every) + select { + case <-ctx.Done(): + timer.Stop() + return + case <-timer.C: + case <-s.nudged: + timer.Stop() + // Let what else is arriving arrive, then compose once for all of it. + select { + case <-ctx.Done(): + return + case <-time.After(s.settle): + } + select { + case <-s.nudged: + default: + } + } + } +} + +func (s *statusSummary) composeOnce(ctx context.Context) { + start := time.Now() + asking, cancel := context.WithTimeout(ctx, s.within) + body, err := s.compose(asking) + cancel() + took := time.Since(start) + s.mu.Lock() + defer s.mu.Unlock() + if err != nil { + s.failed, s.failedAt = err.Error(), time.Now() + fmt.Printf("status could not be composed (after %s): %v — `status` answers the last summary, "+ + "with its age\n", took.Round(time.Millisecond), err) + return + } + s.body, s.composedAt, s.took, s.failed = body, start, took, "" +} + +// answer is what the `status` verb answers: the last summary at once, the same document `status +// --json` prints, with when it was composed. Before the first composition has ended it waits for it, +// but never past the caller's window; a controller that has none says so and why, rather than +// answering an empty mesh as a well one. +func (s *statusSummary) answer(ctx context.Context) (any, error) { + wait := time.NewTimer(link.AnswerWithin - time.Second) + defer wait.Stop() + select { + case <-s.first: + case <-wait.C: + case <-ctx.Done(): + } + s.mu.Lock() + defer s.mu.Unlock() + if s.body == nil { + why := "its first composition has not finished" + if s.failed != "" { + why = "it could not be composed: " + s.failed + } + return nil, fmt.Errorf("this controller started %s ago and has no status to answer yet — %s. "+ + "Ask again shortly", time.Since(s.started).Round(time.Second), why) + } + var parsed any + _ = json.Unmarshal(s.body, &parsed) + out := map[string]any{ + "output": string(s.body), "ok": true, "answer": parsed, + "composed": s.composedAt.UTC().Format(time.RFC3339), + "age": time.Since(s.composedAt).Round(time.Second).String(), + "composedIn": s.took.Round(time.Millisecond).String(), + "note": "composed by the serving controller at its start, after each report, build or act, and " + + "every minute; answered at once from the last composition", + } + if s.failed != "" && s.failedAt.After(s.composedAt) { + out["lastAttemptFailed"] = fmt.Sprintf("%s: %s — this summary is the last that could be composed", + s.failedAt.UTC().Format(time.RFC3339), s.failed) + } + return out, nil +} + +// composeStatus is `status --json`, composed in this process against its stores. +func composeStatus(open *stores) func(context.Context) ([]byte, error) { + return func(ctx context.Context) ([]byte, error) { + asked, err := theThreeQuestions(ctx, open) + if err != nil { + return nil, err + } + return statusAsJSON(asked) + } +} + +// readingVerbs are the verbs that only read; after any other, what `status` says may have changed, so the summary is +// composed again; a verb that only reads leaves it alone, or a console polling `nodes` would keep the +// controller composing for ever. +var readingVerbs = map[string]bool{ + "tools": true, "calls": true, "status": true, "nodes": true, "node": true, "modules": true, + "seats": true, "builds": true, "plan": true, "queue": true, "durations": true, "hand-acts": true, +} + +// nudgingListener is the enrolment, nudging the summary when a machine said something new. +type nudgingListener struct { + link.Enrolment + summary *statusSummary +} + +func (l nudgingListener) Heard(ctx context.Context, report link.Report) (bool, error) { + news, err := l.Enrolment.Heard(ctx, report) + if news { + l.summary.nudge() + } + return news, err +} diff --git a/cmd/mesh-controller/status_summary_test.go b/cmd/mesh-controller/status_summary_test.go new file mode 100644 index 0000000..6f8a698 --- /dev/null +++ b/cmd/mesh-controller/status_summary_test.go @@ -0,0 +1,152 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/link" +) + +// quickly shortens the summary's clocks for one test. +func quickly(t *testing.T) { + t.Helper() + every, settle := statusEvery, statusSettle + statusEvery, statusSettle = time.Hour, 10*time.Millisecond + t.Cleanup(func() { statusEvery, statusSettle = every, settle }) +} + +// **`status` answers in full within ten seconds, five times in a row** (novox/hq to-be 45 Phase 0, +// D9) — however long composing it takes. On 2026-10-06 composing took eighteen seconds and the +// verb answered "still running"; from the summary it answers at once, in full, saying when. +func TestStatusAnswersAtOnceHoweverLongComposingTakes(t *testing.T) { + quickly(t) + var composed atomic.Int32 + slow := make(chan struct{}) + s := newStatusSummary(func(ctx context.Context) ([]byte, error) { + if composed.Add(1) > 1 { + <-slow // every composition after the first outlasts any caller + } + return []byte(`{"wrong":[],"machines":4}`), nil + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + defer close(slow) + go s.keep(ctx) + for i := 0; i < 5; i++ { + s.nudge() // a composition is under way and does not finish + start := time.Now() + got, err := s.answer(ctx) + if err != nil { + t.Fatal(err) + } + if took := time.Since(start); took > time.Second { + t.Fatalf("answer %d took %s", i+1, took) + } + m := got.(map[string]any) + if m["answer"].(map[string]any)["machines"] != float64(4) || m["composed"] == "" || m["ok"] != true { + t.Fatalf("answer %d was not in full: %v", i+1, m) + } + } +} + +// The first composition is waited for, never past the caller's window; a controller with none yet +// says so rather than answering an empty mesh as a well one. +func TestStatusBeforeItsFirstCompositionSaysSo(t *testing.T) { + quickly(t) + was := link.AnswerWithin + link.AnswerWithin = 1100 * time.Millisecond + t.Cleanup(func() { link.AnswerWithin = was }) + never := make(chan struct{}) + defer close(never) + s := newStatusSummary(func(context.Context) ([]byte, error) { <-never; return nil, nil }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go s.keep(ctx) + start := time.Now() + _, err := s.answer(ctx) + if err == nil || !strings.Contains(err.Error(), "has no status to answer yet") { + t.Fatalf("answered %v", err) + } + if took := time.Since(start); took > link.AnswerWithin { + t.Fatalf("waited %s, past the caller's window", took) + } +} + +// A nudge composes it again, once for several close together; a failed composition leaves the last +// summary standing and says it is the last that could be composed. +func TestANudgeComposesAgainAndAFailureKeepsTheLastSummary(t *testing.T) { + quickly(t) + var composed atomic.Int32 + fail := atomic.Bool{} + s := newStatusSummary(func(context.Context) ([]byte, error) { + n := composed.Add(1) + if fail.Load() { + return nil, errors.New("the store did not answer") + } + body, _ := json.Marshal(map[string]any{"n": n}) + return body, nil + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go s.keep(ctx) + if _, err := s.answer(ctx); err != nil { + t.Fatal(err) + } + s.nudge() + s.nudge() + s.nudge() + waitFor(t, func() bool { return composed.Load() == 2 }) + time.Sleep(50 * time.Millisecond) + if n := composed.Load(); n != 2 { + t.Fatalf("three nudges together composed %d times after the first", n-1) + } + fail.Store(true) + s.nudge() + waitFor(t, func() bool { return composed.Load() == 3 }) + waitFor(t, func() bool { + got, err := s.answer(ctx) + if err != nil { + t.Fatal(err) + } + m := got.(map[string]any) + return m["lastAttemptFailed"] != nil && m["answer"].(map[string]any)["n"] == float64(2) + }) +} + +// The seat's `status` answers from the summary when this process keeps one. +func TestTheStatusVerbAnswersFromTheSummary(t *testing.T) { + quickly(t) + s := newStatusSummary(func(context.Context) ([]byte, error) { return []byte(`{"from":"summary"}`), nil }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go s.keep(ctx) + statusFrom = s + t.Cleanup(func() { statusFrom = nil }) + handlers, _, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + got, err := handlers["status"](ctx, json.RawMessage(`{}`)) + if err != nil { + t.Fatal(err) + } + if got.(map[string]any)["answer"].(map[string]any)["from"] != "summary" { + t.Fatalf("answered %v", got) + } +} + +func waitFor(t *testing.T, ok func() bool) { + t.Helper() + deadline := time.Now().Add(3 * time.Second) + for !ok() { + if time.Now().After(deadline) { + t.Fatal("never happened") + } + time.Sleep(5 * time.Millisecond) + } +} diff --git a/cmd/mesh-controller/supersede_test.go b/cmd/mesh-controller/supersede_test.go index 2c269fb..824aaca 100644 --- a/cmd/mesh-controller/supersede_test.go +++ b/cmd/mesh-controller/supersede_test.go @@ -75,21 +75,25 @@ func TestAPersonClosesAStuckPlan(t *testing.T) { if err := open.inventory.SavePlan(ctx, stuck); err != nil { t.Fatal(err) } - if err := plansCommand(ctx, []string{"close", stuck.ID}); err != nil { + if err := plansCommand(ctx, []string{"close", stuck.ID}); err == nil || !strings.Contains(err.Error(), "--why") { + t.Fatalf("a plan was closed by hand without saying why: %v", err) + } + if err := plansCommand(ctx, []string{"close", stuck.ID, "--why", "its report will not come"}); err != nil { t.Fatal(err) } closed, err := open.inventory.PlanByID(ctx, stuck.ID) if err != nil { t.Fatal(err) } - if closed.State != inventory.PlanFailed || !strings.Contains(closed.Note, "closed by hand") { + if closed.State != inventory.PlanFailed || !strings.Contains(closed.Note, "closed by hand") || + !strings.Contains(closed.Note, "its report will not come") { t.Fatalf("the plan was left %s: %q", closed.State, closed.Note) } - if err := plansCommand(ctx, []string{"close", stuck.ID}); err == nil { + if err := plansCommand(ctx, []string{"close", stuck.ID, "--why", "again"}); err == nil { t.Fatal("a plan already closed was closed again") } - if argv, err := argvFor("plans", map[string]any{"close": stuck.ID}); err != nil || - !reflect.DeepEqual(argv, []string{"plans", "close", stuck.ID}) { + if argv, err := argvFor("plans", map[string]any{"close": stuck.ID, "why": "w"}); err != nil || + !reflect.DeepEqual(argv, []string{"plans", "close", stuck.ID, "--why", "w"}) { t.Fatalf("the seat's verb does not close a plan: %v %v", argv, err) } } diff --git a/internal/broker/controller_buckets.go b/internal/broker/controller_buckets.go new file mode 100644 index 0000000..4b18354 --- /dev/null +++ b/internal/broker/controller_buckets.go @@ -0,0 +1,102 @@ +package broker + +import ( + "context" + "fmt" + "time" + + "github.com/nats-io/nats.go/jetstream" +) + +// The controller's own key-value buckets (novox/hq to-be 45 §1, §6, §7). +// +// **What the controller must remember across its own restart, it keeps on the bus.** A call's +// outcome lived in the memory of the process that served it (novox/hq issue 265), so a controller +// replaced while a push ran answered "no such call" for the one thing its caller had been told to +// ask about. The bus already outlives the controller and is the shape ADR 0201 gives a module's +// current state: one value per key, written by one owner, read by anybody granted it. These are the +// controller's, written by it alone — the writers table of to-be 45 §1 — and asserted on every start +// like the streams, so a bus raised from nothing has them before the first call is served. + +// CallsBucket keeps every call of the mesh's own verbs and what came of it; HandActsBucket every act +// a person did by hand, with why. +var ( + CallsBucket = BucketName(ControllerSeat, "calls") + HandActsBucket = BucketName(ControllerSeat, "hand-acts") +) + +// The bounds to-be 45 §6 sets for calls: the last thousand, or fourteen days, whichever is fewer. +// A call is two keys — its record, and its answer apart so a listing does not read every answer — +// so the stream holds twice as many messages as it keeps calls. +const ( + KeptCallsDurably = 1000 + CallsKeptFor = 14 * 24 * time.Hour + // CallAnswerBytes is the most of one answer kept: a whole declaration is far smaller, and an + // answer larger is cut and says so. + CallAnswerBytes = 64 << 10 + // HandActsKeptFor is as long as a condition's history (to-be 45 §2): an act by hand is read + // back beside what it addressed. + HandActsKeptFor = 90 * 24 * time.Hour +) + +// IsControllerBucket says a bucket is the controller's own, not a module's state nothing declares. +func IsControllerBucket(bucket string) bool { + return bucket == CallsBucket || bucket == HandActsBucket +} + +// ControllerBucketsAsserter is what raising the controller's buckets needs of a connection. +type ControllerBucketsAsserter interface { + EnsureControllerBuckets() error +} + +// EnsureControllerBuckets creates the controller's buckets if absent and brings their options to +// match. An update, never a delete: what they hold is the record of what the mesh was asked. +func (j *JetStream) EnsureControllerBuckets() error { + js, err := jetstream.New(j.conn) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: CallsBucket, + Description: "the calls of the mesh's own verbs and what came of each (novox/hq to-be 45 §6, issue " + + "265): written by the controller alone, read through `calls`; the last thousand, or fourteen days", + History: 1, + TTL: CallsKeptFor, + MaxValueSize: CallAnswerBytes + 4<<10, + MaxBytes: 2 * KeptCallsDurably * (CallAnswerBytes + 4<<10), + Storage: jetstream.FileStorage, + }); err != nil { + return fmt.Errorf("asserting bucket %s: %w", CallsBucket, err) + } + // **The count, on the stream under the bucket.** A bucket has an age and a size and no count; + // the stream it is made of does, and with one value per key the oldest message is the oldest + // call. Asserted after the bucket, every time, because asserting the bucket writes the stream's + // configuration whole and puts the count back to none. + stream, err := js.Stream(ctx, "KV_"+CallsBucket) + if err != nil { + return fmt.Errorf("reading the stream under %s: %w", CallsBucket, err) + } + cfg := stream.CachedInfo().Config + if cfg.MaxMsgs != 2*KeptCallsDurably { + cfg.MaxMsgs = 2 * KeptCallsDurably + cfg.Discard = jetstream.DiscardOld + if _, err := js.UpdateStream(ctx, cfg); err != nil { + return fmt.Errorf("bounding %s to the last %d calls: %w", CallsBucket, KeptCallsDurably, err) + } + } + if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: HandActsBucket, + Description: "every act a person did by hand, with why (novox/hq to-be 45 §7): written by the " + + "controller's repairing verbs and `hand-act record`, read through `hand-acts`", + History: 1, + TTL: HandActsKeptFor, + MaxValueSize: 16 << 10, + MaxBytes: 64 << 20, + Storage: jetstream.FileStorage, + }); err != nil { + return fmt.Errorf("asserting bucket %s: %w", HandActsBucket, err) + } + return nil +} diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 50ebab5..560bdd6 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -288,6 +288,11 @@ func PermissionsFor(p Principal) (Permissions, error) { // enrolments and can reach nothing else. pub = append(pub, "_INBOX."+enrolmentPrefix+".>") + // **And its own buckets** (novox/hq to-be 45 §1): the calls it served and the acts done by + // hand, which it alone writes. A put is a publish to the bucket's subject, which `$JS.API.>` + // does not cover; each bucket named, not `$KV.>`, which would let it write any module's state. + pub = append(pub, "$KV."+CallsBucket+".>", "$KV."+HandActsBucket+".>") + case KindPerson: // Tools, and nothing else. Every subject a person may publish is a tool call; a person // who could publish an event would be able to claim a module said something. diff --git a/internal/broker/state.go b/internal/broker/state.go index 9c1e594..033ee4c 100644 --- a/internal/broker/state.go +++ b/internal/broker/state.go @@ -160,8 +160,9 @@ func RaiseBuckets(a BucketAsserter, buckets []Bucket) (undeclared []string, err return nil, fmt.Errorf("listing the bus's state: %w", err) } for _, n := range names { - // A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state. - if !declared[n] && !IsCancelledSet(n) { + // A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state; so are + // the controller's own buckets (novox/hq to-be 45 §1). + if !declared[n] && !IsCancelledSet(n) && !IsControllerBucket(n) { undeclared = append(undeclared, n) } } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 7eaf3a6..db5831c 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "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.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_hand-acts.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "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.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] } subscribe: { allow: ["$JS.API.>", "$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"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 6d30eb8..6ab9fd2 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -110,6 +110,8 @@ var ControllerVerbs = []Verb{ "paths": "with repository: the files the merge would change, comma-separated, from the repository's root", "modules": "with repository: or the modules it would change, comma-separated", "limit": "how many plans to list (default 10); only when listing", + "why": "with stop or close: why it is ended by hand — required, and recorded in the hand-act log (novox/hq to-be 45 §7)", + "cause": "with stop or close: the cause in a word, or a condition's kind (optional)", }, nil)}, {Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " + "or, with files, the files it would be given.", @@ -137,11 +139,14 @@ var ControllerVerbs = []Verb{ {Name: "push", Description: "Send one machine everything it should be. With no machine named it is a push of the " + "WHOLE mesh — every machine that is behind — and the answer says so first; behind says that outright. " + "Answers at once that it is running, with a call id: `calls` with that id says what it sent " + - "(a push can reload the bus, which then refuses any answer still to come).", + "(a push can reload the bus, which then refuses any answer still to come). A push by hand is a repair, " + + "and says why: recorded in the hand-act log (novox/hq to-be 45 §7).", Input: schema(map[string]string{ "node": "the machine's name; without it, every machine that is behind", "behind": "\"true\": every machine that is behind, the whole mesh — the same as naming none, said outright; not with node", - }, nil, "behind")}, + "why": "why this is pushed by hand: recorded in the hand-act log", + "cause": "the cause in a word, or a condition's kind — the word a second push for the same reason uses (optional)", + }, []string{"why"}, "behind")}, {Name: "rotate", Description: "Replace a credential. A pair credential, by provision (and a consuming machine and module, " + "else every holder): both ends are re-sent together. Or a module's own secret, by machine, module and " + "name: made anew and the machine sent, so the module starts again on it — only for a secret its " + @@ -206,6 +211,26 @@ var ControllerVerbs = []Verb{ Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)}, {Name: "resume", Description: "The build seat's holder on one machine — or every holder — takes builds again.", Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)}, + // Acts done by hand, and what the bounds are set from (novox/hq to-be 45 §7, Phase 0). + {Name: "hand-act", Description: "Record an act done by hand outside the mesh — a container restarted, a file " + + "edited, a service started on a machine — with why and its cause, in the hand-act log beside the pushes and " + + "plans ended by hand (novox/hq to-be 45 §7). A cause recorded twice in a fortnight is a healer wanted.", + Input: schema(map[string]string{ + "what": "what was done, in a line", + "why": "why it had to be done by hand", + "cause": "the cause in a word, or a condition's kind — the word a second act for the same reason uses", + "condition": "the key of the condition it addressed, if any (optional)", + }, []string{"what", "why", "cause"})}, + {Name: "hand-acts", Description: "What was done by hand lately — pushes, plans ended, consumers re-made, acts " + + "recorded — who, why and the cause of each, and which causes repeat: each repeat is a healer the mesh lacks.", + Input: schema(map[string]string{"days": "how many days back (default 14)"}, nil)}, + {Name: "durations", Description: "How long things take, as the controller measured them: a send to its machine's " + + "report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build — per machine, " + + "repository or module, with median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).", + Input: schema(map[string]string{ + "kind": "one kind: apply, heartbeat-gap, plan-tier or build; every kind when absent", + "days": "how many days back (default 14)", + }, nil)}, {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/inventory/durations.go b/internal/inventory/durations.go new file mode 100644 index 0000000..df14b74 --- /dev/null +++ b/internal/inventory/durations.go @@ -0,0 +1,146 @@ +package inventory + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" +) + +// The durations the core's bounds are set from (novox/hq to-be 45 Phase 0): see migration 0066. + +// The kinds of duration recorded. +const ( + DurationApply = "apply" + DurationHeartbeatGap = "heartbeat-gap" + DurationPlanTier = "plan-tier" + DurationBuild = "build" +) + +// DurationKinds are every kind, in the order `durations` shows them. +var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild} + +// DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it +// in Phase 1's first live week. +const DurationsKeptFor = 30 * 24 * time.Hour + +// Duration is one measurement. +type Duration struct { + Kind string `json:"kind"` + Subject string `json:"subject"` + Node string `json:"node,omitempty"` + Ref string `json:"ref"` + Started time.Time `json:"started"` + Took time.Duration `json:"took"` + Detail string `json:"detail,omitempty"` +} + +// RecordDuration keeps one measurement, once: the same thing measured again is not a second row. +func (i *Inventory) RecordDuration(ctx context.Context, d Duration) error { + if d.Took < 0 { + return nil // a clock that went back measures nothing + } + _, err := i.store.Pool().Exec(ctx, + `insert into duration (kind, subject, node, ref, started, took_ms, detail) + values ($1, $2, $3, $4, $5, $6, $7) on conflict do nothing`, + d.Kind, d.Subject, d.Node, d.Ref, d.Started, d.Took.Milliseconds(), d.Detail) + return err +} + +// RecordApplyDuration measures a machine's report of the declaration it was last sent: from the send +// to the first report of it. A report of anything else, or of a send already measured, measures +// nothing — a machine reconciling reports the same declaration every few minutes. +func (i *Inventory) RecordApplyDuration(ctx context.Context, node, declared, outcome string) error { + if declared == "" { + return nil + } + _, err := i.store.Pool().Exec(ctx, + `insert into duration (kind, subject, node, ref, started, took_ms, detail) + select $1, n.name, n.name, n.sent || '@' || to_char(n.sent_at at time zone 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US'), + n.sent_at, (extract(epoch from (now() - n.sent_at)) * 1000)::bigint, $4 + from node n + where n.name = $2 and n.sent = $3 and n.sent_at is not null + on conflict do nothing`, + DurationApply, node, declared, outcome) + return err +} + +// RecordHeartbeatGap measures the silence before a machine's word: from the last word heard before +// it, which the caller read before recording this one. +func (i *Inventory) RecordHeartbeatGap(ctx context.Context, node string, before time.Time) error { + if before.IsZero() { + return nil + } + now := time.Now() + return i.RecordDuration(ctx, Duration{Kind: DurationHeartbeatGap, Subject: node, Node: node, + Ref: before.UTC().Format(time.RFC3339Nano), Started: before, Took: now.Sub(before)}) +} + +// Durations is every measurement of a kind since a moment, oldest first; every kind when kind is empty. +func (i *Inventory) Durations(ctx context.Context, kind string, since time.Time) ([]Duration, error) { + rows, err := i.store.Pool().Query(ctx, + `select kind, subject, node, ref, started, took_ms, detail from duration + where ($1 = '' or kind = $1) and recorded >= $2 order by recorded`, kind, since) + if err != nil { + return nil, err + } + defer rows.Close() + var out []Duration + for rows.Next() { + var d Duration + var ms int64 + if err := rows.Scan(&d.Kind, &d.Subject, &d.Node, &d.Ref, &d.Started, &ms, &d.Detail); err != nil { + return nil, err + } + d.Took = time.Duration(ms) * time.Millisecond + out = append(out, d) + } + return out, rows.Err() +} + +// ForgetOldDurations removes what is older than DurationsKeptFor, and says how many. +func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) { + tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`, + time.Now().Add(-DurationsKeptFor)) + if err != nil { + return 0, err + } + return tag.RowsAffected(), nil +} + +// planTierLeft is the measurement a plan's save makes when it leaves a tier: the tier moved on, or +// the plan ended. Read in the save's own transaction, so two saves cannot both measure one tier. +func planTierLeft(ctx context.Context, tx pgx.Tx, p Plan, now time.Time) (entered time.Time, err error) { + var oldTier int + var oldState string + var since time.Time + err = tx.QueryRow(ctx, + `select tier, state, coalesce(tier_entered, created) from release_plan where id = $1 for update`, + p.ID).Scan(&oldTier, &oldState, &since) + if errors.Is(err, pgx.ErrNoRows) { + return now, nil // a new plan enters its first tier now + } + if err != nil { + return time.Time{}, err + } + wasOpen := oldState == PlanBuilding || oldState == PlanRolling + if !wasOpen || (oldTier == p.Tier && p.Open()) { + return since, nil // still in the tier, or already ended + } + var modules []string + if oldTier >= 0 && oldTier < len(p.Tiers) { + modules = p.Tiers[oldTier] + } + detail := fmt.Sprintf("plan %s, tier %d of %d (%v), left %s", p.ID, oldTier, len(p.Tiers), modules, p.State) + if p.Open() { + detail = fmt.Sprintf("plan %s, tier %d of %d (%v), moved on to tier %d", p.ID, oldTier, len(p.Tiers), modules, p.Tier) + } + _, err = tx.Exec(ctx, + `insert into duration (kind, subject, node, ref, started, took_ms, detail) + values ($1, $2, '', $3, $4, $5, $6) on conflict do nothing`, + DurationPlanTier, p.Repository, fmt.Sprintf("%s/tier-%d", p.ID, oldTier), since, + now.Sub(since).Milliseconds(), detail) + return now, err +} diff --git a/internal/inventory/durations_test.go b/internal/inventory/durations_test.go new file mode 100644 index 0000000..c377ac4 --- /dev/null +++ b/internal/inventory/durations_test.go @@ -0,0 +1,97 @@ +package inventory + +import ( + "testing" + "time" +) + +// **The durations the bounds are set from are recorded once each** (novox/hq to-be 45 Phase 0): a +// send measured at its first report and not again at every reconcile that repeats it; a send of +// something else measures nothing; a new send is a new measurement. +func TestAnApplyIsMeasuredOncePerSend(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + node, err := inv.AddNode(ctx, "anchor") + if err != nil { + t.Fatal(err) + } + if err := inv.RecordSent(ctx, node.ID, "d1", nil); err != nil { + t.Fatal(err) + } + for range 3 { // the report, then two reconciles saying the same + if err := inv.RecordApplyDuration(ctx, "anchor", "d1", OutcomeApplied); err != nil { + t.Fatal(err) + } + } + if err := inv.RecordApplyDuration(ctx, "anchor", "d0", OutcomeApplied); err != nil { + t.Fatal(err) + } + ds, err := inv.Durations(ctx, DurationApply, time.Now().Add(-time.Hour)) + if err != nil || len(ds) != 1 || ds[0].Subject != "anchor" || ds[0].Detail != OutcomeApplied || ds[0].Took < 0 { + t.Fatalf("%v %+v", err, ds) + } + time.Sleep(10 * time.Millisecond) + if err := inv.RecordSent(ctx, node.ID, "d2", nil); err != nil { + t.Fatal(err) + } + if err := inv.RecordApplyDuration(ctx, "anchor", "d2", OutcomeFailed); err != nil { + t.Fatal(err) + } + if ds, _ = inv.Durations(ctx, "", time.Now().Add(-time.Hour)); len(ds) != 2 { + t.Fatalf("a second send was not measured: %+v", ds) + } +} + +// A machine's silence is the time since its last word, once per word. +func TestASilenceIsMeasuredFromTheLastWord(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + before := time.Now().Add(-90 * time.Second) + for range 2 { + if err := inv.RecordHeartbeatGap(ctx, "anchor", before); err != nil { + t.Fatal(err) + } + } + if err := inv.RecordHeartbeatGap(ctx, "anchor", time.Time{}); err != nil { + t.Fatal(err) + } + ds, err := inv.Durations(ctx, DurationHeartbeatGap, time.Now().Add(-time.Hour)) + if err != nil || len(ds) != 1 || ds[0].Took < 90*time.Second || ds[0].Took > 2*time.Minute { + t.Fatalf("%v %+v", err, ds) + } +} + +// A plan's tier is measured when the plan leaves it — moving on, or ending — and only then. +func TestAPlansTierIsMeasuredWhenItIsLeft(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + p := Plan{ID: "plan-1", Repository: "novox/mesh-tools", Commit: "abc", Created: time.Now().UTC(), + State: PlanBuilding, Tiers: [][]string{{"mesh-tools"}, {"builder"}}, + Modules: map[string]*PlanModule{"mesh-tools": {}, "builder": {}}} + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + p.State = PlanRolling // the same tier, saved again + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + if ds, _ := inv.Durations(ctx, DurationPlanTier, time.Now().Add(-time.Hour)); len(ds) != 0 { + t.Fatalf("a tier not left was measured: %+v", ds) + } + p.Tier = 1 + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + p.State = PlanDone + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + if err := inv.SavePlan(ctx, p); err != nil { // saved again once ended: nothing more to measure + t.Fatal(err) + } + ds, err := inv.Durations(ctx, DurationPlanTier, time.Now().Add(-time.Hour)) + if err != nil || len(ds) != 2 || ds[0].Subject != "novox/mesh-tools" || ds[0].Ref != "plan-1/tier-0" || + ds[1].Ref != "plan-1/tier-1" { + t.Fatalf("%v %+v", err, ds) + } +} diff --git a/internal/inventory/migrations/0066-the-durations-the-bounds-are-set-from.sql b/internal/inventory/migrations/0066-the-durations-the-bounds-are-set-from.sql new file mode 100644 index 0000000..2b5e2e3 --- /dev/null +++ b/internal/inventory/migrations/0066-the-durations-the-bounds-are-set-from.sql @@ -0,0 +1,31 @@ +-- The durations the core's bounds are set from (novox/hq to-be 45 Phase 0). +-- +-- Every watchdog of the core's signals table has a bound — how long a machine may take to report +-- after a send, how long it may be silent, how long a plan's tier or a build may take — and a bound +-- guessed is a condition that cries wolf or one that never fires. So the controller records each +-- duration as it is observed, and Phase 1 sets the bounds from what was recorded: +-- +-- apply a declaration sent → the machine's first report of that declaration, per machine +-- heartbeat-gap one word from a machine → the next, per machine +-- plan-tier a plan entering a tier → leaving it, per repository +-- build a build asked → its outcome heard, per module +-- +-- One row per thing measured: `ref` names it (the send, the earlier word, the plan's tier, the +-- build), so a report repeated by a reconcile is not a second measurement. Kept a month; read +-- through `durations`. +create table duration ( + kind text not null, + subject text not null, + node text not null default '', + ref text not null, + started timestamptz not null, + took_ms bigint not null, + detail text not null default '', + recorded timestamptz not null default now(), + primary key (kind, subject, ref) +); + +create index duration_by_kind on duration (kind, recorded); + +-- When a plan entered the tier it is at, so leaving it measures the tier. +alter table release_plan add column tier_entered timestamptz; diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index 50e8b89..cb8d875 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -80,13 +80,29 @@ func (i *Inventory) SavePlan(ctx context.Context, p Plan) error { if err != nil { return err } - _, err = i.store.Pool().Exec(ctx, - `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch) - values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10) + // **And how long the tier it left took** (novox/hq to-be 45 Phase 0): measured here, where the + // plan moves, in the same transaction as the move, so no save can move a tier unmeasured or + // measure one twice. + tx, err := i.store.Pool().Begin(ctx) + if err != nil { + return err + } + defer func() { _ = tx.Rollback(ctx) }() + entered, err := planTierLeft(ctx, tx, p, time.Now()) + if err != nil { + return err + } + _, err = tx.Exec(ctx, + `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch, tier_entered) + values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11) on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, - tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch`, - p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch) - return err + tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch, + tier_entered = excluded.tier_entered`, + p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered) + if err != nil { + return err + } + return tx.Commit(ctx) } // OpenPlans is every plan still being worked, oldest first. diff --git a/internal/inventory/seats.go b/internal/inventory/seats.go index f744528..a0065f3 100644 --- a/internal/inventory/seats.go +++ b/internal/inventory/seats.go @@ -105,14 +105,27 @@ func (i *Inventory) widenProtocol(ctx context.Context, s catalogue.Seat) error { changed := false row.Accepts, changed = union(row.Accepts, s.Accepts, changed) row.Emits, changed = union(row.Emits, s.Emits, changed) - have := map[string]bool{} - for _, v := range row.Serves { - have[v.Name] = true + have := map[string]int{} + for n, v := range row.Serves { + have[v.Name] = n } for _, v := range s.Serves { - if !have[v.Name] { + at, kept := have[v.Name] + if !kept { row.Serves = append(row.Serves, v) changed = true + continue + } + // **The controller's own verbs are described by the binary that runs them** (novox/hq + // issue 244, to-be 45 §7). Its seat's verbs are not an operator's to reshape: each is a + // command line this binary composes from the arguments its own table declares, and the + // console judges a call against the row. A row kept from an older build described `push` + // without the `behind` and `why` the binary takes, so the console refused an argument the verb + // needs. So a verb this binary defines takes this binary's definition; a verb only the row has + // — a newer build's, during a roll-out (ADR 0185) — is left as it is. + if s.Name == catalogue.ControllerSeatName && !sameVerb(row.Serves[at], v) { + row.Serves[at] = v + changed = true } } if !changed { @@ -282,3 +295,18 @@ func (i *Inventory) Holdings(ctx context.Context) ([]catalogue.Held, error) { } return out, rows.Err() } + +// sameVerb is whether two definitions of a verb say the same, read as the row stores them. +func sameVerb(a, b catalogue.Verb) bool { + ja, errA := json.Marshal(a) + jb, errB := json.Marshal(b) + if errA != nil || errB != nil { + return false + } + var ra, rb any + _ = json.Unmarshal(ja, &ra) + _ = json.Unmarshal(jb, &rb) + ca, _ := json.Marshal(ra) + cb, _ := json.Marshal(rb) + return string(ca) == string(cb) +} diff --git a/internal/inventory/seats_test.go b/internal/inventory/seats_test.go index e8a7135..ab51567 100644 --- a/internal/inventory/seats_test.go +++ b/internal/inventory/seats_test.go @@ -113,3 +113,46 @@ func TestRenameSeatKeepsTheFormerNameAsAnAlias(t *testing.T) { t.Fatal("renaming a seat to its own name was accepted") } } + +// **The controller's verbs in the row are the binary's** (novox/hq issue 244, to-be 45 §7): a row +// seeded by an older build describes `push` without the arguments this one takes, and the console +// judges a call against the row — so re-seeding brings the controller's verbs to this binary's +// definition, and keeps a verb only the row has, which a newer build added (ADR 0185). +func TestTheControllersVerbsInTheRowAreTheBinarys(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + if _, err := inv.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil { + t.Fatal(err) + } + old := `[{"name":"push","description":"an older push","input":{"type":"object","properties":{"node":{"type":"string"}}}}, + {"name":"newer","description":"a verb of a newer build"}]` + if _, err := inv.store.Pool().Exec(ctx, `update seat set serves = $1 where name = $2`, + []byte(old), catalogue.ControllerSeatName); err != nil { + t.Fatal(err) + } + if _, err := inv.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil { + t.Fatal(err) + } + seats, err := inv.Seats(ctx) + if err != nil { + t.Fatal(err) + } + verbs := map[string]catalogue.Verb{} + for _, s := range seats { + if s.Name == catalogue.ControllerSeatName { + for _, v := range s.Serves { + verbs[v.Name] = v + } + } + } + props, _ := verbs["push"].Input["properties"].(map[string]any) + if _, takesWhy := props["why"]; !takesWhy || verbs["push"].Description == "an older push" { + t.Fatalf("push in the row is still the older build's: %+v", verbs["push"]) + } + if _, kept := verbs["newer"]; !kept { + t.Fatal("a verb only the row has was dropped") + } + if _, added := verbs["durations"]; !added { + t.Fatal("a verb this binary adds was not added") + } +} diff --git a/internal/link/calls.go b/internal/link/calls.go index f094d26..51f996f 100644 --- a/internal/link/calls.go +++ b/internal/link/calls.go @@ -7,7 +7,9 @@ import ( "fmt" "log" "regexp" + "sort" "strconv" + "strings" "sync" "time" @@ -33,8 +35,9 @@ import ( // either. A variable so a test need not wait. var AnswerWithin = 10 * time.Second -// KeptCalls is how many calls are kept, newest first; keptAnswer the largest answer kept of a call -// whose caller was sent it. One its caller never had is kept whole. +// KeptCalls is how many calls this process keeps in memory, newest first — to match a refusal to its +// call, and to answer at once; the bus keeps the last thousand (Durably). keptAnswer is the largest +// answer kept of a call whose caller was sent it. One its caller never had is kept whole. const ( KeptCalls = 100 keptAnswer = 64 << 10 @@ -47,6 +50,10 @@ const ( // CallFinishedAfter is a call that finished after its caller was told it was still running: its // answer is here and nowhere else. CallFinishedAfter = "finished after its caller was answered" + // CallAbandoned is a call whose controller stopped before it finished (novox/hq to-be 45 §6): a + // controller starting finds it running under another and says so, rather than leaving it running + // for ever in the record. + CallAbandoned = "abandoned: the controller running it stopped before it finished" ) // Call is one call of a role's tool, as `calls` shows it. @@ -65,10 +72,27 @@ type Call struct { // Refused is the bus refusing the answer this holder sent: the caller got nothing, and this is // the only place that says what it would have. Refused string `json:"answer refused by the bus,omitempty"` + // Caller is the bus principal that asked, read from the inbox its answer went to — every principal + // is granted only its own (novox/hq to-be 45 §7: a hand act says who). + Caller string `json:"caller,omitempty"` + // Holder is the controller process that served it, so one starting can tell its own running + // calls from those a stopped one left. + Holder string `json:"holder,omitempty"` reply string // the subject the answer went to, which the bus names when it refuses it } +// CallKeeper keeps calls where the process serving them does not: the controller's bucket on the +// bus (novox/hq to-be 45 §6). A call is kept whole on every change — begun, finished, refused — so a +// controller replaced at any moment leaves the last word on each. +type CallKeeper interface { + Keep(ctx context.Context, c Call) error + // Kept is one call with its whole answer. + Kept(ctx context.Context, id string) (Call, bool, error) + // Recent is the kept calls, newest first, without their answers. + Recent(ctx context.Context) ([]Call, error) +} + // CallLog keeps the latest calls a holder served. type CallLog struct { mu sync.Mutex @@ -78,6 +102,15 @@ type CallLog struct { // Follow is the tool that reads this log back, named in a running answer — set by a holder that // serves one (the controller's `calls`); without it, the answer points at the holder's journal. Follow string + + // keeper keeps every call beyond this process, when Durably was given one; writes go through + // one goroutine, in order, so a call's last state is the one kept. + keeper CallKeeper + holder string + writes chan Call + logger *log.Logger + lost int // writes the keeper could not take, said once each + keepErr error } // Calls is this process's log: one holder process serves its seats on one connection. @@ -85,16 +118,130 @@ var Calls = NewCallLog() func NewCallLog() *CallLog { return &CallLog{now: time.Now} } +// keepTries is how many times one call's state is offered to the keeper before it is said lost: the +// bus reloading its user list refuses for a moment, and that is exactly when a push runs. +const keepTries = 5 + +// Durably keeps every call from now on with keeper as well as in memory, under this process's name, +// and marks running the calls a controller before this one left running: it stopped, so they cannot +// finish (novox/hq to-be 45 §6). Said, naming each. +func (l *CallLog) Durably(ctx context.Context, keeper CallKeeper, holder string, logger *log.Logger) error { + l.mu.Lock() + l.keeper, l.holder, l.logger = keeper, holder, logger + if l.writes == nil { + l.writes = make(chan Call, 256) + go l.keepWrites() + } + l.mu.Unlock() + kept, err := keeper.Recent(ctx) + if err != nil { + return fmt.Errorf("reading the calls kept on the bus: %w", err) + } + for _, c := range kept { + if c.State != CallRunning || c.Holder == holder { + continue + } + whole, found, err := keeper.Kept(ctx, c.ID) + if err != nil || !found { + whole = c + } + whole.State = CallAbandoned + if err := keeper.Keep(ctx, whole); err != nil { + return fmt.Errorf("marking %s abandoned: %w", c.ID, err) + } + if logger != nil { + logger.Printf("%s (%s.%s, asked %s by %s) was running under %s, which stopped: marked abandoned — "+ + "it may have done part of what it was asked, and nothing will finish it", c.ID, c.Seat, c.Verb, + c.Started.Format(time.RFC3339), orSomebody(c.Caller), orSomebody(c.Holder)) + } + } + return nil +} + +func orSomebody(s string) string { + if s == "" { + return "an unnamed caller" + } + return s +} + +// keep queues one call's state for the keeper. Never blocks a call: a queue that is full is a keeper +// that is not taking writes, and that is said rather than waited on. +func (l *CallLog) keep(c Call) { + if l.writes == nil { + return + } + select { + case l.writes <- c: + default: + l.lost++ + if l.logger != nil { + l.logger.Printf("%s (%s.%s) is kept in memory only: the bus is not taking calls' records (%d not kept)", + c.ID, c.Seat, c.Verb, l.lost) + } + } +} + +func (l *CallLog) keepWrites() { + for c := range l.writes { + var err error + for try := 0; try < keepTries; try++ { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + err = l.keeper.Keep(ctx, c) + cancel() + if err == nil { + break + } + time.Sleep(time.Duration(try+1) * time.Second) + } + if err != nil && l.logger != nil { + l.logger.Printf("%s (%s.%s, %s) could not be kept on the bus after %d tries: %v — `calls` "+ + "answers it from memory until this controller stops", c.ID, c.Seat, c.Verb, c.State, keepTries, err) + } + } +} + +// callerOf is the bus principal an answer goes to: every principal's inbox is `_INBOX..` +// followed by the client's own random token, and a user may itself hold dots. +func callerOf(reply string) string { + rest, ok := strings.CutPrefix(reply, "_INBOX.") + if !ok { + return "" + } + tokens := strings.Split(rest, ".") + for i, t := range tokens { + if i > 0 && isNUID(t) { + return strings.Join(tokens[:i], ".") + } + } + return "" +} + +// isNUID is the client library's random inbox token: twenty-two letters and digits. +func isNUID(t string) bool { + if len(t) != 22 { + return false + } + for _, r := range t { + if !(r >= '0' && r <= '9' || r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z') { + return false + } + } + return true +} + func (l *CallLog) begin(seat, verb string, args json.RawMessage, reply string) *Call { l.mu.Lock() defer l.mu.Unlock() l.next++ c := &Call{ID: "call-" + strconv.FormatInt(l.now().UnixNano(), 10) + "-" + strconv.FormatUint(l.next, 10), - Seat: seat, Verb: verb, Args: args, Started: l.now(), State: CallRunning, reply: reply} + Seat: seat, Verb: verb, Args: args, Started: l.now(), State: CallRunning, reply: reply, + Caller: callerOf(reply), Holder: l.holder} l.calls = append(l.calls, c) if len(l.calls) > KeptCalls { l.calls = l.calls[len(l.calls)-KeptCalls:] } + l.keep(*c) return c } @@ -117,29 +264,67 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) { } else { c.State = CallAnswered } + l.keep(*c) } -// Recent is the kept calls, newest first, as copies. -func (l *CallLog) Recent() []Call { +// Recent is the kept calls, newest first, as copies: this process's from memory, and, when they are +// kept durably, every other the bus holds — a call a controller before this one served included. +// Memory wins for a call in both, being the newer word on it. A bus that cannot be read is said in +// the error beside what memory holds, never answered as no calls. +func (l *CallLog) Recent() ([]Call, error) { l.mu.Lock() - defer l.mu.Unlock() out := make([]Call, 0, len(l.calls)) + seen := map[string]bool{} for i := len(l.calls) - 1; i >= 0; i-- { out = append(out, *l.calls[i]) + seen[l.calls[i].ID] = true } - return out -} - -// Get is one kept call. -func (l *CallLog) Get(id string) (Call, bool) { - l.mu.Lock() - defer l.mu.Unlock() - for _, c := range l.calls { - if c.ID == id { - return *c, true + keeper := l.keeper + l.mu.Unlock() + if keeper == nil { + return out, nil + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + kept, err := keeper.Recent(ctx) + if err != nil { + return out, fmt.Errorf("the calls kept on the bus could not be read, so only this controller's own "+ + "are listed: %w", err) + } + for _, c := range kept { + if !seen[c.ID] { + out = append(out, c) } } - return Call{}, false + sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) }) + return out, nil +} + +// Get is one kept call: from memory, or from the bus when this process did not serve it. +func (l *CallLog) Get(id string) (Call, bool, error) { + l.mu.Lock() + for _, c := range l.calls { + if c.ID == id { + found := *c + l.mu.Unlock() + return found, true, nil + } + } + keeper := l.keeper + l.mu.Unlock() + if keeper == nil { + return Call{}, false, nil + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + return keeper.Kept(ctx, id) +} + +// IsDurable says whether calls outlive this process. +func (l *CallLog) IsDurable() bool { + l.mu.Lock() + defer l.mu.Unlock() + return l.keeper != nil } // kept is what a call's arguments are kept as: each argument by name, a value only when it is short @@ -195,6 +380,7 @@ func (l *CallLog) Refusal(err error, logger *log.Logger) bool { at := l.now() hit.Refused = fmt.Sprintf("%s: %s", at.Format(time.RFC3339), err) id, verb, took := hit.ID, hit.Seat+"."+hit.Verb, at.Sub(hit.Started).Round(time.Second) + l.keep(*hit) l.mu.Unlock() if logger != nil { // Said in the mesh's words, beside the library's own line: which call, and where its answer is. @@ -277,6 +463,8 @@ func (l *CallLog) serveCall(seat, verb string, args json.RawMessage, reply strin args = json.RawMessage(`{}`) } c := l.begin(seat, verb, kept(args), reply) + // Who asked travels with the call, so an act it does by hand says so (novox/hq to-be 45 §7). + ctx = context.WithValue(ctx, callerKey{}, c.Caller) acknowledged := make(chan struct{}) var once sync.Once ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) }) diff --git a/internal/link/calls_kept.go b/internal/link/calls_kept.go new file mode 100644 index 0000000..06e4fff --- /dev/null +++ b/internal/link/calls_kept.go @@ -0,0 +1,115 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "sort" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// Calls kept on the bus (novox/hq to-be 45 §6). +// +// **Two keys a call**: `` holds the record — verb, arguments as kept, caller, state, times — and +// `.answer` the answer, bounded. A listing watches the records alone, so reading the last thousand +// calls reads the last thousand small records and not a thousand answers; one call asked by id reads +// both. A call's id is one token, so the record's key never holds a dot and `*` matches records only. + +// BusCalls keeps calls in the controller's calls bucket. +type BusCalls struct { + kv jetstream.KeyValue +} + +// CallsOnTheBus opens the calls bucket the controller asserts at its start. +func CallsOnTheBus(ctx context.Context, conn *nats.Conn) (*BusCalls, error) { + api, err := jetstream.New(conn) + if err != nil { + return nil, err + } + kv, err := api.KeyValue(ctx, broker.CallsBucket) + if err != nil { + return nil, fmt.Errorf("the calls bucket %s is not on the bus — the controller asserts it at its "+ + "start, so one older than this has not: %w", broker.CallsBucket, err) + } + return &BusCalls{kv: kv}, nil +} + +const answerKey = ".answer" + +// Keep writes a call's record, and its answer when it has one. +func (b *BusCalls) Keep(ctx context.Context, c Call) error { + answer := c.Answer + c.Answer = nil + if len(answer) > 0 { + if len(answer) > broker.CallAnswerBytes { + answer, _ = json.Marshal(map[string]any{"cut": fmt.Sprintf("an answer of %d bytes; the first %d "+ + "are kept", len(answer), broker.CallAnswerBytes), "start": string(answer[:broker.CallAnswerBytes])}) + } + // The answer before the record, so a record saying a call finished never points at an answer + // not yet written. + if _, err := b.kv.Put(ctx, c.ID+answerKey, answer); err != nil { + return err + } + } + record, err := json.Marshal(c) + if err != nil { + return err + } + _, err = b.kv.Put(ctx, c.ID, record) + return err +} + +// Kept is one call, with its answer. +func (b *BusCalls) Kept(ctx context.Context, id string) (Call, bool, error) { + entry, err := b.kv.Get(ctx, id) + if errors.Is(err, jetstream.ErrKeyNotFound) || errors.Is(err, jetstream.ErrInvalidKey) { + return Call{}, false, nil + } + if err != nil { + return Call{}, false, err + } + var c Call + if err := json.Unmarshal(entry.Value(), &c); err != nil { + return Call{}, false, fmt.Errorf("the record of %s on the bus is not a call: %w", id, err) + } + answer, err := b.kv.Get(ctx, id+answerKey) + switch { + case err == nil: + c.Answer = json.RawMessage(answer.Value()) + case !errors.Is(err, jetstream.ErrKeyNotFound): + return Call{}, false, err + } + return c, true, nil +} + +// Recent is every call the bucket holds, newest first, without answers: read once through a watch +// of the records, which hands over the current value of each and then says it has. +func (b *BusCalls) Recent(ctx context.Context) ([]Call, error) { + w, err := b.kv.Watch(ctx, "*", jetstream.IgnoreDeletes()) + if err != nil { + return nil, err + } + defer func() { _ = w.Stop() }() + var out []Call + for { + select { + case <-ctx.Done(): + return nil, fmt.Errorf("reading the calls bucket: %w", ctx.Err()) + case entry := <-w.Updates(): + if entry == nil { + // Every current value handed over. + sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) }) + return out, nil + } + var c Call + if json.Unmarshal(entry.Value(), &c) == nil && c.ID != "" { + out = append(out, c) + } + } + } +} diff --git a/internal/link/calls_kept_test.go b/internal/link/calls_kept_test.go new file mode 100644 index 0000000..84a23be --- /dev/null +++ b/internal/link/calls_kept_test.go @@ -0,0 +1,152 @@ +package link + +import ( + "bytes" + "context" + "encoding/json" + "log" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The calls bucket against a real server (novox/hq to-be 45 §6): whether a call's outcome is +// readable by id from a controller that did not serve it is a claim about what the bus keeps. + +func callsBucket(t *testing.T) *BusCalls { + t.Helper() + js := aBus(t) + api, err := jetstream.New(js.Conn()) + if err != nil { + t.Fatal(err) + } + _ = api.DeleteKeyValue(t.Context(), broker.CallsBucket) + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + keeper, err := CallsOnTheBus(t.Context(), js.Conn()) + if err != nil { + t.Fatal(err) + } + return keeper +} + +// waitKept waits until the keeper holds a call in a state, the writes being queued. +func waitKept(t *testing.T, keeper CallKeeper, id, state string) Call { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + for { + c, found, err := keeper.Kept(context.Background(), id) + if err == nil && found && c.State == state { + return c + } + if time.Now().After(deadline) { + t.Fatalf("%s never kept as %q: %+v %v %v", id, state, c, found, err) + } + time.Sleep(20 * time.Millisecond) + } +} + +// **A controller restart keeps every call's outcome** (to-be 45 Phase 0): a call one controller +// answered is read whole by id from the next, and a call it left running is said abandoned rather +// than running for ever. +func TestNatsACallsOutcomeOutlivesItsController(t *testing.T) { + keeper := callsBucket(t) + first := NewCallLog() + if err := first.Durably(t.Context(), keeper, "controller@one", nil); err != nil { + t.Fatal(err) + } + a := newAnswers(t) + first.serveCall("mesh-controller", "push", json.RawMessage(`{"node":"anchor","values":"secret"}`), + "_INBOX.node-tools.g14.Pe3sGzAtv8jUBYQ6sKSvKz.1", + func(context.Context, json.RawMessage) (any, error) { return "anchor told", nil }, a.respond, nil) + recent, _ := first.Recent() + finished := waitKept(t, keeper, recent[0].ID, CallAnswered) + // A second call, still running when its controller stops. + left := first.begin("mesh-controller", "plans", json.RawMessage(`{}`), "") + waitKept(t, keeper, left.ID, CallRunning) + + var said bytes.Buffer + second := NewCallLog() + if err := second.Durably(t.Context(), keeper, "controller@two", log.New(&said, "", 0)); err != nil { + t.Fatal(err) + } + c, found, err := second.Get(finished.ID) + if err != nil || !found { + t.Fatalf("the next controller cannot read %s: %v %v", finished.ID, found, err) + } + if !strings.Contains(string(c.Answer), "anchor told") || c.Caller != "node-tools.g14" || c.Holder != "controller@one" { + t.Fatalf("kept %+v", c) + } + if strings.Contains(string(c.Args), "secret") { + t.Fatalf("a setting was kept: %s", c.Args) + } + abandoned, _, _ := second.Get(left.ID) + if abandoned.State != CallAbandoned || !strings.Contains(said.String(), left.ID) { + t.Fatalf("a call left running reads %q, and was said: %q", abandoned.State, said.String()) + } + listed, err := second.Recent() + if err != nil || len(listed) != 2 { + t.Fatalf("the next controller lists %d calls: %v", len(listed), err) + } + // Its own running calls are not its predecessor's: a controller starting again under the same + // name leaves them alone. + mine := second.begin("mesh-controller", "status", nil, "") + waitKept(t, keeper, mine.ID, CallRunning) + if err := second.Durably(t.Context(), keeper, "controller@two", nil); err != nil { + t.Fatal(err) + } + if c, _, _ := keeper.Kept(t.Context(), mine.ID); c.State != CallRunning { + t.Fatalf("a controller marked its own running call %q", c.State) + } +} + +// The bucket holds the last thousand calls or fourteen days: a call is two keys, so the stream under +// it holds twice as many messages, and asserting it again keeps the count. +func TestNatsTheCallsBucketIsBounded(t *testing.T) { + callsBucket(t) + js := aBus(t) + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + info, err := js.Context().StreamInfo("KV_" + broker.CallsBucket) + if err != nil { + t.Fatal(err) + } + if info.Config.MaxMsgs != 2*broker.KeptCallsDurably || info.Config.MaxAge != broker.CallsKeptFor { + t.Fatalf("kept %d messages for %s", info.Config.MaxMsgs, info.Config.MaxAge) + } +} + +// An answer larger than the bound is cut and says so, and the record is still written. +func TestNatsAnAnswerLargerThanTheBoundIsCut(t *testing.T) { + keeper := callsBucket(t) + big, _ := json.Marshal(map[string]string{"result": strings.Repeat("x", broker.CallAnswerBytes+10)}) + c := Call{ID: "call-1-1", Seat: "mesh-controller", Verb: "plan", State: CallAnswered, Answer: big} + if err := keeper.Keep(t.Context(), c); err != nil { + t.Fatal(err) + } + got, found, err := keeper.Kept(t.Context(), "call-1-1") + if err != nil || !found || !strings.Contains(string(got.Answer), "the first") { + t.Fatalf("%v %v %.200s", found, err, got.Answer) + } +} + +// Who asked is read from the inbox the answer goes to. +func TestTheCallerIsTheInboxsPrincipal(t *testing.T) { + for reply, want := range map[string]string{ + "_INBOX.node-tools.g14.Pe3sGzAtv8jUBYQ6sKSvKz.1": "node-tools.g14", + "_INBOX.jochen.Pe3sGzAtv8jUBYQ6sKSvKz": "jochen", + "_INBOX.x.1": "", + "mesh.control.one": "", + "": "", + } { + if got := callerOf(reply); got != want { + t.Errorf("%q: %q, want %q", reply, got, want) + } + } +} diff --git a/internal/link/calls_test.go b/internal/link/calls_test.go index 0d1f9b7..1291dee 100644 --- a/internal/link/calls_test.go +++ b/internal/link/calls_test.go @@ -62,7 +62,7 @@ func TestACallThatFinishesInTimeAnswersInFull(t *testing.T) { if got := a.only()["result"]; got != "all well" { t.Fatalf("answered %v", got) } - recent := l.Recent() + recent, _ := l.Recent() if len(recent) != 1 || recent[0].State != CallAnswered || string(recent[0].Args) != "{}" { t.Fatalf("kept %+v", recent) } @@ -90,12 +90,12 @@ func TestACallThatOutlastsTheWindowSaysItIsRunningAndKeepsItsAnswer(t *testing.T if got["running"] != true || id == "" || !strings.Contains(got["output"].(string), id) { t.Fatalf("the running answer does not name its call: %v", got) } - if c, _ := l.Get(id); c.State != CallRunning { + if c, _, _ := l.Get(id); c.State != CallRunning { t.Fatalf("while it runs it is kept as %q", c.State) } close(release) <-finished - c, ok := l.Get(id) + c, ok, _ := l.Get(id) if !ok || c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "anchor told") { t.Fatalf("its answer was not kept: %+v", c) } @@ -125,7 +125,7 @@ func TestAnAcknowledgedCallIsAnsweredBeforeItGoesOn(t *testing.T) { if got := a.only()["result"].(map[string]any); got["running"] != true { t.Fatalf("an acknowledged call answered %v", got) } - if c := l.Recent()[0]; c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "sent") { + if c := recentOf(l)[0]; c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "sent") { t.Fatalf("kept %+v", c) } } @@ -143,7 +143,7 @@ func TestARefusedAnswerIsKeptAgainstItsCall(t *testing.T) { if !l.Refusal(refusal, log.New(&logged, "", 0)) { t.Fatal("the refusal of a kept call's answer was not recognised") } - c := l.Recent()[0] + c := recentOf(l)[0] if c.Refused == "" || !strings.Contains(logged.String(), c.ID) { t.Fatalf("the refusal is not kept or not said: %+v / %q", c, logged.String()) } @@ -167,8 +167,13 @@ func TestTheLogKeepsTheNewest(t *testing.T) { for i := 0; i < KeptCalls+5; i++ { l.begin("s", "v", nil, "") } - recent := l.Recent() + recent, _ := l.Recent() if len(recent) != KeptCalls || !strings.HasSuffix(recent[0].ID, "-105") || !strings.HasSuffix(recent[KeptCalls-1].ID, "-6") { t.Fatalf("kept %d, newest %s, oldest %s", len(recent), recent[0].ID, recent[len(recent)-1].ID) } } + +func recentOf(l *CallLog) []Call { + out, _ := l.Recent() + return out +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index 0f59178..463fbe5 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -380,6 +380,10 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err if report.Applied == nil && report.Refused == "" && len(report.Failed) == 0 { if report.Superseded != "" { log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded) + } else if err := e.Inventory.RecordHeartbeatGap(ctx, report.Node, node.LastSeen); err != nil { + // The silence before this word, which the machine-silent bound will be set from (novox/hq + // to-be 45 Phase 0). A measurement lost is said and costs the report nothing. + log.Printf("the silence before %s's word could not be recorded: %v", report.Node, err) } return false, e.Inventory.Seen(ctx, node.ID) } @@ -435,6 +439,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err if err != nil { return false, err } + // And how long the machine took, from the send to this first account of it (novox/hq to-be 45 + // Phase 0): what the sent-not-reported bound will be set from. Said if lost, never a failure. + if err := e.Inventory.RecordApplyDuration(ctx, report.Node, report.Declared, doing.Outcome); err != nil { + log.Printf("how long %s took to apply could not be recorded: %v", report.Node, err) + } // A refusal, a failure, or a bare word that the node is there — none of them is an account of // what the machine holds, so each moves last_seen and nothing else. Recording a partial list diff --git a/internal/link/handacts.go b/internal/link/handacts.go new file mode 100644 index 0000000..ed437e7 --- /dev/null +++ b/internal/link/handacts.go @@ -0,0 +1,162 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "os" + "sort" + "strconv" + "strings" + "sync/atomic" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The hand-act log (novox/hq to-be 45 §7). +// +// **A repair a person makes by hand is the record of a healer the mesh does not have yet.** Tonight's +// pushes of one machine, plans closed because a report never came, a consumer re-made because it fell +// a week behind — each was done, and the only trace was a chat. So every verb that repairs by hand +// asks why, and writes one entry: who, which verb and arguments, why, when, the condition it +// addresses if one is named, and a cause — a word that, recorded twice in a fortnight, says a healer +// is wanted (S15, from Phase 3). `hand-act record` is the same entry for an act done outside the mesh. +// The controller is the bucket's only writer; its verbs are the way in. + +// HandAct is one entry. +type HandAct struct { + ID string `json:"id"` + At time.Time `json:"at"` + // By is who: the bus principal a seat call came from, or the account and machine at a shell. + By string `json:"by"` + // Verb and Args are the act as given: `push`, `plans close`, `broker consumer-reset`, or + // `hand-act record` with what was done outside the mesh. + Verb string `json:"verb"` + Args []string `json:"arguments,omitempty"` + Why string `json:"why"` + // Cause is the condition kind, or a word the person gives; the verb's own name when neither. + Cause string `json:"cause"` + // Condition is the condition's key the act addresses, when it names one. + Condition string `json:"condition,omitempty"` +} + +// CallerVar carries a seat call's caller to the command the controller runs for it, so an act done +// through the console says who asked rather than "the controller". +const CallerVar = "MESH_CALLER" + +// Caller is who is acting in this process: the seat call's caller when the controller ran it for +// one, otherwise the account and machine at the shell. +func Caller() string { + if c := strings.TrimSpace(os.Getenv(CallerVar)); c != "" { + return c + } + user := os.Getenv("USER") + if user == "" { + user = "an unnamed account" + } + host, _ := os.Hostname() + return fmt.Sprintf("%s at a shell on %s", user, host) +} + +type callerKey struct{} + +// CallerIn is the caller of the seat call ctx belongs to, empty outside one. +func CallerIn(ctx context.Context) string { + c, _ := ctx.Value(callerKey{}).(string) + return c +} + +var handActSeq atomic.Uint64 + +// RecordHandAct writes one entry. Its key is its time and a sequence, so the bucket lists in order. +func RecordHandAct(ctx context.Context, conn *nats.Conn, act HandAct) (HandAct, error) { + if strings.TrimSpace(act.Why) == "" { + return act, errors.New("an act by hand says why: --why ") + } + if act.At.IsZero() { + act.At = time.Now().UTC() + } + if act.ID == "" { + act.ID = "act-" + strconv.FormatInt(act.At.UnixNano(), 10) + "-" + strconv.FormatUint(handActSeq.Add(1), 10) + } + if act.By == "" { + act.By = Caller() + } + if act.Cause == "" { + act.Cause = act.Verb + } + kv, err := handActs(ctx, conn) + if err != nil { + return act, err + } + body, err := json.Marshal(act) + if err != nil { + return act, err + } + _, err = kv.Put(ctx, act.ID, body) + return act, err +} + +// HandActs is every entry since a moment, oldest first. +func HandActs(ctx context.Context, conn *nats.Conn, since time.Time) ([]HandAct, error) { + kv, err := handActs(ctx, conn) + if err != nil { + return nil, err + } + w, err := kv.WatchAll(ctx, jetstream.IgnoreDeletes()) + if err != nil { + return nil, err + } + defer func() { _ = w.Stop() }() + var out []HandAct + for { + select { + case <-ctx.Done(): + return nil, fmt.Errorf("reading the hand-act log: %w", ctx.Err()) + case entry := <-w.Updates(): + if entry == nil { + sort.SliceStable(out, func(i, j int) bool { return out[i].At.Before(out[j].At) }) + return out, nil + } + var a HandAct + if json.Unmarshal(entry.Value(), &a) == nil && !a.At.Before(since) { + out = append(out, a) + } + } + } +} + +// RepeatedCauses are the causes recorded more than once within the fortnight before now, with how +// often: each is a repair done by hand again, which is what S15 will raise as a healer wanted. +func RepeatedCauses(acts []HandAct, now time.Time) map[string]int { + counts := map[string]int{} + for _, a := range acts { + if now.Sub(a.At) <= 14*24*time.Hour { + counts[a.Cause]++ + } + } + for c, n := range counts { + if n < 2 { + delete(counts, c) + } + } + return counts +} + +func handActs(ctx context.Context, conn *nats.Conn) (jetstream.KeyValue, error) { + api, err := jetstream.New(conn) + if err != nil { + return nil, err + } + kv, err := api.KeyValue(ctx, broker.HandActsBucket) + if err != nil { + return nil, fmt.Errorf("the hand-act log %s is not on the bus — the controller asserts it at its "+ + "start, so one older than this has not: %w", broker.HandActsBucket, err) + } + return kv, nil +} diff --git a/internal/link/handacts_test.go b/internal/link/handacts_test.go new file mode 100644 index 0000000..059e552 --- /dev/null +++ b/internal/link/handacts_test.go @@ -0,0 +1,59 @@ +package link + +import ( + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The hand-act log against a real server (novox/hq to-be 45 §7): an act is written with who, why and +// its cause, read back in order, and a cause recorded twice within a fortnight is found. +func TestNatsAnActByHandIsKeptWithWhyAndARepeatIsFound(t *testing.T) { + js := aBus(t) + api, err := jetstream.New(js.Conn()) + if err != nil { + t.Fatal(err) + } + _ = api.DeleteKeyValue(t.Context(), broker.HandActsBucket) + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + t.Setenv(CallerVar, "node-tools.g14, through the mesh-controller seat") + if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"anchor"}}); err == nil { + t.Fatal("an act without why was written") + } + first, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"anchor"}, + Why: "it never reported the send", Cause: "sent-not-reported"}) + if err != nil { + t.Fatal(err) + } + if first.By != "node-tools.g14, through the mesh-controller seat" || first.ID == "" { + t.Fatalf("written as %+v", first) + } + if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "plans close", Args: []string{"plan-1"}, + Why: "waiting on the same report"}); err != nil { + t.Fatal(err) + } + if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"ace"}, + Why: "again", Cause: "sent-not-reported", At: time.Now().UTC().Add(time.Second)}); err != nil { + t.Fatal(err) + } + acts, err := HandActs(t.Context(), js.Conn(), time.Now().Add(-time.Hour)) + if err != nil || len(acts) != 3 || acts[0].ID != first.ID || acts[1].Cause != "plans close" { + t.Fatalf("%v %+v", err, acts) + } + repeated := RepeatedCauses(acts, time.Now()) + if len(repeated) != 1 || repeated["sent-not-reported"] != 2 { + t.Fatalf("repeated %v", repeated) + } + if old, _ := HandActs(t.Context(), js.Conn(), time.Now().Add(time.Hour)); len(old) != 0 { + t.Fatalf("acts before the moment asked were listed: %+v", old) + } + if !strings.HasPrefix(acts[2].ID, "act-") { + t.Fatalf("%q", acts[2].ID) + } +}