Author SHA1 Message Date
mesh-admin 8170fc58a3 Merge pull request 'Keep a passed gate's verdict when the first machine's later reports go quiet (hq issue 335)' (#165) from fix/335-a-passed-gate-is-not-judged-again into main 2026-10-08 19:48:03 +00:00
mesh-admin efcdd5dd7d Merge pull request 'Serve the read verbs on the serving controller's own connection, and name every connection (hq issue 327, ADR 0265)' (#161) from fix/327-a-verb-reads-on-the-serving-connection into main 2026-10-08 19:32:11 +00:00
jochen 5b7e6ff453 Answer dead-letters on the lent serving connection, and keep one clip helper
mesh/delivery delivered
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
Two handles to the same serving connection, and two copies of one helper,
would drift (review of hq issues 327 and 330).
2026-10-08 21:09:08 +02:00
jochen cc7fb99f29 Answer a panicking verb with an error, say flag errors in the answer, and leave refused logins out of D15
Read verbs now run in the serving process, where a panic would end every
call; refused logins are nobody's reconnect loop (review of hq issue 327).
2026-10-08 21:08:34 +02:00
jochen 1e04670052 Serve the read verbs on the serving controller's own connection, and name every connection
Each verb ran as a process that dialled the bus, so hundreds of short
connections an hour, all named mesh-controller, hid any client reconnecting
in a loop (hq issue 327). D15 now says a user whose connections keep dropping.
2026-10-08 21:08:34 +02:00
jochen 81e5458cbf Test that a carried module takes its lead's pass before its step is read (hq issue 335 review)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
Without the reorder the passed build's walk stopped on the first machine's
later silence; the guard that caught it is now said to be one.
2026-10-08 21:08:12 +02:00
jochen fb74e24c9e Keep a passed gate's verdict when the first machine's later reports go quiet (hq issue 335)
A build that passed on its first machine was judged again from that
machine's next reports while its send to the rest waited; another walk's
unreported send there then failed the passed build at the wait's bound
and put it back.
2026-10-08 21:08:12 +02:00
mesh-admin ca09a07fdf Merge pull request 'Keep what a consumer gives up on until a person delivers it again or drops it (hq issue 330, ADR 0264)' (#159) from fix/330-a-message-given-up-on-is-kept into main 2026-10-08 19:07:33 +00:00
mesh-admin 9b028b4212 Merge pull request 'Add the node-nfs-server and node-mounts seats (hq ADR 0263)' (#163) from feat/mounts-module into main 2026-10-08 18:39:52 +00:00
jochen d7fab82a89 Say the adopt switch is the string "true" and that reload only has the kernel reread its exports, as the holders do
mesh/delivery delivered
mesh/delivery-group group feat/mounts-module delivered: every member is delivered
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
2026-10-08 20:11:33 +02:00
mesh-admin 7f65e62743 Merge pull request 'Keep a left-out module's provisions and backups, and refuse more identity keys (hq ADR 0262)' (#162) from feat/left-out-keeps-what-it-provides into main 2026-10-08 16:37:42 +00:00
jochen 2b01f8786e Let node-nfs-server.test take a client's address: the server knows the range, not the mesh's node names
mesh/delivery-group group feat/mounts-module rejected: a member's own check failed
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
2026-10-08 18:34:38 +02:00
jochen 909062e729 Deliver a dead letter again only where it is received, and never stop serving for the notices
mesh/delivery delivered
mesh/delivery-group group fix/330-a-message-given-up-on-is-kept delivered: every member is delivered
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
A dead letter was let go as delivered even when its consumer did not filter
its again subject; seat asks needed a grant over every seat's queue and left
the original stuck; a notices bind failure stopped the controller (review).
2026-10-08 18:32:58 +02:00
jochen 826dcb91b1 Keep what a consumer gives up on until a person delivers it again or drops it
Design 25 promised a dead-letter stream that did not exist: a message a
consumer gave up on stayed only in its source, which drops it after a week,
and its condition cleared when the advisories stopped (hq issue 330, ADR 0264).
2026-10-08 18:32:58 +02:00
jochen 193168e086 Place a left-out module's backup lines best effort, and refuse more identity keys (hq ADR 0262 review)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
An unplaceable line of a left-out module, such as an access nobody placed, failed the whole machine's
declaration. Say it among what could not be placed instead, never copy the definition's path past a
placement that does not read, and accept a removal only when the decoder is past it.
2026-10-08 18:27:50 +02:00
jochen 59fdffb979 Add the node-nfs-server and node-mounts seats, so a share and a mount have a role the mesh defines (hq ADR 0263)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery-group group feat/mounts-module ready: every member ready, and composed together they pass
mesh/delivery superseded: a newer head of the same pull request
A machine sharing folders and a machine mounting them each need one holder
per machine, with verbs an agent calls instead of exportfs, fstab edits or
zfs set. Both seats deliver nothing: nfs-share is provided at the mesh's
scope. The two adopt verbs are dry runs unless confirmed.
2026-10-08 18:24:38 +02:00
jochen 5d47e0bfd6 Keep a left-out module's provisions and backups, and refuse more identity keys (hq ADR 0262)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
A module left out for an unknown key inside an entry lost its whole manifest, so every consumer of what
it provides was refused and its data stopped being copied. Read past only the unknown key, keep its
backup lines, and say in the condition what stops.
2026-10-08 18:03:46 +02:00
mesh-admin e3ec15f707 Merge pull request 'Fill a preference's ${setting:} from its manifest default (hq ADR 0262)' (#153) from feat/setting-defaults into main 2026-10-08 15:54:34 +00:00
mesh-admin ac91357a53 Merge pull request 'Promise unlink-dangling on the service manager, optional (hq issue 332)' (#160) from feat/service-manager-unlink-dangling into main 2026-10-08 15:52:53 +00:00
jochen b5f2c3b961 Promise unlink-dangling on the service manager, optional (hq issue 332)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer delivery to the same trunk took over its walk
disable cannot remove an enable link whose unit file is gone, so a
leftover unit stays wanted at every login with no verb to end it. Optional
until the systemd module serves it (ADR 0246 step 1).
2026-10-08 17:39:08 +02:00
jochen af63b233db Leave out a module whose stored manifest has an unknown field, and raise it (hq ADR 0262 review)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group feat/setting-defaults failed: a member failed
A key dropped silently ran a module without what its manifest says, and a key inside a block still
failed the whole catalogue. Judge a key by what it is about, and narrow the listing to one machine.
2026-10-08 17:34:53 +02:00
mesh-admin 0cf5a3a5ba Merge pull request 'Promise reset-failed and wanted-by on the service manager, optional (hq issue 332)' (#158) from feat/service-manager-reset-failed-why-started into main 2026-10-08 15:31:41 +00:00
jochen f5680ba8da List every module's preferences in the settings verb (hq ADR 0262)
One verb is the interface to every preference, so no module builds a settings tool of its own: each
key, its default and why, and every assigned machine's value with its source.
2026-10-08 17:24:32 +02:00
jochen 76babaea52 Read stored manifests leniently and mark the defaults layer (hq ADR 0262 review)
A strict read of the stored catalogue fails every plan and send once a manifest uses a field an older
controller lacks; registration stays strict. A node named default lost its layer to the name check.
Judge the operator's keys by whole words, and scan a default under any key.
2026-10-08 17:24:32 +02:00
jochen e41b78cd77 Fill a preference's ${setting:} from its manifest default (hq ADR 0262)
Without a default, a running module could never gain a setting: the file asking for it
failed to compose until set, and the key was refused as stray until a file asked for it.
Defaults sit under the mesh's and the node's settings, never merge into a JSON file, are
refused for the operator's own values, and settings shows each value's source.
2026-10-08 17:24:32 +02:00
jochen 1acce7132e Promise reset-failed and wanted-by on the service manager, optional (hq issue 332)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
A failed unit whose file is gone stays raised until its record is reset,
and nothing could say which unit or enable link still asks for it. Both
verbs are optional until the systemd module serves them (ADR 0246 step 1).
2026-10-08 17:21:16 +02:00
mesh-admin 3f68a495f1 Merge pull request 'Name what a held release holds, and where it is released (ADR 0258)' (#155) from fix/release-held-says-what-waits into main 2026-10-08 15:09:29 +00:00
jochen 298ec06ae0 Say held updates wait because a walk failed, in the glossary's words
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
The backlog is held after any failed walk, not only a release, and a check is a pull
request's status; the words said the last release failed its check.
2026-10-08 17:02:30 +02:00
jochen e33da2dc1c Name what a held release holds, and where it is released
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
The release-held words said "release them, or leave them held" without the modules,
the machines or the mesh MCP server, so the operator could neither tell what waited
nor where to act (ADR 0258). The controller's restart needs missed the same suffix.
A test now holds every need that opens with a verb only the mesh MCP server performs
to name it, so a new kind cannot miss it.
2026-10-08 16:44:02 +02:00
74 changed files with 4418 additions and 168 deletions
+5 -1
View File
@@ -27,6 +27,8 @@ import (
"syscall"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/builder"
"github.com/novox/mesh-controller/internal/link"
@@ -175,7 +177,9 @@ func dialFor(credential Credential) (*broker.JetStream, string, error) {
if !credential.onTheNewBus() {
return nil, "", fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL)
}
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
// Named for what it is, not the controller whose code dials it (novox/hq issue 327).
host, _ := os.Hostname()
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint, nats.Name("build agent on "+host))
if err != nil {
return nil, "", err
}
+3 -1
View File
@@ -8,6 +8,7 @@ import (
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
@@ -187,7 +188,8 @@ func (a *actor) serveUnderTheLease(ctx context.Context, inv *inventory.Inventory
a.mu.Lock()
a.serving = true
a.mu.Unlock()
js, err := broker.Dial(address)
// Its own connection, held as long as the lease, and named so (novox/hq issue 327).
js, err := broker.Dial(address, nats.Name(broker.ConnectionName+" lease"))
if err != nil {
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it to take the "+
"lease: %w", broker.BareAddress(address), err)
+5 -5
View File
@@ -879,6 +879,10 @@ func heldBy(ctx context.Context) map[string]string {
// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus
// being built it dials, because a build request is a one-shot and holds nothing else.
func askOverOn(seat string) (link.Builders, error) {
// The serving controller asks on its own connection (novox/hq issue 327).
if serving := servingBus.Load(); serving != nil {
return link.BuildsOn(serving, seat), nil
}
address, err := broker.BusAddress()
if err != nil {
return nil, err
@@ -937,11 +941,7 @@ func buildSeatAmong(entries []inventory.Entry) string {
// with a consumer of its own that is gone when this returns, so nothing accumulates in the server
// for the reading, and filtered by subject, so one build's lines are all that travel.
func buildLog(ctx context.Context, id string) error {
address, err := broker.BusAddress()
if err != nil {
return err
}
js, err := broker.Dial(address)
js, err := aBus()
if err != nil {
return fmt.Errorf("cannot reach the bus to read a build's log: %w", err)
}
+29 -25
View File
@@ -5,6 +5,7 @@ import (
"errors"
"flag"
"fmt"
"io"
"os"
"slices"
"strconv"
@@ -107,19 +108,20 @@ func conditionsCommand(ctx context.Context, args []string) error {
}
switch sub {
case "list":
return listConditions(ctx, args)
return listConditions(ctx, args, os.Stdout)
case "show":
return showCondition(ctx, args)
return showCondition(ctx, args, os.Stdout)
case "silence":
return silenceCondition(ctx, args)
case "history":
return conditionHistory(ctx, args)
return conditionHistory(ctx, args, os.Stdout)
}
return errors.New(conditionsUsage)
}
func listConditions(ctx context.Context, args []string) error {
func listConditions(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions", flag.ContinueOnError)
usageTo(set, w)
scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", "))
severity := set.String("severity", "", "only urgent, or only warning")
machine := set.String("machine", "", "only those about this machine")
@@ -144,20 +146,20 @@ func listConditions(ctx context.Context, args []string) error {
}
}
if *asJSON {
return printJSON(map[string]any{"conditions": inBrief(out), "open": len(open), "counted": counted(out),
return printJSONTo(w, map[string]any{"conditions": inBrief(out), "open": len(open), "counted": counted(out),
"note": "urgent first, then oldest first; a condition clears when observation says so, never by hand; " +
"each with its newest evidence — `conditions key=<key>` gives one whole"})
}
if len(out) == 0 {
if len(open) == 0 {
fmt.Println("no open conditions")
fmt.Fprintln(w, "no open conditions")
} else {
fmt.Printf("none of the %d open condition(s) is about that\n", len(open))
fmt.Fprintf(w, "none of the %d open condition(s) is about that\n", len(open))
}
return nil
}
for _, line := range conditionLines(out, time.Now()) {
fmt.Println(line)
fmt.Fprintln(w, line)
}
return nil
}
@@ -233,8 +235,9 @@ func conditionLines(list []conditions.Condition, now time.Time) []string {
return out
}
func showCondition(ctx context.Context, args []string) error {
func showCondition(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions show", flag.ContinueOnError)
usageTo(set, w)
asJSON := set.Bool("json", false, "as data")
rest, err := parseAround(set, args)
if err != nil {
@@ -266,34 +269,34 @@ func showCondition(ctx context.Context, args []string) error {
"history --key %s` what became of it", key, key)
}
if *asJSON {
return printJSON(c)
return printJSONTo(w, c)
}
now := time.Now()
fmt.Printf("%s %s\n %s\n\n", strings.ToUpper(string(c.Severity)), c.Key, c.Summary)
fmt.Printf(" kind %s\n about %s %s", c.Kind, c.Subject.Scope, c.Subject.ID)
fmt.Fprintf(w, "%s %s\n %s\n\n", strings.ToUpper(string(c.Severity)), c.Key, c.Summary)
fmt.Fprintf(w, " kind %s\n about %s %s", c.Kind, c.Subject.Scope, c.Subject.ID)
if c.Subject.Machine != "" {
fmt.Printf(", on %s", c.Subject.Machine)
fmt.Fprintf(w, ", on %s", c.Subject.Machine)
}
fmt.Printf("\n raised by %s\n since %s (%s ago), observed %d time(s), last %s ago\n",
fmt.Fprintf(w, "\n raised by %s\n since %s (%s ago), observed %d time(s), last %s ago\n",
c.Source, c.Raised.Local().Format("2006-01-02 15:04:05"), roughly(now.Sub(c.Raised)), c.Observations,
now.Sub(c.LastObserved).Round(time.Second))
if c.Count > 1 {
fmt.Printf(" raised %d times, each within ten minutes of clearing\n", c.Count)
fmt.Fprintf(w, " raised %d times, each within ten minutes of clearing\n", c.Count)
}
fmt.Printf(" resolved by %s\n", resolverWords(c.Resolver))
fmt.Fprintf(w, " resolved by %s\n", resolverWords(c.Resolver))
if c.Silenced != nil {
fmt.Printf(" silenced until %s by %s: %s\n", c.Silenced.Until.Local().Format("2006-01-02 15:04"),
fmt.Fprintf(w, " silenced until %s by %s: %s\n", c.Silenced.Until.Local().Format("2006-01-02 15:04"),
c.Silenced.By, c.Silenced.Why)
}
if len(c.Tried) > 0 {
fmt.Println("\n tried:")
fmt.Fprintln(w, "\n tried:")
for _, t := range c.Tried {
fmt.Printf(" %s %s — %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), orHealer(t.By), t.What, t.Outcome)
fmt.Fprintf(w, " %s %s — %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), orHealer(t.By), t.What, t.Outcome)
}
}
fmt.Println("\n evidence, newest first:")
fmt.Fprintln(w, "\n evidence, newest first:")
for _, e := range c.Evidence {
fmt.Printf(" %s %s\n", e.At.Local().Format("2006-01-02 15:04:05"), e.Said)
fmt.Fprintf(w, " %s %s\n", e.At.Local().Format("2006-01-02 15:04:05"), e.Said)
}
return nil
}
@@ -378,8 +381,9 @@ func parseFor(s string) (time.Duration, error) {
return d, nil
}
func conditionHistory(ctx context.Context, args []string) error {
func conditionHistory(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions history", flag.ContinueOnError)
usageTo(set, w)
days := set.Int("days", 7, "how many days back, at most 90")
key := set.String("key", "", "only this condition")
asJSON := set.Bool("json", false, "as data")
@@ -420,10 +424,10 @@ func conditionHistory(ctx context.Context, args []string) error {
if out == nil {
out = []conditions.Event{}
}
return printJSON(map[string]any{"history": out, "days": *days})
return printJSONTo(w, map[string]any{"history": out, "days": *days})
}
if len(out) == 0 {
fmt.Printf("nothing was raised, changed or cleared in the last %d day(s)\n", *days)
fmt.Fprintf(w, "nothing was raised, changed or cleared in the last %d day(s)\n", *days)
return nil
}
for _, e := range out {
@@ -440,7 +444,7 @@ func conditionHistory(ctx context.Context, args []string) error {
line += " — " + e.Why
}
}
fmt.Println(line)
fmt.Fprintln(w, line)
}
return nil
}
+162
View File
@@ -0,0 +1,162 @@
package main
import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
// What a consumer gave up on, answered by the serving controller (novox/hq issue 330).
//
// **In this process, on its own connection**: DEAD_LETTERS is read and changed on the bus, and the
// serving controller is already on it. A verb run as a fresh process would open a connection of its own
// for each call (novox/hq issue 327).
// busHandles are the serving controller's connection and JetStream handle.
type busHandles struct {
conn *nats.Conn
js nats.JetStreamContext
}
// defaultDeadLetters is how many the list says when not asked for more.
const defaultDeadLetters = 50
// causeDeadLetter is the cause a delivery or drop gives when the caller gives none.
const causeDeadLetter = "dead-letter"
// deadLettersAnswer is what `dead-letters` answers: the list, one whole, or what came of delivering one
// again or dropping it.
func deadLettersAnswer(ctx context.Context, a *verbArguments) (any, error) {
serving := servingBus.Load()
if serving == nil {
return nil, errors.New("this controller is not serving, so it does not read DEAD_LETTERS: ask again, " +
"and the serving controller answers")
}
on := &busHandles{conn: serving.Conn(), js: serving.Context()}
// The shape first, and every argument it reads; one given beside it is refused before anything is
// done, as every verb refuses what it would pass over (novox/hq issue 244).
var deliver, drop, why, cause, idText, consumer, limit string
switch {
case a.given["deliver"] != "" || a.given["drop"] != "":
deliver, drop, why, cause = a.str("deliver"), a.str("drop"), a.str("why"), a.str("cause")
case a.given["id"] != "":
idText = a.str("id")
default:
consumer, limit = a.str("consumer"), a.str("limit")
}
if unused := a.unused(); len(unused) > 0 {
return nil, fmt.Errorf("dead-letters did not use %s together with %s, and an argument a verb would pass "+
"over is refused: nothing was done", quoteAll(unused), quoteAll(a.usedGiven()))
}
switch {
case deliver != "" && drop != "":
return nil, errors.New("dead-letters delivers one again or drops one, not both. Nothing was done")
case deliver != "" || drop != "":
act, text := "deliver", deliver
if drop != "" {
act, text = "drop", drop
}
if strings.TrimSpace(why) == "" {
return nil, fmt.Errorf("dead-letters %s is a hand act, and says why: why is required and recorded in "+
"the hand-act log (novox/hq to-be 45 §7). Nothing was done", act)
}
id, err := deadLetterID(text)
if err != nil {
return nil, err
}
return actOnDeadLetter(ctx, on, act, id, why, cause)
case idText != "":
id, err := deadLetterID(idText)
if err != nil {
return nil, err
}
return link.DeadLetterNamed(on.js, id)
}
most := defaultDeadLetters
if limit != "" {
n, err := strconv.Atoi(limit)
if err != nil || n <= 0 {
return nil, fmt.Errorf("limit is a number of dead letters, not %q", limit)
}
most = n
}
held, total, err := link.DeadLetters(on.js, consumer, most)
if err != nil {
return nil, err
}
answer := map[string]any{"dead_letters": held, "held": total,
"note": "newest first; with id, one whole; deliver or drop one with why"}
if total == 0 {
answer["note"] = "no consumer gave up on a message that is still kept"
}
return answer, nil
}
// deadLetterID is a dead letter's id as a caller wrote it.
func deadLetterID(text string) (uint64, error) {
id, err := strconv.ParseUint(strings.TrimSpace(text), 10, 64)
if err != nil || id == 0 {
return 0, fmt.Errorf("a dead letter's id is its number in %s, as dead-letters lists it, not %q",
broker.DeadLettersStream, text)
}
return id, nil
}
// actOnDeadLetter delivers one again or drops it, recorded in the hand-act log before it is done. A log
// that cannot be written is said, and the act still happens: the log is never the reason a person's act
// is refused (handacts.go).
func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64, why, cause string) (any, error) {
d, err := link.DeadLetterNamed(on.js, id)
if err != nil {
return nil, err
}
if act == "deliver" {
// Refused before it is recorded: an act that cannot be done is not an act.
if _, err := link.AgainTo(on.js, d); err != nil {
return nil, err
}
}
if cause == "" {
cause = causeDeadLetter
}
by := link.CallerIn(ctx)
if by == "" {
by = "a seat call whose caller the bus did not name"
}
answer := map[string]any{"dead_letter": d.ID, "consumer": d.Who, "subject": d.Subject}
recorded, logErr := link.RecordHandAct(ctx, on.conn, link.HandAct{Verb: "dead-letters " + act,
Args: []string{strconv.FormatUint(id, 10), d.Stream + "." + d.Consumer}, Why: why, Cause: cause,
Condition: conditions.Key(conditions.ScopeBus, d.Stream+"."+d.Consumer, link.AdvisoryMaxDeliveries),
By: by + ", through the " + catalogue.ControllerSeatName + " seat"})
if logErr != nil {
answer["unrecorded"] = "the hand-act log could not be written, and the act was done all the same: " + logErr.Error()
} else {
answer["recorded"] = recorded.ID
}
switch act {
case "deliver":
_, to, err := link.DeliverAgain(on.js, id)
if err != nil {
return nil, err
}
answer["delivered_on"] = to
answer["done"] = fmt.Sprintf("dead letter %d was delivered again to %s, and nobody else; it is no longer kept",
id, consumerWho(d.Stream, d.Consumer))
case "drop":
if _, err := link.DropDeadLetter(on.js, id); err != nil {
return nil, err
}
answer["done"] = fmt.Sprintf("dead letter %d, which %s gave up on, was dropped for good", id,
consumerWho(d.Stream, d.Consumer))
}
return answer, nil
}
+132
View File
@@ -0,0 +1,132 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
"github.com/novox/mesh-controller/internal/testbus"
)
// The serving controller's bus, with the mesh's streams, for the verb to read and act on.
func servingDeadLetters(t *testing.T) *broker.JetStream {
t.Helper()
js, err := broker.Dial(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
if err := broker.AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
if err := js.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
before := servingBus.Load()
servingBus.Store(js)
t.Cleanup(func() { servingBus.Store(before) })
return js
}
func askDeadLetters(t *testing.T, args map[string]any) (any, error) {
t.Helper()
a, err := readArguments("dead-letters", args)
if err != nil {
return nil, err
}
return deadLettersAnswer(context.Background(), a)
}
func TestDeadLettersListsDropsAndRecordsWhy(t *testing.T) {
js := servingDeadLetters(t)
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{"n":1}`)); err != nil {
t.Fatal(err)
}
d, err := link.KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":1,"deliveries":5}`))
if err != nil {
t.Fatal(err)
}
answer, err := askDeadLetters(t, map[string]any{})
if err != nil {
t.Fatal(err)
}
listed := answer.(map[string]any)
if listed["held"] != 1 || len(listed["dead_letters"].([]link.DeadLetter)) != 1 {
t.Fatalf("listed %v", listed)
}
// Refused before anything is done: an act without why, an argument the shape passes over, both acts.
for _, args := range []map[string]any{
{"drop": "1"},
{"drop": "1", "why": "x", "limit": "3"},
{"drop": "1", "deliver": "1", "why": "x"},
{"why": "x"},
{"id": "nought"},
} {
if _, err := askDeadLetters(t, args); err == nil {
t.Errorf("%v was done", args)
}
}
if _, total, _ := link.DeadLetters(js.Context(), "", 0); total != 1 {
t.Fatalf("a refused call changed what is kept: %d left", total)
}
done, err := askDeadLetters(t, map[string]any{"drop": "1", "why": "the media server took the download in by hand"})
if err != nil {
t.Fatal(err)
}
if said := done.(map[string]any); said["recorded"] == nil || !strings.Contains(said["done"].(string), "dropped") {
t.Fatalf("answered %v", said)
}
acts, err := link.HandActs(context.Background(), js.Conn(), time.Now().Add(-time.Minute))
if err != nil || len(acts) != 1 || acts[0].Verb != "dead-letters drop" || acts[0].Cause != causeDeadLetter ||
!strings.Contains(acts[0].Condition, "media_sonarr") {
t.Fatalf("recorded %+v (%v)", acts, err)
}
if !personsDecision(acts[0]) {
t.Error("dropping a dead letter is counted as a repair, so S15 would want a healer for it")
}
if _, err := askDeadLetters(t, map[string]any{"id": "1"}); err == nil {
t.Errorf("dead letter %d is still answered after it was dropped", d.ID)
}
}
// Open while DEAD_LETTERS holds a message for the consumer, in words the operator reads in one pass:
// what is held, and where to act.
func TestAConsumersDeadLettersAreSaidUntilActedOn(t *testing.T) {
f := &signalFacts{now: time.Now(), deadLetters: map[string]int{"EVENTS.media_sonarr": 4, "EVENTS.controller": 1}}
found := watchDeadLetters(f)
if len(found) != 2 {
t.Fatalf("said %d conditions", len(found))
}
for _, o := range found {
if o.Kind != "max-deliveries" || o.Severity != conditions.Warning || o.Needs == "" {
t.Errorf("%+v", o)
}
if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Needs: o.Needs,
Explanation: o.Explanation, Resolved: o.Resolved}, "media"); !ok {
t.Errorf("%q is not plain: %s", o.Headline, why)
}
if !strings.Contains(o.Needs, "mesh MCP server") {
t.Errorf("does not say where to act: %q", o.Needs)
}
}
sonarr := found[1]
if sonarr.ID != "EVENTS.media_sonarr" || sonarr.Machine != "media" ||
sonarr.Headline != "Sonarr on media could not handle 4 messages" ||
!strings.Contains(sonarr.Summary, "DEAD_LETTERS") {
t.Errorf("%+v", sonarr)
}
if found[0].Headline != "The controller could not handle a message" {
t.Errorf("%q", found[0].Headline)
}
// None held, none said: it clears when they are delivered again or dropped.
if left := watchDeadLetters(&signalFacts{now: time.Now()}); len(left) != 0 {
t.Fatalf("%v", left)
}
}
+5
View File
@@ -116,6 +116,11 @@ var probeRegistry = []probe{
{ID: probeDeliveriesID, Asserts: "no delivery is held past its state's bound unsaid: mesh-delivery's " +
"`stalled`, each with the transition its table lets healer H2 take", From: "ADR 0239",
Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries},
// A client of the bus reconnecting in a loop (novox/hq issue 327), from the server's record of closed
// connections, which the bus's own module reads.
{ID: probeReconnectsID, Asserts: "no user of the bus had its connection dropped more than twelve times in the " +
"last hour: the bus module's nats_closed_connections", From: "issue 327", Kind: kindBusReconnects,
Phase: 1, run: probeReconnects},
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
// The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is
+20 -13
View File
@@ -6,6 +6,7 @@ import (
"errors"
"flag"
"fmt"
"io"
"os"
"slices"
"sort"
@@ -14,7 +15,6 @@ import (
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
@@ -82,6 +82,11 @@ var handActVerbs = []handActVerb{
{Verb: "retire approve", Decision: "nothing is retired past its bound without a person (ADR 0230)"},
{Verb: "retire reject", Decision: "keeping a consumer active is a person's word (ADR 0230)"},
{Verb: "cleanup delete", Decision: "nothing retired is deleted without a person (ADR 0230)"},
// What becomes of a message a consumer gave up on (novox/hq issue 330): kept until a person says.
{Verb: "dead-letters deliver", Decision: "a message a consumer gave up on is delivered again only on a " +
"person's word (issue 330)"},
{Verb: "dead-letters drop", Decision: "a message a consumer gave up on is let go only on a person's word " +
"(issue 330)"},
// The sweep run on a person's word rather than after a build: the same decision the records make, at
// a moment the person chose (ADR 0251) — never a repair.
{Verb: "collect", Decision: "letting the store go of what the records keep for no reason, now rather " +
@@ -149,14 +154,10 @@ func onTheBus(f func(*nats.Conn) error) error {
if handActConn != nil {
return f(handActConn)
}
address, err := broker.BusAddress()
js, err := aBus()
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())
}
@@ -247,7 +248,13 @@ func handActCommand(ctx context.Context, args []string) error {
if len(args) > 0 && args[0] == "list" {
args = args[1:]
}
return listHandActs(ctx, args, os.Stdout)
}
// listHandActs is `hand-acts`: what was done by hand lately, and the causes done more than once.
func listHandActs(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("hand-acts", flag.ContinueOnError)
usageTo(set, w)
days := set.Int("days", 14, "how many days back")
asJSON := set.Bool("json", false, "as data")
if _, err := parseAround(set, args); err != nil {
@@ -265,27 +272,27 @@ func handActCommand(ctx context.Context, args []string) error {
if err != nil {
return err
}
fmt.Println(string(body))
fmt.Fprintln(w, string(body))
return nil
}
if len(acts) == 0 {
fmt.Printf("nothing was done by hand in the last %d day(s)\n", *days)
fmt.Fprintf(w, "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,
fmt.Fprintf(w, "%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.Fprintf(w, ", condition %s", a.Condition)
}
fmt.Println(")")
fmt.Fprintln(w, ")")
if pushedRecorded(a) {
carried := strings.Join(a.Carried, "; ")
if carried == "" {
carried = recordedBefore[a.ID]
}
fmt.Printf(" a push of recorded builds, no repair: %s\n", carried)
fmt.Fprintf(w, " a push of recorded builds, no repair: %s\n", carried)
}
}
if len(repeated) > 0 {
@@ -294,7 +301,7 @@ func handActCommand(ctx context.Context, args []string) error {
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",
fmt.Fprintf(w, "\ndone by hand more than once in a fortnight — a healer is wanted (to-be 45 S15): %s\n",
strings.Join(causes, ", "))
}
return nil
+19
View File
@@ -15,6 +15,7 @@ import (
"os/signal"
"syscall"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/identity"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/licences"
@@ -53,6 +54,9 @@ func run() error {
return fmt.Errorf("no command given")
}
// Every connection this process dials says what it is (novox/hq issue 327).
broker.ConnectionName = connectionName(args[0], os.Getenv(verbVar))
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
// Whatever this process holds of the controller's lease is given back as it ends (novox/hq to-be
@@ -416,3 +420,18 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error {
statusFrom.nudge()
return nil
}
// verbVar carries the seat verb a command runs for, from the serving controller to the process it starts.
const verbVar = "MESH_VERB"
// connectionName is what this process's connections say they are in the bus's list (novox/hq issue 327):
// the serving controller, a verb's own process and which verb, or a command run at a shell and which.
func connectionName(command, verb string) string {
switch {
case command == "serve":
return "mesh-controller serving"
case verb != "":
return "mesh-controller verb " + verb
}
return "mesh-controller command " + command
}
+163 -5
View File
@@ -8,6 +8,7 @@ import (
"fmt"
"net"
"os"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/broker"
@@ -396,7 +397,8 @@ func assignCommand(ctx context.Context, verb string, args []string) error {
func settingsCommand(ctx context.Context, args []string) error {
if len(args) == 0 {
return errors.New("settings show <module> [--node <node>] [--history], settings set <module> <file> " +
"[--node <node>] [--replace], or settings clear <module> [--node <node>]")
"[--node <node>] [--replace], settings clear <module> [--node <node>], or settings preferences " +
"[<module>] [--node <node>]")
}
open, err := openStores(ctx)
if err != nil {
@@ -502,13 +504,92 @@ func settingsCommand(ctx context.Context, args []string) error {
}
if !has {
fmt.Printf("%s on %s: no layer — the module's definition says\n", positionals[0], where)
return nil
} else {
shown, err := json.MarshalIndent(values, "", " ")
if err != nil {
return err
}
fmt.Println(string(shown))
}
shown, err := json.MarshalIndent(values, "", " ")
// Every value the module gives a default or a layer sets, and where it came from (novox/hq
// ADR 0262): the default, the mesh's layer, or this node's. Said after the layer, which stays
// the first thing printed because a caller reads it before replacing it (ADR 0217).
known, err := inv.Catalogue(ctx)
if err != nil {
return err
}
fmt.Println(string(shown))
m, ok := known[positionals[0]]
if !ok || len(m.Settings) == 0 {
return nil
}
var layers []catalogue.Layer
if *node != "" {
if mesh, has, err := inv.Layer(ctx, "", positionals[0]); err != nil {
return err
} else if has {
layers = append(layers, catalogue.Layer{From: catalogue.MeshWideLayer, Values: mesh})
}
if has {
layers = append(layers, catalogue.Layer{From: *node, Values: values})
}
} else if has {
layers = append(layers, catalogue.Layer{From: catalogue.MeshWideLayer, Values: values})
}
fmt.Print(describeEffective(positionals[0], where, catalogue.Effective(m, layers)))
return nil
case "preferences":
// Every module's preferences, and each machine's value with its source (novox/hq ADR 0262).
if len(positionals) > 1 {
return errors.New("settings preferences [<module>] [--node <node>]")
}
only := ""
if len(positionals) == 1 {
only = positionals[0]
}
if *node != "" {
// A machine the mesh does not know is refused, never answered with an empty listing.
if _, err := inv.NodeByName(ctx, *node); err != nil {
return err
}
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
var listed []preferencesOf
for _, e := range entries {
m := e.Manifest
if len(m.Settings) == 0 || (only != "" && m.Module != only) {
continue
}
if *node != "" && !containsString(e.On, *node) {
// Asked for one machine: a module not on it has no value there to say.
continue
}
p := preferencesOf{Manifest: m, On: map[string][]catalogue.SettingSource{}}
for _, n := range e.On {
if *node != "" && n != *node {
continue
}
p.Nodes = append(p.Nodes, n)
layers, err := inv.SettingsFor(ctx, n, m.Module)
if err != nil {
return err
}
p.On[n] = catalogue.Effective(m, layers)
}
listed = append(listed, p)
}
if only != "" && len(listed) == 0 {
fmt.Printf("%s declares no preferences%s\n", only, onNode(*node))
return nil
}
if len(listed) == 0 && *node != "" {
fmt.Printf("no module on %s declares a preference\n", *node)
return nil
}
fmt.Print(describePreferences(listed))
return nil
case "clear":
@@ -522,10 +603,87 @@ func settingsCommand(ctx context.Context, args []string) error {
return nil
default:
return fmt.Errorf("settings has no %q; it has show, set and clear", args[0])
return fmt.Errorf("settings has no %q; it has show, set, clear and preferences", args[0])
}
}
// describeEffective says each setting's value on a machine or the whole mesh, where it came from, and
// the module's default when a layer overrides it.
func describeEffective(module, where string, values []catalogue.SettingSource) string {
var b strings.Builder
fmt.Fprintf(&b, "%s on %s, every value and where it comes from:\n", module, where)
for _, v := range values {
value, _ := json.Marshal(v.Value)
fmt.Fprintf(&b, " %s = %s (%s", v.Key, value, v.From)
if v.HasDefault && !v.FromDefault {
d, _ := json.Marshal(v.Default)
fmt.Fprintf(&b, "; the default is %s", d)
}
b.WriteString(")\n")
}
return b.String()
}
// onNode is ` on <node>` for one machine, nothing for the whole mesh.
func onNode(node string) string {
if node == "" {
return ""
}
return " on " + node
}
// preferencesOf is one module's preferences and its value on each machine it is assigned to.
type preferencesOf struct {
Manifest catalogue.Manifest
Nodes []string
On map[string][]catalogue.SettingSource
}
// describePreferences lists each module's preferences — key, default and why — and, per machine it is
// assigned to, the value and where it comes from (novox/hq ADR 0262).
func describePreferences(modules []preferencesOf) string {
if len(modules) == 0 {
return "no module declares a preference\n"
}
var b strings.Builder
for i, p := range modules {
if i > 0 {
b.WriteString("\n")
}
on := "assigned nowhere"
if len(p.Nodes) > 0 {
on = "on " + strings.Join(p.Nodes, ", ")
}
fmt.Fprintf(&b, "%s (%s)\n", p.Manifest.Module, on)
keys := make([]string, 0, len(p.Manifest.Settings))
for k := range p.Manifest.Settings {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
d := p.Manifest.Settings[k]
def, _ := json.Marshal(d.Default)
fmt.Fprintf(&b, " %s, default %s: %s\n", k, def, d.Why)
for _, n := range p.Nodes {
for _, s := range p.On[n] {
if s.Key != k {
continue
}
v, _ := json.Marshal(s.Value)
from := s.From
if s.FromDefault {
from = catalogue.DefaultLayer
} else if s.From != catalogue.MeshWideLayer {
from = "the node"
}
fmt.Fprintf(&b, " %s: %s (%s)\n", n, v, from)
}
}
}
}
return b.String()
}
// nodeFlag is ` --node <node>` for a machine's layer, nothing for the whole mesh's.
func nodeFlag(node string) string {
if node == "" {
+7 -2
View File
@@ -7,7 +7,9 @@ import (
"errors"
"flag"
"fmt"
"io"
"net"
"os"
"strings"
"time"
@@ -684,11 +686,14 @@ func orNotReported(s string) string {
}
// printJSON prints a value as indented JSON, the shape every `--json` answers in.
func printJSON(v any) error {
func printJSON(v any) error { return printJSONTo(os.Stdout, v) }
// printJSONTo is printJSON to a writer of the caller's.
func printJSONTo(w io.Writer, v any) error {
body, err := json.MarshalIndent(v, "", " ")
if err != nil {
return err
}
fmt.Println(string(body))
fmt.Fprintln(w, string(body))
return nil
}
+3
View File
@@ -478,6 +478,9 @@ func settlingPending(ctx context.Context, open *stores) {
for _, line := range settlePending(ctx, open, time.Now()) {
fmt.Println(line)
}
for _, line := range raiseUnknownFields(ctx, open.inventory) {
fmt.Println(line)
}
select {
case <-ctx.Done():
return
+41 -10
View File
@@ -387,10 +387,7 @@ var plainWordings = map[string]func(conditions.Observation) words{
Resolved: "Resolved: " + m + " runs a good build again"}
}),
"release-held": worded(func(o conditions.Observation) words {
return words{Headline: "Updates wait for your release",
Needs: "release them, or leave them held.",
Explanation: "Some module updates wait for a person to release them, and are not delivered until then.",
Resolved: "Resolved: the held updates are released"}
return releaseHeldWords(nil, o.Also)
}),
"facts-stale": worded(func(o conditions.Observation) words {
return words{Headline: "Merge checks use outdated facts",
@@ -402,21 +399,21 @@ var plainWordings = map[string]func(conditions.Observation) words{
// The controller and the core.
"controller-deaf": worded(func(o conditions.Observation) words {
return words{Headline: "The controller stopped listening",
Needs: "restart the controller if this stays.",
Needs: "restart the controller if this stays, " + FromMeshMCPServer,
Explanation: "The controller, which coordinates the mesh, has taken no messages for minutes while some " +
"wait. Changes and repairs do not happen until it recovers.",
Resolved: "The controller listens again"}
}),
"self-check-silent": worded(func(o conditions.Observation) words {
return words{Headline: "The mesh's self-check stopped",
Needs: "restart the controller if this stays.",
Needs: "restart the controller if this stays, " + FromMeshMCPServer,
Explanation: "The self-check, which looks over the whole mesh every few minutes, has not finished a run. " +
"Problems may go unnoticed until it runs again.",
Resolved: "The self-check runs again"}
}),
"watchdogs-silent": worded(func(o conditions.Observation) words {
return words{Headline: "The mesh's watchdogs stopped",
Needs: "restart the controller if this stays.",
Needs: "restart the controller if this stays, " + FromMeshMCPServer,
Explanation: "The watchdogs, which notice when something expected does not happen, have not run, so " +
"missed signals are not noticed.",
Resolved: "The watchdogs run again"}
@@ -546,10 +543,18 @@ var plainWordings = map[string]func(conditions.Observation) words{
Explanation: "The bus reports a listener too slow to keep up, so messages to it are late.",
Resolved: "The listener keeps up again"}
}),
kindBusReconnects: worded(func(o conditions.Observation) words {
return words{Headline: "A client keeps losing the bus",
Explanation: "One of the mesh's clients lost its connection to the bus again and again in the last hour. " +
"While it reconnects, what it says and what it is asked waits.",
Resolved: "Resolved: the client stays connected"}
}),
"max-deliveries": worded(func(o conditions.Observation) words {
return words{Headline: "A message could not be handled",
Explanation: "The bus gave up on a message after trying to hand it over too many times.",
Resolved: "Resolved: messages are handled again"}
return words{Headline: "A listener gave up on messages",
Needs: "deliver them again or drop them, from the mesh MCP server.",
Explanation: "A listener on the bus could not handle messages after several tries, so what they asked " +
"for was not done. The mesh keeps them until you deliver them again or drop them.",
Resolved: "Resolved: the messages given up on were delivered again or dropped"}
}),
"refused": worded(func(o conditions.Observation) words {
return words{Headline: "The bus refuses some messages",
@@ -683,6 +688,32 @@ func moduleNeeds(node string, rs []inventory.ResourceHealth) string {
// authorised (to-be 46 phases 5 and 6).
const FromMeshMCPServer = "from the mesh MCP server; this notification cannot do it."
// releaseHeldWords are the plain words of updates held after a walk failed: which modules wait,
// on which machines, and where the operator releases them (ADR 0258: a release is not an acknowledgement, so
// no notification gives it). Raised with the modules and machines (backlogObservation); the kind's fallback
// knows neither.
func releaseHeldWords(modules, machines []string) words {
what := "Some module updates"
headline := "Updates wait for your release"
if len(modules) > 0 {
what = "Updates of " + namesWords(modules, 3)
if h := what + " wait for your release"; len(h) <= conditions.HeadlineMax {
headline = h
} else if h := fmt.Sprintf("%d module updates wait for your release", len(modules)); len(h) <= conditions.HeadlineMax {
headline = h
}
}
where := ""
if len(machines) > 0 {
where = " on " + namesWords(machines, 4)
}
return words{Headline: headline,
Needs: "release them " + FromMeshMCPServer,
Explanation: fmt.Sprintf("%s%s wait for a person to release them, because the last walk failed."+
" They are not delivered until then, and stay held if you leave them.", what, where),
Resolved: "Resolved: the held updates are released"}
}
// reloginNeeds is what an account waiting for its groups needs (ADR 0252).
func reloginNeeds(node string) string {
return fmt.Sprintf("log out of %s completely and log in again, or restart it.", node)
+87
View File
@@ -1,6 +1,7 @@
package main
import (
"regexp"
"strings"
"testing"
"time"
@@ -222,3 +223,89 @@ func TestDataLossOffersNoSilence(t *testing.T) {
}
}
}
// **Updates held after a failed release** (2026-10-08): the popup read "Needs you: release them, or leave
// them held." — naming neither what waits nor where it is released. It names the modules and machines, and
// the mesh MCP server.
func TestUpdatesHeldNameWhatWaitsAndWhereItIsReleased(t *testing.T) {
saved := backlogNow
t.Cleanup(func() { backlogNow = saved })
backlogNow.held = "release-1791457717307061152 failed (failed its gate on g14); what waits is released again by a person"
backlogNow.waiting = map[string][]inventory.CarriedMove{
"shanks": {{Module: "openrazer"}},
"g14": {{Module: "openrazer"}, {Module: "sensors"}},
}
got := backlogObservation()
if len(got) != 1 {
t.Fatalf("raised %d", len(got))
}
plainExample(t, got[0], "Updates of openrazer and sensors wait for your release",
"Needs you: release them from the mesh MCP server; this notification cannot do it. Updates of openrazer and "+
"sensors on g14 and shanks wait for a person to release them, because the last walk failed. "+
"They are not delivered until then, and stay held if you leave them.")
}
// **What held them is said in the glossary's words** (2026-10-08 review): the backlog is held after any
// walk failed, not a release, and a "check" is a pull request's status. So the words say the walk failed,
// and never that a release failed or a check did — with the modules and machines named or not.
func TestUpdatesHeldSayTheWalkFailed(t *testing.T) {
for _, w := range []words{
releaseHeldWords([]string{"openrazer"}, []string{"g14"}),
releaseHeldWords(nil, nil),
plainWordings["release-held"](conditions.Observation{Scope: conditions.ScopeMesh, ID: "release"}),
} {
if !strings.Contains(w.Explanation, "because the last walk failed.") {
t.Errorf("does not say the walk failed: %q", w.Explanation)
}
for _, wrong := range []string{"release failed", "check"} {
if strings.Contains(w.Explanation, wrong) {
t.Errorf("says %q: %q", wrong, w.Explanation)
}
}
}
}
// mcpVerb is a need that opens with a verb only the mesh MCP server performs (ADR 0258 §1): release, stop,
// start, restart, and a restore.
var mcpVerb = regexp.MustCompile(`^(release|stop|start|restart|restore)\b`)
// **A need no notification can answer says where it is answered** (ADR 0258 §1): every wording whose need
// opens with a verb the mesh MCP server performs, and offers no action, ends with FromMeshMCPServer — the
// kinds worded here for every subject shape, and those worded where they are raised. A new kind that misses
// it fails here; release-held did (2026-10-08).
func TestANeedNoNotificationAnswersNamesTheMeshMCPServer(t *testing.T) {
check := func(what string, w words) {
t.Helper()
if w.Needs == "" || len(w.Actions) > 0 || !mcpVerb.MatchString(w.Needs) {
return
}
if !strings.HasSuffix(w.Needs, FromMeshMCPServer) {
t.Errorf("%s needs %q without %q", what, w.Needs, FromMeshMCPServer)
}
}
subjects := []conditions.Observation{
{Scope: conditions.ScopeMachine, ID: "ace", Machine: "ace"},
{Scope: conditions.ScopeMachine, ID: "ace.immich.library", Machine: "ace"},
{Scope: conditions.ScopeModule, ID: "openrazer.g14", Machine: "g14"},
{Scope: conditions.ScopeDelivery, ID: "novox/hq@055550802096"},
{Scope: conditions.ScopeCore, ID: "controller.anchor", Machine: "anchor"},
{Scope: conditions.ScopeMesh, ID: "release", Also: []string{"g14"}},
}
for kind, fn := range plainWordings {
for _, s := range subjects {
for _, sev := range []conditions.Severity{conditions.Warning, conditions.Urgent} {
s.Kind, s.Severity, s.Resolver = kind, sev, conditions.ResolverOperator
check(kind+" about "+s.ID, fn(s))
}
}
}
check("a walk waiting", words{Needs: waitingNeeds(conditions.Urgent)})
check("a module's failed service", words{Needs: moduleNeeds("g14",
[]inventory.ResourceHealth{{Kind: link.KindUnit, Target: "openrazer-daemon.service"}})})
for _, state := range []string{"held", "ready", "failing"} {
_, _, _, needs, actions := stalledWords(stalledLine{ID: "novox/hq@055550802096", State: state, For: "36h"},
conditions.Observation{Resolver: conditions.ResolverOperator})
check("a delivery "+state, words{Needs: needs, Actions: actions})
}
check("updates held", releaseHeldWords([]string{"openrazer"}, []string{"g14"}))
}
+3 -4
View File
@@ -515,9 +515,8 @@ func sortedKeysOf(m map[string]string) []string {
// 0163, rule 6), one line each: the machine is told everything else, and is told it was left out.
func reportLeftOut(node string, declared sendable) {
for _, m := range declared.LeftOut {
fmt.Printf("%s: %s left out — a setting stored for it cannot compose with its definition; "+
"what the machine holds for it is kept and its containers are untouched. %s\n",
node, m, declared.leftOutWhy[m])
fmt.Printf("%s: %s left out — what the machine holds for it is kept and its containers are "+
"untouched. %s\n", 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.
@@ -1425,7 +1424,7 @@ func servedOnNode(ctx context.Context, inv *inventory.Inventory, node string,
}
serves := catalogue.ServedOn(m, provision, ports)
if len(serves) > 0 {
serves, err = catalogue.Settle(serves, layers)
serves, err = catalogue.Settle(serves, catalogue.WithDefaults(m, layers))
if err != nil {
return nil, err
}
+3 -6
View File
@@ -9,7 +9,6 @@ import (
"strings"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
@@ -98,11 +97,9 @@ func buildSeatPause(ctx context.Context, inv *inventory.Inventory, plans []inven
if len(holders) == 0 {
return pauseView{}
}
address, err := broker.BusAddress()
if err != nil {
return pauseView{}
}
js, err := broker.Dial(address)
// On the serving controller's own connection when this is it: a watchdog tick while a walk waits for a
// build dialled one every 30 seconds (novox/hq issue 327).
js, err := aBus()
if err != nil {
fmt.Fprintf(os.Stderr, "could not reach the bus to read whether the build seat is paused: %v\n", err)
return pauseView{}
+10
View File
@@ -875,6 +875,16 @@ func streamDiffers(want broker.Stream, have nats.StreamConfig) string {
if perSubject != 0 && have.MaxMsgsPerSubject != perSubject {
differs = append(differs, fmt.Sprintf("keeps %d per subject, defined %d", have.MaxMsgsPerSubject, perSubject))
}
if want.MaxBytes > 0 && have.MaxBytes != want.MaxBytes {
differs = append(differs, fmt.Sprintf("holds up to %d bytes, defined %d", have.MaxBytes, want.MaxBytes))
}
if want.DuplicatesSeconds > 0 && have.Duplicates != time.Duration(want.DuplicatesSeconds)*time.Second {
differs = append(differs, fmt.Sprintf("keeps one of a message id for %s, defined %s", have.Duplicates,
time.Duration(want.DuplicatesSeconds)*time.Second))
}
if want.DiscardNew && have.Discard != nats.DiscardNew {
differs = append(differs, "drops what it holds when full, defined to refuse what comes next")
}
return strings.Join(differs, "; ")
}
+4
View File
@@ -219,6 +219,10 @@ func serve(ctx context.Context) (err error) {
}
// The hand-act log is counted for `status` on this connection rather than a new one a minute.
handActConn = bus.Conn
// And everything else this controller does on the bus for a moment (novox/hq issue 327).
servingBus.Store(server.JetStream())
// No longer serving: nothing is lent, and dead-letters says it is not read here (novox/hq issue 330).
defer servingBus.Store(nil)
// And says when it replaced a value given by hand (novox/hq ADR 0228).
givenEvents = bus
// And a pull request's merge check, asked when the forge announces its head and said when judged
+12 -12
View File
@@ -6,6 +6,7 @@ import (
"errors"
"flag"
"fmt"
"io"
"os"
"sort"
"strings"
@@ -37,22 +38,21 @@ const (
killAnswer = 75 * time.Second
)
// dialTheBus opens the controller's own connection, for a command that reads or changes the queue.
// dialTheBus is the controller's connection, for a command that reads or changes the queue: the serving
// controller's own, lent, when this process is it (novox/hq issue 327).
func dialTheBus() (*broker.JetStream, error) {
address, err := broker.BusAddress()
if err != nil {
return nil, err
}
js, err := broker.Dial(address)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus: %w", err)
}
return js, nil
return aBus()
}
// queueCommand prints every ask in the build seat's work queue.
func queueCommand(ctx context.Context, args []string) error {
return listQueue(ctx, args, os.Stdout)
}
// listQueue is `queue`: the build queue, as a person reads it or as JSON.
func listQueue(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("queue", flag.ContinueOnError)
usageTo(set, w)
asJSON := set.Bool("json", false, "the queue as JSON")
if _, err := parseAround(set, args); err != nil {
return err
@@ -72,10 +72,10 @@ func queueCommand(ctx context.Context, args []string) error {
if err != nil {
return err
}
fmt.Println(string(body))
fmt.Fprintln(w, string(body))
return nil
}
fmt.Print(queueText(q, time.Now()))
fmt.Fprint(w, queueText(q, time.Now()))
return nil
}
+173
View File
@@ -0,0 +1,173 @@
package main
import (
"context"
"fmt"
"sort"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
// D15: no client of the bus reconnects in a loop (novox/hq issue 327).
//
// The server's connection total was the one number that would show a client reconnecting, and every verb
// the controller served opened and closed a connection of its own, hundreds an hour, so a loop was
// invisible in it. The bus's own module reads the server's record of closed connections
// (`nats_closed_connections`); this asks it every run, and says each user whose connections were dropped —
// closed by anything but the client itself: a read or write error, a stale connection, a slow consumer, a
// refused login — more often than reconnectBound in the last hour.
const (
probeReconnectsID = "D15"
kindBusReconnects = "bus-reconnects"
// reconnectBound is how many dropped connections in an hour one user may have before it is said:
// a client that loses its connection every five minutes. Provisional.
reconnectBound = 12
// busModule is the module that is the bus, and closedTool its tool that reads closed connections.
busModule = "nats"
closedTool = "nats_closed_connections"
closedAsk = 20 * time.Second
)
// closedConnections is what the bus's module answers.
type closedConnections struct {
Hours float64 `json:"hours"`
Reaches bool `json:"reaches"`
Users []struct {
User string `json:"user"`
Closed int `json:"closed"`
Dropped int `json:"dropped"`
DroppedPerHour float64 `json:"dropped_per_hour"`
Names []struct {
Name string `json:"name"`
Closed int `json:"closed"`
Dropped int `json:"dropped"`
Reasons map[string]int `json:"reasons"`
} `json:"names"`
} `json:"users"`
}
// probeReconnects is D15.
func probeReconnects(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
if d.js == nil {
return nil, fmt.Errorf("no bus to ask the bus's module over")
}
on, err := d.open.inventory.Running(ctx, busModule)
if err != nil {
return nil, err
}
if len(on) == 0 {
return nil, nil // no bus module assigned: a mesh whose bus is not the mesh's module
}
return reconnectsOn(ctx, d.js.Conn(), on[0])
}
// reconnectsOn asks the bus module on its machine and says who reconnects in a loop.
func reconnectsOn(ctx context.Context, conn *nats.Conn, node string) ([]conditions.Observation, error) {
read, err := askClosed(ctx, conn, node)
if isNothingServes(err) {
// A bus module older than its tool has nothing it can say, which is not a failure of the probe.
return nil, nil
}
if err != nil {
return nil, err
}
return reconnecting(read), nil
}
// askClosed asks the bus's module, on its machine, who closed connections in the last hour.
func askClosed(ctx context.Context, conn *nats.Conn, node string) (closedConnections, error) {
var read closedConnections
answer, err := link.AskModuleToolOn(ctx, conn, busModule, closedTool, node, map[string]any{"hours": 1}, closedAsk)
if err != nil {
return read, err
}
if answer.Error != "" {
return read, fmt.Errorf("%s on %s answered %s with an error: %s", busModule, node, closedTool, answer.Error)
}
if err := unmarshalAnswer(answer, &read); err != nil {
return read, fmt.Errorf("%s on %s answered %s with something unreadable: %w", busModule, node, closedTool, err)
}
return read, nil
}
// reconnecting is one observation per user whose connections were dropped more than reconnectBound times
// an hour.
func reconnecting(read closedConnections) []conditions.Observation {
hours := read.Hours
if hours <= 0 {
hours = 1
}
var out []conditions.Observation
for _, u := range read.Users {
if u.User == "" || strings.HasPrefix(u.User, "(") {
// Refused before it logged in: no client of the mesh's, so nobody's reconnect loop. A login
// refused again and again is a question of its own, not this probe's.
continue
}
perHour := float64(u.Dropped) / hours
if perHour <= reconnectBound {
continue
}
reasons := map[string]int{}
var names []string
for _, n := range u.Names {
if n.Dropped == 0 {
continue
}
names = append(names, fmt.Sprintf("%q ×%d", n.Name, n.Dropped))
for r, c := range n.Reasons {
if r != "Client Closed" {
reasons[r] += c
}
}
}
var why []string
for r, c := range reasons {
why = append(why, fmt.Sprintf("%s ×%d", r, c))
}
sort.Strings(why)
partial := ""
if !read.Reaches {
partial = " (at least: the server's record of closed connections does not reach back the whole hour)"
}
who, machine := busUserWords(u.User)
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: u.User, Kind: kindBusReconnects,
Machine: machine, Severity: conditions.Warning,
Summary: fmt.Sprintf("the bus dropped %s's connection %d times in the last hour%s, more than %d: a "+
"client reconnecting in a loop", u.User, u.Dropped, partial, reconnectBound),
Said: fmt.Sprintf("%d of %d closed connections dropped in %.0f h; by name %s; why %s", u.Dropped,
u.Closed, hours, strings.Join(names, ", "), strings.Join(why, ", ")),
Headline: clip(conditions.Capital(who)+" keeps losing the bus", 60),
Explanation: conditions.Capital(fmt.Sprintf("%s lost its connection to the bus %d times in the last hour "+
"and connected again each time. While it reconnects, what it says and what it is asked waits.", who,
u.Dropped)),
Resolved: "Resolved: " + who + " stays connected"})
}
return out
}
// busUserWords is a bus user as the operator says it, and the machine it is on: `node.<machine>` the
// node-engine, `<machine>.node-tools` the tool runner, `controller` the controller, `<machine>.<module>` a
// module.
func busUserWords(user string) (string, string) {
switch {
case user == "controller":
return "the controller", ""
case strings.HasPrefix(user, "node."):
m := strings.TrimPrefix(user, "node.")
return "the node-engine on " + m, m
}
if m, module, ok := strings.Cut(user, "."); ok {
if module == "node-tools" {
return "the tool runner on " + m, m
}
return module + " on " + m, m
}
return "a client of the bus", ""
}
+189
View File
@@ -0,0 +1,189 @@
package main
import (
"context"
"encoding/json"
"strings"
"testing"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
"github.com/novox/mesh-controller/internal/testbus"
)
// The serving controller's own connection serves what it does for a moment: no connection is opened,
// and closing what it was lent leaves the serving one open (novox/hq issue 327).
func TestTheServingControllerLendsItsOwnConnection(t *testing.T) {
bus := testbus.Start(t)
serving, err := broker.Dial(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
defer serving.Close()
if err := serving.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
was := servingBus.Load()
servingBus.Store(serving)
defer servingBus.Store(was)
t.Setenv(broker.NATSVar, bus.ClientURL())
before, _ := bus.Varz(nil)
lent, err := aBus()
if err != nil {
t.Fatal(err)
}
lent.Close()
if !serving.Conn().IsConnected() {
t.Fatal("closing a lent connection closed the serving controller's")
}
if err := onTheBus(func(*nats.Conn) error { return nil }); err != nil {
t.Fatal(err)
}
handlers, _, err := seatToolHandlers()
if err != nil {
t.Fatal(err)
}
answer, err := handlers["hand-acts"](context.Background(), json.RawMessage(`{}`))
if err != nil {
t.Fatal(err)
}
if a := answer.(verbAnswer); !a.OK || a.Answer == nil {
t.Fatalf("hand-acts answered %+v", a)
}
_, _ = handlers["queue"](context.Background(), json.RawMessage(`{}`))
after, _ := bus.Varz(nil)
if opened := after.TotalConnections - before.TotalConnections; opened != 0 {
t.Fatalf("the serving controller opened %d connection(s) of its own", opened)
}
}
// A verb that still runs as a process of its own says which in the bus's list of connections.
func TestAVerbsOwnProcessNamesItsConnection(t *testing.T) {
bus := testbus.Start(t)
setup, err := broker.Dial(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
if err := setup.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
setup.Close()
asAProcess(t, bus.ClientURL())
handlers, _, err := seatToolHandlers()
if err != nil {
t.Fatal(err)
}
// A silence acts, so it is still its own process; it dials, and finds no such condition.
_, _ = handlers["conditions"](context.Background(),
json.RawMessage(`{"silence":"bus.nothing.here","for":"1h","why":"a test"}`))
names := closedNames(t, bus)
found := false
for _, n := range names {
found = found || n == "mesh-controller verb conditions"
}
if !found {
t.Fatalf("the verb's connection was named %q", names)
}
}
func TestEachProcessNamesItsConnectionsForWhatItIs(t *testing.T) {
for _, c := range []struct{ command, verb, want string }{
{"serve", "", "mesh-controller serving"},
{"conditions", "conditions", "mesh-controller verb conditions"},
{"delivery", "delivery-check", "mesh-controller verb delivery-check"},
{"push", "", "mesh-controller command push"},
} {
if got := connectionName(c.command, c.verb); got != c.want {
t.Errorf("%s/%s: %q", c.command, c.verb, got)
}
}
}
// D15: a user whose connections are dropped more than twelve times an hour is said, in plain words; one
// whose client closes its own short connections is not.
func TestAClientReconnectingInALoopIsSaid(t *testing.T) {
var read closedConnections
if err := json.Unmarshal([]byte(`{"hours":1,"reaches":true,"users":[
{"user":"ace.node-tools","closed":40,"dropped":40,"dropped_per_hour":40,"names":[
{"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]},
{"user":"controller","closed":300,"dropped":0,"names":[
{"name":"mesh-controller verb delivery-check","closed":300,"dropped":0,"reasons":{"Client Closed":300}}]},
{"user":"(no user: refused before it logged in)","closed":90,"dropped":90,"names":[{"name":"","closed":90,
"dropped":90,"reasons":{"Authentication Failure":90}}]},
{"user":"node.anchor","closed":5,"dropped":5,"names":[{"name":"mesh-host/anchor","closed":5,"dropped":5,
"reasons":{"Read Error":5}}]}]}`), &read); err != nil {
t.Fatal(err)
}
found := reconnecting(read)
if len(found) != 1 {
t.Fatalf("%+v", found)
}
o := found[0]
if o.ID != "ace.node-tools" || o.Machine != "ace" || o.Kind != kindBusReconnects ||
o.Headline != "The tool runner on ace keeps losing the bus" || !strings.Contains(o.Said, "Stale Connection ×30") {
t.Fatalf("%+v", o)
}
if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Explanation: o.Explanation,
Resolved: o.Resolved}, "ace"); !ok {
t.Fatalf("not plain: %s", why)
}
}
// The bus's module is asked on its machine, over the bus.
func TestTheBusModuleIsAskedWhoClosedConnections(t *testing.T) {
bus := testbus.Start(t)
conn, err := nats.Connect(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
defer conn.Close()
sub, err := conn.Subscribe(link.ModuleToolOn("nats", "nats_closed_connections", "anchor"), func(m *nats.Msg) {
_ = m.Respond([]byte(`{"result":{"hours":1,"reaches":true,"users":[{"user":"controller","closed":2,"dropped":0}]}}`))
})
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
read, err := askClosed(context.Background(), conn, "anchor")
if err != nil {
t.Fatal(err)
}
if len(read.Users) != 1 || read.Users[0].Closed != 2 {
t.Fatalf("%+v", read)
}
}
// A bus module older than the tool answers nothing to the question: the probe passes over it quietly.
func TestAnOlderBusModuleIsPassedOverQuietly(t *testing.T) {
bus := testbus.Start(t)
conn, err := nats.Connect(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
defer conn.Close()
found, err := reconnectsOn(context.Background(), conn, "anchor")
if err != nil || len(found) != 0 {
t.Fatalf("a bus module without the tool: %v %v", found, err)
}
}
// A verb answered in the serving controller says its flag errors in its answer, not in the controller's log.
func TestAFlagErrorIsSaidInTheAnswer(t *testing.T) {
bus := testbus.Start(t)
serving, err := broker.Dial(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
defer serving.Close()
was := servingBus.Load()
servingBus.Store(serving)
defer servingBus.Store(was)
answer, read := readHere(context.Background(), []string{"hand-acts", "--bogus"})
if !read || answer.OK || !strings.Contains(answer.Output, "flag provided but not defined") {
t.Fatalf("answered %+v", answer)
}
}
+15 -3
View File
@@ -710,15 +710,27 @@ func backlogObservation() []conditions.Observation {
}
n := 0
var machines []string
seen := map[string]bool{}
var modules []string
for node, moves := range backlogNow.waiting {
n += len(moves)
machines = append(machines, node)
for _, mv := range moves {
if !seen[mv.Module] {
seen[mv.Module] = true
modules = append(modules, mv.Module)
}
}
}
sort.Strings(machines)
return []conditions.Observation{{Scope: conditions.ScopeMesh, ID: "release", Token: "held", Kind: "release-held",
Severity: conditions.Warning, Resolver: conditions.ResolverOperator,
sort.Strings(modules)
o := conditions.Observation{Scope: conditions.ScopeMesh, ID: "release", Token: "held", Kind: "release-held",
Severity: conditions.Warning, Resolver: conditions.ResolverOperator, Also: machines,
Summary: fmt.Sprintf("%d build move(s) on %s wait for a gate and are not released: %s — `upgrade backlog` lists "+
"them, `upgrade release-backlog --why …` releases them", n, strings.Join(machines, ", "), backlogNow.held)}}
"them, `upgrade release-backlog --why …` releases them", n, strings.Join(machines, ", "), backlogNow.held)}
w := releaseHeldWords(modules, machines)
o.Headline, o.Explanation, o.Resolved, o.Needs = w.Headline, w.Explanation, w.Resolved, w.Needs
return []conditions.Observation{o}
}
// backlogCommand is `upgrade backlog`, read-only, and `upgrade release-backlog --why`.
+25 -1
View File
@@ -701,7 +701,6 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
}
}
now := time.Now().UTC()
step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound)
// **Sent with others, judged with them** (issue 281): the gate of the send that carried it is its
// verdict on its first machine. A failure there stopped the plan already.
if state.GatedBy != "" {
@@ -715,6 +714,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
}
passedWith(m, state, state.GatedBy, lead.Gate)
}
// Read after its pass is taken over from the send that carried it, so a passed gate is never judged
// again from the first machine's later reports (novox/hq issue 335).
step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound)
switch {
case step.failed != "":
// The first machine refused or failed what it was sent, or never said: the gate failed, and
@@ -974,6 +976,19 @@ func failFirstSend(ctx context.Context, open *stores, p *inventory.Plan, module
fmt.Printf("%s: %s\n", p.ID, p.Note)
}()
g := state.Gate
if g != nil && g.Verdict == inventory.GatePassed {
// **A guard: a build that passed its gate is never put back for what came after** (novox/hq issue 335).
// Not reached while nextRollout answers a passed gate with the rest to send, and advanceOnce reads it
// after a carried module takes over its lead's pass; it is here so that a path added later cannot
// overturn a verdict. If it is reached, the plan stops and says why, and the build is not marked failed.
// The module is left a stopped rollout (its Why said, sent first and not to the rest), which
// `plans retry` takes: it sends the first machine again, and the passed gate then sends the rest.
state.Why = "passed its gate; its walk then stopped: " + why
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s passed its gate on %s in tier %d (%s) and is kept; its walk stopped after: %s",
module, strings.Join(g.Machines, ", "), p.Tier, g.Why, why)
return
}
if g != nil && slices.Contains(g.Returned, module) {
// Put back at once when it broke: its rollback was made then, and is not made again.
state.Why = "put back when it broke; its send failed: " + g.Why
@@ -1061,6 +1076,15 @@ func nextRollout(s inventory.PlanModule, running []string, together bool, report
rest = append(rest, n)
}
}
// **A passed gate is the first machine's verdict, and its later reports are not** (novox/hq issue 335).
// Once the build passed there, what that machine reports next is about whatever it was sent after —
// another walk's send, a push — and says nothing of this build. On 2026-10-08 a build passed on the
// laptop, its send to the rest waited on another walk, that walk sent the laptop a new declaration it did
// not report for half an hour, and the plan read the silence as the first machine never applying the
// passed build: it marked it failed at its gate and put it back. The rest are sent, as the pass said.
if s.Gate != nil && s.Gate.Verdict == inventory.GatePassed {
return rolloutStep{send: rest}
}
var waiting, failed []string
for _, n := range s.First {
r, said := byNode[n]
+106
View File
@@ -0,0 +1,106 @@
package main
import (
"context"
"encoding/json"
"os/exec"
"path/filepath"
"testing"
"time"
"github.com/nats-io/nats-server/v2/server"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/testbus"
)
// novox/hq issue 327, replayed with only what the controller had before its fix, so it can be laid over
// the older commit. On 2026-10-08 the bus's connection total rose by 5.2 a minute while the same 16
// connections stayed open, and three calls of `conditions` in one second added three: each verb ran as a
// process of the controller's own binary, which dialled the bus, logged in and left. The operator's
// channel reads `conditions` at least once a minute. A verb that only reads, served by the serving
// controller, opens no connection of its own.
func TestReplay327(t *testing.T) {
bus := testbus.Start(t)
serving, err := broker.Dial(bus.ClientURL())
if err != nil {
t.Fatal(err)
}
defer serving.Close()
if err := serving.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
keeper, err := keeperOn(context.Background(), serving.Conn())
if err != nil {
t.Fatal(err)
}
// The serving controller: its keeper and its connection, as serve sets them.
keptBefore, connBefore := conditionsFrom, handActConn
conditionsFrom, handActConn = keeper, serving.Conn()
defer func() { conditionsFrom, handActConn = keptBefore, connBefore }()
// A verb that runs as a process of its own runs this controller's binary, on this bus.
asAProcess(t, bus.ClientURL())
handlers, _, err := seatToolHandlers()
if err != nil {
t.Fatal(err)
}
accepted := func() uint64 {
v, err := bus.Varz(nil)
if err != nil {
t.Fatal(err)
}
return v.TotalConnections
}
before := accepted()
for i := 0; i < 3; i++ {
answer, err := handlers["conditions"](context.Background(), json.RawMessage(`{}`))
if err != nil {
t.Fatal(err)
}
if a, ok := answer.(verbAnswer); !ok || !a.OK {
t.Fatalf("conditions answered %+v", answer)
}
}
if opened := accepted() - before; opened != 0 {
t.Fatalf("three conditions calls opened %d connection(s) to the bus; the serving controller is on it already",
opened)
}
}
// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url:
// built into the test's own directory, which goes with the test.
func asAProcess(t *testing.T, url string) {
t.Helper()
path := filepath.Join(t.TempDir(), "mesh-controller")
if out, err := exec.Command("go", "build", "-o", path, ".").CombinedOutput(); err != nil {
t.Fatalf("the controller could not be built to run a verb as its own process: %v: %s", err, out)
}
was := ownImage
ownImage = func() string { return path }
t.Cleanup(func() { ownImage = was })
t.Setenv(broker.NATSVar, url)
t.Setenv(broker.CertificateVar, "")
}
// closedNames are the names of the connections the bus saw closed.
func closedNames(t *testing.T, bus *server.Server) []string {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
var names []string
for time.Now().Before(deadline) {
connz, err := bus.Connz(&server.ConnzOptions{State: server.ConnClosed})
if err != nil {
t.Fatal(err)
}
names = names[:0]
for _, c := range connz.Conns {
names = append(names, c.Name)
}
if len(names) > 0 {
return names
}
time.Sleep(50 * time.Millisecond)
}
return names
}
+132
View File
@@ -0,0 +1,132 @@
package main
import (
"reflect"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// TestReplay335 replays novox/hq issue 335 (2026-10-08): dunst's new build passed its gate on the laptop
// (healthy 3 times over 2m5s); its send to the rest was refused while another walk's build waited on the
// workstation; that walk then sent the laptop a declaration the laptop did not report on; and thirty minutes
// after the first send the plan read the laptop's silence as the passed build never applied, marked it failed
// at its gate with the pass's own words, and put it back. Written with only what the controller had before
// its fix, so it is laid over the commit before.
func TestReplay335(t *testing.T) {
t.Run("the first machine's later reports", testAPassedGateIsNotJudgedAgainFromTheFirstMachinesLaterReports)
t.Run("a walk stopped after the pass", testAWalkStoppedAfterItsGatePassedKeepsThePass)
}
// novox/hq issue 335: a build that passed its gate on its first machine is not judged again from that
// machine's later reports. On 2026-10-08 the send to the rest waited on another walk, that walk sent the
// first machine a declaration it did not report for half an hour, and the plan failed the passed build at
// the wait's bound and put it back.
func testAPassedGateIsNotJudgedAgainFromTheFirstMachinesLaterReports(t *testing.T) {
sentAt := time.Date(2026, 10, 8, 16, 44, 45, 0, time.UTC)
judged := sentAt.Add(2 * time.Minute)
later := sentAt.Add(5 * time.Minute)
state := inventory.PlanModule{First: []string{"laptop"}, FirstAt: &sentAt,
Gate: &inventory.PlanGate{Machines: []string{"laptop"}, Verdict: inventory.GatePassed,
Why: "healthy 3 times over 2m5s", JudgedAt: &judged, Kept: true}}
running := []string{"laptop", "workstation"}
now := sentAt.Add(30*time.Minute + 9*time.Second)
for what, reports := range map[string][]inventory.Reported{
"no report about what it was sent since": {{Node: "laptop", At: &later, Outcome: inventory.OutcomeApplied, Current: false}},
"another send failed there since": {{Node: "laptop", At: &later, Outcome: inventory.OutcomeFailed, Current: true}},
"no report at all": nil,
} {
step := nextRollout(state, running, false, reports, now, 30*time.Minute)
if step.failed != "" || step.waiting != "" || !reflect.DeepEqual(step.send, []string{"workstation"}) {
t.Errorf("%s: %+v, want the rest sent as the pass said", what, step)
}
}
}
// novox/hq issue 335: whatever stops a walk after its gate passed, the passed build is not marked failed at
// its gate, nor put back.
func testAWalkStoppedAfterItsGatePassedKeepsThePass(t *testing.T) {
sentAt := time.Date(2026, 10, 8, 16, 44, 45, 0, time.UTC)
g := &inventory.PlanGate{Machines: []string{"laptop"}, Verdict: inventory.GatePassed, Why: "healthy 3 times over 2m5s"}
state := &inventory.PlanModule{Build: "build-1", First: []string{"laptop"}, FirstAt: &sentAt, Gate: g}
p := &inventory.Plan{ID: "plan-1", Modules: map[string]*inventory.PlanModule{"dunst": state}}
// No stores: a walk that keeps the pass touches none, and one that reaches for them is putting it back.
defer func() {
if r := recover(); r != nil {
t.Fatalf("the passed build was taken to be failed and put back: %v", r)
}
}()
failFirstSend(t.Context(), nil, p, "dunst", state, []string{"laptop"}, "laptop did not report it applied within 30m0s",
[]string{"workstation"})
if g.Verdict != inventory.GatePassed || g.Rollback != "" {
t.Fatalf("the passed gate became %q, rollback %q", g.Verdict, g.Rollback)
}
if strings.Contains(p.Note, "failed its gate") || strings.Contains(p.Note, "put back") ||
!strings.Contains(p.Note, "did not report it applied") {
t.Fatalf("the note reads %q", p.Note)
}
}
// novox/hq issue 335 review: **a module carried by its lead's send takes over the lead's pass before its own
// step is read.** The state is the one a step leaves when the lead's gate passed and the step ended before the
// carried module's turn (an error read before it, kept with the plan): the lead passed, the carried module has
// no gate of its own yet. Since then another send reached the first machine and it has not reported on it, and
// the first send is older than the wait for a first machine's report. Read before the pass is taken over, the
// carried module's step said the first machine never applied it, and the passed build was put back.
func TestACarriedModuleTakesItsLeadsPassBeforeItsStepIsRead(t *testing.T) {
tm := aTierMesh(t, "m01", "m02")
ctx := t.Context()
inv := tm.open.inventory
advancePlans(ctx, tm.open)
if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}}) {
t.Fatalf("sent %v: the first machine once, for the tier", tm.sent)
}
p := tm.plan(t)
lead, carried := p.Modules["m01"], p.Modules["m02"]
if lead.Gate == nil || carried.GatedBy != "m01" {
t.Fatalf("m02 is not carried by m01's send: lead %+v, carried %+v", lead.Gate, carried)
}
// The lead passed; the carried module's turn did not come. The first send is past the wait's bound.
sent := time.Now().UTC().Add(-planWaitBound - time.Minute)
judged := sent.Add(2 * time.Minute)
lead.FirstAt, carried.FirstAt, lead.Gate.Since = &sent, &sent, &sent
lead.Gate.Verdict, lead.Gate.Why, lead.Gate.JudgedAt, lead.Gate.Kept = inventory.GatePassed,
"healthy 3 times over 2m5s", &judged, true
carried.Gate = nil
if err := inv.SavePlan(ctx, &p); err != nil {
t.Fatal(err)
}
// Another walk's send reached anchor, which has not reported on it.
if err := inv.RecordSent(ctx, nodeID(t, tm.open, "anchor"), "d-anchor-elsewhere",
map[string]string{"m01": "c2", "m02": "c2"}); err != nil {
t.Fatal(err)
}
advancePlans(ctx, tm.open)
p = tm.plan(t)
if p.State == inventory.PlanFailed {
t.Fatalf("the plan failed after its gate passed: %s", p.Note)
}
if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}, {"laptop"}}) {
t.Fatalf("sent %v: the rest once, as the pass said", tm.sent)
}
current, err := inv.CurrentBuilds(ctx)
if err != nil {
t.Fatal(err)
}
for _, m := range []string{"m01", "m02"} {
if current[m].Commit != "c2" {
t.Errorf("%s was put back to %s after its gate passed", m, current[m].Commit)
}
if failed, _ := inv.GateFailed(ctx, "build-"+m+"-2"); failed {
t.Errorf("%s's build was marked failed at its gate after the gate passed", m)
}
}
if g := p.Modules["m02"].Gate; g == nil || g.Verdict != inventory.GatePassed {
t.Errorf("m02 did not take over m01's pass: %+v", g)
}
}
+4 -1
View File
@@ -131,7 +131,10 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines
// **Dialled the way the mesh dials it** — credential and pin — because a bare connect to a
// bus that requires TLS and a user fails at the handshake, and the check then reported a
// standing server as absent (seen live, 2026-09-28).
if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil {
// The serving controller's own connection answers it without a second (novox/hq issue 327).
if serving := servingBus.Load(); serving != nil && serving.Conn().IsConnected() {
state.ServerStanding = true
} else if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil {
state.ServerStanding = true
js.Close()
}
+106 -13
View File
@@ -7,6 +7,7 @@ import (
"errors"
"fmt"
"github.com/nats-io/nats.go/micro"
"io"
"os"
"os/exec"
"slices"
@@ -28,6 +29,9 @@ import (
// answer. It also means a refusal is the same refusal in the same words, because it is the same
// output.
// verbKey carries the verb a call is for, to the process it runs.
type verbKey struct{}
// verbAnswer is what a verb answers: what the command printed, whether it succeeded, and — where the
// command speaks JSON — the same as data.
type verbAnswer struct {
@@ -250,6 +254,8 @@ func (a *verbArguments) commandLine() ([]string, error) {
return argv, nil
case "tools":
return nil, errors.New("tools is answered from the records, not by a command")
case "dead-letters":
return nil, errors.New("dead-letters is answered by the serving controller, on its own connection, not by a command")
case "status":
return []string{"status", "--json"}, nil
case "nodes":
@@ -761,6 +767,28 @@ func (a *verbArguments) commandLine() ([]string, error) {
}
return argv, nil
case "settings":
// Every module's preferences, their defaults and each machine's value (novox/hq ADR 0262):
// the one interface for them, so no module builds a settings tool of its own. Asked for by
// name, or by naming no module, since a layer is always some module's.
if str("module") == "" && str("values") != "" {
return nil, errors.New("settings: a module is needed to set values; name it with module")
}
if str("module") == "" && on("clear") {
return nil, errors.New("settings: a module is needed to clear a layer; name it with module")
}
if list := str("list"); list != "" || str("module") == "" {
if list != "" && list != "preferences" {
return nil, fmt.Errorf("settings lists %q only; %q is not a listing", "preferences", list)
}
argv := []string{"settings", "preferences"}
if m := str("module"); m != "" {
argv = append(argv, m)
}
if n := str("node"); n != "" {
argv = append(argv, "--node", n)
}
return argv, nil
}
// `settings set|clear` at a shell (novox/hq issue 198). The values travel as an argument
// because a tool has no file to hand the command; the command reads either.
if err := need("module"); err != nil {
@@ -897,6 +925,12 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) {
caller = "a seat call whose caller the bus did not name"
}
cmd.Env = append(cmd.Env, link.CallerVar+"="+caller+", through the "+catalogue.ControllerSeatName+" seat")
// And which verb, so the connection it dials says so in the bus's list (novox/hq issue 327).
verb, _ := ctx.Value(verbKey{}).(string)
if verb == "" {
verb = argv[0]
}
cmd.Env = append(cmd.Env, verbVar+"="+verb)
// Two buffers, one answer. What the command *says* is both streams, in the order a person at
// a shell would read them; what it *answers as data* is standard output alone — `status --json`
// prints its warnings beside the document, and a JSON parsed from the two together parsed
@@ -905,18 +939,7 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) {
cmd.Stdout = &stdout
cmd.Stderr = &stderr
runErr := cmd.Run()
answer := verbAnswer{Output: stdout.String() + stderr.String(), OK: runErr == nil}
if jsonVerbs[argv[0]] && runErr == nil {
var parsed any
if json.Unmarshal(bytes.TrimSpace(stdout.Bytes()), &parsed) == nil {
answer.Answer = parsed
if overviewVerbs[argv[0]] {
// Once, as data: the same document again as text doubled an answer that already
// outgrew one message of the bus (novox/hq issue 314).
answer.Output = stderr.String() + "its answer, as data, is `answer`\n"
}
}
}
answer := answerOf(argv, stdout.Bytes(), stderr.String(), runErr == nil)
var exit *exec.ExitError
if runErr != nil && !errors.As(runErr, &exit) {
// Not the command refusing — the command not running at all, which is this process's fault.
@@ -931,6 +954,66 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) {
return answer, nil
}
// answerOf is what a command said, as a verb answers it: both streams as text, and standard output as data
// where the command speaks JSON.
func answerOf(argv []string, stdout []byte, stderr string, ok bool) verbAnswer {
answer := verbAnswer{Output: string(stdout) + stderr, OK: ok}
if jsonVerbs[argv[0]] && ok {
var parsed any
if json.Unmarshal(bytes.TrimSpace(stdout), &parsed) == nil {
answer.Answer = parsed
if overviewVerbs[argv[0]] {
// Once, as data: the same document again as text doubled an answer that already
// outgrew one message of the bus (novox/hq issue 314).
answer.Output = stderr + "its answer, as data, is `answer`\n"
}
}
}
return answer
}
// readHere answers a verb that only reads the bus in the serving controller itself, on its own connection
// and keeper (novox/hq issue 327): the same command, writing to the answer rather than to a process's
// output, so the answer is the one the command prints. False for any other command line, which runs as a
// command of its own. Every `conditions` call — the operator's channel reads it at least once a minute —
// was a process that dialled the bus, logged in and left.
func readHere(ctx context.Context, argv []string) (verbAnswer, bool) {
var read func(context.Context, []string, io.Writer) error
args := argv[1:]
// Each where this process holds what it reads: the serving keeper, the hand-act log's connection, the
// serving connection.
switch argv[0] {
case "conditions":
sub := "list"
if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
sub, args = args[0], args[1:]
}
if conditionsFrom != nil {
read = map[string]func(context.Context, []string, io.Writer) error{
"list": listConditions, "show": showCondition, "history": conditionHistory}[sub]
}
case "hand-acts":
if handActConn != nil || servingBus.Load() != nil {
read = listHandActs
}
case "queue":
if servingBus.Load() != nil {
read = listQueue
}
}
if read == nil {
return verbAnswer{}, false
}
var out bytes.Buffer
stderr := ""
err := read(ctx, args, &out)
if err != nil {
// As the command says it when it fails (main).
stderr = "mesh-controller: " + err.Error() + "\n"
}
return answerOf(argv, out.Bytes(), stderr, err == nil), true
}
// seatToolHandlers are the handlers for every verb the mesh-controller seat declares, from the
// store's row, so a verb the row does not carry is not served. A verb it carries that this binary
// cannot run is named at start and answers the reason when called — never a refusal to serve, which
@@ -960,6 +1043,9 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
if verb == "calls" {
return callsAnswer(link.Calls, a.given["call"])
}
if verb == "dead-letters" {
return deadLettersAnswer(ctx, a)
}
if verb == "doctor" {
// From the serving controller, which runs the self-check and hears the signals
// (novox/hq to-be 45 §4): the last verdict at once, or a run now.
@@ -1011,6 +1097,7 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
if err != nil {
return nil, err
}
ctx = context.WithValue(ctx, verbKey{}, verb)
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)
@@ -1019,6 +1106,9 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
// Whatever it did, `status` is composed again once it has.
defer statusFrom.nudge()
}
if answer, read := readHere(ctx, argv); read {
return answer, nil
}
if answersFirst(argv) {
// Before anything is sent: a push sends the bus's own machine first, and a broker
// reloading its user list forgets the answer it was about to permit (novox/hq issue 265).
@@ -1042,7 +1132,10 @@ func actsOnAPlan(args map[string]any) bool {
// inProcess are the verbs answered by this process rather than by a command it runs: `tools` from
// the records, `calls` from what this process served.
var inProcess = map[string]bool{"tools": true, "calls": true, "doctor": true}
var inProcess = map[string]bool{"tools": true, "calls": true, "doctor": true,
// What a consumer gave up on, read and changed on the serving controller's own connection (novox/hq
// issue 330).
"dead-letters": true}
// answersFirst is a command line whose caller is answered before it runs: a push, by its verb or
// through `command`. A push sends the machine holding the bus first when its user list changed, the
+49
View File
@@ -0,0 +1,49 @@
package main
import (
"flag"
"fmt"
"io"
"os"
"sync/atomic"
"github.com/novox/mesh-controller/internal/broker"
)
// The serving controller's own connection, lent to whatever it does for a moment (novox/hq issue 327).
//
// Every place that needed the bus for a moment dialled it: a verb's own process, and in the serving
// controller a watchdog tick reading whether the build seat paused, a walk's step reading readiness, a
// queue read. Each paid a connection, a TLS handshake and a login on the control node, and hundreds an
// hour hid in the server's connection total the one thing it would show: a client reconnecting in a loop.
// The serving controller is on the bus already; what it does is done on that connection.
// servingBus is the serving controller's connection, set when it starts serving; nil in every other
// process, which dials its own.
var servingBus atomic.Pointer[broker.JetStream]
// aBus is a connection for something done for a moment: the serving controller's own, lent — so its
// Close closes nothing — when this process is it, and otherwise one dialled for it, named for this
// process (broker.ConnectionName), which its Close closes.
func aBus() (*broker.JetStream, error) {
if serving := servingBus.Load(); serving != nil {
return broker.Borrow(serving), nil
}
address, err := broker.BusAddress()
if err != nil {
return nil, err
}
js, err := broker.Dial(address)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus: %w", err)
}
return js, nil
}
// usageTo sends a command's flag errors and usage to where its answer goes when that is not this process's
// output: a verb answered in the serving controller says them in its answer, not in the controller's log.
func usageTo(set *flag.FlagSet, w io.Writer) {
if w != os.Stdout {
set.SetOutput(w)
}
}
@@ -0,0 +1,109 @@
package main
import (
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// `settings` says each value and where it came from, and the default a layer overrides (novox/hq ADR 0262).
func TestSettingsSayWhereEachValueComesFrom(t *testing.T) {
got := describeEffective("dunst", "laptop", []catalogue.SettingSource{
{Key: "font-size", Value: float64(13), From: "laptop", Default: float64(10), HasDefault: true},
{Key: "width", Value: float64(250), From: catalogue.DefaultLayer, FromDefault: true, Default: float64(250), HasDefault: true},
})
want := "dunst on laptop, every value and where it comes from:\n" +
" font-size = 13 (laptop; the default is 10)\n" +
" width = 250 (default)\n"
if got != want {
t.Fatalf("said\n%s\nwant\n%s", got, want)
}
}
// `settings` with list "preferences" is the one listing of every module's preferences; module and node
// narrow it, and no other listing is taken.
func TestSettingsListPreferences(t *testing.T) {
for _, c := range []struct {
args map[string]any
want string
}{
{map[string]any{"list": "preferences"}, "settings preferences"},
{map[string]any{"list": "preferences", "module": "dunst"}, "settings preferences dunst"},
{map[string]any{"list": "preferences", "node": "laptop"}, "settings preferences --node laptop"},
} {
argv, err := argvFor("settings", c.args)
if err != nil || strings.Join(argv, " ") != c.want {
t.Errorf("%v: %v %v, want %s", c.args, argv, err, c.want)
}
}
for args, want := range map[string]map[string]any{
"a module is needed to set values": {"values": `{"width": 300}`},
"a module is needed to clear a layer": {"clear": "true"},
} {
if _, err := argvFor("settings", want); err == nil || !strings.Contains(err.Error(), args) {
t.Errorf("%v: %v, want %q", want, err, args)
}
}
if _, err := argvFor("settings", map[string]any{"list": "everything"}); err == nil {
t.Error("a listing other than preferences was taken")
}
if argv, err := argvFor("settings", map[string]any{}); err != nil || strings.Join(argv, " ") != "settings preferences" {
t.Errorf("settings naming no module is the listing: %v %v", argv, err)
}
}
func TestPreferencesSayEachMachinesValueAndItsSource(t *testing.T) {
m := catalogue.Manifest{Module: "dunst", Settings: map[string]catalogue.SettingDeclaration{
"font-size": {Kind: catalogue.KindPreference, Default: float64(10), Why: "readable at 100 DPI"},
"width": {Kind: catalogue.KindPreference, Default: float64(250), Why: "forty characters"},
}}
on := map[string][]catalogue.SettingSource{
"laptop": catalogue.Effective(m, []catalogue.Layer{{From: "laptop", Values: map[string]any{"font-size": float64(16)}}}),
"desk": catalogue.Effective(m, []catalogue.Layer{{From: catalogue.MeshWideLayer, Values: map[string]any{"width": float64(300)}}}),
}
got := describePreferences([]preferencesOf{{Manifest: m, Nodes: []string{"desk", "laptop"}, On: on}})
want := "dunst (on desk, laptop)\n" +
" font-size, default 10: readable at 100 DPI\n" +
" desk: 10 (default)\n" +
" laptop: 16 (the node)\n" +
" width, default 250: forty characters\n" +
" desk: 300 (the mesh)\n" +
" laptop: 250 (default)\n"
if got != want {
t.Fatalf("said\n%s\nwant\n%s", got, want)
}
if describePreferences(nil) != "no module declares a preference\n" {
t.Fatal("an empty listing")
}
}
// The listing over the real stores: each machine's value with its source; a machine names only the
// modules on it; a machine the mesh does not know is refused (novox/hq ADR 0262).
func TestPreferencesListedFromTheStores(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
register(t, open, catalogue.Manifest{Module: "notes", Version: "1",
Settings: map[string]catalogue.SettingDeclaration{
"font-size": {Kind: catalogue.KindPreference, Default: float64(10), Why: "readable at 100 DPI"},
},
Resources: []map[string]any{{"id": "rc", "type": "file", "path": "/etc/notes.conf", "mode": "0644",
"content": "font = ${setting:font-size}\n"}}})
if _, err := assign(ctx, open, "laptop", "notes"); err != nil {
t.Fatal(err)
}
if err := open.inventory.SetSettings(ctx, "laptop", "notes", map[string]any{"font-size": float64(16)}); err != nil {
t.Fatal(err)
}
all := stdoutOf(t, func() error { return settingsCommand(ctx, []string{"preferences"}) })
if !strings.Contains(all, "notes (on laptop)") || !strings.Contains(all, "laptop: 16 (the node)") ||
!strings.Contains(all, "font-size, default 10: readable at 100 DPI") {
t.Fatalf("the listing:\n%s", all)
}
if got := stdoutOf(t, func() error { return settingsCommand(ctx, []string{"preferences", "--node", "anchor"}) }); got != "no module on anchor declares a preference\n" {
t.Fatalf("a machine without the module:\n%s", got)
}
if err := settingsCommand(ctx, []string{"preferences", "--node", "nowhere"}); err == nil {
t.Fatal("a machine the mesh does not know was answered")
}
}
+87 -4
View File
@@ -154,8 +154,9 @@ var signalsTable = []signalRow{
}},
{Row: "S9", Signal: "bus advisories: maximum deliveries, consumer deleted; the controller's own slow " +
"consumer and refused subjects", Emitter: "bus server's advisory subjects; the controller's connection",
Trigger: "any", Bound: "any occurrence; clears after an hour without another, and a deleted consumer " +
"once it exists again or the mesh no longer expects it",
Trigger: "any", Bound: "any occurrence; clears after an hour without another, a deleted consumer " +
"once it exists again or the mesh no longer expects it, and a message given up on once DEAD_LETTERS " +
"no longer holds it (novox/hq issue 330)",
Kind: "slow-consumer, max-deliveries, refused, consumer-lost", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.advisoriesErr }, watch: watchAdvisories,
newest: func(f *signalFacts) time.Time {
@@ -485,12 +486,94 @@ func watchAdvisories(f *signalFacts) []conditions.Observation {
if a.ID == "controller" {
machine = f.host
}
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind,
Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said})
o := conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind, Token: a.Token,
Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said}
if a.Kind == link.AdvisoryMaxDeliveries && a.Token == link.AdvisoryNotKept {
o.Machine = consumerMachine(a.Stream, a.Consumer)
o.Headline = clip(conditions.Capital(fmt.Sprintf("%s gave up on a message, not kept",
consumerWho(a.Stream, a.Consumer))), 60)
o.Explanation = "A listener on the bus could not handle a message, and the mesh could not keep it " +
"for you yet. It tries again every minute."
o.Resolved = "Resolved: the message is kept"
}
out = append(out, o)
}
return append(out, watchDeadLetters(f)...)
}
// watchDeadLetters says each consumer that DEAD_LETTERS holds a message for (novox/hq issue 330): open
// while it holds any, so it clears when they are delivered again or dropped, never because the server
// stopped saying it.
func watchDeadLetters(f *signalFacts) []conditions.Observation {
keys := make([]string, 0, len(f.deadLetters))
for k := range f.deadLetters {
keys = append(keys, k)
}
sort.Strings(keys)
var out []conditions.Observation
for _, key := range keys {
n := f.deadLetters[key]
stream, consumer, _ := strings.Cut(key, ".")
messages, them := "a message", "it"
if n > 1 {
messages, them = fmt.Sprintf("%d messages", n), "them"
}
who := consumerWho(stream, consumer)
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: key, Kind: link.AdvisoryMaxDeliveries,
Machine: consumerMachine(stream, consumer), Severity: conditions.Warning,
Summary: fmt.Sprintf("%s gave up on %s; %s kept in %s until delivered again or dropped, with why, "+
"through the controller's dead-letters verb", link.ConsumerInWords(stream, consumer), messages,
map[bool]string{true: "they are", false: "it is"}[n > 1], broker.DeadLettersStream),
Said: fmt.Sprintf("%d held for %s", n, key),
Headline: clip(conditions.Capital(fmt.Sprintf("%s could not handle %s", who, messages)), 60),
Needs: fmt.Sprintf("deliver %s again or drop %s, from the mesh MCP server.", them, them),
Explanation: conditions.Capital(fmt.Sprintf("%s was handed %s several times and gave up, so what %s "+
"asked for was not done. The mesh keeps %s until you deliver %s again or drop %s.", who, messages,
them, them, them, them)),
Resolved: "Resolved: the messages it gave up on were delivered again or dropped"})
}
return out
}
// consumerWho is a durable consumer's holder as the operator says it: a module on its machine, the
// controller, or a seat's holders.
func consumerWho(stream, consumer string) string {
switch {
case consumer == broker.ControllerName:
return "the controller"
case strings.HasPrefix(stream, "SEAT_") && strings.HasSuffix(consumer, "_worker"):
seat := strings.ToLower(strings.ReplaceAll(strings.TrimSuffix(strings.TrimPrefix(consumer, "SEAT_"), "_worker"), "_", "-"))
return "the holder of " + seat
case stream == broker.EventsStream:
if node, module, ok := strings.Cut(consumer, "_"); ok {
return module + " on " + node
}
}
return "a listener on the bus"
}
// consumerMachine is the machine a module's consumer is on; empty for the others.
func consumerMachine(stream, consumer string) string {
if stream == broker.EventsStream {
if node, _, ok := strings.Cut(consumer, "_"); ok {
return node
}
}
return ""
}
// clip is words at most n characters long, cut at a word.
func clip(s string, n int) string {
if len(s) <= n {
return s
}
cut := s[:n]
if i := strings.LastIndex(cut, " "); i > 0 {
cut = cut[:i]
}
return cut
}
func watchSelfCheck(f *signalFacts) []conditions.Observation {
every := f.selfCheck.every
if every <= 0 {
+64
View File
@@ -0,0 +1,64 @@
package main
import (
"context"
"fmt"
"sort"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// A stored manifest with a key this controller does not know (novox/hq ADR 0262). The module is left out
// of every machine's declaration by name; this is the loud half: a condition per module until the
// controller is updated, or the module is registered again in a shape this controller reads.
const (
sourceUnknownFields = "the catalogue"
kindUnknownField = "unknown-field"
)
// unknownFieldObservations is one condition for each module of the catalogue whose stored manifest has
// a key this controller does not know.
func unknownFieldObservations(known map[string]catalogue.Manifest) []conditions.Observation {
names := make([]string, 0, len(known))
for name, m := range known {
if m.UnknownField() != "" {
names = append(names, name)
}
}
sort.Strings(names)
var out []conditions.Observation
for _, name := range names {
m := known[name]
out = append(out, conditions.Observation{
Scope: conditions.ScopeMesh, ID: name, Kind: kindUnknownField, Severity: conditions.Warning,
Resolver: conditions.ResolverOperator, Source: sourceUnknownFields,
Summary: catalogue.UnknownFieldReason(m),
Said: m.UnknownField(),
Headline: name + " is left out until the controller is updated",
Explanation: name + " uses a field this controller does not know. Until the controller is updated, nothing of it changes on its machines, and what it adds to other modules and the ports opened for it stop. Its data is still backed up as this controller reads it, which may not be what its newer version asks.",
Needs: "update the controller, or register " + name + " again at a version this controller knows.",
Resolved: "the controller reads " + name + " again",
})
}
return out
}
// raiseUnknownFields raises those conditions and clears the ones no longer true, on the controller's
// tick. A catalogue that could not be read raises and clears nothing: "none" is not said for "could not
// tell" (ADR 0227 rule 4).
func raiseUnknownFields(ctx context.Context, inv *inventory.Inventory) []string {
if conditionsFrom == nil {
return nil
}
known, err := inv.Catalogue(ctx)
if err != nil {
return []string{fmt.Sprintf("the catalogue could not be read to say which modules it cannot read: %v", err)}
}
if err := conditionsFrom.Reconcile(ctx, sourceUnknownFields, unknownFieldObservations(known)); err != nil {
return []string{fmt.Sprintf("the modules with a field this controller does not know could not be kept as conditions: %v", err)}
}
return nil
}
@@ -0,0 +1,50 @@
package main
import (
"encoding/json"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
)
// A module whose stored manifest has a key this controller does not know is a condition, in plain
// words, until it is read again; every other module raises nothing (novox/hq ADR 0262).
func TestAModuleWithAnUnknownFieldIsACondition(t *testing.T) {
var later, now catalogue.Manifest
if err := json.Unmarshal([]byte(`{"module": "dunst", "version": "2", "settings": {}, "a-field-from-later": 1}`), &later); err != nil {
t.Fatal(err)
}
if err := json.Unmarshal([]byte(`{"module": "xorg", "version": "1"}`), &now); err != nil {
t.Fatal(err)
}
observed := unknownFieldObservations(map[string]catalogue.Manifest{"dunst": later, "xorg": now})
if len(observed) != 1 || observed[0].ID != "dunst" || observed[0].Kind != kindUnknownField {
t.Fatalf("observed: %+v", observed)
}
o := observed[0]
if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Explanation: o.Explanation,
Resolved: o.Resolved, Needs: o.Needs}); !ok {
t.Fatalf("not plain: %s", why)
}
k, _ := withConditionsInMemory(t)
if err := k.Reconcile(t.Context(), sourceUnknownFields, observed); err != nil {
t.Fatal(err)
}
if _, open, _ := k.Get(t.Context(), o.Key()); !open {
t.Fatal("not raised")
}
if err := k.Reconcile(t.Context(), sourceUnknownFields, unknownFieldObservations(map[string]catalogue.Manifest{"xorg": now})); err != nil {
t.Fatal(err)
}
still, err := k.Open(t.Context())
if err != nil {
t.Fatal(err)
}
for _, c := range still {
if c.Key == o.Key() {
t.Fatalf("not cleared once read again: %+v", c)
}
}
}
+6
View File
@@ -75,6 +75,9 @@ type signalFacts struct {
advisories []link.Advisory
lostConsumers map[string]bool
// deadLetters are how many messages DEAD_LETTERS holds per consumer, by `<stream>.<consumer>`
// (novox/hq issue 330): each consumer's max-deliveries condition is open while it holds any.
deadLetters map[string]int
advisoriesErr error
selfCheck selfCheckFacts
@@ -332,6 +335,9 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
}
f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet))
f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories)
if f.advisoriesErr == nil && w.js != nil {
f.deadLetters, f.advisoriesErr = link.HeldDeadLetters(w.js.Context())
}
f.handActs, f.handActsErr = w.gatherHandActs(ctx, now)
f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last()
return f
+7 -3
View File
@@ -38,7 +38,8 @@ type Consumer struct {
Push bool
// AckWaitSeconds before an unacknowledged delivery is redelivered.
AckWaitSeconds int
// MaxDeliver before the message is dead-lettered; zero for the mesh's default.
// MaxDeliver is how often a message is handed over before the consumer gives it up; zero for no
// bound. What a consumer gives up is kept in DEAD_LETTERS by the controller (novox/hq issue 330).
MaxDeliver int
// MaxAckPending is how many deliveries the server lets stand unacknowledged at once; zero for
// the server's default, which is many. **One, for a consumer handled one at a time**
@@ -145,13 +146,16 @@ func ConsumerFor(p Principal) (Consumer, bool) {
return Consumer{}, false
}
sort.Strings(filters)
// And the events given up on and delivered again to this consumer alone (novox/hq issue 330).
filters = append(filters, AgainFilter(consumerDurable(p)))
return Consumer{
Name: consumerDurable(p),
Stream: consumerStream(p),
Filters: filters,
AckWaitSeconds: 30,
MaxDeliver: 5,
Why: "what " + p.Module + " declared it consumes; after max-deliver it dead-letters",
Why: "what " + p.Module + " declared it consumes; after max-deliver it gives an event up, and the " +
"controller keeps it in DEAD_LETTERS until a person delivers it again or drops it",
}, true
}
@@ -234,7 +238,7 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
//
// **No max-deliver, and a long ack wait.** A declaration is settled only after the node has applied
// it and reported, which is minutes on a machine pulling images; and a declaration the mesh cannot
// get a node to accept is not one to dead-letter, because the stream keeps only the newest per node
// get a node to accept is not one to give up on, because the stream keeps only the newest per node
// anyway — so there is exactly one message per node to redeliver, for as long as that node is away.
func NodeConsumer(node string) Consumer {
return Consumer{
+3 -2
View File
@@ -61,8 +61,9 @@ func TestAModuleGetsOneConsumerCarryingEveryFilter(t *testing.T) {
if !ok {
t.Fatal("a module that consumes got no consumer")
}
if len(c.Filters) != 2 {
t.Fatalf("expected both subjects as filters, got %v", c.Filters)
// Both, and its own share of what is delivered again (novox/hq issue 330).
if len(c.Filters) != 3 || c.Filters[2] != "mesh.again.one_audit.>" {
t.Fatalf("expected both subjects and its own again filter, got %v", c.Filters)
}
perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
+30 -2
View File
@@ -28,6 +28,8 @@ import (
type JetStream struct {
conn *nats.Conn
js nats.JetStreamContext
// borrowed is a connection lent by its owner (Borrow): closing it is the owner's.
borrowed bool
// Note is how this says something it decided not to fail over. Nil is silent, which is only
// right for a caller that has no way to report; the controller sets it.
Note func(string, ...any)
@@ -40,9 +42,17 @@ func (j *JetStream) note(format string, args ...any) {
}
}
// ConnectionName is what a connection this process dials says it is, in the server's list of
// connections: the controller's role and what it is doing (novox/hq issue 327). Every connection was
// named `mesh-controller` — the serving controller's, each verb's own process, the build agents' — so the
// server's list could not say which was which. The process sets it once, at its start; an option a
// caller passes to Dial names one connection otherwise.
var ConnectionName = "mesh-controller"
// Dial connects and returns the controller's JetStream handle.
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
// The name first, so a name the caller gives is the one that stands.
opts = append([]nats.Option{nats.Name(ConnectionName), nats.Timeout(10 * time.Second)}, opts...)
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
// checks nothing else, and so does this). Without this, the first connection failed with
@@ -130,11 +140,20 @@ func (j *JetStream) Conn() *nats.Conn { return j.conn }
func (j *JetStream) Context() nats.JetStreamContext { return j.js }
func (j *JetStream) Close() {
if j.conn != nil {
if j.conn != nil && !j.borrowed {
j.conn.Close()
}
}
// Borrow is the same connection for a caller that will close what it was handed when it is done: its
// Close closes nothing, and the connection stays its owner's (novox/hq issue 327). How the serving
// controller lends its own connection to work that would otherwise dial one of its own.
func Borrow(j *JetStream) *JetStream {
lent := *j
lent.borrowed = true
return &lent
}
// EnsureStream creates the stream if it is absent and brings it to match if it is present.
//
// **Idempotent, because the controller asserts on every start** rather than creating once at
@@ -154,6 +173,15 @@ func (j *JetStream) EnsureStream(s Stream) error {
Description: s.Why,
}
want.AllowDirect = s.Direct
if s.MaxBytes > 0 {
want.MaxBytes = s.MaxBytes
}
if s.DiscardNew {
want.Discard = nats.DiscardNew
}
if s.DuplicatesSeconds > 0 {
want.Duplicates = time.Duration(s.DuplicatesSeconds) * time.Second
}
if s.Retention == RetentionLastPerSubject {
// Last-per-subject is a limits stream with one message kept per subject, not a
// retention policy of its own — the state shape, spelled the way the server spells it.
+31
View File
@@ -3,6 +3,8 @@ package broker
import (
"testing"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/testbus"
)
@@ -66,3 +68,32 @@ func TestAgainstARealServer(t *testing.T) {
}
})
}
// A connection says what it is in the server's list (novox/hq issue 327): the process's name, unless the
// caller names this one; and a lent connection's Close leaves its owner's open.
func TestAConnectionIsNamedAndALentOneIsNotClosed(t *testing.T) {
url := testbus.URL(t)
was := ConnectionName
ConnectionName = "mesh-controller verb conditions"
defer func() { ConnectionName = was }()
named, err := Dial(url)
if err != nil {
t.Fatal(err)
}
defer named.Close()
if got := named.Conn().Opts.Name; got != "mesh-controller verb conditions" {
t.Errorf("named %q", got)
}
lease, err := Dial(url, nats.Name("mesh-controller serving lease"))
if err != nil {
t.Fatal(err)
}
defer lease.Close()
if got := lease.Conn().Opts.Name; got != "mesh-controller serving lease" {
t.Errorf("a name the caller gave became %q", got)
}
Borrow(named).Close()
if !named.Conn().IsConnected() {
t.Error("closing a lent connection closed its owner's")
}
}
+5
View File
@@ -412,6 +412,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// both in the mesh's own account; the controller says each as a condition in the mesh's words.
// Named, not `$JS.EVENT.>`: the other advisories are every API call the mesh makes.
sub = append(sub, BusAdvisories...)
// **And what a consumer gave up on, kept and delivered again** (novox/hq issue 330): the notice
// acknowledged once the message is copied, the copy kept, and an event delivered again to the one
// consumer that gave it up. An ask to a seat is not delivered again, so no seat's queue is granted.
pub = append(pub, "$JS.ACK."+DeadLetterNoticesStream+"."+ControllerName+".>", deadLetterPrefix+">",
againPrefix+">")
case KindPerson:
// Tools, and nothing else. Every subject a person may publish is a tool call; a person
+134 -6
View File
@@ -3,11 +3,13 @@ package broker
import (
"fmt"
"sort"
"strings"
)
// The mesh's own streams.
//
// **These four and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118). An earlier
// **These and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118; the two that keep what a
// consumer gave up on added for issue 330). An earlier
// reading had the controller create *every* stream at genesis, from a fixed set. That is only the
// mesh's own half: a seat's streams are created when the module declaring it is registered, and a
// module's durable consumers when it is assigned — neither of which has happened at genesis. What
@@ -50,11 +52,80 @@ type Stream struct {
// Direct lets a client read a subject's last message without a consumer, which is how a
// runtime reads its own membership with no JetStream API beyond one request (ADR 0160).
Direct bool
// MaxBytes bounds the stream's size, zero for unbounded. With DiscardNew a full stream refuses
// what comes next rather than dropping what it holds: the publisher is told, and says so.
MaxBytes int64
DiscardNew bool
// DuplicatesSeconds is the window in which a message id published twice is kept once; zero for the
// server's default (two minutes).
DuplicatesSeconds int
}
// AssignmentsStream holds every assignment's membership, the newest per subject.
const AssignmentsStream = "ASSIGNMENTS"
// What a durable consumer gave up on is kept (novox/hq issue 330, design 25 §3).
//
// **The server says it and keeps it; the controller copies it.** A consumer that handed a message over
// as often as it may stops offering it and publishes `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES` with
// the stream and the message's sequence. DeadLetterNoticesStream captures those advisories as the
// server publishes them, so one said while no controller listens is still there when one starts. The
// controller's consumer on it fetches the given-up message by its sequence, while the source stream
// still holds it, and keeps a copy in DeadLettersStream under DeadLetterSubject, with the consumer,
// the subject, how often it was handed over and when it was given up. It stays there until a person
// delivers it again or drops it, with why; a condition is open for as long as it does.
const (
DeadLetterNoticesStream = "DEAD_LETTER_NOTICES"
DeadLettersStream = "DEAD_LETTERS"
// MaxDeliveriesAdvisories is the subject the server says a given-up message on.
MaxDeliveriesAdvisories = "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>"
deadLetterPrefix = "mesh.events.dead."
againPrefix = "mesh.again."
// DeadLettersBytes bounds the kept copies. Full, the stream refuses the next copy, which is said
// as the consumer's condition; it never drops one it holds.
DeadLettersBytes = 256 << 20
)
// DeadLetterSubject is where a message one consumer gave up on is kept: `mesh.events.dead.<stream>.<consumer>`.
func DeadLetterSubject(stream, consumer string) string {
return deadLetterPrefix + stream + "." + consumer
}
// DeadLetterOf is the stream and consumer a kept message's subject names; false for any other subject.
func DeadLetterOf(subject string) (stream, consumer string, ok bool) {
rest, found := strings.CutPrefix(subject, deadLetterPrefix)
if !found {
return "", "", false
}
stream, consumer, ok = strings.Cut(rest, ".")
return stream, consumer, ok && stream != "" && consumer != "" && !strings.Contains(consumer, ".")
}
// AgainSubject is where an event given up on is delivered again to the one consumer that gave it up,
// and to nobody else: `mesh.again.<consumer>.` and the original subject without its `mesh.`. Every
// consumer on EVENTS filters its own (AgainFilter), and a module's runtime reads the event's key from
// the tokens around `.event.`, so the handler sees the same key it saw the first time.
func AgainSubject(consumer, original string) string {
return againPrefix + consumer + "." + strings.TrimPrefix(original, "mesh.")
}
// AgainFilter is the one consumer's share of the subjects events are delivered again on.
func AgainFilter(consumer string) string { return againPrefix + consumer + ".>" }
// OriginalOfAgain is the subject an event delivered again was first published on; false for a subject
// that is not one delivered again.
func OriginalOfAgain(subject string) (string, bool) {
rest, found := strings.CutPrefix(subject, againPrefix)
if !found {
return "", false
}
_, original, ok := strings.Cut(rest, ".")
if !ok || original == "" {
return "", false
}
return "mesh." + original, true
}
// MeshStreams is the foundation set, in the order a person reads it.
//
// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live
@@ -97,7 +168,9 @@ func MeshStreams() []Stream {
// A seat's own events ride here too: they are 1:many like any event, and the
// `event` token keeps them clear of both the seat's work queue (`accept`) and its
// tools (`tool`), which must not be persisted.
Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>"},
// And an event given up on, delivered again to the one consumer that gave it up
// (novox/hq issue 330): under `mesh.again.<consumer>.`, which only that consumer filters.
Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>", againPrefix + ">"},
Retention: RetentionLimits,
MaxAge: 7 * 24 * 60 * 60,
MaxMsgsPerSubject: 10000,
@@ -105,6 +178,26 @@ func MeshStreams() []Stream {
"excluded by the event token; per-subject caps keep a noisy emitter from " +
"evicting a quiet one without splitting the stream",
},
{
Name: DeadLetterNoticesStream,
Subjects: []string{MaxDeliveriesAdvisories},
Retention: RetentionWorkQueue,
MaxAge: 7 * 24 * 60 * 60,
Why: "the server's word that a consumer gave up on a message, kept until the controller has " +
"copied the message into DEAD_LETTERS (novox/hq issue 330); a week, the longest the source " +
"streams keep what they are about",
},
{
Name: DeadLettersStream,
Subjects: []string{deadLetterPrefix + ">"},
Retention: RetentionLimits,
MaxBytes: DeadLettersBytes,
DiscardNew: true,
DuplicatesSeconds: 24 * 60 * 60,
Why: "every message a consumer gave up on, with its consumer, subject, deliveries and when, kept " +
"until a person delivers it again or drops it with why (novox/hq issue 330); no age, and full " +
"it refuses the next copy rather than drop one it holds",
},
}
}
@@ -298,7 +391,7 @@ const EventsStream = "EVENTS"
//
// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's,
// not the server's (window.go): a message is held with a nak-and-delay until the controller either
// takes it or gives up and says so. A max-deliver here would dead-letter a push that was being
// takes it or gives up and says so. A max-deliver here would give up on a push that was being
// held through a store restart — the exact message the stream exists to protect — some minutes
// before the controller had finished deciding about it.
func MeshConsumers() []Consumer {
@@ -314,7 +407,7 @@ func MeshConsumers() []Consumer {
{
Name: ControllerName,
Stream: "EVENTS",
Filters: ControllerFollows,
Filters: append(append([]string(nil), ControllerFollows...), AgainFilter(ControllerName)),
Push: true,
AckWaitSeconds: 30,
MaxDeliver: 5,
@@ -331,12 +424,47 @@ func MeshConsumers() []Consumer {
Resettable: "what it drops is caught up: merges by the catch-up pass (issue 266), build outcomes " +
"from the build records (issue 214), a provider's failing word said again (ADR 0224)",
Why: "the events the mesh's own controller reacts to, one at a time; after " +
"max-deliver it dead-letters, because an announcement it cannot act on will not " +
"become actionable",
"max-deliver it gives the event up, and the controller keeps it in DEAD_LETTERS until " +
"a person delivers it again or drops it",
},
// What the server said a consumer gave up on (novox/hq issue 330): copied into DEAD_LETTERS and
// acknowledged. No max-deliver: a notice the controller could not copy is offered again, and said.
{
Name: ControllerName,
Stream: DeadLetterNoticesStream,
Push: true,
AckWaitSeconds: 30,
Why: "the controller copies each message a consumer gave up on into DEAD_LETTERS; no max-deliver, " +
"because a notice it gave up on would lose the message it is about",
},
}
}
// NoticesConsumer is the controller's consumer on DEAD_LETTER_NOTICES (novox/hq issue 330).
func NoticesConsumer() Consumer {
for _, c := range MeshConsumers() {
if c.Stream == DeadLetterNoticesStream {
return c
}
}
panic("the mesh's consumers carry none on " + DeadLetterNoticesStream)
}
// AssertServingConsumers are the controller's own consumers its serving cannot go without: all but the
// one on DEAD_LETTER_NOTICES, which the keeper of dead letters asserts and retries by itself, so a fault
// there never stops the controller serving (novox/hq issue 330).
func AssertServingConsumers(e Ensurer) error {
for _, c := range MeshConsumers() {
if c.Stream == DeadLetterNoticesStream {
continue
}
if err := e.EnsureConsumer(c); err != nil {
return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err)
}
}
return nil
}
// Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason
// Asserter is.
type Ensurer interface {
+4
View File
@@ -127,6 +127,10 @@ func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) {
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
"ASSIGNMENTS": RetentionLastPerSubject,
// What a consumer gave up on (novox/hq issue 330): the server's notice taken once, the message
// kept until somebody acts.
"DEAD_LETTER_NOTICES": RetentionWorkQueue,
"DEAD_LETTERS": RetentionLimits,
}
got := map[string]Retention{}
for _, s := range MeshStreams() {
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "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.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] }
subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "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" }
} }
+5
View File
@@ -133,6 +133,11 @@ var WritersTable = []WriterRow{
Others: "—"},
{State: "the facts snapshot", Writer: "controller", KeptIn: "the artifact store, facts/latest",
Others: "the build seat reads"},
// What a consumer gave up on (novox/hq issue 330): copied by the controller from the server's notice,
// and delivered again or dropped only through its verb, with why.
{State: "a message a consumer gave up on", Writer: "controller", KeptIn: "the bus, the stream " + DeadLettersStream,
Others: "read, delivered again or dropped through the controller's dead-letters verb",
Subjects: []string{deadLetterPrefix + ">", againPrefix + ">"}, Writes: isController},
}
// CheckWriters refuses a grant that lets a principal publish on a subject the writers table gives
+1
View File
@@ -30,6 +30,7 @@ var designRows = []string{
"a provider's standing",
"the operator-channel's open messages",
"the facts snapshot",
"a message a consumer gave up on",
}
func TestTheWritersTableIsTheDesigns(t *testing.T) {
+1 -1
View File
@@ -238,7 +238,7 @@ func (r Resolution) derivedFor(provision, as, consumer, local string, settings S
"different things and nothing would compare them (novox/hq ADR 0202)",
consumer, local, provision, m.Module, orNothing(sortedAnyKeys(names)))
}
settled, err := Settle(names, settings[m.Module])
settled, err := Settle(names, WithDefaults(m, settings[m.Module]))
if err != nil {
return nil, fmt.Errorf("%s serving %s: %w", m.Module, provision, err)
}
+1 -1
View File
@@ -159,7 +159,7 @@ func (b *DataBackup) UnmarshalJSON(raw []byte) error {
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(&full); err != nil {
return fmt.Errorf("a data item's backup is \"copy\", \"none\" or {dump, into}: %w", err)
return fmt.Errorf("a data item's backup is \"copy\", \"none\" or {dump, into}: %w", typedUnknown(err))
}
*b = DataBackup{Dump: full.Dump, Into: full.Into}
return nil
+14 -4
View File
@@ -318,6 +318,10 @@ func (e *NotMadeError) Error() string {
func (r Resolution) LeftOut(settings SettingsBy, adopted bool) map[string]string {
out := map[string]string{}
for _, m := range r.Modules {
if why := UnknownFieldReason(m); why != "" {
out[m.Module] = why
continue
}
if err := JudgeSettings(m, settings[m.Module], adopted); err != nil {
out[m.Module] = err.Error()
}
@@ -365,11 +369,16 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// Before placing, because a placement is a setting too.
left := r.LeftOut(with.Settings, with.Adopted)
kept := make([]Manifest, 0, len(r.Modules))
var stillBackedUp []Manifest
for _, m := range r.Modules {
if why, isLeft := left[m.Module]; isLeft {
if leftOut != nil {
leftOut[m.Module] = why
}
// **Its data is still copied** (novox/hq ADR 0262): a module left out runs nothing new, and
// the data it already holds on the machine is the reason to keep copying it. Only its data,
// as the backup holder's lines are derived from it, and the directories they name.
stillBackedUp = append(stillBackedUp, backupView(m))
continue
}
kept = append(kept, m)
@@ -907,7 +916,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// **An operator's value, from the assignment** (novox/hq ADR 0112, ADR 0155): what a
// definition may not carry because it is true of one installation only. Filled from
// the same layers a mergeable file takes, and refused when no layer set it.
if err := settingInto(copied, with.Settings[m.Module], m.Module); err != nil {
if err := settingInto(copied, WithDefaults(m, with.Settings[m.Module]), m.Module); err != nil {
return nil, err
}
// **Placed before anything reads a path.** A pathless directory receives the path
@@ -975,7 +984,8 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// seat that places them (novox/hq ADR 0203, ADR 0204). Gathered from every module on
// the node, as the jails are, and **last of every placeholder pass**: shell code is a
// shell's own syntax, full of `${…}` no pass above should ever be shown.
if err := contributionsInto(copied, m, r.Modules, thisMachine, with, r.Capabilities, unplaced); err != nil {
contributing := append(append([]Manifest(nil), r.Modules...), stillBackedUp...)
if err := contributionsInto(copied, m, contributing, thisMachine, with, r.Capabilities, unplaced); err != nil {
return nil, err
}
copied["id"] = m.Module + "." + fmt.Sprint(resource["id"])
@@ -1495,7 +1505,7 @@ func (r Resolution) composed(m Manifest, to string, raw map[string]any, layers [
// Overridden, not merged: a setting changes a key the contribution declares and adds none.
// The provider reads the contribution as a contract, and a setting made for one of this
// module's files is no part of it (novox/hq 04-ISSUES/173).
values, err := overridden(raw, layers, what)
values, err := overridden(raw, WithDefaults(m, layers), what)
if err != nil {
return nil, fmt.Errorf("%s: %w", what, err)
}
@@ -1926,7 +1936,7 @@ func (r Resolution) servedOnThisMachine(provision string, with Rendering) (map[s
// one, so keep looking rather than concluding from the first.
continue
}
settled, err := Settle(serves, with.Settings[m.Module])
settled, err := Settle(serves, WithDefaults(m, with.Settings[m.Module]))
if err != nil {
return nil, false, fmt.Errorf("%s serving %s: %w", m.Module, provision, err)
}
+4 -1
View File
@@ -153,7 +153,10 @@ func walk(node any, at string, meant map[string]bool, visit func(at, value strin
}
sort.Strings(keys)
for _, k := range keys {
if prose[k] || k == NamesOnPurpose || (at == "" && k == "module") {
// Prose is a string a person reads. A key that is called `why` or `description` and holds
// anything else — a setting of that name, whose default the mesh writes — is walked like any
// other (novox/hq ADR 0262).
if _, isString := v[k].(string); (prose[k] && isString) || k == NamesOnPurpose || (at == "" && k == "module") {
continue
}
child := at + "." + k
+61 -6
View File
@@ -197,7 +197,7 @@ func (i *OfferIdentity) UnmarshalJSON(raw []byte) error {
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)
return fmt.Errorf("an offer's identity is false or {max, in}: %w", typedUnknown(err))
}
*i = OfferIdentity{Max: full.Max, In: full.In}
return nil
@@ -316,7 +316,7 @@ func (o *Offer) UnmarshalJSON(raw []byte) error {
dec.DisallowUnknownFields()
if err := dec.Decode(&full); err != nil {
return fmt.Errorf("a provided name is either a string or {name, scope, credential, reach, identity, "+
"keeps-consumer-data}: %w", err)
"keeps-consumer-data}: %w", typedUnknown(err))
}
o.Name, o.Scope, o.Credential, o.Reach, o.Identity = full.Name, full.Scope, full.Credential, full.Reach, full.Identity
o.KeepsConsumerData = full.Keeps
@@ -458,6 +458,24 @@ type Manifest struct {
// (novox/hq ADR 0201). Not history — that is an event — and never a secret, sealed or not.
State []StateDeclaration `json:"state,omitempty"`
// Settings are the defaults this module gives its settings (novox/hq ADR 0262): each key a file,
// a contribution or a served fact asks for as `${setting:<key>}`, its default, and why that
// default. Only a preference has one — a font size, a width, a number of workers — and a value
// that is inherently the operator's (a domain, an identity, a secret) is declared nowhere here and
// stays refused by name until a layer sets it. The mesh's layer, then the node's, override it.
Settings map[string]SettingDeclaration `json:"settings,omitempty"`
// unknown is the first key this manifest has that this controller does not know, at any depth, when
// it was read from the store (novox/hq ADR 0262): a manifest a newer controller registered.
// ParseManifest refuses it; the store's catalogue still loads, and the module is left out of every
// machine's declaration by name until the controller is updated (LeftOut).
unknown string
// bestEffort marks the view of a left-out module that only its backup lines are made from
// (backupView, novox/hq ADR 0262): a line of it that cannot be placed is said in the plan's list of
// what could not be placed, never an error that would cost the whole machine its declaration.
bestEffort bool
// Data is every kind of data this module keeps — its own, by directory, and what it keeps for
// its consumers, by provision — each with a class the mesh protects and watches it by (novox/hq
// ADR 0233). One list: the backup holder's lines, the bindings that do not move, what an
@@ -1178,7 +1196,13 @@ func ReceivedID(requirement string) string { return "received-" + requirement }
type manifestFields Manifest
// UnmarshalJSON reads `secrets` in both of its shapes — a path, or an object of local names to
// paths (ADR 0094) — and everything else exactly as the fields declare, unknown keys refused.
// paths (ADR 0094) — and everything else exactly as the fields declare.
//
// **An unknown key is kept aside, not refused here** (novox/hq ADR 0262). Registration and the module
// check refuse it, through ParseManifest. Reading the catalogue the store already holds does not: a
// manifest registered under a newer controller carries a field an older one does not know, and a
// strict read there failed the whole catalogue, and with it every plan and every send, the moment a
// controller was rolled back. UnknownField says what was set aside.
func (m *Manifest) UnmarshalJSON(raw []byte) error {
var keys map[string]json.RawMessage
if err := json.Unmarshal(raw, &keys); err != nil {
@@ -1259,10 +1283,30 @@ func (m *Manifest) UnmarshalJSON(raw []byte) error {
decoder := json.NewDecoder(bytes.NewReader(rest))
decoder.DisallowUnknownFields()
var fields manifestFields
unknown := ""
if err := decoder.Decode(&fields); err != nil {
return err
if asUnknownField(err) == nil {
return err
}
// Read without it, at whatever depth it is (prunedFields): the module is left out of every
// declaration by name (LeftOut), and still provides what it provides, and still has its data
// copied, so nothing that requires it is refused and nothing it holds goes uncopied.
unknown = err.Error()
var pruned []string
var perr error
if fields, pruned, perr = prunedFields(keys); perr != nil {
// Not read past: the manifest keeps its name and version alone. It is left out and raised
// all the same, and one manifest never fails the whole catalogue.
fields = manifestFields{}
_ = json.Unmarshal(keys["module"], &fields.Module)
_ = json.Unmarshal(keys["version"], &fields.Version)
unknown += " (read no further: " + perr.Error() + ")"
} else {
unknown += " (read without " + strings.Join(pruned, ", ") + ")"
}
}
*m = Manifest(fields)
m.unknown = unknown
if len(plain) > 0 {
m.Secrets = plain
}
@@ -1366,6 +1410,10 @@ func SecretLocal(to, local string) string {
return local
}
// UnknownField is the first key a leniently read manifest had that this controller does not know, as
// the JSON decoder words it, or "" when it had none (novox/hq ADR 0262).
func (m Manifest) UnknownField() string { return m.unknown }
func ParseManifest(raw []byte) (Manifest, error) {
var m Manifest
// Strictly. **An unknown key is refused**, which is the discipline the host's declaration
@@ -1377,7 +1425,12 @@ func ParseManifest(raw []byte) (Manifest, error) {
// checking whether something is restricted will find that it is, and be wrong.
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&m); err != nil {
err := decoder.Decode(&m)
if err == nil && m.unknown != "" {
// Kept aside by UnmarshalJSON for the stored catalogue's sake; registration refuses it.
err = errors.New(m.unknown)
}
if err != nil {
// A key that used to mean something says what it became. Refusing a renamed field with
// "unknown field" is correct and unhelpful: whoever wrote it knew what they meant, and
// the mesh knows what it is called now.
@@ -1541,6 +1594,8 @@ func ParseManifest(raw []byte) (Manifest, error) {
problems = append(problems, EventProblems(m)...)
// And what it may call its state, and whose it may read (state.go, novox/hq ADR 0201).
problems = append(problems, StateProblems(m)...)
// And the defaults it gives its settings (setting_defaults.go, novox/hq ADR 0262).
problems = append(problems, SettingProblems(m)...)
wellFormed := true
for _, c := range m.Claims {
if !name.MatchString(c.Name) {
@@ -2441,7 +2496,7 @@ func (o *OwnSecrets) UnmarshalJSON(raw []byte) error {
dec := json.NewDecoder(bytes.NewReader(body))
dec.DisallowUnknownFields()
if err := dec.Decode(&long); err != nil {
return fmt.Errorf("own-secrets.%s: a path, or {\"path\", \"taken\", \"issued-by\"}: %w", name, err)
return fmt.Errorf("own-secrets.%s: a path, or {\"path\", \"taken\", \"issued-by\"}: %w", name, typedUnknown(err))
}
out[name] = OwnSecret{Path: long.Path, Taken: long.Taken, IssuedBy: long.IssuedBy}
}
+68 -2
View File
@@ -1,6 +1,7 @@
package catalogue
import (
"slices"
"strings"
"testing"
)
@@ -42,7 +43,7 @@ func TestTheServiceManagerSeatServesTheUnitVerbs(t *testing.T) {
if seat.Scope != ScopeNode {
t.Fatalf("the service manager is a role each machine has once, and the seat is %s-scoped", seat.Scope)
}
want := []string{"units", "status", "start", "stop", "restart", "enable", "disable", "journal", "failed"}
want := []string{"units", "status", "start", "stop", "restart", "enable", "disable", "journal", "failed", "reset-failed", "wanted-by", "unlink-dangling"}
var got []string
for _, v := range seat.Serves {
got = append(got, v.Name)
@@ -80,13 +81,78 @@ func TestFailedIsAnOptionalVerbOfTheServiceManager(t *testing.T) {
if err := CanHold(holder(eight[1:]...), seat); err == nil || !strings.Contains(err.Error(), "does not serve units") {
t.Fatalf("a holder missing a required verb was accepted: %v", err)
}
optional := map[string]bool{"failed": true, "reset-failed": true, "wanted-by": true, "unlink-dangling": true}
for _, v := range seat.Serves {
if v.Optional != (v.Name == "failed") {
if v.Optional != optional[v.Name] {
t.Errorf("%s optional: %v", v.Name, v.Optional)
}
}
}
// **`reset-failed` and `wanted-by` join the seat optional** (novox/hq issue 332, ADR 0246 step 1): the
// holder running today, serving the nine verbs, still holds the seat; a holder serving all eleven is not
// refused; each replaces the shell command it names; and read back from the store's row they stay optional.
func TestResetFailedAndWantedByAreOptionalVerbsOfTheServiceManager(t *testing.T) {
defer UseSeats(DefaultSeats())
nine := []string{"units", "status", "start", "stop", "restart", "enable", "disable", "journal", "failed"}
holder := func(serves ...string) Manifest {
return Manifest{Module: "systemd", Version: "1", Claims: []Claim{{Name: ServiceManagerSeat, Scope: ScopeNode, Serves: serves}}}
}
seat, _ := SeatNamed(ServiceManagerSeat)
if err := CanHold(holder(nine...), seat); err != nil {
t.Fatalf("today's holder, without the two verbs, is refused: %v", err)
}
if err := CanHold(holder(append(nine, "reset-failed", "wanted-by")...), seat); err != nil {
t.Fatalf("a holder serving both is refused: %v", err)
}
if err := CanHold(holder(append(nine, "reset-failed", "wanted-by", "unlink-dangling")...), seat); err != nil {
t.Fatalf("a holder serving unlink-dangling too is refused: %v", err)
}
if err := CanHold(holder(append(nine, "why-started")...), seat); err == nil || !strings.Contains(err.Error(), "does not promise") {
t.Fatalf("a verb the seat does not promise was accepted: %v", err)
}
replaces := map[string]string{"reset-failed": "systemctl reset-failed", "wanted-by": "systemctl list-dependencies --reverse",
"unlink-dangling": "rm ~/.config/systemd/user/*.wants/<unit>"}
for _, v := range seat.Serves {
want, ok := replaces[v.Name]
if !ok {
continue
}
if !slices.Contains(v.Replaces, want) {
t.Errorf("%s does not say it replaces %q: %v", v.Name, want, v.Replaces)
}
props, _ := v.Input["properties"].(map[string]any)
if _, has := props["unit"]; !has {
t.Errorf("%s takes no unit", v.Name)
}
if v.Name == "unlink-dangling" {
// Removing is the person's act: their reason is required, and kept with the record.
if required, _ := v.Input["required"].([]string); !slices.Contains(required, "why") {
t.Errorf("unlink-dangling does not require why: %v", v.Input["required"])
}
}
delete(replaces, v.Name)
}
if len(replaces) > 0 {
t.Fatalf("the seat does not promise %v", replaces)
}
// The row as seeding stores it: every verb, the mark not a column.
var rows []Seat
for _, s := range DefaultSeats() {
row := s
row.Serves = nil
for _, v := range s.Serves {
row.Serves = append(row.Serves, Verb{Name: v.Name, Description: v.Description, Input: v.Input})
}
rows = append(rows, row)
}
UseSeats(rows)
seat, _ = SeatNamed(ServiceManagerSeat)
if err := CanHold(holder(nine...), seat); err != nil {
t.Fatalf("read back from the row, the two verbs are required: %v", err)
}
}
// The optional mark is not stored, so a seat set read back from the store's rows takes it from the
// compiled seat: otherwise `failed`, seeded into the row, would come back required.
func TestAnOptionalVerbStaysOptionalInASetReadFromTheStore(t *testing.T) {
+33 -11
View File
@@ -212,9 +212,20 @@ func seatContributions(modules []Manifest, holder Manifest, placeholder string,
return shapedContributions(modules, holder, facts["name"], s, r, where, caps)
}
var failed error
var unplacedLines []string
var b strings.Builder
for _, m := range inModuleOrder(modules) {
named := false
// A left-out module's lines are best effort (novox/hq ADR 0262): what cannot be placed is said,
// and the machine is declared without it. A placement setting that does not read is not
// replaced by the definition's own path, which may not be where the data is.
if m.bestEffort && r.Dirs {
if _, err := Places(m, with.Settings[m.Module]); err != nil {
unplacedLines = append(unplacedLines, fmt.Sprintf("%s's %s for %s, kept while it is left out, "+
"is not placed: its placement setting does not read (%v)", m.Module, kind, s.Name, err))
continue
}
}
for _, c := range m.allContributions() {
if c.Kind != kind || !capable(c, caps) {
continue
@@ -222,15 +233,12 @@ func seatContributions(modules []Manifest, holder Manifest, placeholder string,
if cs, known := SeatNamed(c.Seat); !known || cs.Name != s.Name {
continue
}
if !named {
fmt.Fprintf(&b, "%s %s\n", r.Comment, m.Module)
named = true
}
content := c.Content
var lineFailed error
if r.Dirs {
filled, err := dirFill(content, dirsFor(m, with), m.Module)
if err != nil && failed == nil {
failed = err
if err != nil && lineFailed == nil {
lineFailed = err
}
// An operator's path the module was given, as an item of data on it (novox/hq ADR 0233).
if accessRef.MatchString(filled) {
@@ -238,15 +246,15 @@ func seatContributions(modules []Manifest, holder Manifest, placeholder string,
if err == nil {
filled, err = accessFill(filled, byID, m.Module)
}
if err != nil && failed == nil {
failed = err
if err != nil && lineFailed == nil {
lineFailed = err
}
}
for _, key := range machineUsed(filled) {
value, has := facts[key]
if !has {
if failed == nil {
failed = fmt.Errorf("%s's %s for %s says ${machine:%s}, and this machine says %s",
if lineFailed == nil {
lineFailed = fmt.Errorf("%s's %s for %s says ${machine:%s}, and this machine says %s",
m.Module, kind, s.Name, key, orNothing(namesOfFacts(facts)))
}
continue
@@ -255,11 +263,25 @@ func seatContributions(modules []Manifest, holder Manifest, placeholder string,
}
content = filled
}
if lineFailed != nil {
if m.bestEffort {
unplacedLines = append(unplacedLines, fmt.Sprintf("%s's %s for %s, kept while it is left "+
"out, is not placed: %v", m.Module, kind, s.Name, lineFailed))
continue
}
if failed == nil {
failed = lineFailed
}
}
if !named {
fmt.Fprintf(&b, "%s %s\n", r.Comment, m.Module)
named = true
}
b.WriteString(content)
if !strings.HasSuffix(content, "\n") {
b.WriteString("\n")
}
}
}
return b.String(), nil, failed
return b.String(), unplacedLines, failed
}
+34 -1
View File
@@ -304,6 +304,11 @@ var defaultSeats = append([]Seat{
// Where it is held, its holder writes the resolver file and the uplink's holder steps back from it
// (node_resolver.go). It knows nothing of any VPN: its verbs route domains to servers over a link.
{Name: ResolverSeat, Scope: ScopeNode, Decision: "novox/hq ADR 0247", Serves: resolverVerbs()},
// A machine's shares and the shares it mounts (novox/hq ADR 0263): the holder of node-nfs-server
// exports a machine's folders to the private network and provides each as `nfs-share`; the holder of
// node-mounts writes a mount and an automount unit per share on a machine that asks (shares.go).
{Name: NFSServerSeat, Scope: ScopeNode, Decision: "novox/hq ADR 0263", Serves: nfsServerVerbs()},
{Name: MountsSeat, Scope: ScopeNode, Decision: "novox/hq ADR 0263", Serves: mountsVerbs()},
},
// The graphical session's roles (novox/hq ADR 0208), last because they are a workstation's.
graphicalSessionSeats()...)
@@ -620,7 +625,8 @@ func SeatsWithAProtocol() []Seat {
// 0177): the units on the machine in both scopes, read and acted on by name. Every verb takes an
// optional scope — "system" when absent, "user" for the operator account's own manager — so a
// caller asks for a user unit the way it asks for a system one; `failed` alone reads both managers
// when none is named, and is optional (Verb.Optional).
// when none is named. `failed`, `reset-failed`, `wanted-by` and `unlink-dangling` are optional
// (Verb.Optional) until every holder serves them (ADR 0246).
func serviceManagerVerbs() []Verb {
scoped := func(more map[string]string, required []string) map[string]any {
props := map[string]string{"scope": "\"system\" (the default) or \"user\": the operator account's own manager"}
@@ -669,6 +675,33 @@ func serviceManagerVerbs() []Verb {
"account's; a manager that does not answer is reported with its error, never as nothing failed.",
Input: schema(map[string]string{"scope": "\"system\" or \"user\": only that manager (both when absent)"}, nil),
Replaces: []string{"systemctl --failed"}},
// **Clearing a failed record, and finding what starts a unit** (novox/hq issue 332): a unit whose
// file is gone stays failed in its manager until the record is reset, and what keeps asking for it
// is a dependency or an enable link no verb could show. Both optional while their holders catch up
// (ADR 0246 step 1); a later change requires them.
{Name: "reset-failed", Optional: true, Description: "Clear one unit's failed record in its manager — " +
"its failed state and its start-limit count — and answer its state after. It removes no file and " +
"starts nothing.",
Input: scoped(unit, []string{"unit"}), Replaces: []string{"systemctl reset-failed"}},
{Name: "wanted-by", Optional: true, Description: "What wants, requires or triggers one unit, read-only: the " +
"manager's reverse dependencies (WantedBy, RequiredBy, UpheldBy, BoundBy, TriggeredBy and the rest) and " +
"every enable link naming it in the scope's configuration directories (*.wants, *.requires, *.upholds), " +
"each with where it points and whether it dangles — the account's own directory included in user scope.",
Input: scoped(unit, []string{"unit"}),
Replaces: []string{"systemctl list-dependencies --reverse", "systemctl show -p WantedBy",
"ls ~/.config/systemd/user/*.wants"}},
// **A dangling enable link removed on the person's word** (novox/hq issue 332): `systemctl disable`
// leaves a link whose unit file is gone, so the manager keeps asking for the unit. Optional until its
// holders serve it (ADR 0246 step 1).
{Name: "unlink-dangling", Optional: true, Description: "Remove the dangling enable links named for one " +
"unit — links in a *.wants, *.requires or *.upholds directory of the account's or the machine's own " +
"configuration whose target does not exist, which `disable` leaves once the unit's file is gone — on " +
"the person's word: `why` is required. Each is recorded (its path, its target, the reason) before it is " +
"removed, and the manager reloaded; a link whose target exists, or one a package, a generator or the " +
"runtime placed, is left and said.",
Input: scoped(map[string]string{"unit": unit["unit"],
"why": "the person's reason for removing the links, kept in the record"}, []string{"unit", "why"}),
Replaces: []string{"rm ~/.config/systemd/user/*.wants/<unit>", "unlink /etc/systemd/system/*.wants/<unit>"}},
}
}
+3 -3
View File
@@ -46,7 +46,7 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
delivered[s.Delivers] = s.Name
}
}
// Thirty-eight with node-resolver (novox/hq ADR 0247); thirty-seven with mesh-delivery (novox/hq ADR 0239); thirty-six since node-resolver-config retired
// Forty with node-nfs-server and node-mounts (novox/hq ADR 0263); thirty-eight with node-resolver (novox/hq ADR 0247); thirty-seven with mesh-delivery (novox/hq ADR 0239); thirty-six since node-resolver-config retired
// into node-uplink (novox/hq ADR 0223); thirty-seven
// since the retired node-dns-resolver went (novox/hq ADR 0220); thirty-eight with
// node-backup (novox/hq ADR 0214); thirty-seven with node-message-bus (novox/hq ADR 0215);
@@ -57,8 +57,8 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
// node-container-runtime (ADR 0207); nineteen with node-environment and node-login-shell (ADR 0203,
// ADR 0204); seventeen with node-build-agent (ADR 0190). One fewer once the retired
// mesh-build-machine row goes, when no registered manifest claims it.
if len(Seats()) != 38 {
t.Errorf("the mesh defines %d seats rather than 38; the set is closed, so a change here is "+
if len(Seats()) != 40 {
t.Errorf("the mesh defines %d seats rather than 40; the set is closed, so a change here is "+
"a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames())
}
}
+234
View File
@@ -0,0 +1,234 @@
package catalogue
import (
"fmt"
"sort"
"strings"
)
// A preference has a default; the operator's own value has none (novox/hq ADR 0262, extending ADR
// 0112 and taking the default half of ADR 0164).
//
// `${setting:<key>}` was refused whenever no layer set the key, and `settings set` refused a key no
// file asked for yet. Together they meant a running module could never gain a setting: the file that
// asks for it fails to compose until somebody sets it, and nobody can set it until the file asks.
// A definition may now give a key a default, with why, when the value is a **preference** — a font
// size, a width, a number of workers: something true of the software that a machine may tune. A
// value that is inherently the operator's — a domain, a public name, an identity, a secret — gets no
// default and stays refused by name until a layer sets it, which is what ADR 0112 and ADR 0155 were
// written for.
//
// The default is the lowest layer. The mesh's layer, then the node's, override it. It fills only
// `${setting:<key>}`: a mergeable file's content already is its defaults, and laying a default over
// it would reach every mergeable file of the module (novox/hq issue 168).
// SettingDeclaration is one key's default, as the definition gives it.
type SettingDeclaration struct {
// Kind says what sort of value this is. `preference` is the only kind with a default; it is
// stated rather than assumed so a reviewer sees the claim being made.
Kind string `json:"kind"`
// Default is the value when no layer sets the key: a string, a number or a boolean.
Default any `json:"default"`
// Why this default: one sentence for whoever wonders whether to change it.
Why string `json:"why"`
}
// KindPreference is the kind of a setting that may have a default.
const KindPreference = "preference"
// DefaultLayer is what the layer of a module's own defaults is called where a value's source is said.
// Only said: the layer is recognised by Layer.Default, never by this name, which a node may also have.
const DefaultLayer = "default"
// meshWords are the settings keys the mesh reads itself; a module declares none of them.
var meshWords = map[string]bool{
PortsSetting: true, ExposeSetting: true, ReachSetting: true, EndpointsSetting: true,
PlacesSetting: true, AccessesSetting: true, NetworksSetting: true,
}
// operatorsOwn are the words that, as what a key's name is about, say its value is the operator's and
// never a preference: a default for one would be the literal ADR 0112 removed from definitions.
var operatorsOwn = map[string]bool{
"domain": true, "host": true, "hostname": true, "servername": true, "fqdn": true, "zone": true,
"realm": true, "tenant": true, "site": true, "timezone": true,
"issuer": true, "url": true, "uri": true, "webhook": true, "origin": true, "dsn": true,
"ip": true, "ipv4": true, "ipv6": true,
"email": true, "mail": true, "phone": true, "address": true,
"identity": true, "login": true, "user": true, "username": true, "account": true, "owner": true,
"uid": true, "gid": true, "puid": true, "pgid": true,
"password": true, "pass": true, "passwd": true, "passphrase": true, "secret": true, "token": true,
"key": true, "apikey": true, "bearer": true, "cert": true, "certificate": true, "credential": true,
"nameserver": true, "gateway": true, "subnet": true, "sender": true, "recipient": true, "contact": true,
"trusted": true, "whitelist": true, "peer": true, "bind": true, "listen": true, "upstream": true,
"proxy": true, "admin": true, "mac": true,
}
// operatorsCompounds are names of two words that are the operator's though neither word alone says so
// at the end of a key: a client's identifier, and a name the world knows a site or server by.
var operatorsCompounds = map[string]bool{
"client-id": true, "site-name": true, "server-name": true, "public-name": true, "smtp-relay": true,
"host-name": true, "user-name": true, "domain-name": true, "dns-server": true,
// Whom a rule lets in or keeps out: a list of addresses or networks.
"allow-from": true, "deny-from": true, "allow-list": true,
}
// aboutAnAmount are first words that make a key about how many or whether, never about whom:
// `max-tokens` is a number, `show-hostname` a switch. Not `allow` or `use`: `allow-from` and
// `use-host` name whom.
var aboutAnAmount = map[string]bool{
"max": true, "min": true, "num": true, "count": true, "show": true, "hide": true, "enable": true,
"disable": true,
}
// operatorsWord is what in a key's name says its value is the operator's, or "". A key is about its
// last word — `url-timeout` is a timeout, `user-agent` an agent, `mail-domain` a domain — or its last two
// as one of operatorsCompounds. A plural is read as its singular.
func operatorsWord(key string) string {
words := strings.FieldsFunc(key, func(r rune) bool { return r == '-' || r == '_' || r == '.' })
if len(words) == 0 || (len(words) > 1 && aboutAnAmount[words[0]]) {
return ""
}
for i, w := range words {
switch {
case operatorsOwn[w]:
case strings.HasSuffix(w, "ies") && operatorsOwn[strings.TrimSuffix(w, "ies")+"y"]:
words[i] = strings.TrimSuffix(w, "ies") + "y"
case strings.HasSuffix(w, "s") && operatorsOwn[strings.TrimSuffix(w, "s")]:
words[i] = strings.TrimSuffix(w, "s")
}
}
if n := len(words); n > 1 {
// Where a secret or an identity is kept is the operator's too: `password-file`, `token-path`.
if (words[n-1] == "file" || words[n-1] == "path") && operatorsOwn[words[n-2]] {
return words[n-2] + "-" + words[n-1]
}
for _, last := range []string{words[n-1], strings.TrimSuffix(words[n-1], "s")} {
if pair := words[n-2] + "-" + last; operatorsCompounds[pair] {
return pair
}
}
}
if last := words[len(words)-1]; operatorsOwn[last] {
return last
}
return ""
}
// SettingProblems is every way a definition's setting defaults are wrong, in its own words.
func SettingProblems(m Manifest) []string {
if len(m.Settings) == 0 {
return nil
}
used := settingKeysUsedBy(m)
var problems []string
for _, key := range sortedSettingKeys(m.Settings) {
d := m.Settings[key]
say := func(format string, args ...any) {
problems = append(problems, fmt.Sprintf("%s's setting %q: ", m.Module, key)+fmt.Sprintf(format, args...))
}
if !settingRef.MatchString("${setting:" + key + "}") {
say("not a usable key: lower-case letters, digits, dots, dashes and underscores")
continue
}
if meshWords[key] {
say("the mesh reads %q itself, and a module gives it no default", key)
continue
}
if d.Kind != KindPreference {
say("kind is %q, and only a %q has a default (novox/hq ADR 0262)", d.Kind, KindPreference)
}
if word := operatorsWord(key); word != "" {
say("a key naming %q is the operator's value, and has no default — it is the "+
"assignment's, never the definition's (novox/hq ADR 0112, ADR 0262)", word)
}
switch v := d.Default.(type) {
case string:
if strings.TrimSpace(v) == "" {
say("an empty default is no default; give the value, or declare nothing")
}
case float64, bool:
case nil:
say("no default: a key without one is the operator's, and is not declared here")
default:
say("a default is a string, a number or a boolean, and this is %T", d.Default)
}
if strings.TrimSpace(d.Why) == "" {
say("no why: say in one sentence why this default")
}
if !used[key] {
say("nothing asks for ${setting:%s}, so the default reaches nothing", key)
}
}
return problems
}
// Defaults is the layer a module's own defaults make, or nothing when it gives none.
func Defaults(m Manifest) (Layer, bool) {
if len(m.Settings) == 0 {
return Layer{}, false
}
values := map[string]any{}
for key, d := range m.Settings {
if d.Default != nil {
values[key] = d.Default
}
}
return Layer{From: DefaultLayer, Values: values, Default: true}, len(values) > 0
}
// WithDefaults is a module's layers with its defaults under them, for filling `${setting:<key>}`.
func WithDefaults(m Manifest, layers []Layer) []Layer {
d, has := Defaults(m)
if !has {
return layers
}
return append([]Layer{d}, layers...)
}
// SettingSource is one key's effective value and the layer it came from.
type SettingSource struct {
Key string
Value any
// From is DefaultLayer, MeshWideLayer or the node's name.
From string
// FromDefault is whether the value is the module's default, whatever From reads.
FromDefault bool
// Default is the module's default, when it gives one.
Default any
HasDefault bool
}
// Effective is every key a module gives a default or a layer sets, with its value and where it came
// from: the default, then the mesh's layer, then the node's — later wins.
func Effective(m Manifest, layers []Layer) []SettingSource {
byKey := map[string]*SettingSource{}
for key, d := range m.Settings {
byKey[key] = &SettingSource{Key: key, Value: d.Default, From: DefaultLayer, FromDefault: true,
Default: d.Default, HasDefault: true}
}
for _, layer := range layers {
for key, v := range layer.Values {
s, ok := byKey[key]
if !ok {
s = &SettingSource{Key: key}
byKey[key] = s
}
s.Value, s.From, s.FromDefault = v, layer.From, layer.Default
}
}
out := make([]SettingSource, 0, len(byKey))
for _, s := range byKey {
out = append(out, *s)
}
sort.Slice(out, func(i, j int) bool { return out[i].Key < out[j].Key })
return out
}
func sortedSettingKeys(in map[string]SettingDeclaration) []string {
keys := make([]string, 0, len(in))
for k := range in {
keys = append(keys, k)
}
sort.Strings(keys)
return keys
}
+450
View File
@@ -0,0 +1,450 @@
package catalogue
import (
"encoding/json"
"strings"
"testing"
)
// A preference has a default in the definition; the mesh's layer, then the node's, override it; a key
// with a default is not stray; and a value that is the operator's still has none (novox/hq ADR 0262).
func notifier() Manifest {
return Manifest{Module: "notifier",
Settings: map[string]SettingDeclaration{
"font-size": {Kind: KindPreference, Default: float64(10), Why: "readable at a scale of one"},
"width": {Kind: KindPreference, Default: float64(250), Why: "fits a title of forty characters"},
},
Resources: []map[string]any{{"id": "configuration", "type": "file", "path": "/x/notifierrc",
"content": "font = Inter ${setting:font-size}\nwidth = ${setting:width}\n"}},
}
}
func parsed(t *testing.T, m Manifest) error {
t.Helper()
raw, err := json.Marshal(m)
if err != nil {
t.Fatal(err)
}
_, err = ParseManifest(raw)
return err
}
func TestAnUnsetPreferenceTakesItsDefault(t *testing.T) {
m := notifier()
for _, layers := range [][]Layer{nil, {{From: MeshWideLayer, Values: map[string]any{}}}} {
file := map[string]any{}
for k, v := range m.Resources[0] {
file[k] = v
}
if err := settingInto(file, WithDefaults(m, layers), m.Module); err != nil {
t.Fatal(err)
}
if file["content"] != "font = Inter 10\nwidth = 250\n" {
t.Fatalf("filled as %q", file["content"])
}
}
if err := JudgeSettings(m, nil, false); err != nil {
t.Fatalf("a module whose every key has a default does not compose with no layer: %v", err)
}
}
func TestTheNodeOverTheMeshOverTheDefault(t *testing.T) {
m := notifier()
layers := []Layer{
{From: MeshWideLayer, Values: map[string]any{"width": float64(300)}},
{From: "laptop", Values: map[string]any{"font-size": float64(13), "width": float64(340)}},
}
file := map[string]any{"type": "file", "content": m.Resources[0]["content"]}
if err := settingInto(file, WithDefaults(m, layers), m.Module); err != nil {
t.Fatal(err)
}
if file["content"] != "font = Inter 13\nwidth = 340\n" {
t.Fatalf("filled as %q", file["content"])
}
file = map[string]any{"type": "file", "content": m.Resources[0]["content"]}
if err := settingInto(file, WithDefaults(m, layers[:1]), m.Module); err != nil {
t.Fatal(err)
}
if file["content"] != "font = Inter 10\nwidth = 300\n" {
t.Fatalf("the mesh's layer over the default filled as %q", file["content"])
}
}
// Composed for a machine, the way a push writes it: the default reaches the file, and the node's
// layer overrides it.
func TestAComposedMachineGetsTheDefaultAndTheNodesValue(t *testing.T) {
r := anAdoptedAnchor()
r.Modules = append(r.Modules, notifier())
with := anchorRendering(false)
composed, err := r.Compose(with)
if err != nil {
t.Fatal(err)
}
if why, left := composed.LeftOut["notifier"]; left {
t.Fatalf("the notifier was left out: %s", why)
}
if c := byID(composed.Resources)["notifier.configuration"]["content"]; c != "font = Inter 10\nwidth = 250\n" {
t.Fatalf("composed with no layer as %q", c)
}
with.Settings["notifier"] = []Layer{{From: "anchor", Values: map[string]any{"font-size": float64(13)}}}
if composed, err = r.Compose(with); err != nil {
t.Fatal(err)
}
if c := byID(composed.Resources)["notifier.configuration"]["content"]; c != "font = Inter 13\nwidth = 250\n" {
t.Fatalf("composed with the node's font size as %q", c)
}
}
// Setting a key the module gives a default is not refused as reaching nothing: it overrides the
// default, which is how a running module gains a setting with no gap between.
func TestAKeyWithADefaultIsNotStray(t *testing.T) {
m := notifier()
m.Resources[0]["content"] = "font = Inter ${setting:font-size}\nwidth = ${setting:width}\n"
stray := UnusedSettings(m, []Layer{{From: "laptop", Values: map[string]any{"font-size": float64(13), "colour": "red"}}})
joined := strings.Join(stray, "; ")
if strings.Contains(joined, `"font-size"`) || !strings.Contains(joined, `"colour"`) {
t.Fatalf("stray: %s", joined)
}
}
// A default fills ${setting:…} only: it never becomes a key of a mergeable file (issue 168).
func TestADefaultIsNotMergedIntoAJSONFile(t *testing.T) {
m := notifier()
m.Resources = append(m.Resources, map[string]any{"id": "other", "type": "file", "path": "/x/other.json",
"merge": MergeJSON, "content": `{"keep": 1}`})
out, err := ApplySettings(m.Resources[1], WithDefaults(m, nil))
if err != nil {
t.Fatal(err)
}
if strings.Contains(out["content"].(string), "font-size") {
t.Fatalf("a default reached a mergeable file: %s", out["content"])
}
}
// A value that is the operator's has no default: no layer setting it is refused by name, as before.
func TestAnOperatorsValueWithoutADefaultIsStillRefused(t *testing.T) {
m := notifier()
m.Resources[0]["content"] = "font = Inter ${setting:font-size}\nwidth = ${setting:width}\nfrom = ${setting:domain}\n"
err := JudgeSettings(m, nil, false)
if err == nil || !strings.Contains(err.Error(), "${setting:domain}") {
t.Fatalf("judged %v", err)
}
if strings.Contains(err.Error(), "font-size") {
t.Fatalf("a default was named as set today: %v", err)
}
}
func TestTheParserTakesAPreferenceAndRefusesTheRest(t *testing.T) {
if err := parsed(t, notifier()); err != nil {
t.Fatalf("a preference with a default was refused: %v", err)
}
for _, c := range []struct {
name string
key string
d SettingDeclaration
refuse string
}{
{"no kind", "font-size", SettingDeclaration{Default: float64(10), Why: "x"}, `only a "preference" has a default`},
{"another kind", "font-size", SettingDeclaration{Kind: "operator", Default: float64(10), Why: "x"}, `only a "preference"`},
{"no default", "font-size", SettingDeclaration{Kind: KindPreference, Why: "x"}, "no default"},
{"an empty default", "font-size", SettingDeclaration{Kind: KindPreference, Default: " ", Why: "x"}, "an empty default"},
{"an object", "font-size", SettingDeclaration{Kind: KindPreference, Default: map[string]any{"a": 1.0}, Why: "x"}, "a string, a number or a boolean"},
{"no why", "font-size", SettingDeclaration{Kind: KindPreference, Default: float64(10)}, "no why"},
{"the operator's", "mail-domain", SettingDeclaration{Kind: KindPreference, Default: "example.tld", Why: "x"}, "is the operator's value"},
{"a secret", "api-token", SettingDeclaration{Kind: KindPreference, Default: "x", Why: "x"}, "is the operator's value"},
{"the mesh's word", PortsSetting, SettingDeclaration{Kind: KindPreference, Default: float64(1), Why: "x"}, "the mesh reads"},
{"read by nothing", "height", SettingDeclaration{Kind: KindPreference, Default: float64(300), Why: "x"}, "reaches nothing"},
} {
m := notifier()
m.Resources[0]["content"] = m.Resources[0]["content"].(string) + "x = ${setting:" + c.key + "}\n"
if c.name == "read by nothing" {
m.Resources[0]["content"] = "font = Inter ${setting:font-size}\nwidth = ${setting:width}\n"
}
m.Settings[c.key] = c.d
err := parsed(t, m)
if err == nil || !strings.Contains(err.Error(), c.refuse) {
t.Errorf("%s: parsed %v, want %q", c.name, err, c.refuse)
}
}
}
func TestEveryValueSaysWhereItCameFrom(t *testing.T) {
m := notifier()
m.Resources[0]["content"] = m.Resources[0]["content"].(string) + "x = ${setting:position}\n"
got := Effective(m, []Layer{
{From: MeshWideLayer, Values: map[string]any{"width": float64(300), "position": "top-right"}},
{From: "laptop", Values: map[string]any{"font-size": float64(13)}},
})
want := map[string]string{"font-size": "laptop", "position": MeshWideLayer, "width": MeshWideLayer}
if len(got) != 3 {
t.Fatalf("effective: %+v", got)
}
for _, s := range got {
if s.From != want[s.Key] {
t.Errorf("%s from %q, want %q", s.Key, s.From, want[s.Key])
}
}
if got[0].Key != "font-size" || got[0].Value != float64(13) || got[0].Default != float64(10) {
t.Fatalf("font-size: %+v", got[0])
}
only := Effective(m, nil)
if only[0].From != DefaultLayer || only[0].Value != float64(10) {
t.Fatalf("with no layer: %+v", only)
}
}
// A node may be called `default`. Its layer is a node's like any other: it overrides the module's
// default, it merges into a mergeable file, and it is named among what is set.
func TestANodeCalledDefaultIsANodesLayer(t *testing.T) {
m := notifier()
node := []Layer{{From: DefaultLayer, Values: map[string]any{"font-size": float64(13)}}}
file := map[string]any{"type": "file", "content": m.Resources[0]["content"]}
if err := settingInto(file, WithDefaults(m, node), m.Module); err != nil {
t.Fatal(err)
}
if file["content"] != "font = Inter 13\nwidth = 250\n" {
t.Fatalf("the node called default was dropped: %q", file["content"])
}
out, err := ApplySettings(map[string]any{"id": "j", "type": "file", "merge": MergeJSON, "content": `{}`}, node)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(out["content"].(string), `"font-size": 13`) {
t.Fatalf("the node called default did not merge: %s", out["content"])
}
if got := Effective(m, node); got[0].FromDefault || got[0].Value != float64(13) {
t.Fatalf("the node's value read as the default: %+v", got[0])
}
}
// The operator's own value is told by what a key's name is about: its last word, or its last two as
// a known compound; a plural as its singular; and never a number or a switch.
func TestAKeyIsTheOperatorsByWhatItIsAbout(t *testing.T) {
for _, key := range []string{"max-tokens", "show-hostname", "ghost-opacity", "users-per-page", "mailbox-size",
"font-size", "width", "keyboard-delay", "ipv6-preferred", "client-width", "user-agent", "url-timeout",
"site-title", "cert-renewal-days", "name", "font-name", "allow-resize", "use-gpu", "disable-sender-check",
"max-recipients", "gateway-timeout", "sender-delay", "upstream-resolvers", "mirror-countries",
"pool-region", "proxy-timeout", "listen-backlog", "peer-keepalive", "admin-theme", "log-file",
"cache-path"} {
if w := operatorsWord(key); w != "" {
t.Errorf("%s read as the operator's (%s)", key, w)
}
}
for key, word := range map[string]string{"mail-domain": "domain", "api-key": "key", "admin-password": "password",
"oauth-client-id": "client-id", "db-dsn": "dsn", "public-ip": "ip", "puid": "puid", "dns-zone": "zone",
"webhook": "webhook", "notify_phone": "phone", "cors.origin": "origin", "smtp-pass": "pass",
"backup-passphrase": "passphrase", "data-owner": "owner", "fqdn": "fqdn", "tenant": "tenant", "host": "host",
"allowed-hosts": "host", "admin-emails": "email", "tokens": "token", "hostname": "hostname",
"apikey": "apikey", "servername": "servername", "tls-cert": "cert", "site": "site", "timezone": "timezone",
"bearer": "bearer", "bind-ipv4": "ipv4", "listen-ipv6": "ipv6", "site-name": "site-name",
"server-name": "server-name", "public-name": "public-name", "smtp-relay": "smtp-relay",
"host-name": "host-name", "user-name": "user-name", "domain-name": "domain-name",
"nameserver": "nameserver", "upstream-nameservers": "nameserver", "dns-server": "dns-server",
"dns-servers": "dns-server", "default-gateway": "gateway", "lan-subnet": "subnet", "sender": "sender",
"notify-recipients": "recipient", "contact": "contact", "tls-certificate": "certificate",
"allow-from": "allow-from", "deny-from": "deny-from", "allow-hosts": "host", "use-host": "host",
"trusted": "trusted", "allow-list": "allow-list", "ip-whitelist": "whitelist", "peer": "peer",
"wireguard-peers": "peer", "bind": "bind", "listen": "listen", "upstream": "upstream", "http-proxy": "proxy",
"trusted-proxies": "proxy", "admin": "admin", "notify-admin": "admin", "wake-mac": "mac",
"password-file": "password-file", "key-file": "key-file", "token-path": "token-path",
"secret-file": "secret-file", "cert-path": "cert-path"} {
if w := operatorsWord(key); w != word {
t.Errorf("%s: read %q, want %q", key, w, word)
}
}
}
// A default is a value the mesh writes, so the installation check reads it, even under a setting
// called `description` or `why`; a setting's own why is prose and is not read.
func TestTheInstallationCheckReadsADefault(t *testing.T) {
m := notifier()
m.Settings["relay"] = SettingDeclaration{Kind: KindPreference, Default: "relay.acme.be", Why: "as at relay.acme.be"}
m.Settings["description"] = SettingDeclaration{Kind: KindPreference, Default: "notes.acme.be", Why: "x"}
problems := strings.Join(InstallationProblems(m), "; ")
for _, want := range []string{"relay.acme.be at settings.relay.default", "notes.acme.be at settings.description.default"} {
if !strings.Contains(problems, want) {
t.Errorf("not reported: %q in %s", want, problems)
}
}
if strings.Contains(problems, "settings.relay.why") {
t.Errorf("a why was read as a value: %s", problems)
}
}
// A default fills ${setting:…} in what a module contributes and serves, and never replaces a value a
// contribution states itself.
func TestADefaultFillsAContributionAndAServedFactButNotALiteral(t *testing.T) {
m := notifier()
m.Settings["site-title"] = SettingDeclaration{Kind: KindPreference, Default: "Notes", Why: "x"}
contribution := map[string]any{"title": "${setting:site-title}", "width": "fixed", "port": float64(8080)}
got, err := overridden(contribution, WithDefaults(m, nil), "a contribution")
if err != nil {
t.Fatal(err)
}
if got["title"] != "Notes" || got["width"] != "fixed" {
t.Fatalf("contribution: %v", got)
}
got, err = overridden(contribution, WithDefaults(m, []Layer{{From: "laptop", Values: map[string]any{"width": "wide"}}}), "a contribution")
if err != nil || got["width"] != "wide" {
t.Fatalf("a node's value did not override a contribution's own key: %v %v", got, err)
}
served, err := Settle(map[string]any{"name": "${setting:site-title}"}, WithDefaults(m, nil))
if err != nil || served["name"] != "Notes" {
t.Fatalf("served: %v %v", served, err)
}
if _, err := Settle(map[string]any{"name": "${setting:site-title}"}, nil); err == nil {
t.Fatal("a served fact without the defaults was filled")
}
}
// A stored manifest with a key this controller does not know, at the top or inside any block, is read
// and its module is left out of every declaration by name; the rest of the catalogue is read; the module
// check refuses it. One case per block whose decoder wraps the decoder's words in its own.
func TestAStoredManifestWithAnUnknownKeyIsLeftOutAndRegistrationRefusesIt(t *testing.T) {
for where, raw := range map[string]string{
"the top": `{"module": "later", "version": "2", "a-field-from-later": {"x": 1}}`,
"a state": `{"module": "later", "version": "2", "state": [{"name": "s", "a-field-from-later": 1}]}`,
"a provided name": `{"module": "later", "version": "2", "provides": [{"name": "p", "a-field-from-later": 1}]}`,
"an offer's identity": `{"module": "later", "version": "2", "provides": [{"name": "p", "identity": {"in": "x", "a-field-from-later": 1}}]}`,
"a backup": `{"module": "later", "version": "2", "data": {"own": [{"id": "d", "path": "${dir:d}", "class": "valuable", "backup": {"dump": "x", "into": "y", "a-field-from-later": 1}}]}}`,
"an own secret": `{"module": "later", "version": "2", "own-secrets": {"s": {"path": "/x", "a-field-from-later": 1}}}`,
"a seat's verb": `{"module": "later", "version": "2", "seats": [{"name": "later-seat", "serves": [{"name": "v", "a-field-from-later": 1}]}]}`,
} {
var m Manifest
if err := json.Unmarshal([]byte(raw), &m); err != nil {
t.Errorf("%s: the stored manifest was not read: %v", where, err)
continue
}
if m.Module != "later" || m.Version != "2" || !strings.Contains(m.UnknownField(), `"a-field-from-later"`) {
t.Errorf("%s: read as %q %q, unknown %q", where, m.Module, m.Version, m.UnknownField())
continue
}
if where != "the top" && !strings.Contains(m.UnknownField(), ": json: unknown field") {
t.Errorf("%s: the block's decoder did not say it: %q", where, m.UnknownField())
}
left := Resolution{Node: "laptop", Modules: []Manifest{m, notifier()}}.LeftOut(nil, false)
if why := left["later"]; !strings.Contains(why, "uses a field this controller does not know") ||
!strings.Contains(why, "a-field-from-later") {
t.Errorf("%s: not left out by name: %v", where, left)
}
if _, notifierLeft := left["notifier"]; notifierLeft {
t.Errorf("%s: another module was left out with it: %v", where, left)
}
if _, err := ParseManifest([]byte(raw)); err == nil || !strings.Contains(err.Error(), "a-field-from-later") {
t.Errorf("%s: the module check took it: %v", where, err)
}
}
var known Manifest
if err := json.Unmarshal([]byte(`{"module": "now", "state": [{"name": "s"}]}`), &known); err != nil || known.UnknownField() != "" {
t.Fatalf("a known manifest: %v %q", err, known.UnknownField())
}
// A malformed manifest is still refused: only an unknown key is read past.
var bad Manifest
if err := json.Unmarshal([]byte(`{"module": "bad", "state": [{"name": 3}]}`), &bad); err == nil {
t.Fatal("a malformed stored manifest was read")
}
}
// A module left out for a key this controller does not know, inside an entry, is read past that key
// alone: it still provides what it provides, its other entries are whole, and a key of the same name
// that another entry knows is kept (novox/hq ADR 0262).
func TestALeftOutModuleStillProvidesWhatItProvides(t *testing.T) {
var m Manifest
raw := `{"module": "later", "version": "2",
"provides": [{"name": "db", "scope": "mesh"}, {"name": "cache", "a-field-from-later": 1}],
"state": [{"name": "s", "history": 3}],
"data": {"own": [{"id": "d", "path": "${dir:d}", "class": "valuable", "backup": {"dump": "x", "into": "d", "class": "later"}}]},
"resources": [{"id": "d", "type": "directory", "mode": "0700"}]}`
if err := json.Unmarshal([]byte(raw), &m); err != nil {
t.Fatal(err)
}
if len(m.Provides) != 2 || m.Provides[0].Name != "db" || m.Provides[1].Name != "cache" {
t.Fatalf("provides: %+v", m.Provides)
}
if len(m.State) != 1 || m.State[0].History != 3 {
t.Fatalf("state: %+v", m.State)
}
if m.Data == nil || len(m.Data.Own) != 1 || m.Data.Own[0].Class != "valuable" {
t.Fatalf("data: the item's own class was taken for the backup's unknown one: %+v", m.Data)
}
for _, want := range []string{"provides[1].a-field-from-later", "data.own[0].backup.class"} {
if !strings.Contains(m.UnknownField(), want) {
t.Errorf("unknown %q does not say it read without %s", m.UnknownField(), want)
}
}
if !strings.Contains(UnknownFieldReason(m), "left out") {
t.Fatal(UnknownFieldReason(m))
}
}
// A left-out module's data is still copied: the backup holder on its machine keeps its lines while every
// other thing of it is left out.
func TestALeftOutModulesDataIsStillBackedUp(t *testing.T) {
var later Manifest
if err := json.Unmarshal([]byte(`{"module": "later", "version": "2", "a-field-from-later": 1,
"resources": [{"id": "d", "type": "directory", "mode": "0700"}, {"id": "rc", "type": "file", "path": "/etc/later.conf", "content": "x\n"}],
"data": {"own": [{"id": "d", "path": "${dir:d}", "class": "valuable"}]}}`), &later); err != nil {
t.Fatal(err)
}
holder := Manifest{Module: "backups", Version: "1",
Claims: []Claim{{Name: BackupSeat, Scope: ScopeNode}},
Resources: []map[string]any{{"id": "list", "type": "file", "path": "/etc/backups.list", "mode": "0644",
"content": "${contribution:" + BackupSeat + ":backup}"}}}
r := anAdoptedAnchor()
r.Modules = append(r.Modules, holder, later)
composed, err := r.Compose(anchorRendering(false))
if err != nil {
t.Fatal(err)
}
if _, left := composed.LeftOut["later"]; !left {
t.Fatalf("not left out: %v", composed.LeftOut)
}
got := byID(composed.Resources)
if _, declared := got["later.rc"]; declared {
t.Fatal("the left-out module's file is still declared")
}
list, _ := got["backups.list"]["content"].(string)
if !strings.Contains(list, "# later") || !strings.Contains(list, "path /var/lib/later/d") {
t.Fatalf("the left-out module's data is no longer backed up:\n%s", list)
}
}
// A left-out module's backup lines are best effort: a line that cannot be placed — an access nobody
// placed, a placement setting that does not read — is said among what could not be placed, and the
// machine is declared (novox/hq ADR 0262).
func TestALeftOutModulesUnplaceableBackupLineCostsOnlyThatLine(t *testing.T) {
holder := Manifest{Module: "backups", Version: "1",
Claims: []Claim{{Name: BackupSeat, Scope: ScopeNode}},
Resources: []map[string]any{{"id": "list", "type": "file", "path": "/etc/backups.list", "mode": "0644",
"content": "${contribution:" + BackupSeat + ":backup}"}}}
media := Manifest{Module: "media", Version: "1",
Accesses: []Access{{ID: "library", Mode: "read"}},
Data: &Data{Own: []DataItem{{ID: "library", Path: "${access:library}", Class: "valuable"}}}}
placedBadly := Manifest{Module: "notes", Version: "1",
Resources: []map[string]any{{"id": "d", "type": "directory", "path": "/srv/notes", "mode": "0700"}},
Data: &Data{Own: []DataItem{{ID: "d", Path: "${dir:d}", Class: "valuable"}}}}
r := anAdoptedAnchor()
r.Modules = append(r.Modules, holder, media, placedBadly)
with := anchorRendering(false)
with.Settings["notes"] = []Layer{{From: "anchor", Values: map[string]any{PlacesSetting: "not a map"}}}
composed, err := r.Compose(with)
if err != nil {
t.Fatalf("an unplaceable line of a left-out module failed the machine: %v", err)
}
for _, m := range []string{"media", "notes"} {
if _, left := composed.LeftOut[m]; !left {
t.Errorf("%s is not left out: %v", m, composed.LeftOut)
}
}
unplaced := strings.Join(composed.Unplaced, "\n")
for _, want := range []string{"media's backup for node-backup", "notes's backup for node-backup", "placement setting does not read"} {
if !strings.Contains(unplaced, want) {
t.Errorf("not said among what could not be placed: %q in\n%s", want, unplaced)
}
}
list, _ := byID(composed.Resources)["backups.list"]["content"].(string)
if strings.Contains(list, "/srv/notes") || strings.Contains(list, "${") {
t.Fatalf("a line was placed from the definition's default or unfilled:\n%s", list)
}
}
+8 -3
View File
@@ -40,8 +40,9 @@ func settingsUsed(content string) []string {
// settingInto fills a file's ${setting:…} placeholders from the layers over a module.
//
// The last layer setting a key wins, which is the node's over the mesh's — the same order settle
// applies to a mergeable file. A value that is not a string is written the way a program would read
// The last layer setting a key wins, which is the node's over the mesh's over the module's own
// default (novox/hq ADR 0262) — the same order settle applies to a mergeable file. The caller lays
// the defaults under the layers with WithDefaults. A value that is not a string is written the way a program would read
// it (a number without a trailing .000000, a boolean as true/false).
func settingInto(resource map[string]any, layers []Layer, module string) error {
if fmt.Sprint(resource["type"]) != "file" {
@@ -56,7 +57,8 @@ func settingInto(resource map[string]any, layers []Layer, module string) error {
if !set {
return fmt.Errorf(
"%s has a file that says ${setting:%s}, and nothing sets %q for it — an operator's "+
"value is the assignment's, never the definition's (novox/hq ADR 0112): "+
"value is the assignment's, never the definition's (novox/hq ADR 0112), and only a "+
"preference has a default in the definition (ADR 0262): "+
"`settings set %s <file>` with {%q: …}%s",
module, key, key, module, key, orNoSettings(layers))
}
@@ -80,6 +82,9 @@ func settingValue(layers []Layer, key string) (any, bool) {
func orNoSettings(layers []Layer) string {
var keys []string
for _, l := range layers {
if l.Default {
continue
}
for k := range l.Values {
keys = append(keys, k)
}
+16 -2
View File
@@ -30,6 +30,9 @@ type Layer struct {
// Where these came from, for saying which layer set a value.
From string
Values map[string]any
// Default marks the layer a module's own defaults make (novox/hq ADR 0262). Marked rather than
// recognised by From, which is a node's name for a node's layer, and a node may be called anything.
Default bool
}
// ApplySettings produces a resource's final content from the module's own and the layers over it.
@@ -114,7 +117,7 @@ func overridden(base map[string]any, layers []Layer, what string) (map[string]an
values[key] = value
}
}
kept = append(kept, Layer{From: layer.From, Values: values})
kept = append(kept, Layer{From: layer.From, Values: values, Default: layer.Default})
}
merged, err := settle(base, kept, nil, what)
if err != nil {
@@ -151,6 +154,12 @@ func settle(base map[string]any, layers []Layer, protected map[string]bool, what
map[string]any, error) {
merged := deepCopy(base)
for _, layer := range layers {
if layer.Default {
// A module's own defaults fill ${setting:<key>} only (novox/hq ADR 0262). A mergeable
// file's content already is its defaults, and a contribution or served fact declares its
// own; laid on here, a default would reach every mergeable file of its module (issue 168).
continue
}
for key, value := range layer.Values {
if key == PortsSetting {
// Where the machine puts a port is the mesh's to apply, not a value for a file or
@@ -224,6 +233,11 @@ func UnusedSettings(m Manifest, layers []Layer) []string {
}
lands := settingKeysUsedBy(m)
// A key the module gives a default reaches what asks for it (novox/hq ADR 0262): setting it
// overrides the default, which is how a running module gains a setting without a gap between.
for key := range m.Settings {
lands[key] = true
}
for _, values := range m.Contributes {
for key := range values {
lands[key] = true
@@ -420,7 +434,7 @@ func JudgeSettings(m Manifest, layers []Layer, adopted bool) error {
for k, v := range settled {
copied[k] = v
}
if err := settingInto(copied, layers, m.Module); err != nil {
if err := settingInto(copied, WithDefaults(m, layers), m.Module); err != nil {
return err
}
}
+96
View File
@@ -0,0 +1,96 @@
package catalogue
// A machine's shares, and the shares a machine mounts (novox/hq ADR 0263).
//
// **Two roles, one per side of the wire.** A machine that shares a directory with the mesh does it
// through the holder of `node-nfs-server`: it writes the machine's export file, runs the NFS service,
// opens its port to the private network only, and provides each share as the provision `nfs-share`. A
// machine that wants the files gets them through the holder of `node-mounts`: it writes a mount unit and
// an automount unit per share, so nothing mounts at boot and nothing can fail a boot, and it says a
// device that comes and goes is absent rather than failed.
//
// **One holder per machine on each side.** Two modules writing one machine's export file, or two writing
// mount units for one mount point, is the conflict a seat exists to refuse. Both seats deliver nothing:
// `nfs-share` is provided at the mesh's scope, by the module holding `node-nfs-server` on the machine that
// shares, and a seat at a node's scope cannot be the answer for a provision at the mesh's.
//
// **The data is the operator's** (ADR 0051). Neither holder creates, chowns or removes anything under a
// shared path; the export maps every client to the path's owner, so no client acts as another account on
// the server. Their verbs read, and the ones that act take over what a person wrote by hand only on a
// person's word: a dataset's export property, an fstab line.
// NFSServerSeat is the role of the machine that shares directories over NFS (novox/hq ADR 0263).
const NFSServerSeat = "node-nfs-server"
// MountsSeat is the role that mounts a machine's shares and occasional sources (novox/hq ADR 0263).
const MountsSeat = "node-mounts"
// nfsServerVerbs is the contract every holder of node-nfs-server serves (novox/hq ADR 0263).
func nfsServerVerbs() []Verb {
return []Verb{
{Name: "exports", Description: "Every share this machine exports: its name, its path, read-write or " +
"read-only, the owner every client is mapped to (uid and gid), the clients it is exported to (the " +
"private network's range), and whether the kernel holds it now. Also the exports found that are " +
"not the mesh's: a dataset's sharenfs property, a line in /etc/exports.",
Input: schema(map[string]string{}, nil),
Replaces: []string{"exportfs -v", "cat /etc/exports", "zfs get sharenfs"}},
{Name: "clients", Description: "Which machines have mounted which share now, as the NFS server " +
"knows its clients.",
Input: schema(map[string]string{}, nil),
Replaces: []string{"ss -tn sport = :2049", "cat /proc/fs/nfsd/clients/*/info"}},
{Name: "test", Description: "Whether this machine exports one share for the mesh now, and whether its " +
"NFS service is up; with an address, also whether that address is inside the private network's " +
"range the share is exported to. What a machine mounting the share asks before it mounts.",
Input: schema(map[string]string{
"share": "the share, by its name",
"address": "a client's address on the private network (optional)",
}, []string{"share"}),
Replaces: []string{"showmount -e"}},
{Name: "reload", Description: "Have the kernel read this machine's export files again now (exportfs " +
"-ra) and answer the shares as it then holds them. The module's own process writes its export " +
"file; this is for after a change made outside it. Changes no shared path.",
Input: schema(map[string]string{}, nil),
Replaces: []string{"exportfs -ra"}},
{Name: "adopt", Description: "Take over an export a person made by hand: clear a dataset's sharenfs " +
"property for a path the mesh's own export now serves — only when that export is live. Without " +
"confirm, says what it would do and changes nothing.",
Input: schema(map[string]string{
"path": "the shared path whose hand-made export is taken over",
"confirm": "\"true\": change it (needs why); anything else is a dry run",
"why": "why, for the record",
}, []string{"path"}, "confirm"),
Replaces: []string{"zfs set sharenfs=off"}},
}
}
// mountsVerbs is the contract every holder of node-mounts serves (novox/hq ADR 0263).
func mountsVerbs() []Verb {
return []Verb{
{Name: "list", Description: "Every mount point on this machine: its fstab line, the mesh's mount and " +
"automount units for it, its state (armed, mounted, absent, unreachable, failed) and since when, " +
"and which settings asked for it. A password in a mount's options is never shown.",
Input: schema(map[string]string{}, nil),
Replaces: []string{"cat /etc/fstab", "findmnt", "systemctl list-units --type=mount"}},
{Name: "test", Description: "Whether one share's or occasional source's server answers from this " +
"machine now, without touching its mount point.",
Input: schema(map[string]string{"name": "the share or source, by its name"}, []string{"name"}),
Replaces: []string{"showmount -e", "ping"}},
{Name: "mount", Description: "Mount one share or occasional source now, rather than at its first " +
"access.",
Input: schema(map[string]string{"name": "the share or source, by its name"}, []string{"name"}),
Replaces: []string{"mount"}},
{Name: "unmount", Description: "Release one share or occasional source now; its automount stays " +
"armed, so the next access mounts it again.",
Input: schema(map[string]string{"name": "the share or source, by its name"}, []string{"name"}),
Replaces: []string{"umount"}},
{Name: "adopt", Description: "Take over a mount a person wrote by hand: comment out the /etc/fstab " +
"line for a mount point the mesh's own units now serve, keeping a copy of the file — only when " +
"the mesh's automount for it is armed. Without confirm, says what it would do and changes nothing.",
Input: schema(map[string]string{
"mountpoint": "the mount point whose fstab line is taken over",
"confirm": "\"true\": change it (needs why); anything else is a dry run",
"why": "why, for the record",
}, []string{"mountpoint"}, "confirm"),
Replaces: []string{"sed -i /etc/fstab"}},
}
}
+161
View File
@@ -0,0 +1,161 @@
package catalogue
import (
"strings"
"testing"
)
// Defends novox/hq ADR 0263: a machine's shares and the shares a machine mounts are two node seats, each
// with its verbs required of every holder, each delivering nothing — `nfs-share` is provided at the mesh's
// scope by the holder of node-nfs-server, and a node seat cannot answer for a provision at the mesh's.
func TestTheSharesAndMountsAreNodeSeatsWithTheirVerbs(t *testing.T) {
for seat, want := range map[string]string{
NFSServerSeat: "exports clients test reload adopt",
MountsSeat: "list test mount unmount adopt",
} {
s, ok := SeatNamed(seat)
if !ok {
t.Fatalf("%s is not in the mesh's set", seat)
}
if s.Scope != ScopeNode || s.Decision != "novox/hq ADR 0263" || s.Delivers != "" || s.Replicated {
t.Errorf("%s is %+v; a node seat under ADR 0263 that delivers nothing", seat, s)
}
var got []string
for _, v := range s.Serves {
got = append(got, v.Name)
if v.Optional {
t.Errorf("%s.%s is optional; its first holder serves it", seat, v.Name)
}
if v.Description == "" || v.Input["type"] != "object" || len(v.Replaces) == 0 {
t.Errorf("%s.%s has no description, no object schema or says it replaces nothing", seat, v.Name)
}
}
if strings.Join(got, " ") != want {
t.Errorf("%s serves %v, not %s", seat, got, want)
}
}
}
// Both acts that take over what a person wrote by hand are dry runs unless confirmed, and name the path.
func TestTheAdoptVerbsAreDryRunsUnlessConfirmed(t *testing.T) {
for seat, key := range map[string]string{NFSServerSeat: "path", MountsSeat: "mountpoint"} {
s, _ := SeatNamed(seat)
for _, v := range s.Serves {
if v.Name != "adopt" {
continue
}
props, _ := v.Input["properties"].(map[string]any)
confirm, _ := props["confirm"].(map[string]any)
if confirm == nil || confirm["enum"] == nil {
t.Errorf("%s.adopt has no confirm switch: %v", seat, v.Input)
}
if req, _ := v.Input["required"].([]string); len(req) != 1 || req[0] != key {
t.Errorf("%s.adopt requires %v, not %q alone", seat, v.Input["required"], key)
}
if !strings.Contains(v.Description, "Without confirm, says what it would do and changes nothing") {
t.Errorf("%s.adopt does not say a call without confirm changes nothing: %s", seat, v.Description)
}
}
}
}
// A holder claiming the seat at a node with every verb holds it; one naming a verb the seat does not
// promise, or leaving one out, is refused — as the catalogue's nfs-server and mounts modules claim them.
func TestAHolderOfTheShareSeatsServesEveryVerbAndNothingElse(t *testing.T) {
cases := []struct {
seat, module string
verbs []string
}{
{NFSServerSeat, "nfs-server", []string{"exports", "clients", "test", "reload", "adopt"}},
{MountsSeat, "mounts", []string{"list", "test", "mount", "unmount", "adopt"}},
}
for _, c := range cases {
seat, _ := SeatNamed(c.seat)
holder := Manifest{Module: c.module, Version: "1",
Claims: []Claim{{Name: c.seat, Scope: ScopeNode, Serves: c.verbs}}}
if c.seat == NFSServerSeat {
holder.Provides = []Offer{{Name: "nfs-share", Scope: ScopeMesh}}
}
if err := CanHold(holder, seat); err != nil {
t.Errorf("%s serving every verb is refused: %v", c.module, err)
}
typo := holder
typo.Claims = []Claim{{Name: c.seat, Scope: ScopeNode, Serves: append(append([]string{}, c.verbs...), "export")}}
if err := CanHold(typo, seat); err == nil || !strings.Contains(err.Error(), "does not promise") {
t.Errorf("%s naming a verb the seat does not promise was accepted: %v", c.module, err)
}
short := holder
short.Claims = []Claim{{Name: c.seat, Scope: ScopeNode, Serves: c.verbs[:len(c.verbs)-1]}}
if err := CanHold(short, seat); err == nil || !strings.Contains(err.Error(), "adopt") {
t.Errorf("%s leaving adopt out was accepted: %v", c.module, err)
}
mesh := holder
mesh.Claims = []Claim{{Name: c.seat, Scope: ScopeMesh, Serves: c.verbs}}
if err := CanHold(mesh, seat); err == nil {
t.Errorf("%s claiming a node seat at the mesh's scope was accepted", c.module)
}
}
}
// What each verb replaces is what an agent would type over ssh to read or change a share by hand.
func TestTheShareVerbsSayWhatTheyReplace(t *testing.T) {
want := map[string]string{
NFSServerSeat + ".exports": "exportfs -v",
NFSServerSeat + ".adopt": "zfs set sharenfs=off",
MountsSeat + ".list": "cat /etc/fstab",
MountsSeat + ".adopt": "sed -i /etc/fstab",
}
for _, s := range DefaultSeats() {
for _, v := range s.Serves {
cmd, ok := want[s.Name+"."+v.Name]
if !ok {
continue
}
delete(want, s.Name+"."+v.Name)
found := false
for _, r := range v.Replaces {
found = found || r == cmd
}
if !found {
t.Errorf("%s.%s does not say it replaces %q: %v", s.Name, v.Name, cmd, v.Replaces)
}
}
}
for verb := range want {
t.Errorf("%s is not a verb of the compiled seats", verb)
}
}
// The confirm switch is the seat's string "true", as every verb's switch is, and its description says so:
// a holder handed the boolean reading of "true to change it" would treat the string as a dry run.
func TestTheConfirmSwitchIsTheStringTrue(t *testing.T) {
for _, seat := range []string{NFSServerSeat, MountsSeat} {
s, _ := SeatNamed(seat)
for _, v := range s.Serves {
props, _ := v.Input["properties"].(map[string]any)
confirm, _ := props["confirm"].(map[string]any)
if confirm == nil {
continue
}
if confirm["type"] != "string" {
t.Errorf("%s.%s confirm is %v, not the string switch", seat, v.Name, confirm["type"])
}
if d, _ := confirm["description"].(string); !strings.Contains(d, `"true"`) {
t.Errorf("%s.%s confirm does not name the string \"true\": %q", seat, v.Name, d)
}
}
}
}
// reload has the kernel read the export files; the module's process writes its file. The description
// must not promise a write it does not do.
func TestReloadSaysWhatItDoes(t *testing.T) {
s, _ := SeatNamed(NFSServerSeat)
for _, v := range s.Serves {
if v.Name == "reload" && (strings.Contains(v.Description, "Write this machine's export file") ||
!strings.Contains(v.Description, "exportfs")) {
t.Errorf("reload's description promises a write or does not name exportfs: %s", v.Description)
}
}
}
+1 -1
View File
@@ -49,7 +49,7 @@ func (s *StateDeclaration) UnmarshalJSON(raw []byte) error {
dec := json.NewDecoder(bytes.NewReader(trimmed))
dec.DisallowUnknownFields()
if err := dec.Decode(&full); err != nil {
return fmt.Errorf("a state is either a name or {name, history, ttl-seconds, per-machine}: %w", err)
return fmt.Errorf("a state is either a name or {name, history, ttl-seconds, per-machine}: %w", typedUnknown(err))
}
*s = StateDeclaration(full)
return nil
+206
View File
@@ -0,0 +1,206 @@
package catalogue
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"regexp"
)
// A key this controller does not know, in a manifest (novox/hq ADR 0262).
//
// Registration refuses it, as it always has. A manifest the store already holds was registered by a newer
// controller, and is read by this one after a rollback: refusing it there failed the whole catalogue, and
// with it every plan and every send. Dropping the key silently would be worse — a module running without
// something its manifest says. So the manifest is read, the module is left out of every machine's
// declaration by name, and the controller raises a condition until it is updated.
// UnknownFieldError is a key a manifest has that this controller does not know, wherever it is: at the
// top of the manifest or inside a block (a state, an offer, a data item's backup, an own secret, a verb).
// Typed, because the blocks' decoders wrap the decoder's words in their own.
type UnknownFieldError struct {
Field string
err error
}
func (e *UnknownFieldError) Error() string { return e.err.Error() }
func (e *UnknownFieldError) Unwrap() error { return e.err }
// unknownFieldText is how the JSON decoder words a key its target does not have.
var unknownFieldText = regexp.MustCompile(`^json: unknown field "([^"]*)"$`)
// typedUnknown is err as an UnknownFieldError when it is the decoder refusing an unknown key, else err.
func typedUnknown(err error) error {
if err == nil {
return nil
}
if m := unknownFieldText.FindStringSubmatch(err.Error()); m != nil {
return &UnknownFieldError{Field: m[1], err: err}
}
return err
}
// asUnknownField is the unknown key err is about, at any depth, or nil.
func asUnknownField(err error) *UnknownFieldError {
var u *UnknownFieldError
if errors.As(typedUnknown(err), &u) {
return u
}
return nil
}
// UnknownFieldReason is why a module whose stored manifest has a key this controller does not know is
// left out of a machine's declaration, or "" when it has none.
func UnknownFieldReason(m Manifest) string {
if m.unknown == "" {
return ""
}
return m.Module + " uses a field this controller does not know (" + m.unknown + "); it is left out " +
"until the controller is updated: nothing of it is changed on its machines and its contributions to " +
"other modules and its open ports stop. Its data is still backed up as this controller reads it, " +
"which may not be what its newer manifest asks, and it still provides what it provides " +
"(novox/hq ADR 0262)"
}
// backupView is what of a left-out module still reaches its machine: its data, so the backup holder
// keeps copying it, and the directories and accesses its data items name. No contribution, shell code
// or environment of its own: those are what leaving it out stops.
func backupView(m Manifest) Manifest {
view := Manifest{Module: m.Module, Data: m.Data, Accesses: m.Accesses, bestEffort: true}
for _, r := range m.Resources {
if fmt.Sprint(r["type"]) == "directory" {
view.Resources = append(view.Resources, r)
}
}
return view
}
// prunedFields is a stored manifest's fields with every key this controller does not know taken out,
// and the keys taken out (novox/hq ADR 0262). A second pass, after the strict one refused: each
// top-level field is decoded alone, and where an entry inside it has an unknown key, the one
// occurrence whose removal moves the decoder past it is removed — never a key of the same name that
// the entry around it knows. So a left-out module still provides what it provides, and its data and
// directories are still read, which its backup lines are made from.
func prunedFields(keys map[string]json.RawMessage) (manifestFields, []string, error) {
var removed []string
tree := map[string]any{}
for k, raw := range keys {
var v any
dec := json.NewDecoder(bytes.NewReader(raw))
dec.UseNumber()
if err := dec.Decode(&v); err != nil {
return manifestFields{}, nil, err
}
tree[k] = v
}
for _, k := range sortedAnyKeys(tree) {
for tries := 0; ; tries++ {
err := decodesAlone(k, tree[k])
if err == nil {
break
}
unknown := asUnknownField(err)
if unknown == nil || tries > 64 {
return manifestFields{}, nil, err
}
if unknown.Field == k {
delete(tree, k)
removed = append(removed, k)
break
}
path, ok := removalThatHelps(k, tree[k], unknown.Field, err.Error())
if !ok {
// The same key unknown in two entries alike: no one removal changes the words.
// Every occurrence goes, and the field is judged again.
if removeEvery(tree[k], unknown.Field) == 0 {
return manifestFields{}, nil, err
}
path = "…." + unknown.Field
}
removed = append(removed, k+path)
}
}
raw, err := json.Marshal(tree)
if err != nil {
return manifestFields{}, nil, err
}
var fields manifestFields
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(&fields); err != nil {
return manifestFields{}, nil, err
}
return fields, removed, nil
}
// decodesAlone is whether one top-level field decodes strictly on its own.
func decodesAlone(k string, v any) error {
raw, err := json.Marshal(map[string]any{k: v})
if err != nil {
return err
}
var fields manifestFields
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
return dec.Decode(&fields)
}
// removalThatHelps removes, from v, the one occurrence of key whose removal changes what the strict
// decoder says of field k, and says where it was. Every other occurrence is left as it was.
func removalThatHelps(k string, v any, key, said string) (string, bool) {
var found bool
var where string
var walk func(node any, at string) bool
walk = func(node any, at string) bool {
switch n := node.(type) {
case map[string]any:
if value, has := n[key]; has {
delete(n, key)
// Accepted only when the decoder is past it: nothing left, or an unknown key said
// elsewhere. A different kind of error means the removal broke the entry; the same
// words mean this was not the occurrence it refused.
err := decodesAlone(k, v)
if err == nil || (asUnknownField(err) != nil && err.Error() != said) {
found, where = true, at+"."+key
return true
}
n[key] = value
}
for _, sub := range sortedAnyKeys(n) {
if walk(n[sub], at+"."+sub) {
return true
}
}
case []any:
for i, item := range n {
if walk(item, fmt.Sprintf("%s[%d]", at, i)) {
return true
}
}
}
return false
}
walk(v, "")
return where, found
}
// removeEvery removes key from every object in v, and says how many it removed.
func removeEvery(v any, key string) int {
n := 0
switch node := v.(type) {
case map[string]any:
if _, has := node[key]; has {
delete(node, key)
n++
}
for _, sub := range node {
n += removeEvery(sub, key)
}
case []any:
for _, item := range node {
n += removeEvery(item, key)
}
}
return n
}
+24 -4
View File
@@ -56,7 +56,7 @@ func (v *Verb) UnmarshalJSON(raw []byte) error {
decoder := json.NewDecoder(bytes.NewReader(trimmed))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&p); err != nil {
return fmt.Errorf("a served verb is a name or {name, description, input, output}: %w", err)
return fmt.Errorf("a served verb is a name or {name, description, input, output}: %w", typedUnknown(err))
}
if p.Name == "" {
return fmt.Errorf("a served verb has no name: %s", trimmed)
@@ -249,19 +249,24 @@ var ControllerVerbs = []Verb{
}, nil, "adopted"),
Replaces: []string{"mesh-controller token issue", "wg set"}},
{Name: "settings", Description: "Read or set what an assignment is configured with: a module's settings for the whole " +
"mesh, or for one machine. Without values or clear, answers the layer as it stands — read it before setting it; " +
"mesh, or for one machine. Without values or clear, answers the layer as it stands — read it before setting it — " +
"then every value the module gives a default or a layer sets, with where it comes from: the module's default, " +
"the mesh, or the machine (novox/hq ADR 0262); with list \"preferences\", or with no module, every " +
"module's preferences and each machine's value; " +
"with history, the layers it replaced. Setting replaces that layer whole and answers each key it adds (+), " +
"changes (~) and removes (-); a set that would remove a key is refused unless replace says it is meant " +
"(novox/hq ADR 0217). Takes effect at the next push. With clear, removes the layer and the module is back to " +
"what its definition says; a cleared or replaced layer is kept in the history.",
Input: schema(map[string]string{
"module": "the module's name",
"module": "the module's name; with list, only that module's",
"values": "the settings as a JSON object, for set",
"node": "one machine; the whole mesh when absent",
"clear": "\"true\" to remove the layer instead of setting it; not with values",
"replace": "\"true\": with values, the set is meant to remove the keys the layer had and it does not name",
"history": "\"true\": without values or clear, the layers this one replaced, the latest first",
}, []string{"module"}, "clear", "replace", "history")},
"list": "\"preferences\": every module's preferences — key, default and why — and the value on each " +
"machine it is assigned to with where it comes from; module and node narrow it (novox/hq ADR 0262)",
}, nil, "clear", "replace", "history")},
{Name: "command", Description: "Run one command line of the controller's own, as you would type it at its " +
"shell — `node account g14 jochen`, `node show ace`, `module list` — and answer what it printed. The " +
"generic verb beside the named ones (novox/hq ADR 0154): everything the binary can do, without a verb " +
@@ -419,6 +424,21 @@ var ControllerVerbs = []Verb{
"why": "with consumer or older-than: why — required, and recorded in the hand-act log",
"cause": "with consumer or older-than: the cause in a word (cleanup-waiting when absent)",
}, nil, "confirm")},
// What a consumer gave up on (novox/hq issue 330): kept in DEAD_LETTERS until a person acts on it.
{Name: "dead-letters", Description: "Every message a consumer on the bus gave up on after handing it over " +
"as often as it may, kept in DEAD_LETTERS: whose consumer, the subject, how often it was handed over and " +
"when it was given up, newest first. With id: that one whole, with what it said. With deliver: hand it " +
"again to the consumer that gave it up, and nobody else; with drop: let it go for good. Delivering and " +
"dropping are hand acts, which say why (novox/hq issue 330).",
Input: schema(map[string]string{
"id": "a dead letter's id, as the list gives it: that one whole, with its message",
"consumer": "a consumer's name, or <stream>.<consumer>: only what that one gave up on",
"limit": "how many to list, newest first (default 50); only when listing",
"deliver": "a dead letter's id: deliver it again to the consumer that gave it up (needs why)",
"drop": "a dead letter's id: drop it for good (needs why)",
"why": "with deliver or drop: why — required, and recorded in the hand-act log",
"cause": "with deliver or drop: the cause in a word (dead-letter when absent)",
}, nil)},
// The data every machine declares (novox/hq ADR 0233).
{Name: "data", Description: "Every item of data every machine declares, as the self-check last measured it: " +
"its class (irreplaceable, rebuildable, cache), where it is, its size, its newest write, its newest good " +
+24
View File
@@ -5,15 +5,37 @@ import (
"encoding/json"
"errors"
"fmt"
"log"
"sort"
"strconv"
"strings"
"sync"
"time"
"github.com/jackc/pgx/v5"
"github.com/novox/mesh-controller/internal/catalogue"
)
// noted is each stored manifest's unknown key already said, so a catalogue read on every plan says it once.
var noted sync.Map
// noteUnknown says, once per module and key, that a stored manifest has a key this controller does not
// know (novox/hq ADR 0262). A manifest registered under a newer controller is read by an older one after
// a rollback; refusing it here failed the whole catalogue. The module is left out of every machine by name
// (catalogue.LeftOut), and the controller's tick raises a condition for it.
func noteUnknown(m catalogue.Manifest) {
u := m.UnknownField()
if u == "" {
return
}
if _, said := noted.LoadOrStore(m.Module+"\x00"+u, true); said {
return
}
log.Printf("the stored manifest of %s has a key this controller does not know (%s): a newer controller "+
"registered it, and %s is left out of every machine until this controller is updated (novox/hq ADR 0262)",
m.Module, u, m.Module)
}
// ErrNoSuchModule is what the mesh says about a module it has never been told about.
var ErrNoSuchModule = errors.New("no module of that name")
@@ -259,6 +281,7 @@ func (i *Inventory) Catalogue(ctx context.Context) (map[string]catalogue.Manifes
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
noteUnknown(m)
out[m.Module] = m
}
return out, rows.Err()
@@ -1217,6 +1240,7 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
noteUnknown(m)
entry := Entry{Manifest: m, Source: source, On: on}
if source.Repository == providedBy {
// It came with the control plane. Not a repository, and showing it as one would have
+34
View File
@@ -0,0 +1,34 @@
package inventory
import (
"strings"
"testing"
)
// A manifest a newer controller registered, with a key this one does not know, does not stop the
// catalogue from loading: it is read, marked, and every other module is read as before (novox/hq ADR 0262).
func TestTheCatalogueLoadsAroundAManifestWithAnUnknownKey(t *testing.T) {
inv := fresh(t)
if err := inv.RegisterModule(t.Context(), manifest("thing", nil, nil), Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(t.Context(),
`insert into module (name, manifest) values ($1, $2)`, "later",
`{"module": "later", "version": "2", "state": [{"name": "s", "a-field-from-later": 1}]}`); err != nil {
t.Fatal(err)
}
known, err := inv.Catalogue(t.Context())
if err != nil {
t.Fatalf("one manifest with an unknown key failed the whole catalogue: %v", err)
}
if known["thing"].Module != "thing" || known["thing"].UnknownField() != "" {
t.Fatalf("the other module: %+v", known["thing"])
}
if !strings.Contains(known["later"].UnknownField(), "a-field-from-later") {
t.Fatalf("the newer manifest is not marked: %q", known["later"].UnknownField())
}
entries, err := inv.Catalogued(t.Context())
if err != nil || len(entries) < 2 {
t.Fatalf("the listing: %v %d", err, len(entries))
}
}
+15 -4
View File
@@ -34,6 +34,9 @@ const (
AdvisoryMaxDeliveries = "max-deliveries"
AdvisoryRefused = "refused"
AdvisoryConsumerLost = "consumer-lost"
// AdvisoryNotKept is the token of a max-deliveries advisory whose message could not be kept in
// DEAD_LETTERS (novox/hq issue 330): said from the log, since the stream does not hold it.
AdvisoryNotKept = "not-kept"
)
// Advisory is one thing the bus said, kept as its newest word and how often it was said.
@@ -43,6 +46,8 @@ type Advisory struct {
ID string
// Stream and Consumer are the consumer it is about, when it is about one.
Stream, Consumer string
// Token is the condition key's last part, when it is not Kind.
Token string
// Said is the newest saying, in the mesh's words.
Said string
First, Last time.Time
@@ -62,7 +67,7 @@ var Advisories = &AdvisoryLog{seen: map[string]*Advisory{}}
func (l *AdvisoryLog) Heard(a Advisory, at time.Time) {
l.mu.Lock()
defer l.mu.Unlock()
key := a.Kind + "/" + a.ID
key := a.Kind + "/" + a.ID + "/" + a.Token
if had, ok := l.seen[key]; ok {
had.Last, had.Said, had.Count = at, a.Said, had.Count+1
return
@@ -107,8 +112,9 @@ func ReadAdvisory(subject string, body []byte) (Advisory, bool) {
switch {
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."):
return Advisory{Kind: AdvisoryMaxDeliveries, ID: id, Stream: a.Stream, Consumer: a.Consumer,
Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: it will not be delivered "+
"again, and what it asked for was not done", who, a.StreamSeq, a.Deliveries)}, true
Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: what it asked for was not "+
"done, and the controller keeps it in %s until it is delivered again or dropped", who, a.StreamSeq,
a.Deliveries, broker.DeadLettersStream)}, true
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.DELETED."):
if !MeshNamed(a.Stream, a.Consumer) {
// A reader's own consumer, gone when it finished — every watch of a bucket and every
@@ -175,7 +181,12 @@ func (s *Server) HearAdvisories(logf func(string, ...any)) (func(), error) {
for _, subject := range broker.BusAdvisories {
sub, err := conn.Subscribe(subject, func(m *nats.Msg) {
if a, ok := ReadAdvisory(m.Subject, m.Data); ok {
Advisories.Heard(a, time.Now())
// A message given up on is said from DEAD_LETTERS, where it is kept, for as long as it is
// kept (novox/hq issue 330) — not for an hour after the server said it; one that could not
// be kept is recorded by whoever tried.
if a.Kind != AdvisoryMaxDeliveries {
Advisories.Heard(a, time.Now())
}
logf("the bus says: %s", a.Said)
}
})
+2 -1
View File
@@ -37,7 +37,8 @@ func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) {
[]byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`))
if !ok || a.Kind != AdvisoryMaxDeliveries || a.ID != "EVENTS.anchor_shop" ||
a.Said != "how shop on anchor hears what it consumes handed message 7 over 5 times and gave up on it: "+
"it will not be delivered again, and what it asked for was not done" {
"what it asked for was not done, and the controller keeps it in DEAD_LETTERS until it is delivered "+
"again or dropped" {
t.Fatalf("%+v", a)
}
}
+6
View File
@@ -47,6 +47,12 @@ func BuildsOverNATSOn(address, seat string) (Builders, error) {
return &natsBuilds{js: js, owned: true, seat: seat}, nil
}
// BuildsOn is the asking side on a connection the caller holds, which Close leaves open: the serving
// controller asks on its own rather than dialling one per build (novox/hq issue 327).
func BuildsOn(js *broker.JetStream, seat string) Builders {
return &natsBuilds{js: js, seat: seat}
}
// role is the seat asked: what the asker was made for, or the current build role for one made
// without saying (a test building the struct by hand).
func (b *natsBuilds) role() string {
+5 -2
View File
@@ -8,13 +8,16 @@ import (
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
// the streams and the controller's consumers are asserted by Raise, before anything is served.
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
logger := newLog()
inbound := Nats(js).(*natsInbound)
inbound.log = logger
return &Server{
inbound: Nats(js),
inbound: inbound,
bus: OverNATS{JS: js.Context(), Conn: js.Conn()},
js: js,
enroller: enroller,
listener: listener,
log: newLog(),
log: logger,
}
}
+14
View File
@@ -7,6 +7,7 @@ import (
"fmt"
"log"
"regexp"
"runtime/debug"
"sort"
"strconv"
"strings"
@@ -575,6 +576,19 @@ func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply
done := make(chan outcome, 1)
go func() {
var body []byte
// **A handler that panics is an answer, not a dead controller.** A verb answered in the serving
// process (novox/hq issue 327: conditions, hand-acts, queue, dead-letters) runs on this goroutine,
// where a panic would take every other call and the controller with it.
defer func() {
if p := recover(); p != nil {
if logger != nil {
logger.Printf("%s.%s: call %s panicked: %v\n%s", seat, verb, c.ID, p, debug.Stack())
}
body, _ = json.Marshal(map[string]any{"error": fmt.Sprintf("%s failed inside the controller and "+
"answered nothing: %v. It is in the controller's log", verb, p)})
done <- outcome{body, true}
}
}()
result, err := handle(ctx, args)
failed := err != nil
if errors.Is(err, ErrHandingOver) {
+30
View File
@@ -11,6 +11,10 @@ import (
"sync"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/testbus"
)
// answers collects what a call answered, and fails a second answer: the bus permits one.
@@ -223,3 +227,29 @@ func TestAJoinTokenIsShownToItsCallerAndNotKept(t *testing.T) {
}
}
}
// A verb whose handler panics answers an error, and the controller serves the next call (review of
// novox/hq issue 327: read verbs are answered in the serving process now).
func TestAHandlerThatPanicsAnswersAnError(t *testing.T) {
conn, err := nats.Connect(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
stop, err := OverNATS{Conn: conn}.ServeSeatTools("panicky", map[string]ToolHandler{
"boom": func(context.Context, json.RawMessage) (any, error) { panic("nil map") },
"fine": func(context.Context, json.RawMessage) (any, error) { return "ok", nil },
}, nil)
if err != nil {
t.Fatal(err)
}
defer stop()
answer, err := AskMeshSeatTool(context.Background(), conn, "panicky", "boom", map[string]any{}, 5*time.Second)
if err != nil || !strings.Contains(answer.Error, "failed inside the controller") {
t.Fatalf("a panic answered %+v (%v)", answer, err)
}
answer, err = AskMeshSeatTool(context.Background(), conn, "panicky", "fine", map[string]any{}, 5*time.Second)
if err != nil || answer.Error != "" {
t.Fatalf("the call after a panic answered %+v (%v)", answer, err)
}
}
+372
View File
@@ -0,0 +1,372 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"slices"
"strconv"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// What a consumer gave up on, kept until a person acts on it (novox/hq issue 330, design 25 §3).
//
// A durable consumer hands a message over as often as it may, gives up on it and says so on
// `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES`. The server's notice is captured in its own stream
// (broker.DeadLetterNoticesStream), so one said while no controller listens waits for the next. The
// serving controller takes each notice, fetches the message it names by its sequence, while the source
// stream still holds it, and keeps a copy in DEAD_LETTERS with what the notice said: the consumer, the
// subject, how often it was handed over and when it was given up. Until then a message given up on was
// kept only by its source, which drops an event after a week, and nothing said which.
//
// The copy's headers are the message's own, but for the server's (`Nats-…`), which describe the first
// publish and would refuse or de-duplicate the copy; they are kept under DeadHeader names.
// The headers a kept message carries beside its own.
const (
DeadHeaderStream = "Mesh-Dead-Stream"
DeadHeaderConsumer = "Mesh-Dead-Consumer"
DeadHeaderSequence = "Mesh-Dead-Sequence"
DeadHeaderSubject = "Mesh-Dead-Subject"
DeadHeaderDeliveries = "Mesh-Dead-Deliveries"
DeadHeaderGaveUp = "Mesh-Dead-Gave-Up"
DeadHeaderPublished = "Mesh-Dead-Published"
DeadHeaderLost = "Mesh-Dead-Lost"
// DeadHeaderMsgID is the message's own de-duplication id, kept, since the copy carries its own.
DeadHeaderMsgID = "Mesh-Dead-Msg-Id"
// AgainHeader marks a message delivered again, with the kept copy's id it came from.
AgainHeader = "Mesh-Delivered-Again"
)
// DeadLetter is one message a consumer gave up on, as DEAD_LETTERS keeps it.
type DeadLetter struct {
// ID is its sequence in DEAD_LETTERS: what the verb names it by.
ID uint64 `json:"id"`
Stream string `json:"stream"`
Consumer string `json:"consumer"`
// Who is the consumer in the mesh's words: whose, and for what.
Who string `json:"who"`
Subject string `json:"subject"`
Sequence uint64 `json:"sequence"`
Deliveries uint64 `json:"deliveries"`
GaveUp time.Time `json:"gave_up"`
Published time.Time `json:"published,omitzero"`
// Lost says why the message itself is not kept: its source no longer held it when the notice was
// taken. The record of it is kept all the same, so it is said and dropped, never silently missing.
Lost string `json:"lost,omitempty"`
Size int `json:"size"`
// Body and Headers are the message itself, given only for one dead letter asked by its id; a header
// with several values keeps them all.
Body string `json:"body,omitempty"`
Headers map[string][]string `json:"headers,omitempty"`
}
// maxDeliveries is the part of the server's notice a kept message is made from.
type maxDeliveries struct {
Stream string `json:"stream"`
Consumer string `json:"consumer"`
StreamSeq uint64 `json:"stream_seq"`
Deliveries uint64 `json:"deliveries"`
Timestamp time.Time `json:"timestamp"`
}
// KeepDeadLetter keeps the message one maximum-deliveries notice is about. An error leaves nothing
// kept, and the notice is to be taken again: the source still holds the message, or a full
// DEAD_LETTERS has room again. A source that no longer holds it is no error: the record is kept with
// why, so it is said.
func KeepDeadLetter(js nats.JetStreamContext, notice []byte) (DeadLetter, error) {
var a maxDeliveries
if err := json.Unmarshal(notice, &a); err != nil || a.Stream == "" || a.Consumer == "" || a.StreamSeq == 0 {
return DeadLetter{}, fmt.Errorf("not a maximum-deliveries notice the mesh can read: %.200s", notice)
}
gaveUp := a.Timestamp
if gaveUp.IsZero() {
gaveUp = time.Now()
}
kept := &nats.Msg{Subject: broker.DeadLetterSubject(a.Stream, a.Consumer), Header: nats.Header{}}
raw, err := js.GetMsg(a.Stream, a.StreamSeq)
switch {
case err == nil:
for k, v := range raw.Header {
if strings.HasPrefix(k, "Nats-") {
continue
}
kept.Header[k] = v
}
if id := raw.Header.Get(nats.MsgIdHdr); id != "" {
kept.Header.Set(DeadHeaderMsgID, id)
}
kept.Header.Set(DeadHeaderSubject, raw.Subject)
kept.Header.Set(DeadHeaderPublished, raw.Time.UTC().Format(time.RFC3339Nano))
kept.Data = raw.Data
case errors.Is(err, nats.ErrMsgNotFound) || errors.Is(err, nats.ErrStreamNotFound):
kept.Header.Set(DeadHeaderLost, fmt.Sprintf("%s no longer held message %d when the notice was taken: %v",
a.Stream, a.StreamSeq, err))
default:
return DeadLetter{}, fmt.Errorf("message %d of %s could not be read: %w", a.StreamSeq, a.Stream, err)
}
kept.Header.Set(DeadHeaderStream, a.Stream)
kept.Header.Set(DeadHeaderConsumer, a.Consumer)
kept.Header.Set(DeadHeaderSequence, strconv.FormatUint(a.StreamSeq, 10))
kept.Header.Set(DeadHeaderDeliveries, strconv.FormatUint(a.Deliveries, 10))
kept.Header.Set(DeadHeaderGaveUp, gaveUp.UTC().Format(time.RFC3339Nano))
// One copy per message given up, however often its notice is taken: a controller that copied and
// stopped before acknowledging, or two controllers each taking it.
kept.Header.Set(nats.MsgIdHdr, fmt.Sprintf("%s.%s.%d", a.Stream, a.Consumer, a.StreamSeq))
ack, err := js.PublishMsg(kept)
if err != nil {
return DeadLetter{}, fmt.Errorf("message %d of %s could not be kept in %s: %w", a.StreamSeq, a.Stream,
broker.DeadLettersStream, err)
}
return deadLetterOf(ack.Sequence, kept.Subject, kept.Header, kept.Data, false), nil
}
// deadLetterOf reads a kept message.
func deadLetterOf(id uint64, subject string, h nats.Header, data []byte, whole bool) DeadLetter {
d := DeadLetter{ID: id, Stream: h.Get(DeadHeaderStream), Consumer: h.Get(DeadHeaderConsumer),
Subject: h.Get(DeadHeaderSubject), Lost: h.Get(DeadHeaderLost), Size: len(data)}
if d.Stream == "" || d.Consumer == "" {
d.Stream, d.Consumer, _ = broker.DeadLetterOf(subject)
}
d.Who = ConsumerInWords(d.Stream, d.Consumer)
d.Sequence, _ = strconv.ParseUint(h.Get(DeadHeaderSequence), 10, 64)
d.Deliveries, _ = strconv.ParseUint(h.Get(DeadHeaderDeliveries), 10, 64)
d.GaveUp, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderGaveUp))
d.Published, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderPublished))
if whole {
d.Body = string(data)
d.Headers = map[string][]string{}
for k, v := range h {
if !strings.HasPrefix(k, "Mesh-Dead-") && !strings.HasPrefix(k, "Nats-") {
d.Headers[k] = append([]string(nil), v...)
}
}
}
return d
}
// HeldDeadLetters is how many messages DEAD_LETTERS holds for each consumer, by `<stream>.<consumer>`.
func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) {
info, err := js.StreamInfo(broker.DeadLettersStream, &nats.StreamInfoRequest{SubjectsFilter: ">"})
if err != nil {
return nil, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err)
}
held := map[string]int{}
for subject, n := range info.State.Subjects {
if stream, consumer, ok := broker.DeadLetterOf(subject); ok && n > 0 {
held[stream+"."+consumer] += int(n)
}
}
return held, nil
}
// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them (all when most is not
// positive); for one consumer when consumer names one (its name, or `<stream>.<consumer>`). The total is
// the stream's own count per consumer, so it is right however few are read.
func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLetter, int, error) {
held, err := HeldDeadLetters(js)
if err != nil {
return nil, 0, err
}
total := 0
for key, n := range held {
if consumer == "" || key == consumer || strings.HasSuffix(key, "."+consumer) {
total += n
}
}
if total == 0 {
return nil, 0, nil
}
info, err := js.StreamInfo(broker.DeadLettersStream)
if err != nil {
return nil, 0, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err)
}
var out []DeadLetter
for seq := info.State.LastSeq; seq >= info.State.FirstSeq && seq > 0; seq-- {
if most > 0 && len(out) >= most || len(out) >= total {
break
}
raw, err := js.GetMsg(broker.DeadLettersStream, seq)
if errors.Is(err, nats.ErrMsgNotFound) {
continue // delivered again or dropped
}
if err != nil {
return out, total, fmt.Errorf("dead letter %d cannot be read: %w", seq, err)
}
d := deadLetterOf(seq, raw.Subject, raw.Header, raw.Data, false)
if consumer != "" && consumer != d.Consumer && consumer != d.Stream+"."+d.Consumer {
continue
}
out = append(out, d)
}
return out, total, nil
}
// ErrNoDeadLetter is a dead letter's id DEAD_LETTERS does not hold.
var ErrNoDeadLetter = errors.New("no such dead letter")
// DeadLetterNamed is one kept message whole: what it said, and the headers it said it with.
func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
raw, err := js.GetMsg(broker.DeadLettersStream, id)
if errors.Is(err, nats.ErrMsgNotFound) {
return DeadLetter{}, fmt.Errorf("%w: %s holds no message %d — delivered again or dropped already, "+
"or never kept; dead-letters lists what it holds", ErrNoDeadLetter, broker.DeadLettersStream, id)
}
if err != nil {
return DeadLetter{}, fmt.Errorf("dead letter %d cannot be read: %w", id, err)
}
return deadLetterOf(id, raw.Subject, raw.Header, raw.Data, true), nil
}
// AgainTo is where a kept message is delivered again so that only the consumer that gave it up gets
// it: an event under that consumer's own again subject on EVENTS — and only when the consumer exists and
// filters that subject, so a message is never let go as delivered while nobody receives it. Any other
// stream's message is refused, with why: an ask given up on by a seat's worker is not delivered again yet
// (novox/hq issue 330's follow-up), since publishing it again leaves the original stuck in the queue.
func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
switch {
case d.Lost != "":
return "", fmt.Errorf("dead letter %d holds no message to deliver: %s. Drop it", d.ID, d.Lost)
case d.Subject == "":
return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+
"delivered again. Drop it", d.ID)
case d.Stream != broker.EventsStream:
return "", fmt.Errorf("dead letter %d is from %s, and only an event is delivered again: publishing it "+
"again would reach every consumer of its subject, or leave the original in its queue. Drop it, and "+
"have its sender say it again", d.ID, d.Stream)
}
info, err := js.ConsumerInfo(d.Stream, d.Consumer)
if errors.Is(err, nats.ErrConsumerNotFound) {
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so nothing would "+
"receive it. Drop it, or deliver it again once the module is assigned there again", d.ID, d.Who)
}
if err != nil {
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
}
want := broker.AgainFilter(d.Consumer)
filters := append([]string{info.Config.FilterSubject}, info.Config.FilterSubjects...)
if !slices.Contains(filters, want) {
return "", fmt.Errorf("%s does not yet take events delivered again (it does not filter %s): the controller "+
"sets that at the next send to its machine. Nothing was done, and dead letter %d is still kept",
d.Who, want, d.ID)
}
return broker.AgainSubject(d.Consumer, d.Subject), nil
}
// DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it
// from DEAD_LETTERS. The message carries its own headers and AgainHeader; its de-duplication id is the
// kept copy's, so asking twice delivers it once.
func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, error) {
d, err := DeadLetterNamed(js, id)
if err != nil {
return d, "", err
}
to, err := AgainTo(js, d)
if err != nil {
return d, "", err
}
again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)}
for k, v := range d.Headers {
again.Header[k] = append([]string(nil), v...)
}
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
if _, err := js.PublishMsg(again); err != nil {
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err)
}
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
return d, to, fmt.Errorf("dead letter %d was delivered again on %s and could not be removed from %s: %w",
id, to, broker.DeadLettersStream, err)
}
return d, to, nil
}
// DropDeadLetter removes a kept message for good.
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
d, err := DeadLetterNamed(js, id)
if err != nil {
return d, err
}
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil {
return d, fmt.Errorf("dead letter %d could not be dropped: %w", id, err)
}
return d, nil
}
// noticeRetry is how long a notice whose message could not be kept waits before it is taken again.
var noticeRetry = time.Minute
// keepingGivenUp keeps taking the notices until ctx ends; the returned function stops it. A subscription
// that cannot be made is said — in the log, and as the max-deliveries condition of the notices themselves
// — and tried again every noticeRetry, so it never stops the controller serving.
func keepingGivenUp(ctx context.Context, bus *broker.JetStream, logger *log.Logger) func() {
ctx, cancel := context.WithCancel(ctx)
done := make(chan struct{})
go func() {
defer close(done)
for {
sub, err := func() (*nats.Subscription, error) {
if err := bus.EnsureConsumer(broker.NoticesConsumer()); err != nil {
return nil, err
}
return keepGivenUp(bus.Context(), logger)
}()
if err == nil {
<-ctx.Done()
_ = sub.Unsubscribe()
return
}
logger.Printf("the notices of messages consumers gave up on cannot be taken, so none is kept until "+
"they can: %v", err)
Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: broker.DeadLetterNoticesStream + "." +
broker.ControllerName, Stream: broker.DeadLetterNoticesStream, Consumer: broker.ControllerName,
Token: AdvisoryNotKept, Said: "the controller cannot take the notices of messages consumers gave " +
"up on, so none is kept: " + err.Error()}, time.Now())
select {
case <-ctx.Done():
return
case <-time.After(noticeRetry):
}
}
}()
return func() {
cancel()
<-done
}
}
// keepGivenUp takes the server's maximum-deliveries notices off their stream and keeps the message
// each is about, for as long as the subscription stands. A notice that could not be kept is offered
// again after noticeRetry, and said: in the log, and as the consumer's max-deliveries condition.
func keepGivenUp(js nats.JetStreamContext, logger *log.Logger) (*nats.Subscription, error) {
return js.Subscribe("", func(m *nats.Msg) {
d, err := KeepDeadLetter(js, m.Data)
if err != nil {
logger.Printf("a message a consumer gave up on could NOT be kept: %v", err)
var a maxDeliveries
if json.Unmarshal(m.Data, &a) == nil && a.Stream != "" && a.Consumer != "" {
Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: a.Stream + "." + a.Consumer,
Stream: a.Stream, Consumer: a.Consumer, Token: AdvisoryNotKept,
Said: fmt.Sprintf("%s gave up on message %d, and it could not be kept: %v",
ConsumerInWords(a.Stream, a.Consumer), a.StreamSeq, err)}, time.Now())
}
_ = m.NakWithDelay(noticeRetry)
return
}
if d.Lost != "" {
logger.Printf("%s gave up on message %d of %s, which its stream no longer held; kept as dead letter %d "+
"without it", d.Who, d.Sequence, d.Stream, d.ID)
} else {
logger.Printf("%s gave up on message %d of %s (%s) after %d deliveries; kept as dead letter %d",
d.Who, d.Sequence, d.Stream, d.Subject, d.Deliveries, d.ID)
}
_ = m.Ack()
}, nats.Bind(broker.DeadLetterNoticesStream, broker.ControllerName), nats.ManualAck())
}
+336
View File
@@ -0,0 +1,336 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// A module's consumer as the controller makes it, on a bus whose streams the controller asserted.
func aModuleConsumer(t *testing.T, js *broker.JetStream, node, module string, consumes ...string) jetstream.Consumer {
t.Helper()
c, ok := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: node, Module: module,
Consumes: consumes, PasswordHash: "x"})
if !ok {
t.Fatal("a module that consumes got no consumer")
}
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
reading, err := api.Consumer(t.Context(), broker.EventsStream, c.Name)
if err != nil {
t.Fatal(err)
}
return reading
}
// next is the one message a consumer hands over within a second, or nil.
func next(t *testing.T, c jetstream.Consumer) jetstream.Msg {
t.Helper()
batch, err := c.Fetch(1, jetstream.FetchMaxWait(time.Second))
if err != nil {
t.Fatal(err)
}
for msg := range batch.Messages() {
return msg
}
return nil
}
// keeping takes the notices as the serving controller does, until the test ends.
func keeping(t *testing.T, js *broker.JetStream) {
t.Helper()
if err := broker.AssertMeshConsumers(js); err != nil {
t.Fatal(err)
}
sub, err := keepGivenUp(js.Context(), log.New(io.Discard, "", 0))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
}
// givenUp hands one event to a consumer as often as it may, failing each time, and waits for it kept.
func givenUp(t *testing.T, js *broker.JetStream, c jetstream.Consumer) DeadLetter {
t.Helper()
for {
msg := next(t, c)
if msg == nil {
break
}
_ = msg.Nak()
}
var kept []DeadLetter
eventually(t, "the event kept", func() bool {
kept, _, _ = DeadLetters(js.Context(), c.CachedInfo().Name, 0)
return len(kept) == 1
})
return kept[0]
}
// **Delivered again to the consumer that gave it up, and nobody else**: another module consuming the
// same event handled it the first time and must not handle it twice.
func TestADeadLetterIsDeliveredAgainToItsConsumerAlone(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
handling := aModuleConsumer(t, js, "media", "radarr", "gitea.pull.merged")
const subject = "mesh.mod.gitea.event.pull.merged"
msg := &nats.Msg{Subject: subject, Data: []byte(`{"n":1}`), Header: nats.Header{}}
msg.Header.Set("x-node", "forge")
msg.Header.Set(nats.MsgIdHdr, "event-1")
if _, err := js.Context().PublishMsg(msg); err != nil {
t.Fatal(err)
}
if m := next(t, handling); m == nil {
t.Fatal("the other module was not handed the event")
} else {
_ = m.Ack()
}
d := givenUp(t, js, failing)
if d.Consumer != "media_sonarr" || d.Stream != "EVENTS" || d.Subject != subject || d.Deliveries != 5 ||
d.GaveUp.IsZero() || d.Published.IsZero() || d.Lost != "" {
t.Fatalf("kept as %+v", d)
}
held, err := HeldDeadLetters(js.Context())
if err != nil || held["EVENTS.media_sonarr"] != 1 {
t.Fatalf("held %v (%v)", held, err)
}
whole, err := DeadLetterNamed(js.Context(), d.ID)
if err != nil || whole.Body != `{"n":1}` || len(whole.Headers["x-node"]) != 1 || whole.Headers["x-node"][0] != "forge" {
t.Fatalf("one asked whole is %+v (%v)", whole, err)
}
_, to, err := DeliverAgain(js.Context(), d.ID)
if err != nil {
t.Fatal(err)
}
if to != "mesh.again.media_sonarr.mod.gitea.event.pull.merged" {
t.Fatalf("delivered again on %s", to)
}
again := next(t, failing)
if again == nil {
t.Fatal("the consumer that gave it up was not handed it again")
}
if string(again.Data()) != `{"n":1}` || again.Headers().Get("x-node") != "forge" ||
again.Headers().Get(AgainHeader) == "" {
t.Fatalf("handed again as %s %v", again.Data(), again.Headers())
}
if original, ok := broker.OriginalOfAgain(again.Subject()); !ok || original != subject {
t.Fatalf("delivered again on %s, which does not say the event's own subject", again.Subject())
}
_ = again.Ack()
if m := next(t, handling); m != nil {
t.Fatalf("the module that handled it was handed it twice: %s", m.Subject())
}
if _, err := DeadLetterNamed(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
t.Fatalf("still kept after it was delivered again: %v", err)
}
held, _ = HeldDeadLetters(js.Context())
if held["EVENTS.media_sonarr"] != 0 {
t.Fatalf("still held: %v", held)
}
// Asked twice, delivered once: the second finds nothing kept.
if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
t.Fatalf("a second delivery answered %v", err)
}
}
func TestADeadLetterDroppedIsGone(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil {
t.Fatal(err)
}
d := givenUp(t, js, failing)
if _, err := DropDeadLetter(js.Context(), d.ID); err != nil {
t.Fatal(err)
}
if left, total, err := DeadLetters(js.Context(), "", 0); err != nil || total != 0 || len(left) != 0 {
t.Fatalf("after the drop: %v %d %v", left, total, err)
}
if m := next(t, failing); m != nil {
t.Fatalf("a dropped event was handed over: %s", m.Subject())
}
}
// A notice whose message its stream no longer holds is kept as a record that says so — said, and
// dropped by a person, never silently missing — and cannot be delivered again.
func TestANoticeWhoseMessageIsGoneIsKeptAndSaysSo(t *testing.T) {
js := aBus(t)
d, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+
`"deliveries":5,"timestamp":"2026-10-06T15:51:00Z"}`))
if err != nil {
t.Fatal(err)
}
if !strings.Contains(d.Lost, "no longer held message 42") || d.Consumer != "media_sonarr" {
t.Fatalf("kept as %+v", d)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "Drop it") {
t.Fatalf("a record without its message was delivered again: %v", err)
}
// The same notice taken twice is kept once.
if _, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+
`"deliveries":5}`)); err != nil {
t.Fatal(err)
}
if _, total, _ := DeadLetters(js.Context(), "media_sonarr", 0); total != 1 {
t.Fatalf("one notice taken twice is kept %d times", total)
}
}
// What only the giving-up consumer filters is its own: a subject delivered again reads back as the
// event's, and a module's runtime reads the same key from it (the tokens around `.event.`).
func TestASubjectDeliveredAgainSaysTheEventsOwn(t *testing.T) {
for _, original := range []string{"mesh.mod.gitea.event.pull.merged", "mesh.seat.node-build-agent.event.built"} {
again := broker.AgainSubject("ace_sonarr", original)
if back, ok := broker.OriginalOfAgain(again); !ok || back != original {
t.Errorf("%s reads back as %s", again, back)
}
key := func(subject string) string {
before, event, _ := strings.Cut(subject, ".event.")
parts := strings.Split(before, ".")
return parts[len(parts)-1] + "." + event
}
if key(again) != key(original) {
t.Errorf("%s reads as key %s, the event as %s", again, key(again), key(original))
}
}
if _, ok := broker.OriginalOfAgain("mesh.mod.gitea.event.pull.merged"); ok {
t.Error("an event's own subject reads as delivered again")
}
}
// Only an event is delivered again, and only to a consumer that is on the bus and filters its again
// subject: a dead letter is never let go as delivered while nobody receives it.
func TestADeadLetterIsDeliveredAgainOnlyWhereItIsReceived(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil {
t.Fatal(err)
}
d := givenUp(t, js, failing)
// A consumer as it was before this fix: its filters do not take what is delivered again.
c, _ := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr",
Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"})
c.Filters = c.Filters[:len(c.Filters)-1]
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "does not yet take") {
t.Fatalf("delivered to a consumer that does not filter its again subject: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// The consumer gone: nothing would receive it.
if err := js.Context().DeleteConsumer(broker.EventsStream, c.Name); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "no longer on the bus") {
t.Fatalf("delivered to a consumer that is gone: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// Not an event: refused, whatever stream it is from.
for _, other := range []DeadLetter{
{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", Subject: "mesh.seat.telegram-sender.accept.send"},
{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"},
} {
if _, err := AgainTo(js.Context(), other); err == nil {
t.Errorf("%s was given a subject to be delivered again on", other.Stream)
}
}
}
// A total says how many are held, however few the list carries.
func TestTheListSaysTheTotalAndStopsAtItsLimit(t *testing.T) {
js := aBus(t)
for seq := 1; seq <= 5; seq++ {
if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_sonarr",`+
`"stream_seq":%d,"deliveries":5}`, seq))); err != nil {
t.Fatal(err)
}
}
list, total, err := DeadLetters(js.Context(), "media_sonarr", 2)
if err != nil || total != 5 || len(list) != 2 || list[0].Sequence != 5 {
t.Fatalf("%d of %d (%v): %+v", len(list), total, err, list)
}
if _, total, _ := DeadLetters(js.Context(), "another_one", 2); total != 0 {
t.Fatalf("another consumer's total is %d", total)
}
}
// An event the controller itself gave up on, delivered again by a person, is acted on as the event it
// was: it arrives under the controller's again subject, which its consumer filters.
func TestTheControllerActsOnAnEventDeliveredAgain(t *testing.T) {
js := aBus(t)
told := &toldAbout{}
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()}
if err := s.Follows(told); err != nil {
t.Fatal(err)
}
if err := s.Answers(replaysWith{}); err != nil {
t.Fatal(err)
}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
eventually(t, "the controller's event consumer being made", func() bool {
_, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName)
return err == nil
})
moved, _ := json.Marshal(Upgraded{Module: "gitea", Commit: "abcdef0123"})
if _, err := js.Context().Publish(broker.AgainSubject(broker.ControllerName, broker.ControllerFollows[0]), moved); err != nil {
t.Fatal(err)
}
eventually(t, "the upgrade delivered again reaching the controller", func() bool { return told.count() == 1 })
}
// The notices cannot be taken (their stream is gone): said as a condition, and the controller serves on.
func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) {
js := aBus(t)
if err := js.Context().DeleteStream(broker.DeadLetterNoticesStream); err != nil {
t.Fatal(err)
}
stop := keepingGivenUp(context.Background(), js, quiet())
defer stop()
eventually(t, "the failure said", func() bool {
for _, a := range Advisories.Since(time.Now().Add(-time.Minute)) {
if a.Token == AdvisoryNotKept && a.Stream == broker.DeadLetterNoticesStream {
return true
}
}
return false
})
held := &counted{}
_, stopServing := servingOn(t, js, held)
defer stopServing()
if _, err := js.Context().Publish("mesh.control.anchor.report", []byte(`{"node":"anchor"}`)); err != nil {
t.Fatal(err)
}
eventually(t, "a report heard while the notices cannot be taken", func() bool { return held.count() == 1 })
}
+22 -1
View File
@@ -41,6 +41,15 @@ type natsInbound struct {
// that restarts loses these and starts the window again, which is correct — it is holding
// nothing, and the messages are all still on the server.
since map[uint64]time.Time
// log is where it says what it could not do; the standard logger when nobody gave one.
log *log.Logger
}
func (n *natsInbound) logger() *log.Logger {
if n.log != nil {
return n.log
}
return log.Default()
}
// Nats is the consume side of the bus being built.
@@ -70,7 +79,7 @@ func (n *natsInbound) Close() {}
// which message is held, and since when — is read and written without a lock because the AMQP loop
// never had two. A second goroutine would make that wrong in a way no test would catch.
func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
if err := broker.AssertMeshConsumers(n.js); err != nil {
if err := broker.AssertServingConsumers(n.js); err != nil {
return err
}
js, conn := n.js.Context(), n.js.Conn()
@@ -94,6 +103,13 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
holding.Store(true)
defer holding.Store(false)
// What a consumer gave up on, kept by the controller acting now (novox/hq issue 330): the server's
// notices wait in their stream for it, so one said while no controller listened is not lost.
// **Never the reason the controller stops serving**: a notice it cannot take yet waits in its stream,
// and the failure is said and tried again every minute.
stopKeeping := keepingGivenUp(ctx, n.js, n.logger())
defer stopKeeping()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
// they are their own guarantee: a lost one is the next one.
beats := make(chan *nats.Msg, Prefetch)
@@ -167,6 +183,11 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control),
msg *nats.Msg, streamed bool) {
// An event the controller gave up on and a person delivered again (novox/hq issue 330) arrives under
// the controller's own again subject, and is the event it was: acted on as first published.
if original, ok := broker.OriginalOfAgain(msg.Subject); ok {
msg.Subject = original
}
kind, known := kindOfSubject(msg.Subject)
if !known {
if streamed {
+93
View File
@@ -0,0 +1,93 @@
package link
import (
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// novox/hq issue 330, replayed with only what the link had before its fix, so it can be laid over the
// older commit. On 2026-10-06 a media server's event consumer gave up on several messages after five
// deliveries each. The controller raised a condition for each run and cleared it within minutes; the
// messages were kept only by EVENTS, which drops an event after a week, and nothing said which they
// were. Design 25 promised a dead-letter stream that does not exist. A consumer that never acknowledges
// an event: once it gives up, the event is kept with its consumer and how often it was handed over.
func TestReplay330(t *testing.T) {
js := aBus(t)
consumer, ok := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr",
Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"})
if !ok {
t.Fatal("a module that consumes got no consumer")
}
if err := js.EnsureConsumer(consumer); err != nil {
t.Fatal(err)
}
_, stop := servingOn(t, js, &counted{})
defer stop()
const subject = "mesh.mod.gitea.event.pull.merged"
if _, err := js.Context().Publish(subject, []byte(`{"repository":"novox/media"}`)); err != nil {
t.Fatal(err)
}
// The module's handler fails every time, as the media server's did.
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
reading, err := api.Consumer(t.Context(), broker.EventsStream, consumer.Name)
if err != nil {
t.Fatal(err)
}
for handed := 0; handed < consumer.MaxDeliver; handed++ {
batch, err := reading.Fetch(1, jetstream.FetchMaxWait(5*time.Second))
if err != nil {
t.Fatal(err)
}
n := 0
for msg := range batch.Messages() {
n++
_ = msg.Nak()
}
if n != 1 {
t.Fatalf("handed over %d times, then nothing: %v", handed, batch.Error())
}
}
// And it goes on reading, as a module's runtime does: the server gives the event up when it would
// hand it over a sixth time.
if batch, err := reading.Fetch(1, jetstream.FetchMaxWait(time.Second)); err == nil {
for msg := range batch.Messages() {
t.Fatalf("handed over a sixth time: %s", msg.Subject())
}
}
kept := "mesh.events.dead." + broker.EventsStream + "." + consumer.Name
var held *nats.RawStreamMsg
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) && held == nil {
if stream, err := js.Context().StreamNameBySubject(kept); err == nil {
held, _ = js.Context().GetLastMsg(stream, kept)
}
if held == nil {
time.Sleep(50 * time.Millisecond)
}
}
if held == nil {
t.Fatalf("%s gave up on the event and nothing on the bus keeps it under %s", consumer.Name, kept)
}
if string(held.Data) != `{"repository":"novox/media"}` {
t.Errorf("kept %q, not the event", held.Data)
}
for header, want := range map[string]string{"Mesh-Dead-Consumer": consumer.Name, "Mesh-Dead-Stream": "EVENTS",
"Mesh-Dead-Subject": subject, "Mesh-Dead-Deliveries": "5"} {
if got := held.Header.Get(header); got != want {
t.Errorf("the kept event's %s is %q, not %q", header, got, want)
}
}
if held.Header.Get("Mesh-Dead-Gave-Up") == "" {
t.Error("the kept event does not say when it was given up")
}
}
+1
View File
@@ -71,6 +71,7 @@
"bus",
"retire",
"cleanup",
"dead-letters",
"data",
"build",
"artifacts",