Author SHA1 Message Date
mesh-admin cf4834a36c Merge pull request 'Keep calls and hand acts on the bus, answer status at once, record durations (hq to-be 45 Phase 0)' (#78) from feat/a-core-that-cannot-fail-silently-phase-0 into main 2026-10-06 07:13:39 +00:00
jochen 9d8cbe7b81 Grant the controller its work queues' cancelled sets (hq issue 269)
A cancel writes the ask's id into the seat's cancelled set before deleting
the ask, and a write is a publish to the bucket's subject, which the
controller was not granted: against a server holding exactly the
controller's composed list, every cancel timed out.
2026-10-06 03:01:19 +02:00
jochen e74c32ed50 Keep calls and hand acts on the bus, answer status at once, record durations (hq to-be 45 Phase 0)
A controller restart lost every call's outcome, `status` composed the mesh
while its caller waited (18.6s live on 2026-10-06, past the 10s window), a
repair by hand left no trace, and the core's bounds had nothing measured to
be set from.

- calls: kept in the controller's bucket mesh-controller_calls (last 1000 or
  14 days, answers bounded to 64 KiB), read by id across a restart; a
  controller starting marks a stopped one's running calls abandoned; each
  call names its caller from the inbox its answer goes to.
- status: the serving controller composes it at start, after news from a
  machine, a build or an acting verb, and every minute; the verb answers the
  last composition at once with when and how long it took. Composing resolves
  each machine once instead of twice.
- hand-act log in mesh-controller_hand-acts: push (required through the seat),
  plans stop/close, broker consumer-reset and the new hand-act record take
  --why/--cause/--condition; `hand-acts` lists them and repeated causes;
  status counts the week's.
- durations (migration 0066): apply (send to first report), heartbeat gap,
  plan tier and build, recorded as heard; `durations` summarises them.
- the controller's seat row takes this binary's definition of its own verbs,
  so the console no longer judges calls against an older build's schema.
- the controller is granted its two buckets' subjects.
2026-10-06 02:59:36 +02:00
mesh-admin 146c48fd96 Merge pull request 'Bound a consumer's identity by the provision it requires (hq issue 263, ADR 0225)' (#76) from fix/263-identity-bound-per-provision into main 2026-10-06 00:28:47 +00:00
mesh-admin 1f3abd3e0e Merge pull request 'rotate: narrow a pair credential to one consuming module (hq issue 268)' (#75) from feat/rotate-one-consuming-module into main 2026-10-06 00:25:31 +00:00
jochen 6d620f77c3 Bound a consumer's identity by the provision it requires (hq issue 263)
The one global 20-character bound made every consumer pay an object
store's key length, even for provisions that keep no name, and a single
overflow refused the provider's whole declaration. An offer now states
its own bound (identity: {max, in} or false); unsaid, a provider told its
consumers keeps 20 and one told nothing keeps none. module check judges
every identity on the longest machine name before merge, and a provider
leaves an overflowing consumer out of its grants and composes, with the
consumer named by push, plan and status (ADR 0225).
2026-10-06 02:16:20 +02:00
jochen f8286c063d rotate: narrow a pair credential to one consuming module (hq issue 268)
A machine runs many consumers of one provision, each with its own
credential. When one module leaks its credential, `rotate <provision>
--consumer <machine>` was the narrowest act and replaced every module's
on that machine, restarting all of them. --module (and the verb's
module argument beside provision) rotates only that module's.
2026-10-06 02:13:48 +02:00
50 changed files with 3207 additions and 120 deletions
+8
View File
@@ -679,6 +679,14 @@ type answers struct {
// failing is every consumer a provider says it keeps failing (novox/hq ADR 0224): a provider's // failing is every consumer a provider says it keeps failing (novox/hq ADR 0224): a provider's
// journal was the only place that said so for a day (04-ISSUES/179). // journal was the only place that said so for a day (04-ISSUES/179).
failing []inventory.ProviderStanding failing []inventory.ProviderStanding
// overflowing is every module whose identity overflows the bound of a provision it requires
// (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. // heldBy is every artifact this mesh has built, for a build that may need one as its base.
+4 -4
View File
@@ -35,10 +35,10 @@ func TestAPushAnswersBeforeItSends(t *testing.T) {
args map[string]any args map[string]any
want bool want bool
}{ }{
{"push", map[string]any{"node": "anchor"}, true}, {"push", map[string]any{"node": "anchor", "why": "w"}, true},
{"push", map[string]any{}, true}, {"push", map[string]any{"why": "w"}, true},
{"command", map[string]any{"command": "push anchor"}, true}, {"command", map[string]any{"command": "push anchor --why w"}, true},
{"command", map[string]any{"command": "push --behind"}, true}, {"command", map[string]any{"command": "push --behind --why=w"}, true},
{"command", map[string]any{"command": "builds"}, false}, {"command", map[string]any{"command": "builds"}, false},
{"status", map[string]any{}, false}, {"status", map[string]any{}, false},
{"assign", map[string]any{"node": "anchor", "module": "m"}, false}, {"assign", map[string]any{"node": "anchor", "module": "m"}, false},
+21 -1
View File
@@ -25,7 +25,17 @@ import (
// mesh seat is judged fully only at registration. A seat another module declares is unknown unless // mesh seat is judged fully only at registration. A seat another module declares is unknown unless
// that module's manifest is passed too. Both are printed as a note, not as a problem — a check that // that module's manifest is passed too. Both are printed as a note, not as a problem — a check that
// refused what it could not see would teach people to ignore it. // refused what it could not see would teach people to ignore it.
//
// **And every identity against every bound it meets** (novox/hq ADR 0225, issue 263): each module's
// identity, on a machine whose name is `longestMachine` characters, against the bound of every
// provision it wants that a manifest given here offers. An overflow is refused in the pull request
// that introduces it — a new requirement, a lowered bound, a longer slug — instead of on the
// provider's machine when a real machine's name first meets the module's.
func moduleCheck(paths []string, out io.Writer) error { func moduleCheck(paths []string, out io.Writer) error {
return moduleCheckFor(paths, catalogue.DefaultLongestMachine, out)
}
func moduleCheckFor(paths []string, longestMachine int, out io.Writer) error {
if len(paths) == 0 { if len(paths) == 0 {
return errors.New("module check <manifest.json>... — one file per module; pass every " + return errors.New("module check <manifest.json>... — one file per module; pass every " +
"manifest of a repository together so the rules between them are checked too") "manifest of a repository together so the rules between them are checked too")
@@ -74,6 +84,15 @@ func moduleCheck(paths []string, out io.Writer) error {
} }
failed += len(problems) failed += len(problems)
// Between the manifests too: an identity against the bounds of the provisions it wants, which
// only the provider's manifest states.
identities := catalogue.IdentityProblems(shelf, longestMachine)
sort.Strings(identities)
for _, p := range identities {
fmt.Fprintln(out, p)
}
failed += len(identities)
var names []string var names []string
for name := range shelf { for name := range shelf {
names = append(names, name) names = append(names, name)
@@ -109,7 +128,8 @@ func moduleCheck(paths []string, out io.Writer) error {
} }
fmt.Fprintf(out, "%d manifest(s) checked. Judged against the seats this binary carries; a claim on "+ fmt.Fprintf(out, "%d manifest(s) checked. Judged against the seats this binary carries; a claim on "+
"one of the mesh's own seats is judged fully at registration, and a seat declared by a "+ "one of the mesh's own seats is judged fully at registration, and a seat declared by a "+
"module not given here reads as unknown\n", len(paths)) "module not given here reads as unknown. Identities judged on a %d-character machine name, "+
"against the bounds of the providers given here\n", len(paths), longestMachine)
return nil return nil
} }
+170
View File
@@ -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):
}
}
}
+197
View File
@@ -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 <text> (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 <what was done> --why <text> [--cause <word>] [--condition <key>]")
}
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 <word>, 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 <what> --why <text> --cause <word> | 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, ""
}
+94
View File
@@ -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)
}
}
+150
View File
@@ -0,0 +1,150 @@
package main
import (
"bytes"
"encoding/json"
"os"
"path/filepath"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// novox/hq ADR 0225, issue 263: a consumer's identity is bounded by the provision it requires, an
// overflow is refused before merge by `module check`, and a provider's machine is never refused for
// one consumer's identity.
// `module check` refuses the pull request that introduces an overflow, naming the module.
func TestModuleCheckRefusesAnIdentityThatOverflowsWhatItRequires(t *testing.T) {
dir := t.TempDir()
write := func(name, body string) string {
p := filepath.Join(dir, name+".json")
if err := os.WriteFile(p, []byte(body), 0o600); err != nil {
t.Fatal(err)
}
return p
}
objects := write("objects", `{"module":"objects","version":"1",
"provides":[{"name":"s3-bucket","scope":"mesh","identity":{"max":20,"in":"an S3 access key"}}],
"receives":{"s3-bucket":"/var/lib/mesh/objects/mesh.json"}}`)
resolver := write("resolver", `{"module":"resolver","version":"1",
"provides":[{"name":"wildcard-resolution","scope":"mesh","identity":false}]}`)
album := write("photoalbum", `{"module":"photoalbum","version":"1","requires":["s3-bucket"]}`)
nm := write("networkmanager", `{"module":"networkmanager","version":"1","requires":["wildcard-resolution"]}`)
var out bytes.Buffer
if err := moduleCheckFor([]string{resolver, nm}, 6, &out); err != nil {
t.Fatalf("a long name requiring a keyless provision was refused (issue 263): %v\n%s", err, out.String())
}
out.Reset()
err := moduleCheckFor([]string{objects, album, resolver, nm}, 6, &out)
if err == nil {
t.Fatalf("an identity overflowing an S3 access key passed:\n%s", out.String())
}
if !strings.Contains(out.String(), "photoalbum wants s3-bucket") ||
!strings.Contains(out.String(), "`slug` of at most 8 characters") ||
strings.Contains(out.String(), "networkmanager wants") {
t.Fatalf("the refusal does not name the one overflowing module and its remedy:\n%s", out.String())
}
}
// Tonight's case, through the commands: networkmanager on a six-character machine requires the
// resolver provision, and a second consumer there overflows an object store's access key. The
// provider's machine still composes; the overflowing consumer is left out of its grants and named,
// by push and by `status`, and the keyless consumer is granted with its long name.
func TestAnOverflowingConsumerNeverRefusesItsProvidersMachine(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
register(t, open, catalogue.Manifest{Module: "objects", Version: "1",
Provides: []catalogue.Offer{{Name: "s3-bucket", Scope: catalogue.ScopeMesh,
Identity: &catalogue.OfferIdentity{Max: 20, In: "an S3 access key"}}},
Receives: map[string]string{"s3-bucket": "/var/lib/mesh/objects/mesh.json"}})
register(t, open, catalogue.Manifest{Module: "resolver", Version: "1",
Provides: []catalogue.Offer{{Name: "wildcard-resolution", Scope: catalogue.ScopeMesh}}})
register(t, open, catalogue.Manifest{Module: "networkmanager", Version: "1",
Requires: []string{"wildcard-resolution"}})
register(t, open, catalogue.Manifest{Module: "photoalbum", Version: "1", Requires: []string{"s3-bucket"}})
register(t, open, catalogue.Manifest{Module: "files", Version: "1", Requires: []string{"s3-bucket"}})
for _, a := range [][2]string{{"anchor", "objects"}, {"anchor", "resolver"},
{"laptop", "networkmanager"}, {"laptop", "photoalbum"}, {"laptop", "files"}} {
if _, err := assign(ctx, open, a[0], a[1]); err != nil {
t.Fatalf("assign %s %s: %v", a[0], a[1], err)
}
}
// The consumer's machine resolves, and says which of its modules no provider will grant.
consumer, _, err := planFor(ctx, open, "laptop")
if err != nil {
t.Fatal(err)
}
over := consumer.Overflowing()
if len(over) != 1 || over[0].Module != "photoalbum" || over[0].Provision != "s3-bucket" {
t.Fatalf("the consumer's side does not name exactly photoalbum: %+v", over)
}
// The provider's machine composes. Under ADR 0049's one bound this was a refusal naming
// networkmanager, and no push to the provider could go through.
plan, settings, err := planFor(ctx, open, "anchor")
if err != nil {
t.Fatal(err)
}
declared, err := declarationFor(ctx, open, "anchor", plan, settings)
if err != nil {
t.Fatalf("one consumer's identity refused its provider's whole machine: %v", err)
}
if len(declared.withheld) != 1 || declared.withheld[0].Identity != "mesh_laptop_photoalbum" {
t.Fatalf("the overflowing consumer is not the one withheld: %+v", declared.withheld)
}
grants, _, err := grantsFor(ctx, open, "anchor")
if err != nil {
t.Fatal(err)
}
var keyless bool
for _, g := range grants {
keyless = keyless || g.Provision == "wildcard-resolution" && g.From == "networkmanager"
}
if !keyless {
t.Fatalf("networkmanager, 26 characters, is not granted the keyless resolver provision: %+v", grants)
}
var granted []string
for _, c := range declared.Received["objects"]["s3-bucket"] {
granted = append(granted, c.From)
}
if strings.Join(granted, ",") != "files" {
t.Fatalf("the object store grants %v; files and only files fit", granted)
}
said := printed(t, func() error { reportLeftOut("anchor", declared); return nil })
if !strings.Contains(said, `photoalbum on laptop requires s3-bucket from anchor`) ||
!strings.Contains(said, "left out of anchor's grants") {
t.Fatalf("the push does not say whom it leaves out:\n%s", said)
}
// And `status` names it, and does not call the mesh well while it stands.
asked, err := theThreeQuestions(ctx, open)
if err != nil {
t.Fatal(err)
}
if len(asked.overflowing) != 1 || asked.overflowing[0].Module != "photoalbum" {
t.Fatalf("status does not carry the overflow: %+v", asked.overflowing)
}
if asked.well() {
t.Fatal("a mesh with a consumer left out of its grants reads as well")
}
shown := printed(t, func() error { return printStatus(asked) })
if !strings.Contains(shown, "identified too long for a provision they require") ||
!strings.Contains(shown, "mesh_laptop_photoalbum") {
t.Fatalf("status does not say it:\n%s", shown)
}
body, err := statusAsJSON(asked)
if err != nil {
t.Fatal(err)
}
var doc struct {
Overflowing []catalogue.Overflow `json:"overflowing"`
}
if err := json.Unmarshal(body, &doc); err != nil || len(doc.Overflowing) != 1 ||
doc.Overflowing[0].Bound.Max != 20 {
t.Fatalf("the document does not carry it: %v\n%s", err, body)
}
}
+22 -2
View File
@@ -145,6 +145,14 @@ func run() error {
return seatCommand(ctx, args[1:]) return seatCommand(ctx, args[1:])
case "status": case "status":
return statusCommand(ctx, args[1:]) 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": case "version":
fmt.Println(version) fmt.Println(version)
return nil return nil
@@ -226,19 +234,26 @@ func usage() {
kill <id> end a build where it runs; recorded failed, killed by hand kill <id> end a build where it runs; recorded failed, killed by hand
pause [<node>] / resume [<node>] the build seat's holder there, or every holder, takes nothing new / again pause [<node>] / resume [<node>] the build seat's holder there, or every holder, takes nothing new / again
plans retry <id> ask a failed plan's failed builds again, and carry the plan on plans retry <id> ask a failed plan's failed builds again, and carry the plan on
plans stop|close <id> --why <text> end a plan by hand; recorded in the hand-act log
hand-act record <what> --why <text> --cause <word> [--condition <key>]
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 collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go
builder issue <name> a broker account for a build machine, scoped to build work, builder issue <name> a broker account for a build machine, scoped to build work,
delivered as the builder module's broker secret (module add it first) delivered as the builder module's broker secret (module add it first)
licence add|list|use|key model access, under the name a person calls it licence add|list|use|key model access, under the name a person calls it
licence manager <name> <node> the node that holds a refreshable licence's refresh token licence manager <name> <node> the node that holds a refreshable licence's refresh token
licence refresh <name> mint a new access token and seal it to every holder licence refresh <name> mint a new access token and seal it to every holder
rotate <provision> [--consumer <n>] a new credential for every holder, both ends at once rotate <provision> [--consumer <n>] [--module <m>] a new credential for every holder, both ends at once
ask <module> <tool> [json] call one of a module's tools over the broker, and print its answer ask <module> <tool> [json] call one of a module's tools over the broker, and print its answer
pin <node> <provision> <from-node> <module> pin <node> <provision> <from-node> <module>
which provider this one gets a provision from: the module, and its node which provider this one gets a provision from: the module, and its node
unpin <node> <provision> put that question back unpin <node> <provision> put that question back
plan <node> [--files|--json] what that node would run, and why plan <node> [--files|--json] what that node would run, and why
push [<node>] [--behind] send a node everything it should be, or only those that need it push [<node>] [--behind] [--why <text>] 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 version what this binary is
Each context reaches its own store through its own credential (novox/hq ADR 0008), named 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 // 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. // (novox/hq 04-ISSUES/219). Zero when the id does not say.
asked, _ := link.BuildAskedAt(result.ID) 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 { switch {
case err != nil && result.Failed != "": case err != nil && result.Failed != "":
fmt.Printf("%s: %v\n", result.ID, err) 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)) result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit))
saysWhenThePolicyActs(ctx, b.inv, manifest.Module) saysWhenThePolicyActs(ctx, b.inv, manifest.Module)
planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked, result.ID) 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 return nil
} }
+13 -2
View File
@@ -59,8 +59,19 @@ func moduleCommand(ctx context.Context, args []string) error {
// `check` needs no mesh, and must not: it is what somebody runs in their own repository before // `check` needs no mesh, and must not: it is what somebody runs in their own repository before
// there is a mesh in reach (novox/hq issue 148). A directory expands to every manifest under it. // there is a mesh in reach (novox/hq issue 148). A directory expands to every manifest under it.
if args[0] == "check" { if args[0] == "check" {
set := flag.NewFlagSet("module check", flag.ContinueOnError)
// The longest machine name an identity must fit on (novox/hq ADR 0225): a mesh passes its own.
longest := set.Int("longest-machine-name", catalogue.DefaultLongestMachine,
"judge each module's identity on a machine name this many characters long")
given, err := parseAround(set, args[1:])
if err != nil {
return err
}
if *longest < 1 {
return errors.New("--longest-machine-name is a length, at least 1")
}
var paths []string var paths []string
for _, a := range args[1:] { for _, a := range given {
if info, err := os.Stat(a); err == nil && info.IsDir() { if info, err := os.Stat(a); err == nil && info.IsDir() {
under, err := manifestsUnder(a) under, err := manifestsUnder(a)
if err != nil { if err != nil {
@@ -71,7 +82,7 @@ func moduleCommand(ctx context.Context, args []string) error {
} }
paths = append(paths, a) paths = append(paths, a)
} }
return moduleCheck(paths, os.Stdout) return moduleCheckFor(paths, *longest, os.Stdout)
} }
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
+15 -4
View File
@@ -387,7 +387,7 @@ func brokerCommand(ctx context.Context, args []string) error {
return busAccounts(ctx, args[1:]) return busAccounts(ctx, args[1:])
} }
if len(args) > 0 && args[0] == "consumer-reset" { if len(args) > 0 && args[0] == "consumer-reset" {
return consumerReset(args[1:]) return consumerReset(ctx, args[1:])
} }
if len(args) == 0 || args[0] != "show" { if len(args) == 0 || args[0] != "show" {
return errors.New("broker show | broker certificate [--check] --into <directory> | broker accounts --into <file> | " + return errors.New("broker show | broker certificate [--check] --into <directory> | broker accounts --into <file> | " +
@@ -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 // 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 // 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. // 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 { func consumerReset(ctx context.Context, args []string) error {
if len(args) != 2 { set := flag.NewFlagSet("broker consumer-reset", flag.ContinueOnError)
return errors.New("broker consumer-reset <stream> <consumer>, e.g. broker consumer-reset EVENTS controller") // 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 <stream> <consumer> --why <text>, 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() address, err := broker.BusAddress()
if err != nil { if err != nil {
return err return err
+50 -16
View File
@@ -399,7 +399,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
} }
out := sendable{Resources: composed.Resources, Adoption: adoption, out := sendable{Resources: composed.Resources, Adoption: adoption,
Received: composed.Received, Mesh: with.Mesh, BusUsers: with.BusUsers, Received: composed.Received, Mesh: with.Mesh, BusUsers: with.BusUsers,
LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut} LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut, withheld: with.Withheld}
// And which build of each module it carries, for the send to record (novox/hq issue 259, ADR // And which build of each module it carries, for the send to record (novox/hq issue 259, ADR
// 0221). Read only on the send path: a question about what would be sent records nothing. // 0221). Read only on the send path: a question about what would be sent records nothing.
if choosing == Allocating { if choosing == Allocating {
@@ -464,6 +464,11 @@ func reportLeftOut(node string, declared sendable) {
"what the machine holds for it is kept and its containers are untouched. %s\n", "what the machine holds for it is kept and its containers are untouched. %s\n",
node, m, declared.leftOutWhy[m]) node, m, declared.leftOutWhy[m])
} }
// And whom it serves nothing, because their identity overflows what the provision keeps (ADR
// 0225): the machine is sent everything else, and the consumer is named.
for _, o := range declared.withheld {
fmt.Printf("%s: %s\n", node, o)
}
} }
// renderingFor is everything a node's declaration is composed with, and the node's record. // renderingFor is everything a node's declaration is composed with, and the node's record.
@@ -471,7 +476,7 @@ func renderingFor(ctx context.Context, open *stores, node string,
plan catalogue.Resolution, settings catalogue.SettingsBy, plan catalogue.Resolution, settings catalogue.SettingsBy,
gens map[string]catalogue.Generator, choosing Choosing) (catalogue.Rendering, inventory.Node, error) { gens map[string]catalogue.Generator, choosing Choosing) (catalogue.Rendering, inventory.Node, error) {
inv := open.inventory inv := open.inventory
grants, err := grantsFor(ctx, open, node) grants, withheld, err := grantsFor(ctx, open, node)
if err != nil { if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err return catalogue.Rendering{}, inventory.Node{}, err
} }
@@ -794,7 +799,7 @@ func renderingFor(ctx context.Context, open *stores, node string,
Suffix: overlay.Suffix(), MeshRange: meshRange, TunnelInterface: overlay.Interface, Accounts: accounts, Foundation: foundation, Suffix: overlay.Suffix(), MeshRange: meshRange, TunnelInterface: overlay.Interface, Accounts: accounts, Foundation: foundation,
Kept: kept, Adopted: record.Adopted, OutwardLinks: outwardLinks, Kept: kept, Adopted: record.Adopted, OutwardLinks: outwardLinks,
Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, SeatReach: reach, Built: built, Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, SeatReach: reach, Built: built,
BusUsers: busUsers, BusUsers: busUsers, Withheld: withheld,
}, record, nil }, record, nil
} }
@@ -930,28 +935,34 @@ func certificateFor(ctx context.Context, open *stores, node string) (string, str
// The mirror of what a consumer is given, and the half that makes the credential real: a password // The mirror of what a consumer is given, and the half that makes the credential real: a password
// nothing was told to create is a password that authenticates nowhere. Sealed to this node, so // nothing was told to create is a password that authenticates nowhere. Sealed to this node, so
// the mesh hands over something it cannot itself use. // the mesh hands over something it cannot itself use.
func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Grant, error) { //
// **One consumer's identity never refuses the provider's machine** (novox/hq ADR 0225, issue 263).
// A consumer whose identity overflows the provision's bound is left out of the grants and returned
// beside them, for push, plan and `status` to say; every other consumer is granted and the provider's
// declaration composes. Refusing here once made a whole machine unpushable for one module elsewhere.
func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Grant, []catalogue.Overflow, error) {
inv := open.inventory inv := open.inventory
issued, err := inv.SecretsFrom(ctx, node) issued, err := inv.SecretsFrom(ctx, node)
if err != nil { if err != nil {
return nil, err return nil, nil, err
} }
// Where each consumer is, so a provider that must reach back to one does not have to know how // Where each consumer is, so a provider that must reach back to one does not have to know how
// the mesh names machines. // the mesh names machines.
shelf, err := inv.Catalogue(ctx) shelf, err := inv.Catalogue(ctx)
if err != nil { if err != nil {
return nil, err return nil, nil, err
} }
onNetwork, err := whereEveryoneIs(ctx, inv, shelf) onNetwork, err := whereEveryoneIs(ctx, inv, shelf)
if err != nil { if err != nil {
return nil, err return nil, nil, err
} }
// What each consumer actually asked for, taken from that machine's own resolution rather than // What each consumer actually asked for, taken from that machine's own resolution rather than
// from a record beside it. A provider told to create a password and not what to create it for // from a record beside it. A provider told to create a password and not what to create it for
// can do nothing with it, and the name a consumer wants is the consumer's to say. // can do nothing with it, and the name a consumer wants is the consumer's to say.
out := make([]catalogue.Grant, 0, len(issued)) out := make([]catalogue.Grant, 0, len(issued))
var withheld []catalogue.Overflow
for _, s := range issued { for _, s := range issued {
plan, settings, err := planFor(ctx, open, s.Consumer) plan, settings, err := planFor(ctx, open, s.Consumer)
switch { switch {
@@ -964,11 +975,11 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran
// The mesh could not be asked what they wanted, which is not the same as their wanting // The mesh could not be asked what they wanted, which is not the same as their wanting
// nothing — and withholding a grant on that reading takes a consumer's access away // nothing — and withholding a grant on that reading takes a consumer's access away
// (novox/hq 04-ISSUES/152). // (novox/hq 04-ISSUES/152).
return nil, fmt.Errorf("what %s asked of %s cannot be read: %w", s.Consumer, s.Name, err) return nil, nil, fmt.Errorf("what %s asked of %s cannot be read: %w", s.Consumer, s.Name, err)
} }
values, asks, err := plan.ContributionsFrom(s.Name, s.ConsumerModule, settings) values, asks, err := plan.ContributionsFrom(s.Name, s.ConsumerModule, settings)
if err != nil { if err != nil {
return nil, err return nil, nil, err
} }
// A port in there is the consumer's software port until this. The consumer is on another // A port in there is the consumer's software port until this. The consumer is on another
// machine, so the assignment that moved it is that machine's — fetched here rather than // machine, so the assignment that moved it is that machine's — fetched here rather than
@@ -976,7 +987,7 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran
// this case (novox/hq 04-ISSUES/038, the cross-node half). // this case (novox/hq 04-ISSUES/038, the cross-node half).
published, err := portsOn(ctx, inv, s.Consumer, s.ConsumerModule) published, err := portsOn(ctx, inv, s.Consumer, s.ConsumerModule)
if err != nil { if err != nil {
return nil, err return nil, nil, err
} }
values = catalogue.AtPublishedPort(values, s.ConsumerModule, published) values = catalogue.AtPublishedPort(values, s.ConsumerModule, published)
from := s.ConsumerModule from := s.ConsumerModule
@@ -987,9 +998,10 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran
from = "" from = ""
} }
// The consumer's identity slug, from its own manifest, carried on the grant so the provider // The consumer's identity slug, from its own manifest, carried on the grant so the provider
// derives the same login the consumer does (novox/hq ADR 0049). Refused here if it still would // derives the same login the consumer does (novox/hq ADR 0049). Judged against the bound of
// not fit the tightest backend — the mesh chose the name, so the mesh refuses it, with the // this provision, as the consumer's resolution states it from the provider's offer (ADR
// remedy a short slug rather than a login a provider silently shortened. // 0225): a consumer it would not fit is left out of the grants and said, rather than a login a
// provider silently shortened — and rather than this whole machine refused for it.
slug := "" slug := ""
for _, mm := range plan.Modules { for _, mm := range plan.Modules {
if mm.Module == s.ConsumerModule { if mm.Module == s.ConsumerModule {
@@ -998,15 +1010,32 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran
} }
} }
if from != "" { if from != "" {
if err := catalogue.CheckIdentity(s.Consumer, catalogue.IdentitySource(slug, s.ConsumerModule)); err != nil { bound := boundOfGrant(plan, s, node)
return nil, err source := catalogue.IdentitySource(slug, s.ConsumerModule)
if catalogue.CheckIdentityWithin(s.Consumer, source, bound) != nil {
withheld = append(withheld, catalogue.Overflow{Provision: s.Name, Provider: node,
Consumer: s.Consumer, Module: s.ConsumerModule,
Identity: catalogue.ConsumerIdentity(s.Consumer, source), Bound: bound})
continue
} }
} }
out = append(out, catalogue.Grant{ out = append(out, catalogue.Grant{
Provision: s.Name, Consumer: s.Consumer, At: onNetwork[s.Consumer], Provision: s.Name, Consumer: s.Consumer, At: onNetwork[s.Consumer],
From: from, Values: values, Slug: slug, Sealed: s.ForProvider, Local: s.Local}) From: from, Values: values, Slug: slug, Sealed: s.ForProvider, Local: s.Local})
} }
return out, nil return out, withheld, nil
}
// boundOfGrant is the identity bound the consumer's own resolution states for the requirement this
// grant answers. A requirement not found there is held to the tightest bound the mesh knows rather
// than to none: what the provider keeps of it is not known here.
func boundOfGrant(consumer catalogue.Resolution, s inventory.Secret, provider string) catalogue.IdentityBound {
for _, n := range consumer.Needs {
if n.Name == s.Name && n.For == s.ConsumerModule && n.From == provider {
return n.Identity
}
}
return catalogue.DefaultIdentityBound
} }
// listensLines is what a person is told about what this module would open, and why — the same // listensLines is what a person is told about what this module would open, and why — the same
@@ -1082,6 +1111,11 @@ func planCommand(ctx context.Context, args []string) error {
left := plan.LeftOut(settings, record.Adopted) left := plan.LeftOut(settings, record.Adopted)
reportLeftOut(args[0], sendable{LeftOut: sortedKeysOf(left), leftOutWhy: left}) reportLeftOut(args[0], sendable{LeftOut: sortedKeysOf(left), leftOutWhy: left})
} }
// And which of its modules no provider will grant, because the identity overflows the bound of
// what it requires (novox/hq ADR 0225) — said on the machine the remedy is for.
for _, o := range plan.Overflowing() {
fmt.Printf("%s: %s\n", args[0], o)
}
// And a setting that reaches nothing — refused where it is stored, and said here for one // And a setting that reaches nothing — refused where it is stored, and said here for one
// stored before its definition moved from under it. // stored before its definition moved from under it.
for _, m := range plan.Modules { for _, m := range plan.Modules {
+60 -4
View File
@@ -105,7 +105,10 @@ func serve(ctx context.Context) error {
work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known, work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known,
OnNATS: true} 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 { if err != nil {
return err 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 // 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. // machines to report moves when they have, and a plan left by a replaced controller resumes.
go planTicker(ctx, open) 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 // 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 // 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). // 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 { if !isNATS {
return errors.New("the mesh's verbs are served over the bus, and this control plane is not on it") 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). // A call that outlasts its caller's patience is followed by `calls` (novox/hq issue 265).
link.Calls.Follow = catalogue.ControllerSeatName + ".calls" 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)) stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags))
if err != nil { if err != nil {
return err 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. // (novox/hq ADR 0010). 0 waits for nothing, which is the old fire-and-forget.
wait := set.Duration("wait", 0, 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)") "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) positionals, err := parseAround(set, args)
if err != nil { if err != nil {
return err return err
@@ -274,6 +295,11 @@ func pushCommand(ctx context.Context, args []string) error {
return errors.New("push <node> or push --behind, not both: one names a machine and the " + return errors.New("push <node> or push --behind, not both: one names a machine and the " +
"other asks which machines need one") "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) open, err := openStores(ctx)
if err != nil { if err != nil {
return err 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. // 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, func wouldSend(ctx context.Context, open *stores,
nodes []inventory.Node) (map[string]string, error) { 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) gens, err := generators(ctx, open)
if err != nil { if err != nil {
@@ -1046,9 +1086,12 @@ func wouldSend(ctx context.Context, open *stores,
} }
out := map[string]string{} out := map[string]string{}
for _, n := range nodes { for _, n := range nodes {
plan, settings, err := planFor(ctx, open, n.Name) known, have := plans[n.Name]
if err != nil { plan, settings := known.plan, known.settings
continue 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) declared, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
if err != nil { 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 { if err := broker.RaiseCancelledSets(js, inventory.MeshSeats()); err != nil {
return err 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 // 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 // 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. // 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)) 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))
}
+10
View File
@@ -78,11 +78,19 @@ type meshStatus struct {
// that machine holds, with the modules that could hold it (novox/hq ADR 0207). Absent when every // 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. // dependency is met. Reported, not refused, until the switch.
Unheld []catalogue.Unheld `json:"unheld,omitempty"` 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 // 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 // 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 // document without this called the mesh well while the identity provider refused every consumer
// for a day (04-ISSUES/179). // for a day (04-ISSUES/179).
Failing []inventory.ProviderStanding `json:"failing,omitempty"` Failing []inventory.ProviderStanding `json:"failing,omitempty"`
// Overflowing is every module whose identity overflows the bound of a provision it requires, and
// so is left out of its provider's grants (novox/hq ADR 0225). Absent when every identity fits.
Overflowing []catalogue.Overflow `json:"overflowing,omitempty"`
} }
// machineFiltered is one rule set on a converged machine that the mesh did not write and that // machineFiltered is one rule set on a converged machine that the mesh did not write and that
@@ -215,7 +223,9 @@ func statusAsJSON(asked answers) ([]byte, error) {
} }
} }
out.Unheld = asked.unheld out.Unheld = asked.unheld
out.HandActsThisWeek, out.HandActsUnread = asked.handActs, asked.handActsUnread
out.Failing = asked.failing out.Failing = asked.failing
out.Overflowing = asked.overflowing
for name := range asked.refused { for name := range asked.refused {
out.Unresolved = append(out.Unresolved, machineUnresolved{ out.Unresolved = append(out.Unresolved, machineUnresolved{
Node: name, Problem: asked.refused[name]}) Node: name, Problem: asked.refused[name]})
+10 -1
View File
@@ -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") 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") 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") 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) positionals, err := parseAround(set, args)
if err != nil { if err != nil {
return err 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) open, err := openStores(ctx)
if err != nil { if err != nil {
return err return err
@@ -1061,8 +1069,9 @@ func plansCommand(ctx context.Context, args []string) error {
if !p.Open() { if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State) return fmt.Errorf("%s is already %s", p.ID, p.State)
} }
why.record(ctx, "plans "+positionals[0], positionals[1:])
p.State = inventory.PlanFailed 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 { sayUnsent(&p, func(m string) bool {
u, err := inv.UpgradeOf(ctx, m) u, err := inv.UpgradeOf(ctx, m)
return err == nil && u.RollOut return err == nil && u.RollOut
+34 -1
View File
@@ -6,6 +6,8 @@ import (
"flag" "flag"
"fmt" "fmt"
"sort" "sort"
"github.com/novox/mesh-controller/internal/inventory"
) )
// rotateCommand replaces a credential and moves both ends together. // rotateCommand replaces a credential and moves both ends together.
@@ -33,12 +35,16 @@ func rotateCommand(ctx context.Context, args []string) error {
// One consumer rather than all of them. Ordinary: a credential is suspected on one machine, // One consumer rather than all of them. Ordinary: a credential is suspected on one machine,
// and rotating the other nine would be a great deal of disruption for one suspicion. // and rotating the other nine would be a great deal of disruption for one suspicion.
only := set.String("consumer", "", "only this machine's credential, rather than every holder's") only := set.String("consumer", "", "only this machine's credential, rather than every holder's")
// One consuming module rather than every module on the machine. A machine runs many consumers
// of one provision, each with its own credential; one module that leaked its credential (novox/hq
// issue 268) is no reason to restart every other one on the machine.
module := set.String("module", "", "only this consuming module's credential")
positionals, err := parseAround(set, args) positionals, err := parseAround(set, args)
if err != nil { if err != nil {
return err return err
} }
if len(positionals) != 1 { if len(positionals) != 1 {
return errors.New("rotate <provision> [--consumer <machine>]") return errors.New("rotate <provision> [--consumer <machine>] [--module <module>]")
} }
provision := positionals[0] provision := positionals[0]
@@ -53,6 +59,12 @@ func rotateCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err return err
} }
holders = ofModule(holders, *module)
if len(holders) == 0 && *module != "" {
return fmt.Errorf(
"no module %s%s holds a credential for %q, so there is nothing to rotate. `plan <machine>` "+
"says what a machine holds", *module, onMachine(*only), provision)
}
if len(holders) == 0 { if len(holders) == 0 {
// Said, not silent. "Nobody holds this" and "this did not run" must never look the same — // Said, not silent. "Nobody holds this" and "this did not run" must never look the same —
// and a rotation somebody believes happened is worse than one they know did not. // and a rotation somebody believes happened is worse than one they know did not.
@@ -116,6 +128,27 @@ func rotateCommand(ctx context.Context, args []string) error {
return nil return nil
} }
// ofModule is the holders whose consuming module is this one; all of them when none is named.
func ofModule(holders []inventory.Holder, module string) []inventory.Holder {
if module == "" {
return holders
}
var out []inventory.Holder
for _, h := range holders {
if h.ConsumerModule == module {
out = append(out, h)
}
}
return out
}
func onMachine(machine string) string {
if machine == "" {
return ""
}
return " on " + machine
}
// asLocal names the credential inside the consumer where it holds several (ADR 0094). // asLocal names the credential inside the consumer where it holds several (ADR 0094).
func asLocal(local string) string { func asLocal(local string) string {
if local == "" { if local == "" {
+27
View File
@@ -0,0 +1,27 @@
package main
import (
"testing"
"github.com/novox/mesh-controller/internal/inventory"
)
// One consuming module's credential, and not its neighbours' on the same machine (novox/hq issue
// 268): a module that leaked its database password is no reason to restart every other consumer.
func TestARotationNarrowedToAModuleTouchesOnlyThatModulesCredential(t *testing.T) {
holders := []inventory.Holder{
{Provision: "postgres-database", Consumer: "ace", ConsumerModule: "letta", Provider: "ace"},
{Provision: "postgres-database", Consumer: "ace", ConsumerModule: "n8n", Provider: "ace"},
{Provision: "postgres-database", Consumer: "ace", ConsumerModule: "letta", Local: "reader", Provider: "ace"},
}
got := ofModule(holders, "letta")
if len(got) != 2 || got[0].ConsumerModule != "letta" || got[1].Local != "reader" {
t.Fatalf("narrowed to letta: %+v", got)
}
if len(ofModule(holders, "")) != 3 {
t.Fatal("no module named narrowed anyway")
}
if len(ofModule(holders, "absent")) != 0 {
t.Fatal("a module holding nothing matched")
}
}
+129 -9
View File
@@ -9,9 +9,11 @@ import (
"github.com/nats-io/nats.go/micro" "github.com/nats-io/nats.go/micro"
"os" "os"
"os/exec" "os/exec"
"slices"
"sort" "sort"
"strings" "strings"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/link"
) )
@@ -239,6 +241,12 @@ func (a *verbArguments) commandLine() ([]string, error) {
if len(argv) == 0 { if len(argv) == 0 {
return nil, errors.New("command names no command") 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 <text> to the command "+
"line (recorded in the hand-act log). Nothing was done", repair)
}
return argv, nil return argv, nil
case "tools": case "tools":
return nil, errors.New("tools is answered from the records, not by a command") 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"} { for _, act := range []string{"stop", "close", "retry"} {
if id := str(act); id != "" { 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 != "" { if id := str("id"); id != "" {
@@ -355,22 +374,59 @@ func (a *verbArguments) commandLine() ([]string, error) {
// Sent and not waited for: the asker reads `status` for what the machine did, which is // 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 // 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. // 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 != "" { if n := str("node"); n != "" {
// behind is not read here: given with a machine, it is refused as passed over — naming // 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 // 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. // 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 // 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. // first, so a caller who meant one machine reads that it was not one.
on("behind") 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": case "rotate":
if p := str("provision"); p != "" { if p := str("provision"); p != "" {
argv := []string{"rotate", p} argv := []string{"rotate", p}
if c := str("consumer"); c != "" { if c := str("consumer"); c != "" {
argv = append(argv, "--consumer", c) argv = append(argv, "--consumer", c)
} }
// With a provision, module narrows to one consuming module (novox/hq issue 268); node
// and secret stay the other shape's, and are refused as passed over.
if m := str("module"); m != "" {
argv = append(argv, "--module", m)
}
return argv, nil return argv, nil
} }
_, node := a.given["node"] _, node := a.given["node"]
@@ -437,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. // 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. // runVerb runs this binary with the given command line and gathers what it said.
func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { func runVerb(ctx context.Context, argv []string) (verbAnswer, error) {
@@ -449,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 // 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. // from a shell in this container would have, because it is that.
cmd.Env = os.Environ() 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 // 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` // 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 // prints its warnings beside the document, and a JSON parsed from the two together parsed
@@ -539,6 +622,14 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
if err != nil { if err != nil {
return nil, err 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) { if answersFirst(argv) {
// Before anything is sent: a push sends the bus's own machine first, and a broker // 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). // reloading its user list forgets the answer it was about to permit (novox/hq issue 265).
@@ -550,6 +641,16 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
return handlers, behind, nil 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 // 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. // the records, `calls` from what this process served.
var inProcess = map[string]bool{"tools": true, "calls": true} var inProcess = map[string]bool{"tools": true, "calls": true}
@@ -563,22 +664,41 @@ func answersFirst(argv []string) bool {
} }
// callsAnswer is what `calls` answers: the kept calls, newest first, without their answers — or // 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) { func callsAnswer(log *link.CallLog, id string) (any, error) {
if id != "" { 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 !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 "+ 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) "answered them, the last %d, and not across a restart — `calls` lists them", id, link.KeptCalls)
} }
return c, nil return c, nil
} }
recent := log.Recent() recent, err := log.Recent()
for i := range recent { for i := range recent {
recent[i].Answer = nil recent[i].Answer = nil
} }
return map[string]any{"calls": recent, "kept": link.KeptCalls, answer := map[string]any{"calls": recent, "note": "newest first; `calls` with a call's id gives its whole answer"}
"note": "newest first; `calls` with a call's id gives its whole answer"}, nil 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 // seatTools is what `tools` answers: every seat with a protocol, and the tools each serves, from the
+10 -5
View File
@@ -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 // 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. // naming one beside behind is refused rather than one of the two guessed.
func TestAPushIsOneMachineOrSaysItIsTheWholeMesh(t *testing.T) { func TestAPushIsOneMachineOrSaysItIsTheWholeMesh(t *testing.T) {
argv, err := argvFor("push", map[string]any{"node": "g1"}) argv, err := argvFor("push", map[string]any{"node": "g1", "why": "w"})
if err != nil || strings.Join(argv, " ") != "push g1 --wait 0" { if err != nil || strings.Join(argv, " ") != "push g1 --wait 0 --why w" {
t.Fatalf("a named push: %v %v", argv, err) 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) 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) 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"`) { !strings.Contains(err.Error(), `"behind"`) {
t.Fatalf("a named push with behind was taken: %v", err) 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"}, "builds": {"n": "=limit"},
"plans": {"n": "=limit", "what-if": "=repository"}, "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 // **Every flag of the command a verb runs is in the verb's schema, or accounted for here.** Derived
+6 -2
View File
@@ -43,6 +43,10 @@ func TestRotateTakesAProvisionOrAnOwnSecret(t *testing.T) {
if strings.Join(argv, " ") != "rotate postgres-database --consumer ace" { if strings.Join(argv, " ") != "rotate postgres-database --consumer ace" {
t.Fatalf("a pair credential: %v", argv) t.Fatalf("a pair credential: %v", argv)
} }
argv, _ = argvFor("rotate", map[string]any{"provision": "postgres-database", "consumer": "ace", "module": "letta"})
if strings.Join(argv, " ") != "rotate postgres-database --consumer ace --module letta" {
t.Fatalf("one consuming module's pair credential: %v", argv)
}
argv, _ = argvFor("rotate", map[string]any{"node": "ace", "module": "nodered", "secret": "api-token"}) argv, _ = argvFor("rotate", map[string]any{"node": "ace", "module": "nodered", "secret": "api-token"})
if strings.Join(argv, " ") != "secret rotate ace nodered api-token" { if strings.Join(argv, " ") != "secret rotate ace nodered api-token" {
t.Fatalf("an own secret: %v", argv) t.Fatalf("an own secret: %v", argv)
@@ -103,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 // 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). // id, for what happened. A repository given as a forge path is said to be one (issue 176).
func TestActsDoNotBlockTheCall(t *testing.T) { func TestActsDoNotBlockTheCall(t *testing.T) {
argv, _ := argvFor("push", map[string]any{"node": "one"}) argv, _ := argvFor("push", map[string]any{"node": "one", "why": "w"})
if strings.Join(argv, " ") != "push one --wait 0" { if strings.Join(argv, " ") != "push one --wait 0 --why w" {
t.Fatalf("push waits: %v", argv) t.Fatalf("push waits: %v", argv)
} }
argv, _ = argvFor("build", map[string]any{"repository": "novox/x", "path": "modules/x"}) argv, _ = argvFor("build", map[string]any{"repository": "novox/x", "path": "modules/x"})
+3
View File
@@ -43,6 +43,9 @@ type sendable struct {
LeftOut []string LeftOut []string
// leftOutWhy is why each was, for push and plan to say; never on the wire. // leftOutWhy is why each was, for push and plan to say; never on the wire.
leftOutWhy map[string]string leftOutWhy map[string]string
// withheld is every consumer this machine's grants leave out, because its identity overflows the
// provision's bound (novox/hq ADR 0225); for push and plan to say, never on the wire.
withheld []catalogue.Overflow
// Builds is the build of each module this declaration carries — module to the commit its build // Builds is the build of each module this declaration carries — module to the commit its build
// was made from — recorded with the send and never on the wire (novox/hq issue 259, ADR 0221). // was made from — recorded with the send and never on the wire (novox/hq issue 259, ADR 0221).
// Composed only on the send path; nil records that it is not known. // Composed only on the send path; nil records that it is not known.
+43 -3
View File
@@ -8,6 +8,7 @@ import (
"strings" "strings"
"time" "time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/overlay" "github.com/novox/mesh-controller/internal/overlay"
) )
@@ -302,6 +303,29 @@ func printStatus(asked answers) error {
fmt.Printf("\n `assign <node> <holder>` meets it; reported until every machine has its holders, then refused\n\n") fmt.Printf("\n `assign <node> <holder>` meets it; reported until every machine has its holders, then refused\n\n")
} }
if len(asked.overflowing) > 0 {
// **Reported, and the provider still pushed** (novox/hq ADR 0225, issue 263). Each is left
// out of its provider's grants, so the module holds a login nothing created; the provider's
// machine is sent everything else rather than refused for one consumer elsewhere.
fmt.Printf("%d module(s) identified too long for a provision they require, and not granted it:\n",
len(asked.overflowing))
for _, o := range asked.overflowing {
fmt.Printf(" %-12s %-20s %-22s %q is %d, %s keeps %d\n", o.Consumer, o.Module, o.Provision,
o.Identity, len(o.Identity), o.Bound.In, o.Bound.Max)
}
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 { if adopted := adoptedNodes(nodes); len(adopted) > 0 {
// Said, because nothing forces the flip: a node left adopted is visible here rather than // 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". // read as converged (novox/hq ADR 0100). Not a fault, so it does not break "all well".
@@ -410,15 +434,21 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
// ADR 0207). Each machine resolved again rather than threaded through whoResolves, whose answer // 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 // 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. // resolve is already in refused, and is passed over here.
plans := map[string]planned{}
for _, n := range out.nodes { for _, n := range out.nodes {
plan, _, err := planFor(ctx, open, n.Name) plan, settings, err := planFor(ctx, open, n.Name)
if err != nil { if err != nil {
if unresolvable(err) { if unresolvable(err) {
continue continue
} }
return answers{}, err return answers{}, err
} }
plans[n.Name] = planned{plan, settings}
out.unheld = append(out.unheld, plan.Unheld...) 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
// resolution, as the provider's composition judges it.
out.overflowing = append(out.overflowing, plan.Overflowing()...)
} }
// And every consumer a provider says it keeps failing (novox/hq ADR 0224). Read from what the // And every consumer a provider says it keeps failing (novox/hq ADR 0224). Read from what the
// providers announced: nothing else in the mesh knows whether a provision is being made. // providers announced: nothing else in the mesh knows whether a provision is being made.
@@ -432,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). // Whether the build seat takes work, for a plan waiting on it (novox/hq ADR 0219).
out.paused = buildSeatPause(ctx, inv, out.plans) 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 // 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, // module being behind its source, one level down: that one says the catalogue is out of date,
@@ -443,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 // 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 // reason is kept and reported as data; every question that does not depend on it is still
// answered. // answered.
would, err := wouldSend(ctx, open, out.nodes) would, err := wouldSendFrom(ctx, open, out.nodes, plans)
if err != nil { if err != nil {
out.network = err.Error() out.network = err.Error()
would = map[string]string{} would = map[string]string{}
@@ -538,7 +578,7 @@ func untakenModules(ctx context.Context, inv *inventory.Inventory, nodes []inven
func (a answers) well() bool { func (a answers) well() bool {
return len(a.wrong) == 0 && len(a.quiet) == 0 && len(a.behind) == 0 && return len(a.wrong) == 0 && len(a.quiet) == 0 && len(a.behind) == 0 &&
len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 && len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 &&
len(a.filtered) == 0 && len(a.unheld) == 0 && len(a.failing) == 0 len(a.filtered) == 0 && len(a.unheld) == 0 && len(a.failing) == 0 && len(a.overflowing) == 0
} }
// hostSplit is which machines report which host version, for every version more than one machine // hostSplit is which machines report which host version, for every version more than one machine
+188
View File
@@ -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
}
+152
View File
@@ -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)
}
}
+9 -5
View File
@@ -75,21 +75,25 @@ func TestAPersonClosesAStuckPlan(t *testing.T) {
if err := open.inventory.SavePlan(ctx, stuck); err != nil { if err := open.inventory.SavePlan(ctx, stuck); err != nil {
t.Fatal(err) 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) t.Fatal(err)
} }
closed, err := open.inventory.PlanByID(ctx, stuck.ID) closed, err := open.inventory.PlanByID(ctx, stuck.ID)
if err != nil { if err != nil {
t.Fatal(err) 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) 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") t.Fatal("a plan already closed was closed again")
} }
if argv, err := argvFor("plans", map[string]any{"close": stuck.ID}); err != nil || if argv, err := argvFor("plans", map[string]any{"close": stuck.ID, "why": "w"}); err != nil ||
!reflect.DeepEqual(argv, []string{"plans", "close", stuck.ID}) { !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) t.Fatalf("the seat's verb does not close a plan: %v %v", argv, err)
} }
} }
+102
View File
@@ -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
}
@@ -0,0 +1,31 @@
package broker
import (
"slices"
"testing"
)
// **The controller may write every bucket it writes** (novox/hq to-be 45 §1, issue 269). Writing a
// key is a publish to the bucket's own subject, which the management interface's grant does not
// cover: the cancelled sets' writes timed out for want of this, and the controller's own buckets
// would have.
func TestTheControllerMayWriteEveryBucketItWrites(t *testing.T) {
p, err := PermissionsFor(Principal{Kind: KindController, PasswordHash: "x"})
if err != nil {
t.Fatal(err)
}
want := []string{"$KV." + CallsBucket + ".>", "$KV." + HandActsBucket + ".>"}
for _, seat := range seatsTheControllerAsks {
if hasCancelledSet(seat) {
want = append(want, "$KV."+CancelledSetName(seat)+".>")
}
}
for _, subject := range want {
if !slices.Contains(p.Publish, subject) {
t.Errorf("the controller may not publish %s, so it cannot write that bucket", subject)
}
}
if slices.Contains(p.Publish, "$KV.>") {
t.Error("the controller may write any bucket, a module's state included")
}
}
+11
View File
@@ -237,6 +237,12 @@ func PermissionsFor(p Principal) (Permissions, error) {
// building, kill it, pause it, resume it. The queue is the controller's to show and to // building, kill it, pause it, resume it. The queue is the controller's to show and to
// change, and what one machine is doing with an ask it took only that machine can say. // change, and what one machine is doing with an ask it took only that machine can say.
pub = append(pub, "mesh.seat."+seat+".tool.>") pub = append(pub, "mesh.seat."+seat+".tool.>")
// **And its cancelled set** (novox/hq ADR 0219, issue 269): a cancel writes the ask's id
// there before it deletes the ask, and a write is a publish to the bucket's subject, which
// `$JS.API.>` does not cover — so every cancel timed out, refused by this list.
if hasCancelledSet(seat) {
pub = append(pub, "$KV."+CancelledSetName(seat)+".>")
}
} }
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own // **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
// facts under the seat it holds, because a role's events belong to the role and keep their // facts under the seat it holds, because a role's events belong to the role and keep their
@@ -288,6 +294,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// enrolments and can reach nothing else. // enrolments and can reach nothing else.
pub = append(pub, "_INBOX."+enrolmentPrefix+".>") 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: case KindPerson:
// Tools, and nothing else. Every subject a person may publish is a tool call; a person // 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. // who could publish an event would be able to claim a module said something.
+3 -2
View File
@@ -160,8 +160,9 @@ func RaiseBuckets(a BucketAsserter, buckets []Bucket) (undeclared []string, err
return nil, fmt.Errorf("listing the bus's state: %w", err) return nil, fmt.Errorf("listing the bus's state: %w", err)
} }
for _, n := range names { for _, n := range names {
// A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state. // A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state; so are
if !declared[n] && !IsCancelledSet(n) { // the controller's own buckets (novox/hq to-be 45 §1).
if !declared[n] && !IsCancelledSet(n) && !IsControllerBucket(n) {
undeclared = append(undeclared, n) undeclared = append(undeclared, n)
} }
} }
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { 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.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$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"] } 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" } allow_responses: { max: 1, ttl: "1m" }
} } } }
@@ -119,3 +119,41 @@ func TestEveryCatalogueStoreSaysHowItIsBackedUp(t *testing.T) {
t.Fatal("no module in the catalogue provides a store, so this proved nothing") t.Fatal("no module in the catalogue provides a store, so this proved nothing")
} }
} }
// TestEveryCatalogueIdentityFitsWhatItRequires is ADR 0225's check over the real catalogue: every
// module's identity, on the longest machine name, fits the bound of every provision it wants. An
// overflow fails here, in the pull request that introduces it, rather than on the provider's machine
// the first time a real machine's name meets the module's (issue 263).
func TestEveryCatalogueIdentityFitsWhatItRequires(t *testing.T) {
root := catalogueRoot(t)
found, err := filepath.Glob(filepath.Join(root, "modules", "*", "module.json"))
if err != nil || len(found) == 0 {
t.Fatalf("no manifests under %s: %v", root, err)
}
shelf := Shelf{}
for _, p := range found {
raw, err := os.ReadFile(p)
if err != nil {
t.Fatalf("%s: %v", p, err)
}
m, err := ParseManifest(raw)
if err != nil {
t.Fatalf("%s: %v", p, err)
}
shelf[m.Module] = m
}
for _, p := range IdentityProblems(shelf, DefaultLongestMachine) {
t.Error(p)
}
// The night it was found: the resolver provision bounds nothing, and the object store still 20.
if dns, ok := shelf["dnsmasq"]; ok {
if b := dns.IdentityBoundOf("wildcard-resolution"); b.Bounded() {
t.Errorf("the resolver provision bounds its consumers' identities: %+v", b)
}
}
if store, ok := shelf["minio"]; ok {
if b := store.IdentityBoundOf("s3-bucket"); b.Max != 20 {
t.Errorf("the object store's access key is not bounded at 20: %+v", b)
}
}
}
+23 -1
View File
@@ -114,7 +114,9 @@ func consumerInto(value, as string) (string, error) {
// asDNSLabel writes a minted identity as a DNS label. // asDNSLabel writes a minted identity as a DNS label.
// //
// The mesh's identities are already lower-case letters, digits and `_` (ConsumerIdentity), and // The mesh's identities are already lower-case letters, digits and `_` (ConsumerIdentity), and
// already short enough for the tightest backend they reach (CheckIdentity, twenty characters). So // already inside the bound of the provision serving them: a provider that serves one as a DNS label
// bounds it at 63 or less (CheckServes, novox/hq ADR 0225), and one that serves it at all without
// saying is held to twenty (IdentityBoundOf). So
// this is the separator and nothing else — no lower-casing of what is already lower case, no // this is the separator and nothing else — no lower-casing of what is already lower case, no
// truncation to a limit the identity is already inside, no padding of a name that is already long // truncation to a limit the identity is already inside, no padding of a name that is already long
// enough. Each of those would be the mesh guessing at a rule it has not been given. // enough. Each of those would be the mesh guessing at a rule it has not been given.
@@ -141,11 +143,31 @@ func CheckServes(m Manifest) []string {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%s serves %s, and the value it serves as %q %s", m.Module, provision, key, err)) "%s serves %s, and the value it serves as %q %s", m.Module, provision, key, err))
} }
// A label longer than DNS keeps is not truncated here (asDNSLabel), so the offer's own
// bound has to keep the identity inside one (ADR 0225).
if strings.Contains(text, "${consumer:as:dns}") {
if b := m.IdentityBoundOf(provision); !b.Bounded() || b.Max > dnsLabelLimit {
problems = append(problems, fmt.Sprintf(
"%s serves %s's consumers their identity as a DNS label in %q, and bounds "+
"that identity at %s: a label keeps %d — state an `identity` of at most %d",
m.Module, provision, key, boundWords(b), dnsLabelLimit, dnsLabelLimit))
}
}
} }
} }
return problems return problems
} }
// dnsLabelLimit is the longest DNS label (RFC 1035 §2.3.4).
const dnsLabelLimit = 63
func boundWords(b IdentityBound) string {
if !b.Bounded() {
return "nothing"
}
return fmt.Sprintf("%d", b.Max)
}
func sortedServes(serves map[string]map[string]any) []string { func sortedServes(serves map[string]map[string]any) []string {
out := make([]string, 0, len(serves)) out := make([]string, 0, len(serves))
for k := range serves { for k := range serves {
+4
View File
@@ -170,6 +170,10 @@ type Rendering struct {
// rather than resolved, because who consumes a node is a fact about the rest of the mesh and // rather than resolved, because who consumes a node is a fact about the rest of the mesh and
// resolution answers questions about one machine. // resolution answers questions about one machine.
Grants []Grant Grants []Grant
// Withheld is every consumer left out of Grants because its identity overflows the provision's
// bound (novox/hq ADR 0225). Composed into nothing; carried so the machine's declaration can say
// whom it does not serve, and why, beside what it does.
Withheld []Overflow
// Ports is where this machine puts what each module needs reachable, by module and by the // Ports is where this machine puts what each module needs reachable, by module and by the
// port the software itself uses (novox/hq ADR 0038). // port the software itself uses (novox/hq ADR 0038).
+148 -9
View File
@@ -3,6 +3,7 @@ package catalogue
import ( import (
"fmt" "fmt"
"regexp" "regexp"
"sort"
"strings" "strings"
) )
@@ -56,13 +57,36 @@ func ConsumerIdentity(node, module string) string {
return IdentityPrefix + clean(node) + "_" + clean(module) return IdentityPrefix + clean(node) + "_" + clean(module)
} }
// identityLimit is the shortest identifier limit among the systems these names reach: an S3 access // IdentityBound is the longest consumer identity one provision's backend keeps, and what keeps it
// key's 20 (novox/hq 04-ISSUES/010). PostgreSQL keeps 63 and MinIO 20, so 20 is the one that binds — // (novox/hq ADR 0049, refined by ADR 0225). Max zero means no bound: the provision keeps no name
// the comment used to name PostgreSQL and was wrong. A name over it is refused, with the remedy a // derived from its consumer, or keeps one in something with no limit the mesh need respect.
// short slug (ADR 0049), not silently cut to fit. type IdentityBound struct {
const identityLimit = 20 // Max is the longest identity that backend keeps, in characters; zero for none.
Max int `json:"max,omitempty"`
// In is what keeps it, in words a refusal can quote: "an S3 access key", "a PostgreSQL role".
In string `json:"in,omitempty"`
}
// CheckIdentity refuses an identity that would not fit the tightest backend a consumer reaches. // Bounded is whether this bound refuses anything.
func (b IdentityBound) Bounded() bool { return b.Max > 0 }
// DefaultIdentityLimit is the bound on a provision whose provider receives its consumers and does
// not say how long a name it keeps: an S3 access key's 20 (novox/hq 04-ISSUES/010, 034), the
// tightest backend the mesh has met. It was the bound on every provision until ADR 0225; it stays
// the bound on any that has not said otherwise, because a provider that is told each consumer's
// identity may create a name from it in a backend nobody has measured.
const DefaultIdentityLimit = 20
// DefaultIdentityBound is DefaultIdentityLimit, said as a bound.
var DefaultIdentityBound = IdentityBound{Max: DefaultIdentityLimit,
In: "a backend that has not said its limit (the tightest known, an S3 access key's)"}
// identityLimit is the bound CheckIdentity applies, for a caller that does not know which provision
// the identity is for.
const identityLimit = DefaultIdentityLimit
// CheckIdentity refuses an identity that would not fit the tightest backend the mesh knows. A caller
// that knows the provision uses CheckIdentityWithin and that provision's own bound (ADR 0225).
// //
// **Truncation is not an error in most of these systems** — a name past the limit is cut to fit and // **Truncation is not an error in most of these systems** — a name past the limit is cut to fit and
// the statement succeeds, so two consumers agreeing for the first N bytes would become one login // the statement succeeds, so two consumers agreeing for the first N bytes would become one login
@@ -70,12 +94,127 @@ const identityLimit = 20
// name and is the only thing that can choose another. The remedy is a first-class one: give the // name and is the only thing that can choose another. The remedy is a first-class one: give the
// module a short `slug` (ADR 0049), or shorten the machine's name. // module a short `slug` (ADR 0049), or shorten the machine's name.
func CheckIdentity(node, module string) error { func CheckIdentity(node, module string) error {
return CheckIdentityWithin(node, module, IdentityBound{Max: identityLimit, In: "a backend (an S3 access key)"})
}
// CheckIdentityWithin refuses an identity that would not fit one provision's bound, and accepts any
// identity for a provision with none (novox/hq ADR 0225).
func CheckIdentityWithin(node, module string, bound IdentityBound) error {
if !bound.Bounded() {
return nil
}
got := ConsumerIdentity(node, module) got := ConsumerIdentity(node, module)
if len(got) <= identityLimit { if len(got) <= bound.Max {
return nil return nil
} }
return fmt.Errorf( return fmt.Errorf(
"%s on %s is identified as %q, %d characters where a backend (an S3 access key) keeps %d — "+ "%s on %s is identified as %q, %d characters where %s keeps %d — "+
"give the module a shorter `slug` or shorten the machine's name", "give the module a shorter `slug` or shorten the machine's name",
module, node, got, len(got), identityLimit) module, node, got, len(got), bound.In, bound.Max)
}
// Overflow is one consumer whose identity does not fit the provision it requires: left out of its
// provider's grants and reported, never a reason to refuse the provider's machine (ADR 0225).
type Overflow struct {
// Provision is what was required, Provider the machine answering it.
Provision string `json:"provision"`
Provider string `json:"provider"`
// Consumer is the machine, Module the module on it that required it.
Consumer string `json:"consumer"`
Module string `json:"module"`
// Identity is the name the mesh derived, and Bound what it overflows.
Identity string `json:"identity"`
Bound IdentityBound `json:"bound"`
}
func (o Overflow) String() string {
return fmt.Sprintf("%s on %s requires %s from %s and is identified as %q, %d characters where %s "+
"keeps %d — left out of %s's grants until the module's `slug` is shorter",
o.Module, o.Consumer, o.Provision, o.Provider, o.Identity, len(o.Identity), o.Bound.In,
o.Bound.Max, o.Provider)
}
// Overflowing is every requirement of this machine's modules whose identity overflows the bound of
// the provision answering it. The same judgement the provider's composition makes before it grants
// (grantsFor), made from the consumer's side so `status` can say it about every machine.
func (r Resolution) Overflowing() []Overflow {
slugs := map[string]string{}
for _, m := range r.Modules {
slugs[m.Module] = m.Slug
}
var out []Overflow
seen := map[string]bool{}
for _, n := range r.Needs {
if n.ByRecord || !n.Identity.Bounded() {
continue
}
source := IdentitySource(slugs[n.For], n.For)
if CheckIdentityWithin(r.Node, source, n.Identity) == nil {
continue
}
key := n.Name + "\x00" + n.From + "\x00" + n.For
if seen[key] {
continue // one line per requirement, however many local names it has
}
seen[key] = true
out = append(out, Overflow{Provision: n.Name, Provider: n.From, Consumer: r.Node, Module: n.For,
Identity: ConsumerIdentity(r.Node, source), Bound: n.Identity})
}
sort.Slice(out, func(i, j int) bool {
if out[i].Module != out[j].Module {
return out[i].Module < out[j].Module
}
return out[i].Provision < out[j].Provision
})
return out
}
// DefaultLongestMachine is the machine name the catalogue check judges identities on when it is not
// told one: the longest name of the mesh this catalogue is written for, so a catalogue that passes
// passes on every machine that mesh has. `module check --longest-machine-name` says another mesh's;
// a mesh that names a longer machine raises this in the same change (ADR 0225).
const DefaultLongestMachine = 6
// IdentityProblems is every module whose identity would overflow a provision it wants, on a machine
// whose name is `longestMachine` characters — judged before merge, over the catalogue alone, so the
// pull request that introduces an overflow is the one refused (novox/hq ADR 0225, issue 263). A
// provision no module in the shelf offers is not judged: its bound is not known here.
func IdentityProblems(shelf Shelf, longestMachine int) []string {
offeredBy := map[string][]string{}
for _, name := range shelfOrder(shelf) {
for _, o := range shelf[name].Offers() {
offeredBy[o] = append(offeredBy[o], name)
}
}
machine := strings.Repeat("n", longestMachine)
var problems []string
for _, name := range shelfOrder(shelf) {
m := shelf[name]
source := IdentitySource(m.Slug, m.Module)
for _, want := range m.Wants() {
// The tightest bound among the modules offering it: whichever one answers on a given
// machine, the identity has to fit it.
tightest, by := IdentityBound{}, ""
for _, provider := range offeredBy[want] {
if provider == name {
continue // a module answering its own requirement is not its own consumer
}
b := shelf[provider].IdentityBoundOf(want)
if b.Bounded() && (!tightest.Bounded() || b.Max < tightest.Max) {
tightest, by = b, provider
}
}
if CheckIdentityWithin(machine, source, tightest) == nil {
continue
}
got := ConsumerIdentity(machine, source)
problems = append(problems, fmt.Sprintf(
"%s wants %s, and %s keeps its consumers' identities in %s of at most %d characters: "+
"on a machine with a %d-character name it is identified as %q, %d — give %s a "+
"`slug` of at most %d characters",
name, want, by, tightest.In, tightest.Max, longestMachine, got, len(got), name,
tightest.Max-len(IdentityPrefix)-longestMachine-1))
}
}
return problems
} }
+190
View File
@@ -0,0 +1,190 @@
package catalogue
import (
"encoding/json"
"strings"
"testing"
)
// Each test names the decision it defends: novox/hq ADR 0225, which refines ADR 0049 after issue 263.
// The provision a module requires sets the bound on its identity, not the tightest backend anywhere.
func TestABoundIsTheProvisionsOwn(t *testing.T) {
store := Manifest{Module: "objects", Provides: []Offer{{Name: "s3-bucket", Scope: ScopeMesh,
Identity: &OfferIdentity{Max: 20, In: "an S3 access key"}}},
Receives: map[string]string{"s3-bucket": "/var/lib/mesh/objects/mesh.json"}}
database := Manifest{Module: "db", Provides: []Offer{{Name: "postgres-database", Scope: ScopeMesh,
Identity: &OfferIdentity{Max: 63, In: "a PostgreSQL role"}}},
Receives: map[string]string{"postgres-database": "/var/lib/mesh/db/mesh.json"}}
if b := store.IdentityBoundOf("s3-bucket"); b.Max != 20 || b.In != "an S3 access key" {
t.Errorf("an object store's stated bound was not taken: %+v", b)
}
if b := database.IdentityBoundOf("postgres-database"); b.Max != 63 {
t.Errorf("a database's stated bound was not taken: %+v", b)
}
// mesh_workstation_keycloak is 25: refused by the object store, accepted by the database.
if CheckIdentityWithin("workstation", "keycloak", store.IdentityBoundOf("s3-bucket")) == nil {
t.Error("a 25-character identity fit a 20-character access key")
}
if err := CheckIdentityWithin("workstation", "keycloak", database.IdentityBoundOf("postgres-database")); err != nil {
t.Errorf("a database consumer paid the object store's limit: %v", err)
}
}
// What an offer leaves unsaid follows from whether its provider can keep a name at all.
func TestAnUnstatedBoundFollowsWhatTheProviderIsTold(t *testing.T) {
// Told nothing about its consumers — no receives, nothing served from their identity: no bound.
// The resolver provision is exactly this, and its consumers paid an object store's limit (263).
resolver := Manifest{Module: "resolver", Provides: []Offer{{Name: "wildcard-resolution", Scope: ScopeMesh}}}
if b := resolver.IdentityBoundOf("wildcard-resolution"); b.Bounded() {
t.Errorf("a provision that is told nothing of its consumers bounds them: %+v", b)
}
// Told each consumer and silent about its backend: the old global bound, not none.
told := Manifest{Module: "told", Provides: []Offer{{Name: "thing", Scope: ScopeMesh}},
Receives: map[string]string{"thing": "/var/lib/mesh/told/mesh.json"}}
if b := told.IdentityBoundOf("thing"); b.Max != DefaultIdentityLimit {
t.Errorf("a provider told its consumers and silent about its backend is not held to %d: %+v",
DefaultIdentityLimit, b)
}
// Serving a value built from the identity is being told it, too.
serving := Manifest{Module: "serving", Provides: []Offer{{Name: "bucket", Scope: ScopeMesh}},
Serves: map[string]map[string]any{"bucket": {"name": "b-${consumer:as:dns}"}}}
if b := serving.IdentityBoundOf("bucket"); b.Max != DefaultIdentityLimit {
t.Errorf("a provider deriving a name from its consumers is not bounded: %+v", b)
}
// And `false` says it outright, even for a provider that receives.
routes := Manifest{Module: "routes", Provides: []Offer{{Name: "route", Scope: ScopeMesh,
Identity: &OfferIdentity{None: true}}},
Receives: map[string]string{"route": "/var/lib/mesh/routes/mesh.json"}}
if b := routes.IdentityBoundOf("route"); b.Bounded() {
t.Errorf("`identity: false` still bounds: %+v", b)
}
}
// The field reads as written and writes back the same, and refuses what says nothing.
func TestAnOffersIdentityIsParsedStrictly(t *testing.T) {
for _, raw := range []string{
`{"name":"s3-bucket","scope":"mesh","identity":{"max":20,"in":"an S3 access key"}}`,
`{"name":"wildcard-resolution","scope":"mesh","identity":false}`,
`{"name":"redis-cache","scope":"mesh","identity":{"in":"a Redis ACL user"}}`,
} {
var o Offer
if err := json.Unmarshal([]byte(raw), &o); err != nil {
t.Fatalf("%s: %v", raw, err)
}
back, err := json.Marshal(o)
if err != nil || string(back) != raw {
t.Errorf("did not round-trip:\n%s\n%s (%v)", raw, back, err)
}
}
for _, raw := range []string{
`{"name":"x","identity":true}`,
`{"name":"x","identity":{"max":20,"in":"y","most":3}}`,
} {
var o Offer
if err := json.Unmarshal([]byte(raw), &o); err == nil {
t.Errorf("accepted %s", raw)
}
}
bad := Manifest{Module: "bad", Version: "1", Provides: []Offer{
{Name: "unsaid", Scope: ScopeMesh, Identity: &OfferIdentity{Max: 20}},
{Name: "tiny", Scope: ScopeMesh, Identity: &OfferIdentity{Max: 4, In: "nothing usable"}},
}}
raw, err := json.Marshal(bad)
if err != nil {
t.Fatal(err)
}
_, err = ParseManifest(raw)
if err == nil || !strings.Contains(err.Error(), "without saying what keeps them") ||
!strings.Contains(err.Error(), "the shortest the mesh makes") {
t.Fatalf("a bound with no `in`, or too short for any identity, was accepted: %v", err)
}
}
// Refused before merge: the catalogue check judges each module's identity, on the longest machine
// name, against the bound of every provision it wants — and names the module and the slug to set.
func TestTheCatalogueCheckRefusesAnIdentityThatOverflowsWhatItRequires(t *testing.T) {
shelf := Shelf{
"objects": {Module: "objects", Provides: []Offer{{Name: "s3-bucket", Scope: ScopeMesh,
Identity: &OfferIdentity{Max: 20, In: "an S3 access key"}}},
Receives: map[string]string{"s3-bucket": "/var/lib/mesh/objects/mesh.json"}},
"photoalbum": {Module: "photoalbum", Requires: []string{"s3-bucket"}},
"files": {Module: "files", Requires: []string{"s3-bucket"}},
}
problems := IdentityProblems(shelf, 6)
if len(problems) != 1 || !strings.Contains(problems[0], "photoalbum wants s3-bucket") ||
!strings.Contains(problems[0], `"mesh_nnnnnn_photoalbum", 22`) ||
!strings.Contains(problems[0], "`slug` of at most 8 characters") {
t.Fatalf("one overflow, named with its remedy, was expected: %q", problems)
}
// A longer machine name refuses more: the check is about the mesh's machines, not one.
if got := IdentityProblems(shelf, 10); len(got) != 2 {
t.Fatalf("on a 10-character name both overflow (mesh_nnnnnnnnnn_files is 21): %q", got)
}
// And a slug is the remedy it names.
album := shelf["photoalbum"]
album.Slug = "album"
shelf["photoalbum"] = album
if got := IdentityProblems(shelf, 6); len(got) != 0 {
t.Fatalf("a slug that fits is still refused: %q", got)
}
}
// Tonight's case (issue 263): networkmanager, no slug, requiring the mesh's resolver provision on a
// machine with a six-character name. Under ADR 0049's one bound it was refused, and its provider's
// whole machine with it; the resolver keeps no name, so it is not refused at all.
func TestARequirementOnAKeylessProvisionComposesWithALongName(t *testing.T) {
shelf := Shelf{
"resolver": {Module: "resolver", Provides: []Offer{{Name: "wildcard-resolution", Scope: ScopeMesh}}},
"networkmanager": {Module: "networkmanager", Requires: []string{"wildcard-resolution"}},
}
if CheckIdentity("laptop", "networkmanager") == nil {
t.Fatal("the regression is not reproduced: mesh_laptop_networkmanager fits the old global bound")
}
if got := IdentityProblems(shelf, 6); len(got) != 0 {
t.Fatalf("a requirement on a keyless provision was refused for its length: %q", got)
}
r := Resolution{Node: "laptop", Modules: []Manifest{shelf["networkmanager"]},
Needs: []Needed{{Name: "wildcard-resolution", From: "anchor", For: "networkmanager",
Identity: shelf["resolver"].IdentityBoundOf("wildcard-resolution")}}}
if got := r.Overflowing(); len(got) != 0 {
t.Fatalf("a keyless requirement is reported as overflowing: %+v", got)
}
}
// The consumer's side of the same judgement the provider's composition makes: what `status` says.
func TestAnOverflowingRequirementIsNamedFromTheConsumersSide(t *testing.T) {
bound := IdentityBound{Max: 20, In: "an S3 access key"}
r := Resolution{Node: "laptop",
Modules: []Manifest{{Module: "photoalbum"}, {Module: "files"}, {Module: "gallery", Slug: "gal"}},
Needs: []Needed{
{Name: "s3-bucket", From: "anchor", For: "photoalbum", Identity: bound},
{Name: "s3-bucket", From: "anchor", For: "photoalbum", Local: "second", Identity: bound},
{Name: "s3-bucket", From: "anchor", For: "files", Identity: bound},
{Name: "s3-bucket", From: "anchor", For: "gallery", Identity: bound},
{Name: "licence", From: "records", For: "photoalbum", ByRecord: true, Identity: bound},
}}
got := r.Overflowing()
if len(got) != 1 || got[0].Module != "photoalbum" || got[0].Identity != "mesh_laptop_photoalbum" ||
got[0].Provider != "anchor" {
t.Fatalf("one overflow, once, was expected: %+v", got)
}
if said := got[0].String(); !strings.Contains(said, "an S3 access key keeps 20") ||
!strings.Contains(said, "slug") {
t.Fatalf("the overflow does not say what keeps it or the remedy: %s", said)
}
}
// An identity served as a DNS label is never truncated, so the offer's bound must keep it in one.
func TestAnIdentityServedAsADNSLabelIsBoundedToOne(t *testing.T) {
serves := map[string]map[string]any{"bucket": {"name": "${consumer:as:dns}"}}
wide := Manifest{Module: "wide", Serves: serves, Provides: []Offer{{Name: "bucket", Scope: ScopeMesh,
Identity: &OfferIdentity{Max: 255, In: "a client id"}}}}
if got := CheckServes(wide); len(got) != 1 || !strings.Contains(got[0], "a label keeps 63") {
t.Fatalf("a 255-character bound on a DNS label passed: %q", got)
}
unsaid := Manifest{Module: "unsaid", Serves: serves, Provides: []Offer{{Name: "bucket", Scope: ScopeMesh}}}
if got := CheckServes(unsaid); len(got) != 0 {
t.Fatalf("the default bound (20) on a DNS label was refused: %q", got)
}
}
+107 -9
View File
@@ -8,6 +8,7 @@ package catalogue
import ( import (
"bytes" "bytes"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"regexp" "regexp"
"sort" "sort"
@@ -153,6 +154,84 @@ type Offer struct {
// a seat's holder gated by the machine's graphical session, and pulling one in for whatever asked // a seat's holder gated by the machine's graphical session, and pulling one in for whatever asked
// is the misassignment research 026 found. Unmet, the requirement is refused naming who could. // is the misassignment research 026 found. Unmet, the requirement is refused naming who could.
Reach string `json:"reach,omitempty"` Reach string `json:"reach,omitempty"`
// Identity is the longest consumer identity this provision's backend keeps (novox/hq ADR 0225):
// `{"max": 63, "in": "a PostgreSQL role"}`, or `false` for a provision that keeps no name derived
// from its consumer. Unsaid, the mesh assumes the tightest backend it knows when the provider is
// told its consumers, and no bound when it is not — see IdentityBoundOf.
Identity *OfferIdentity `json:"identity,omitempty"`
}
// OfferIdentity is what an offer says about the names its backend keeps for its consumers.
type OfferIdentity struct {
// None is set by `"identity": false`: the provision keeps no name derived from its consumer.
None bool
// Max is the longest identity kept, in characters; zero with In set means no limit worth stating.
Max int
// In is what keeps it, for a refusal to quote.
In string
}
// UnmarshalJSON accepts `false` or `{"max": N, "in": "..."}`.
func (i *OfferIdentity) UnmarshalJSON(raw []byte) error {
var flag bool
if err := json.Unmarshal(raw, &flag); err == nil {
if flag {
return errors.New("an offer's identity is false (it keeps no name) or {max, in}; true says nothing")
}
*i = OfferIdentity{None: true}
return nil
}
var full struct {
Max int `json:"max,omitempty"`
In string `json:"in"`
}
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(&full); err != nil {
return fmt.Errorf("an offer's identity is false or {max, in}: %w", err)
}
*i = OfferIdentity{Max: full.Max, In: full.In}
return nil
}
// MarshalJSON writes it back in the form it was written.
func (i OfferIdentity) MarshalJSON() ([]byte, error) {
if i.None {
return []byte("false"), nil
}
return json.Marshal(struct {
Max int `json:"max,omitempty"`
In string `json:"in"`
}{i.Max, i.In})
}
// IdentityBoundOf is the bound this module's offer of a provision puts on its consumers' identities
// (novox/hq ADR 0225).
//
// **What an offer says, it gets.** Where it says nothing, the bound follows from whether the
// provider can keep a name at all: a provider that receives the provision, or serves its consumers
// a value built from their identity, is told who each consumer is and may create a name from it in
// a backend nobody measured — so it keeps the old global bound, DefaultIdentityBound. One that does
// neither is told nothing about its consumers and keeps nothing of them: no bound. The resolver
// provision is that case, and its consumers paid an object store's limit until this (issue 263).
func (m Manifest) IdentityBoundOf(provision string) IdentityBound {
for _, o := range m.Provides {
if o.Name != provision || o.Identity == nil {
continue
}
if o.Identity.None {
return IdentityBound{}
}
return IdentityBound{Max: o.Identity.Max, In: o.Identity.In}
}
if _, receives := m.Receives[provision]; receives {
return DefaultIdentityBound
}
if served, err := json.Marshal(m.Serves[provision]); err == nil &&
bytes.Contains(served, []byte("${consumer:as")) {
return DefaultIdentityBound
}
return IdentityBound{}
} }
// MachineReach is whether a provision is usable only on its provider's own machine. // MachineReach is whether a provision is usable only on its provider's own machine.
@@ -206,20 +285,21 @@ func (o *Offer) UnmarshalJSON(raw []byte) error {
Scope string `json:"scope,omitempty"` Scope string `json:"scope,omitempty"`
Credential *OfferCredential `json:"credential,omitempty"` Credential *OfferCredential `json:"credential,omitempty"`
Reach string `json:"reach,omitempty"` Reach string `json:"reach,omitempty"`
Identity *OfferIdentity `json:"identity,omitempty"`
} }
dec := json.NewDecoder(bytes.NewReader(raw)) dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields() dec.DisallowUnknownFields()
if err := dec.Decode(&full); err != nil { if err := dec.Decode(&full); err != nil {
return fmt.Errorf("a provided name is either a string or {name, scope, credential, reach}: %w", err) return fmt.Errorf("a provided name is either a string or {name, scope, credential, reach, identity}: %w", err)
} }
o.Name, o.Scope, o.Credential, o.Reach = full.Name, full.Scope, full.Credential, full.Reach o.Name, o.Scope, o.Credential, o.Reach, o.Identity = full.Name, full.Scope, full.Credential, full.Reach, full.Identity
return nil return nil
} }
// MarshalJSON writes back the short form when there is nothing else to say, so a manifest that // MarshalJSON writes back the short form when there is nothing else to say, so a manifest that
// went through the mesh comes out looking like the one that went in. // went through the mesh comes out looking like the one that went in.
func (o Offer) MarshalJSON() ([]byte, error) { func (o Offer) MarshalJSON() ([]byte, error) {
if o.Scope == "" && o.Credential == nil && o.Reach == "" { if o.Scope == "" && o.Credential == nil && o.Reach == "" && o.Identity == nil {
return json.Marshal(o.Name) return json.Marshal(o.Name)
} }
return json.Marshal(struct { return json.Marshal(struct {
@@ -227,7 +307,8 @@ func (o Offer) MarshalJSON() ([]byte, error) {
Scope string `json:"scope,omitempty"` Scope string `json:"scope,omitempty"`
Credential *OfferCredential `json:"credential,omitempty"` Credential *OfferCredential `json:"credential,omitempty"`
Reach string `json:"reach,omitempty"` Reach string `json:"reach,omitempty"`
}{o.Name, o.Scope, o.Credential, o.Reach}) Identity *OfferIdentity `json:"identity,omitempty"`
}{o.Name, o.Scope, o.Credential, o.Reach, o.Identity})
} }
// Manifest is everything a module says about itself. // Manifest is everything a module says about itself.
@@ -237,9 +318,10 @@ type Manifest struct {
// Slug is a short identifier the mesh uses in place of the module name when it derives a // Slug is a short identifier the mesh uses in place of the module name when it derives a
// consumer's login (novox/hq ADR 0049). Optional: a module with a short name needs none. It // consumer's login (novox/hq ADR 0049). Optional: a module with a short name needs none. It
// exists because `mesh_<node>_<module>` must fit the tightest backend a consumer reaches — an S3 // exists because `mesh_<node>_<module>` must fit the bound of every provision the module
// access key is 20 characters — and a long module name would overflow it. A person choosing // requires — an S3 access key's 20 characters is the tightest — and a long module name would
// `kc` for keycloak keeps the identity legible where a hash would not. // overflow it (ADR 0225: the bound is the provision's own, and a keyless one has none). A person
// choosing `kc` for keycloak keeps the identity legible where a hash would not.
Slug string `json:"slug,omitempty"` Slug string `json:"slug,omitempty"`
// Provides are the names other modules may require. A module always provides its own name; // Provides are the names other modules may require. A module always provides its own name;
@@ -1265,8 +1347,10 @@ func ParseManifest(raw []byte) (Manifest, error) {
"%q is not a usable module name: lower-case letters, digits, dashes and dots", m.Module)) "%q is not a usable module name: lower-case letters, digits, dashes and dots", m.Module))
} }
// A slug is a short identifier the mesh derives a login from (novox/hq ADR 0049). The same // A slug is a short identifier the mesh derives a login from (novox/hq ADR 0049). The same
// charset as a name; its length is checked against a backend's limit at assignment, where the // charset as a name; its length is judged against the bound of each provision it wants on the
// node it joins is known — a slug that is fine on one machine's short name can overflow another's. // longest machine name by the catalogue check (IdentityProblems, ADR 0225), and against the
// machine it is on when its provider grants it — a slug that is fine on one machine's short name
// can overflow another's.
if m.Slug != "" && !name.MatchString(m.Slug) { if m.Slug != "" && !name.MatchString(m.Slug) {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%q is not a usable slug: lower-case letters, digits, dashes and dots", m.Slug)) "%q is not a usable slug: lower-case letters, digits, dashes and dots", m.Slug))
@@ -1303,6 +1387,20 @@ func ParseManifest(raw []byte) (Manifest, error) {
if !name.MatchString(p) { if !name.MatchString(p) {
problems = append(problems, fmt.Sprintf("%q is not a usable name to provide", p)) problems = append(problems, fmt.Sprintf("%q is not a usable name to provide", p))
} }
// A bound says what keeps the name, so the refusal it causes can say it (ADR 0225); and it
// leaves room for the shortest identity the mesh makes, or it would refuse every consumer.
if id := offer.Identity; id != nil && !id.None {
if strings.TrimSpace(id.In) == "" {
problems = append(problems, fmt.Sprintf(
"%s bounds the identities of %s's consumers without saying what keeps them: `in`",
m.Module, p))
}
if shortest := len(IdentityPrefix) + 3; id.Max < 0 || id.Max > 0 && id.Max < shortest {
problems = append(problems, fmt.Sprintf(
"%s bounds %s's consumers' identities at %d characters, and the shortest the mesh "+
"makes is %d", m.Module, p, id.Max, shortest))
}
}
if offer.Credential != nil { if offer.Credential != nil {
own, declared := m.OwnSecrets[offer.Credential.Own] own, declared := m.OwnSecrets[offer.Credential.Own]
switch { switch {
+11 -2
View File
@@ -194,6 +194,11 @@ type Needed struct {
// state, not a consumer missing its key. Set by the plan, which is the only layer that knows a // state, not a consumer missing its key. Set by the plan, which is the only layer that knows a
// licence's manager; empty for every consumer. // licence's manager; empty for every consumer.
Manager bool Manager bool
// Identity is the longest consumer identity the answering provision keeps (novox/hq ADR 0225),
// from the provider's own offer: what the mesh judges this consumer's identity against, on the
// consumer's side for `status` and on the provider's before it grants. No bound for a provision
// answered by a record, which keeps no name of anybody's.
Identity IdentityBound
} }
// Refusal is why a set of assignments cannot become a declaration. // Refusal is why a set of assignments cannot become a declaration.
@@ -403,7 +408,7 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
needs = append(needs, Needed{ needs = append(needs, Needed{
Name: want, From: node.Name, At: at, Name: want, From: node.Name, At: at,
Serves: servedByOne(by, want), For: because[want], Serves: servedByOne(by, want), For: because[want],
SharedOwn: sharedByOne(by, want)}) SharedOwn: sharedByOne(by, want), Identity: by.IdentityBoundOf(want)})
} else if served := servedByOne(by, want); len(served) > 0 { } else if served := servedByOne(by, want); len(served) > 0 {
// Answered here with no credential to mint, but the provider serves facts the // Answered here with no credential to mint, but the provider serves facts the
// consumer cannot guess — a port, a model name — and so still needs a binding. // consumer cannot guess — a port, a model name — and so still needs a binding.
@@ -447,11 +452,15 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
return return
} }
shared := "" shared := ""
// A provider whose definition is not in hand is held to the tightest bound the
// mesh knows, not to none: what it keeps is not known here (ADR 0225).
bound := DefaultIdentityBound
if pm, known := catalogue[p.Module]; known { if pm, known := catalogue[p.Module]; known {
shared, _ = pm.SharedCredentialOf(want) shared, _ = pm.SharedCredentialOf(want)
bound = pm.IdentityBoundOf(want)
} }
needs = append(needs, Needed{Name: want, From: p.Node, At: p.At, needs = append(needs, Needed{Name: want, From: p.Node, At: p.At,
Serves: p.Serves, For: because[want], SharedOwn: shared}) Serves: p.Serves, For: because[want], SharedOwn: shared, Identity: bound})
} }
switch { switch {
case world.Unchecked: case world.Unchecked:
+29 -4
View File
@@ -110,6 +110,8 @@ var ControllerVerbs = []Verb{
"paths": "with repository: the files the merge would change, comma-separated, from the repository's root", "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", "modules": "with repository: or the modules it would change, comma-separated",
"limit": "how many plans to list (default 10); only when listing", "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)}, }, nil)},
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " + {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.", "or, with files, the files it would be given.",
@@ -137,12 +139,15 @@ var ControllerVerbs = []Verb{
{Name: "push", Description: "Send one machine everything it should be. With no machine named it is a push of the " + {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. " + "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 " + "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{ Input: schema(map[string]string{
"node": "the machine's name; without it, every machine that is behind", "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", "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",
{Name: "rotate", Description: "Replace a credential. A pair credential, by provision (and a consuming machine, " + "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 " + "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 " + "name: made anew and the machine sent, so the module starts again on it — only for a secret its " +
"definition says it reads at start; a value given to the mesh, or one the module applies to a backend, is refused with the reason.", "definition says it reads at start; a value given to the mesh, or one the module applies to a backend, is refused with the reason.",
@@ -150,7 +155,7 @@ var ControllerVerbs = []Verb{
"provision": "a pair credential: the provision whose credential to replace", "provision": "a pair credential: the provision whose credential to replace",
"consumer": "with provision: only the holder on this machine (optional)", "consumer": "with provision: only the holder on this machine (optional)",
"node": "an own secret: the machine", "node": "an own secret: the machine",
"module": "an own secret: the module", "module": "an own secret: the module; with provision: only this consuming module's credential (optional)",
"secret": "an own secret: its name in the module's definition", "secret": "an own secret: its name in the module's definition",
}, nil)}, }, nil)},
{Name: "issue", Description: "Give a module on a machine its account on the bus: minted, and sealed to the " + {Name: "issue", Description: "Give a module on a machine its account on the bus: minted, and sealed to the " +
@@ -206,6 +211,26 @@ var ControllerVerbs = []Verb{
Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)}, 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.", {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)}, 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: " + {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.", "`builds` with that id follows it line by line, and the module is registered when the outcome comes.",
Input: schema(map[string]string{ Input: schema(map[string]string{
+146
View File
@@ -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
}
+97
View File
@@ -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)
}
}
@@ -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;
+22 -6
View File
@@ -80,13 +80,29 @@ func (i *Inventory) SavePlan(ctx context.Context, p Plan) error {
if err != nil { if err != nil {
return err return err
} }
_, err = i.store.Pool().Exec(ctx, // **And how long the tier it left took** (novox/hq to-be 45 Phase 0): measured here, where the
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch) // plan moves, in the same transaction as the move, so no save can move a tier unmeasured or
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10) // 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, 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`, 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) tier_entered = excluded.tier_entered`,
return err 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. // OpenPlans is every plan still being worked, oldest first.
+32 -4
View File
@@ -105,14 +105,27 @@ func (i *Inventory) widenProtocol(ctx context.Context, s catalogue.Seat) error {
changed := false changed := false
row.Accepts, changed = union(row.Accepts, s.Accepts, changed) row.Accepts, changed = union(row.Accepts, s.Accepts, changed)
row.Emits, changed = union(row.Emits, s.Emits, changed) row.Emits, changed = union(row.Emits, s.Emits, changed)
have := map[string]bool{} have := map[string]int{}
for _, v := range row.Serves { for n, v := range row.Serves {
have[v.Name] = true have[v.Name] = n
} }
for _, v := range s.Serves { for _, v := range s.Serves {
if !have[v.Name] { at, kept := have[v.Name]
if !kept {
row.Serves = append(row.Serves, v) row.Serves = append(row.Serves, v)
changed = true 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 { if !changed {
@@ -282,3 +295,18 @@ func (i *Inventory) Holdings(ctx context.Context) ([]catalogue.Held, error) {
} }
return out, rows.Err() 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)
}
+43
View File
@@ -113,3 +113,46 @@ func TestRenameSeatKeepsTheFormerNameAsAnAlias(t *testing.T) {
t.Fatal("renaming a seat to its own name was accepted") 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")
}
}
+205 -17
View File
@@ -7,7 +7,9 @@ import (
"fmt" "fmt"
"log" "log"
"regexp" "regexp"
"sort"
"strconv" "strconv"
"strings"
"sync" "sync"
"time" "time"
@@ -33,8 +35,9 @@ import (
// either. A variable so a test need not wait. // either. A variable so a test need not wait.
var AnswerWithin = 10 * time.Second var AnswerWithin = 10 * time.Second
// KeptCalls is how many calls are kept, newest first; keptAnswer the largest answer kept of a call // KeptCalls is how many calls this process keeps in memory, newest first — to match a refusal to its
// whose caller was sent it. One its caller never had is kept whole. // 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 ( const (
KeptCalls = 100 KeptCalls = 100
keptAnswer = 64 << 10 keptAnswer = 64 << 10
@@ -47,6 +50,10 @@ const (
// CallFinishedAfter is a call that finished after its caller was told it was still running: its // CallFinishedAfter is a call that finished after its caller was told it was still running: its
// answer is here and nowhere else. // answer is here and nowhere else.
CallFinishedAfter = "finished after its caller was answered" 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. // 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 // 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. // the only place that says what it would have.
Refused string `json:"answer refused by the bus,omitempty"` 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 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. // CallLog keeps the latest calls a holder served.
type CallLog struct { type CallLog struct {
mu sync.Mutex 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 // 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. // serves one (the controller's `calls`); without it, the answer points at the holder's journal.
Follow string 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. // 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} } 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.<its user>.`
// 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 { func (l *CallLog) begin(seat, verb string, args json.RawMessage, reply string) *Call {
l.mu.Lock() l.mu.Lock()
defer l.mu.Unlock() defer l.mu.Unlock()
l.next++ l.next++
c := &Call{ID: "call-" + strconv.FormatInt(l.now().UnixNano(), 10) + "-" + strconv.FormatUint(l.next, 10), 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) l.calls = append(l.calls, c)
if len(l.calls) > KeptCalls { if len(l.calls) > KeptCalls {
l.calls = l.calls[len(l.calls)-KeptCalls:] l.calls = l.calls[len(l.calls)-KeptCalls:]
} }
l.keep(*c)
return c return c
} }
@@ -117,29 +264,67 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) {
} else { } else {
c.State = CallAnswered c.State = CallAnswered
} }
l.keep(*c)
} }
// Recent is the kept calls, newest first, as copies. // Recent is the kept calls, newest first, as copies: this process's from memory, and, when they are
func (l *CallLog) Recent() []Call { // 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() l.mu.Lock()
defer l.mu.Unlock()
out := make([]Call, 0, len(l.calls)) out := make([]Call, 0, len(l.calls))
seen := map[string]bool{}
for i := len(l.calls) - 1; i >= 0; i-- { for i := len(l.calls) - 1; i >= 0; i-- {
out = append(out, *l.calls[i]) out = append(out, *l.calls[i])
seen[l.calls[i].ID] = true
} }
return out keeper := l.keeper
} l.mu.Unlock()
if keeper == nil {
// Get is one kept call. return out, nil
func (l *CallLog) Get(id string) (Call, bool) { }
l.mu.Lock() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer l.mu.Unlock() defer cancel()
for _, c := range l.calls { kept, err := keeper.Recent(ctx)
if c.ID == id { if err != nil {
return *c, true 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 // 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() at := l.now()
hit.Refused = fmt.Sprintf("%s: %s", at.Format(time.RFC3339), err) 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) id, verb, took := hit.ID, hit.Seat+"."+hit.Verb, at.Sub(hit.Started).Round(time.Second)
l.keep(*hit)
l.mu.Unlock() l.mu.Unlock()
if logger != nil { if logger != nil {
// Said in the mesh's words, beside the library's own line: which call, and where its answer is. // 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(`{}`) args = json.RawMessage(`{}`)
} }
c := l.begin(seat, verb, kept(args), reply) 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{}) acknowledged := make(chan struct{})
var once sync.Once var once sync.Once
ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) }) ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) })
+115
View File
@@ -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**: `<id>` holds the record — verb, arguments as kept, caller, state, times — and
// `<id>.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)
}
}
}
}
+152
View File
@@ -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)
}
}
}
+11 -6
View File
@@ -62,7 +62,7 @@ func TestACallThatFinishesInTimeAnswersInFull(t *testing.T) {
if got := a.only()["result"]; got != "all well" { if got := a.only()["result"]; got != "all well" {
t.Fatalf("answered %v", got) t.Fatalf("answered %v", got)
} }
recent := l.Recent() recent, _ := l.Recent()
if len(recent) != 1 || recent[0].State != CallAnswered || string(recent[0].Args) != "{}" { if len(recent) != 1 || recent[0].State != CallAnswered || string(recent[0].Args) != "{}" {
t.Fatalf("kept %+v", recent) 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) { if got["running"] != true || id == "" || !strings.Contains(got["output"].(string), id) {
t.Fatalf("the running answer does not name its call: %v", got) 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) t.Fatalf("while it runs it is kept as %q", c.State)
} }
close(release) close(release)
<-finished <-finished
c, ok := l.Get(id) c, ok, _ := l.Get(id)
if !ok || c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "anchor told") { if !ok || c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "anchor told") {
t.Fatalf("its answer was not kept: %+v", c) 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 { if got := a.only()["result"].(map[string]any); got["running"] != true {
t.Fatalf("an acknowledged call answered %v", got) 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) t.Fatalf("kept %+v", c)
} }
} }
@@ -143,7 +143,7 @@ func TestARefusedAnswerIsKeptAgainstItsCall(t *testing.T) {
if !l.Refusal(refusal, log.New(&logged, "", 0)) { if !l.Refusal(refusal, log.New(&logged, "", 0)) {
t.Fatal("the refusal of a kept call's answer was not recognised") 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) { if c.Refused == "" || !strings.Contains(logged.String(), c.ID) {
t.Fatalf("the refusal is not kept or not said: %+v / %q", c, logged.String()) 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++ { for i := 0; i < KeptCalls+5; i++ {
l.begin("s", "v", nil, "") 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") { 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) 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
}
+9
View File
@@ -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.Applied == nil && report.Refused == "" && len(report.Failed) == 0 {
if report.Superseded != "" { if report.Superseded != "" {
log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, 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) 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 { if err != nil {
return false, err 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 // 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 // what the machine holds, so each moves last_seen and nothing else. Recording a partial list
+162
View File
@@ -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 <text>")
}
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
}
+59
View File
@@ -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)
}
}