Compare commits

...
Author SHA1 Message Date
jschoubben a330c564ba Say a module assigned and left out of its machine's composition as a condition (issue 380)
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 ready: it delivers once merged
A push leaves out a module whose settings do not compose and sends the rest,
which is right, but only plan said so: nfs-server ran nowhere for days while
the operator believed it ran. The self-check now raises needs-operator for a
setting nobody gave, naming the setting and the command, and left-out for any
other cause; both warnings, cleared once the module composes or is unassigned.
status and node show list the same modules as assigned, not applied.
2026-10-11 04:43:58 +02:00
mesh-admin 0a2e58c070 Merge pull request 'Deliver again an ask the controller made, removing the original from its queue (issue 334, part)' (#219) from fix/334-a-given-up-ask-can-be-delivered-again into main 2026-10-11 02:36:26 +00:00
jschoubben ef551fdfb6 Remove an ask's original only when its sequence still holds it, and only with a worker (issue 334)
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 seat's queue made again numbers from one, so an old dead letter's sequence
can name a live ask; deleting by number alone would drop it silently. An ask
with no worker would wait unseen while counted as delivered.
2026-10-11 03:24:52 +02:00
jschoubben f510b46319 Deliver again an ask the controller made, removing the original from its queue (issue 334)
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 build ask its seat's worker gave up on could only be dropped, though the
controller already holds the publish on that seat's accepts and the stream
API to remove the original. Asks to other seats stay refused: delivering them
needs a grant ADR 0264 withholds, in the asker's name ADR 0259 protects.
2026-10-11 03:18:34 +02:00
17 changed files with 958 additions and 19 deletions
+3
View File
@@ -835,6 +835,9 @@ type answers struct {
// public name on the machine went dark. The holds were correct; they were recorded only in the
// machine's own state file, and the one visible symptom was a count that did not add up.
untaken map[string]map[string]int
// leftOut is, per machine, every module of its set its composition leaves out, and why (novox/hq issue
// 380): assigned and not applied, which every push said only in passing. Not well while there is any.
leftOut map[string][]leftOutModule
// unheld is every module on a machine whose resources are applied through a seat nothing on
// that machine holds (novox/hq ADR 0207), with the modules that could hold it. Reported, not
// refused, until the switch — and while there is any, the mesh is not all well: the order the
+5 -1
View File
@@ -144,11 +144,15 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64,
}
switch act {
case "deliver":
_, to, err := link.DeliverAgain(on.js, id)
delivered, to, err := link.DeliverAgain(on.js, id)
if err != nil {
return nil, err
}
answer["delivered_on"] = to
if delivered.Original != "" {
// What became of the ask in its seat's queue (novox/hq issue 334): removed, or left, and why.
answer["original"] = delivered.Original
}
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":
+5
View File
@@ -136,6 +136,11 @@ var probeRegistry = []probe{
{ID: "D-root", Asserts: "no agent can become root without a person on a machine where the router or a channel " +
"proving its sender runs: not by its own account, and not through a tool that runs its command as an account " +
"that can", From: "ADR 0259 §8", Kind: kindRootNotFree, Phase: 2, run: probeAgentRoot},
// A module assigned and left out of its machine's composition (novox/hq issue 380): said on that machine, the
// operator's, never urgent; cleared once it composes again or is unassigned.
{ID: probeLeftOutID, Asserts: "no module assigned to a machine is left out of its composition unsaid: each " +
"raises needs-operator for a setting nobody gave, left-out for any other cause", From: "issue 380",
Kind: kindLeftOut, Raises: []string{kindNeedsOperator}, Phase: 1, run: probeLeftOut},
{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
+226
View File
@@ -0,0 +1,226 @@
package main
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
)
// A module assigned to a machine and left out of its composition is said (novox/hq issue 380).
//
// **An omission is a finding, never a refusal of the push** (ADR 0163, rule 6): the machine is sent everything
// else, and the module's held things are kept. Until issue 380 the only place it showed was `plan`: nfs-server was
// assigned to the home server for days, every push left it out for a setting nobody gave, and the operator believed
// it ran. So the self-check composes every machine as the next push would (Resolution.LeftOutBecause, the push's
// own judgement) and raises a condition on that machine for each module it leaves out:
//
// - a setting nobody gave (catalogue.UnsetSettingError) is the operator's to give: kind needs-operator, naming
// the module, the setting and the command that sets it;
// - any other cause is said as a module not working is: kind left-out, a warning with the reason as evidence.
//
// Never urgent and never escalated by age (as ADR 0283 decision 5): nothing the machine ran was undone. It clears
// on the first run that no longer finds it — the module composed again, or no longer assigned. `status` and
// `node show` list the same modules under "assigned, not applied" with the same reason.
//
// A module left out because its stored manifest has a key this controller does not know is not raised here: the
// catalogue's own unknown-field condition says it once for the whole mesh (ADR 0262); it is still listed.
const (
// probeLeftOutID is the self-check's probe that raises and clears these conditions.
probeLeftOutID = "D-left-out"
// kindLeftOut is a module left out for any cause but a setting nobody gave; it is also every such condition's
// token, whichever its kind, so a cause that changes is the same condition said anew.
kindLeftOut = "left-out"
)
// leftOutModule is one module of a machine's set that its composition leaves out, and why.
type leftOutModule struct {
Module string
// Setting is the setting nobody gave, when that is the cause.
Setting string
// Why is the composition's own reason, whole.
Why string
// Unread is a stored manifest this controller cannot read whole (said by the catalogue's condition).
Unread bool
}
// leftOutOf is every module of a machine's resolution that a push would leave out, sorted. Pure: the judgement
// the push makes (catalogue.Resolution.LeftOutBecause), nothing allocated.
func leftOutOf(plan catalogue.Resolution, settings catalogue.SettingsBy, adopted bool) []leftOutModule {
because := plan.LeftOutBecause(settings, adopted)
out := make([]leftOutModule, 0, len(because))
for module, why := range because {
l := leftOutModule{Module: module, Why: oneLine(why.Error())}
var unset *catalogue.UnsetSettingError
var unread *catalogue.UnreadManifestError
switch {
case errors.As(why, &unset):
l.Setting = unset.Setting
case errors.As(why, &unread):
l.Unread = true
}
out = append(out, l)
}
sort.Slice(out, func(i, j int) bool { return out[i].Module < out[j].Module })
return out
}
// settingCommand is the command that gives a module's setting on one machine.
func settingCommand(module, setting, node string) string {
return fmt.Sprintf("`settings set %s '{%q: …}' --node %s`", module, setting, node)
}
// reason is why a module is left out, as `status`, `node show` and the condition's summary say it: for a setting
// nobody gave, the setting and the command that gives it; otherwise the composition's own words.
func (l leftOutModule) reason(node string) string {
if l.Setting != "" {
return fmt.Sprintf("nothing sets its setting %q, so every push leaves it out — %s sets it", l.Setting,
settingCommand(l.Module, l.Setting, node))
}
return "every push leaves it out: " + l.Why
}
// leftOutObservations are the conditions a machine's left-out modules raise, one each.
func leftOutObservations(node string, left []leftOutModule) []conditions.Observation {
var out []conditions.Observation
for _, l := range left {
if l.Unread {
continue
}
out = append(out, leftOutObservation(node, l))
}
return out
}
// leftOutObservation is one module left out of one machine's composition: needs-operator for a setting nobody
// gave, left-out otherwise; a warning either way, the operator's to resolve.
func leftOutObservation(node string, l leftOutModule) conditions.Observation {
o := conditions.Observation{Scope: conditions.ScopeModule, ID: l.Module + "." + node, Token: kindLeftOut,
Kind: kindLeftOut, Machine: node, Severity: conditions.Warning, Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("%s is assigned to %s and not applied: %s", l.Module, node, l.reason(node)),
Said: l.Why}
w := leftOutWords(l.Module, node, l.Setting)
if l.Setting != "" {
o.Kind = kindNeedsOperator
} else {
// The composition's words may name a path or an address, which the operator's channel withholds: the
// summary sends them to the evidence, and `plan` says them in full.
o.Summary = fmt.Sprintf("%s is assigned to %s and not applied: every push leaves it out, because what is "+
"set for it does not compose with its definition — the evidence and `plan %s` say why", l.Module, node, node)
}
o.Headline, o.Explanation, o.Needs, o.Resolved = w.Headline, w.Explanation, w.Needs, w.Resolved
return o
}
// leftOutWords is what the operator reads of a module left out (ADR 0253): plain, the act named.
func leftOutWords(module, node, setting string) words {
w := words{
Headline: fmt.Sprintf("%s is not applied on %s", module, node),
Explanation: fmt.Sprintf("%s is assigned to %s, and every update of %s leaves it out because what is set "+
"for it does not fit its definition. Nothing of it changes there; the rest of %s is updated as usual.",
module, node, node, node),
Needs: fmt.Sprintf("read why in the details, then change what is set for %s or unassign it.", module),
Resolved: fmt.Sprintf("%s on %s is no longer left out", module, node),
}
if setting != "" {
w.Explanation = fmt.Sprintf("%s is assigned to %s, and every update of %s leaves it out because nothing "+
"sets its setting %s. Nothing of it runs there until it is set; the rest of %s is updated as usual.",
module, node, node, setting, node)
w.Needs = fmt.Sprintf("set %s for %s on %s, or approve it when it is proposed to you.", setting, module, node)
// A setting's name that is not plain (a dotted key) is in the summary instead.
if _, ok := conditions.PlainWords(w, node); !ok {
w.Explanation = fmt.Sprintf("%s is assigned to %s, and every update of %s leaves it out because a "+
"setting it needs is not set. Nothing of it runs there until it is set; the details name it.",
module, node, node)
w.Needs = fmt.Sprintf("set what %s needs on %s; the details name the setting.", module, node)
}
}
return w
}
// probeLeftOut is the self-check's probe of issue 380: every machine's set is judged as its next push would
// judge it, and each module left out raises its condition. A machine that does not resolve is passed over — D1
// says it — and a store that cannot be read is an error, never "nothing left out" (ADR 0227 rule 4).
func probeLeftOut(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
nodes, err := d.open.inventory.Nodes(ctx)
if err != nil {
return nil, err
}
var out []conditions.Observation
for _, n := range nodes {
plan, settings, err := planFor(ctx, d.open, n.Name)
if err != nil {
if unresolvable(err) {
continue
}
return nil, fmt.Errorf("%s cannot be worked out: %w", n.Name, err)
}
out = append(out, leftOutObservations(n.Name, leftOutOf(plan, settings, n.Adopted))...)
if ctx.Err() != nil {
return nil, ctx.Err()
}
}
return out, nil
}
// notAppliedLines is a machine's "assigned, not applied" as `node show` prints it: one line a module, with the
// reason its condition says.
func notAppliedLines(node string, left []leftOutModule) []string {
if len(left) == 0 {
return nil
}
lines := []string{"", " assigned, not applied:"}
for _, l := range left {
lines = append(lines, fmt.Sprintf(" %-22s %s", l.Module, l.reason(node)))
}
return lines
}
// machineNotApplied is one module assigned to a machine and left out of its composition, in `status --json`.
type machineNotApplied struct {
Node string `json:"node"`
Module string `json:"module"`
Setting string `json:"setting,omitempty"`
Reason string `json:"reason"`
}
// notApplied is every machine's left-out modules, in a stated order, as `status --json` carries them.
func notApplied(left map[string][]leftOutModule) []machineNotApplied {
names := make([]string, 0, len(left))
for name := range left {
names = append(names, name)
}
sort.Strings(names)
var out []machineNotApplied
for _, name := range names {
for _, l := range left[name] {
out = append(out, machineNotApplied{Node: name, Module: l.Module, Setting: l.Setting, Reason: l.reason(name)})
}
}
return out
}
// printNotApplied is status's "assigned, not applied", every machine's.
func printNotApplied(left map[string][]leftOutModule) {
rows := notApplied(left)
if len(rows) == 0 {
return
}
fmt.Printf("%d module(s) assigned, not applied — every push leaves them out, and each raises a condition:\n",
len(rows))
for _, r := range rows {
fmt.Printf(" %-12s %-22s %s\n", r.Node, r.Module, r.Reason)
}
fmt.Println()
}
// isLeftOutCondition is whether a condition is this probe's, which the machine's health statements neither
// raise nor clear.
func isLeftOutCondition(c conditions.Condition) bool {
return c.Source == probeLeftOutID || strings.HasSuffix(c.Key, "."+kindLeftOut)
}
+203
View File
@@ -0,0 +1,203 @@
package main
import (
"context"
"encoding/json"
"os"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 380: nfs-server was assigned to the home server for days and every push left it out — "nfs-server
// has a file that says ${setting:shares}, and nothing sets shares for it" — and nothing but `plan` said so. The
// manifest is the catalogue's own at the commit that added it (mesh-catalog 48fba44), which still says
// ${setting:shares}; the reason below is the one the live push printed.
// leftOutNFS is the resolution of a machine assigned that nfs-server, with the settings given.
func leftOutNFS(t *testing.T) catalogue.Resolution {
t.Helper()
raw, err := os.ReadFile("testdata/left-out/nfs-server.json")
if err != nil {
t.Fatal(err)
}
m, err := catalogue.ParseManifest(raw)
if err != nil {
t.Fatal(err)
}
return catalogue.Resolution{Modules: []catalogue.Manifest{m}}
}
// sharesGiven is the operator's setting for it on the machine.
func sharesGiven(value string) catalogue.SettingsBy {
return catalogue.SettingsBy{"nfs-server": {{From: "anchor", Values: map[string]any{"shares": value}}}}
}
func TestAModuleLeftOutForASettingNobodyGaveNeedsTheOperatorNamingTheSettingAndTheCommand(t *testing.T) {
left := leftOutOf(leftOutNFS(t), nil, false)
if len(left) != 1 || left[0].Module != "nfs-server" || left[0].Setting != "shares" {
t.Fatalf("left out: %+v", left)
}
if !strings.Contains(left[0].Why, `nothing sets "shares" for it`) {
t.Fatalf("the reason is not the push's own: %s", left[0].Why)
}
obs := leftOutObservations("anchor", left)
if len(obs) != 1 {
t.Fatalf("raised %+v", obs)
}
o := obs[0]
if o.Key() != "module.nfs-server.anchor.left-out" || o.Kind != kindNeedsOperator || o.Machine != "anchor" ||
o.Severity != conditions.Warning || o.Resolver != conditions.ResolverOperator {
t.Fatalf("raised %s as %s, %s, by %s, on %q", o.Key(), o.Kind, o.Severity, o.Resolver, o.Machine)
}
for _, want := range []string{"nfs-server", "anchor", `"shares"`, "`settings set nfs-server '{\"shares\": …}' --node anchor`"} {
if !strings.Contains(o.Summary, want) {
t.Errorf("the summary does not name %s: %s", want, o.Summary)
}
}
if o.Said != left[0].Why {
t.Errorf("the evidence is not the composition's reason: %s", o.Said)
}
plainExample(t, o, "nfs-server is not applied on anchor",
"Needs you: set shares for nfs-server on anchor, or approve it when it is proposed to you. nfs-server is assigned to "+
"anchor, and every update of anchor leaves it out because nothing sets its setting shares. Nothing of it "+
"runs there until it is set; the rest of anchor is updated as usual.")
}
func TestAModuleLeftOutForAnotherCauseIsAWarningWithTheReasonAsEvidence(t *testing.T) {
// A setting stored that its definition can no longer compose: one with a line break (issue 339).
left := leftOutOf(leftOutNFS(t), sharesGiven("library=/srv/library\nmedia=/srv/media"), false)
if len(left) != 1 || left[0].Setting != "" {
t.Fatalf("left out: %+v", left)
}
obs := leftOutObservations("anchor", left)
if len(obs) != 1 {
t.Fatalf("raised %+v", obs)
}
o := obs[0]
if o.Key() != "module.nfs-server.anchor.left-out" || o.Kind != kindLeftOut || o.Severity != conditions.Warning {
t.Fatalf("raised %s as %s, %s", o.Key(), o.Kind, o.Severity)
}
if !strings.Contains(o.Said, "holds a line break") || strings.Contains(o.Summary, "/srv/") {
t.Fatalf("the reason is not the evidence, or the summary carries it to the channel:\n%s\n%s", o.Summary, o.Said)
}
plainExample(t, o, "nfs-server is not applied on anchor",
"Needs you: read why in the details, then change what is set for nfs-server or unassign it. nfs-server is assigned to "+
"anchor, and every update of anchor leaves it out because what is set for it does not fit its definition. "+
"Nothing of it changes there; the rest of anchor is updated as usual.")
}
// The self-check raises it, keeps it a warning however long it stands, and clears it when the module composes
// again or is no longer assigned; a machine's health statement neither clears nor raises it.
func TestALeftOutModulesConditionClearsWhenItComposesAgainOrIsUnassigned(t *testing.T) {
plan, settings := leftOutNFS(t), catalogue.SettingsBy(nil)
withProbes(t, probe{ID: probeLeftOutID, Asserts: "the test's", Kind: kindLeftOut, Phase: 1,
run: func(context.Context, *doctor) ([]conditions.Observation, error) {
return leftOutObservations("anchor", leftOutOf(plan, settings, false)), nil
}})
store := conditions.NewInMemory()
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store})
defer k.Close(context.Background())
d := &doctor{keeper: k, teller: &conditions.Told{}, host: "anchor"}
key := "module.nfs-server.anchor.left-out"
openOnes := func() map[string]conditions.Condition {
t.Helper()
open, err := k.Open(t.Context())
if err != nil {
t.Fatal(err)
}
out := map[string]conditions.Condition{}
for _, c := range open {
out[c.Key] = c
}
return out
}
d.runOnce(t.Context(), "a test")
c, raised := openOnes()[key]
if !raised || c.Kind != kindNeedsOperator || c.Severity != conditions.Warning {
t.Fatalf("a module left out raised %+v", openOnes())
}
// A statement from the machine that says nothing of it — the module runs nothing there — leaves it open.
if err := judgeModuleHealth(t.Context(), nil, k, "anchor", map[string][]inventory.ResourceHealth{}, nil,
time.Now().Add(72*time.Hour)); err != nil {
t.Fatal(err)
}
if _, still := openOnes()[key]; !still {
t.Fatal("a health statement that says nothing of the module cleared its left-out condition")
}
d.runOnce(t.Context(), "a test")
if c := openOnes()[key]; c.Severity != conditions.Warning {
t.Fatalf("after a health statement and days, it is %+v", openOnes())
}
// The setting given: it composes, and the condition clears.
settings = sharesGiven("library=/srv/library")
d.runOnce(t.Context(), "a test")
if _, still := openOnes()[key]; still {
t.Fatalf("composed again, still open: %+v", openOnes())
}
// Left out again, then unassigned: no longer in the machine's set, and it clears.
settings = nil
d.runOnce(t.Context(), "a test")
if _, raised := openOnes()[key]; !raised {
t.Fatal("left out again and not raised")
}
plan = catalogue.Resolution{}
d.runOnce(t.Context(), "a test")
if _, still := openOnes()[key]; still {
t.Fatalf("unassigned, still open: %+v", openOnes())
}
}
func TestTheLeftOutProbeIsInTheRegistry(t *testing.T) {
for _, p := range probeRegistry {
if p.ID == probeLeftOutID {
if p.run == nil || p.Kind != kindLeftOut {
t.Fatalf("%+v", p)
}
return
}
}
t.Fatalf("no probe %s: a module left out is said nowhere", probeLeftOutID)
}
// status and node show list it under "assigned, not applied" with the reason the condition says, and status is
// not well while there is one.
func TestStatusAndNodeListAModuleAssignedAndNotApplied(t *testing.T) {
left := map[string][]leftOutModule{"anchor": leftOutOf(leftOutNFS(t), nil, false)}
reason := leftOutObservations("anchor", left["anchor"])[0].Summary
asked := answers{leftOut: left}
if asked.well() {
t.Fatal("a mesh with a module assigned and not applied is called well")
}
shown := printed(t, func() error { return printStatus(asked) })
if !strings.Contains(shown, "1 module(s) assigned, not applied") || !strings.Contains(shown, "nfs-server") ||
!strings.Contains(shown, "`settings set nfs-server '{\"shares\": …}' --node anchor` sets it") {
t.Fatalf("status says:\n%s", shown)
}
if !strings.Contains(reason, left["anchor"][0].reason("anchor")) {
t.Fatalf("status and the condition say different reasons:\n%s\n%s", shown, reason)
}
body, err := statusAsJSON(asked)
if err != nil {
t.Fatal(err)
}
var doc struct {
NotApplied []machineNotApplied `json:"not-applied"`
}
if err := json.Unmarshal(body, &doc); err != nil || len(doc.NotApplied) != 1 ||
doc.NotApplied[0].Setting != "shares" || doc.NotApplied[0].Node != "anchor" {
t.Fatalf("status --json: %v %s", err, body)
}
lines := strings.Join(notAppliedLines("anchor", left["anchor"]), "\n")
if !strings.Contains(lines, "assigned, not applied:") || !strings.Contains(lines, "nfs-server") ||
!strings.Contains(lines, `nothing sets its setting "shares"`) {
t.Fatalf("node show says:\n%s", lines)
}
}
+3 -1
View File
@@ -147,9 +147,11 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
}
standing := map[string]conditions.Condition{}
for _, c := range open {
// A module left out of the composition is the self-check's to raise and clear, never a statement's
// (novox/hq issue 380): its needs-operator is not cleared for not being in what the machine runs.
if (c.Kind == kindModuleUnhealthy || c.Kind == kindReloginNeeded || c.Kind == kindUsedAsFound ||
c.Kind == kindNeedsOperator) &&
c.Subject.Machine == node {
c.Subject.Machine == node && !isLeftOutCondition(c) {
standing[c.Key] = c
}
}
+14 -2
View File
@@ -44,7 +44,7 @@ func nodeCommand(ctx context.Context, args []string) error {
if len(args) != 2 {
return errors.New("node show <name>")
}
return showNode(ctx, inv, args[1])
return showNode(ctx, open, args[1])
case "add":
return addNode(ctx, inv, args[1:])
@@ -787,7 +787,8 @@ func roughly(d time.Duration) string {
//
// It is also where "what should it be configured as" is read. The same line that gates an
// assignment carries `card1-DP-1`, and a person composing settings for that machine needs it.
func showNode(ctx context.Context, inv *inventory.Inventory, name string) error {
func showNode(ctx context.Context, stored *stores, name string) error {
inv := stored.inventory
node, err := inv.NodeByName(ctx, name)
if err != nil {
return err
@@ -889,6 +890,17 @@ func showNode(ctx context.Context, inv *inventory.Inventory, name string) error
if len(assigned) > 0 {
fmt.Printf("\n assigned: %s\n", strings.Join(assigned, ", "))
}
// And which of them a push leaves out, and why (novox/hq issue 380): judged as the push judges it, the same
// reason its condition says. Not computable is said, never read as "all applied".
plan, settings, err := planFor(ctx, stored, name)
switch {
case err != nil:
fmt.Printf("\n whether a push leaves any assigned module out is NOT known: %s\n", oneLine(err.Error()))
default:
for _, line := range notAppliedLines(name, leftOutOf(plan, settings, node.Adopted)) {
fmt.Println(line)
}
}
return nil
}
+8
View File
@@ -231,6 +231,14 @@ var plainWordings = map[string]func(conditions.Observation) words{
w := needsOperatorWords(orModule(module), node, nil)
return w
}),
kindLeftOut: worded(func(o conditions.Observation) words {
// The observation carries its own words (novox/hq issue 380); these are its kind's alone.
module := ""
if o.Scope == conditions.ScopeModule && o.Machine != "" {
module = strings.TrimSuffix(o.ID, "."+o.Machine)
}
return leftOutWords(orModule(module), machineOr(o, "a machine"), "")
}),
kindProviderFailing: worded(func(o conditions.Observation) words {
thing, consumer := conditions.ThingWords(o), idPart(o, 2)
if consumer == "" {
+4
View File
@@ -69,6 +69,9 @@ type meshStatus struct {
// **A document without this said an outage was a well mesh.** Read from what each machine
// reported, so it is the machine's account and not the mesh's take-time listing.
Untaken []machineUntaken `json:"untaken,omitempty"`
// NotApplied is every module assigned to a machine and left out of its composition, with why (novox/hq
// issue 380). Absent when every module composes.
NotApplied []machineNotApplied `json:"not-applied,omitempty"`
// Filtered is every converged machine that is not filtered by the mesh alone (novox/hq ADR
// 0168), one entry per rule set the mesh did not write — the found firewall in force again,
// or a chain nobody speaks for. Absent when every converged machine is filtered by the mesh
@@ -251,6 +254,7 @@ func statusAsJSON(asked answers) ([]byte, error) {
}
}
out.Unheld = asked.unheld
out.NotApplied = notApplied(asked.leftOut)
out.HandActsThisWeek, out.HandActsUnread = asked.handActs, asked.handActsUnread
out.HealsThisWeek, out.HealsUnread = asked.heals, asked.healsUnread
// In brief, as `conditions` lists them: status leads with every open condition, and their whole
+1 -1
View File
@@ -66,7 +66,7 @@ func TestAProviderFailingAConsumerBreaksAllWellUntilItRecovers(t *testing.T) {
}
// Both machines' `node show` name it: where the provider runs, and where the consumer is.
for _, node := range []string{"anchor", "laptop"} {
shown := printed(t, func() error { return showNode(ctx, open.inventory, node) })
shown := printed(t, func() error { return showNode(ctx, open, node) })
if !strings.Contains(shown, "open condition(s) about this machine") || !strings.Contains(shown, "mesh_laptop_dashboard") {
t.Fatalf("node show %s does not name it:\n%s", node, shown)
}
+12 -1
View File
@@ -289,6 +289,10 @@ func printStatus(asked answers) error {
fmt.Printf("\n `take <node> <module>` compares what runs against what it declares, and runs it\n\n")
}
// Assigned and not applied (novox/hq issue 380): before what is merely reported, because it reads like work
// finished and is none.
printNotApplied(asked.leftOut)
if len(asked.unheld) > 0 {
// **Reported, and not refused yet** (novox/hq ADR 0207 §4). Each machine still resolves and
// is sent what it would be; this says which of its modules depend on a seat nothing there
@@ -456,6 +460,13 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
return answers{}, err
}
plans[n.Name] = planned{plan, settings}
// And which of its modules a push leaves out (novox/hq issue 380), judged as the push judges it.
if left := leftOutOf(plan, settings, n.Adopted); len(left) > 0 {
if out.leftOut == nil {
out.leftOut = map[string][]leftOutModule{}
}
out.leftOut[n.Name] = left
}
out.unheld = append(out.unheld, plan.Unheld...)
// And which of its modules a provider leaves out of its grants, for an identity too long
// for what the provision keeps (novox/hq ADR 0225) — judged from the consumer's own
@@ -604,7 +615,7 @@ func untakenModules(ctx context.Context, inv *inventory.Inventory, nodes []inven
// read as success for the whole of the edge cut-over outage (novox/hq 04-ISSUES/125).
func (a answers) well() bool {
return len(a.wrong) == 0 && len(a.quiet) == 0 && len(a.behind) == 0 &&
len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 &&
len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 && len(a.leftOut) == 0 &&
len(a.filtered) == 0 && len(a.unheld) == 0 && len(a.overflowing) == 0 &&
len(a.conditions) == 0 && a.conditionsUnread == "" && a.pendingUnread == "" && !pendingOpen(a.pending)
}
+142
View File
@@ -0,0 +1,142 @@
{
"module": "nfs-server",
"version": "1",
"upgrade": {
"policy": "record",
"why": "the folders other machines mount: a build that breaks the exports leaves every client's mount hanging or refused, and the gate on this one machine does not see the clients (hq ADR 0236, ADR 0263)"
},
"capabilities": [
"package-manager",
"service-manager"
],
"provides": [
{
"name": "nfs-share",
"scope": "mesh",
"identity": false
}
],
"data": {
"consumers": {
"nfs-share": {
"class": "none",
"why": "the shared folders are the operator's data (hq ADR 0051): what a client writes lands in them, and they are protected where the operator declares them, never by this module, which keeps nothing of a consumer's"
}
}
},
"claims": [
{
"name": "node-nfs-server",
"scope": "node",
"serves": [
"exports",
"clients",
"test",
"reload",
"adopt"
]
}
],
"state": [
{
"name": "exports",
"ttl-seconds": 120,
"per-machine": true
}
],
"tools": [
"nfs_health"
],
"listens": [
{
"name": "nfs",
"port": 2049,
"protocol": "tcp",
"from": "mesh",
"fixed": true,
"why": "the shares, to the mesh's machines only (hq ADR 0263): NFS version 4 alone, which needs no other port, and never the home network, where a device that is not a node could claim any user id"
}
],
"resources": [
{
"id": "package",
"type": "package",
"package": "nfs-utils"
},
{
"id": "nfs-conf",
"type": "file",
"path": "/etc/nfs.conf.d/50-mesh.conf",
"mode": "0644",
"content": "# Written by the mesh (module nfs-server, novox/hq ADR 0263). Replaced on every push; a drop-in of\n# the operator's that sorts after this one overrides it, and is theirs.\n#\n# NFS version 4 only: a client needs port 2049 and nothing else, so the module opens nothing more\n# than that, to the private network. Version 3 needs rpcbind and mountd, on ports the mesh does not open.\n[nfsd]\nvers2=n\nvers3=n\nvers4=y\nvers4.0=n\nvers4.1=y\nvers4.2=y\n"
},
{
"id": "config-dir",
"type": "directory",
"path": "/etc/nfs-server",
"mode": "0755"
},
{
"id": "config",
"type": "file",
"path": "/etc/nfs-server/shares.conf",
"mode": "0644",
"content": "# Written by the mesh (module nfs-server, novox/hq ADR 0263) from this machine's assignment.\n# Replaced on every push; change the `shares` setting, never this file.\n#\n# The shares: name=folder, or name=folder:ro, one share per folder. The module's process exports each\n# to the private network's range below, every client mapped to the folder's owner.\nshares=${setting:shares}\nrange=${machine:mesh-range}\n"
},
{
"id": "run-dir",
"type": "directory",
"path": "/run/nfs-server",
"mode": "0755"
},
{
"id": "server",
"type": "service",
"unit": "nfs-server.service",
"state": "running",
"boot": "enabled",
"restart-on": [
"nfs-conf"
],
"health": {
"kind": "unit"
}
},
{
"id": "exports",
"type": "process",
"name": "nfs-server-exports",
"artifact": "tools",
"run": [
"./nfs-server",
"exports"
],
"restart-on": [
"config"
],
"health": {
"kind": "tool",
"tool": "nfs_health",
"interval": "60s",
"timeout": "10s",
"looks": 2,
"grace": "90s"
}
}
],
"build": {
"artifacts": [
{
"name": "tools",
"kind": "bundle",
"language": "go",
"system": "arch",
"from": "cmd/nfs-server",
"binary": "nfs-server",
"loads": [
"nfs-server"
]
}
]
}
}
+15
View File
@@ -163,6 +163,21 @@ type Principal struct {
// goes with the retired seat row.
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
// TheControllersAsk says whether a message of stream, published on subject, is an ask the controller
// itself makes: one on the accept subject of a seat in seatsTheControllerAsks, kept in that seat's own
// work queue. Such an ask, given up on by the seat's worker, can be delivered again with the authority
// the controller already holds — the publish on that seat's accepts and the stream API to remove the
// original — and no other can (novox/hq issue 334, ADR 0264's consequences: no grant over the seats'
// queues).
func TheControllersAsk(stream, subject string) bool {
for _, seat := range seatsTheControllerAsks {
if stream == seatStreamName(seat) && strings.HasPrefix(subject, "mesh.seat."+seat+".accept.") {
return true
}
}
return false
}
// SeatVerb is one verb of one seat, on every machine holding it.
type SeatVerb struct{ Seat, Verb string }
+22 -2
View File
@@ -342,18 +342,38 @@ func (e *NotMadeError) Error() string {
// compose. Empty when every module composes. The same judgement SetSettings makes before storing.
func (r Resolution) LeftOut(settings SettingsBy, adopted bool) map[string]string {
out := map[string]string{}
for module, why := range r.LeftOutBecause(settings, adopted) {
out[module] = why.Error()
}
return out
}
// LeftOutBecause is LeftOut with each reason as the error it was, so a reader can tell a setting nobody
// gave (an *UnsetSettingError) from any other cause and say it as the operator's to give (novox/hq issue
// 380). An UnreadManifestError for a stored manifest this controller cannot read whole.
func (r Resolution) LeftOutBecause(settings SettingsBy, adopted bool) map[string]error {
out := map[string]error{}
for _, m := range r.Modules {
if why := UnknownFieldReason(m); why != "" {
out[m.Module] = why
out[m.Module] = &UnreadManifestError{Module: m.Module, said: why}
continue
}
if err := JudgeSettings(m, settings[m.Module], adopted); err != nil {
out[m.Module] = err.Error()
out[m.Module] = err
}
}
return out
}
// UnreadManifestError is a module left out because its stored manifest has a key this controller does not know
// (novox/hq ADR 0262): said once for the whole mesh, by the catalogue's own condition, not per machine.
type UnreadManifestError struct {
Module string
said string
}
func (e *UnreadManifestError) Error() string { return e.said }
// Compose is Declaration with the owner of every resource said.
func (r Resolution) Compose(with Rendering) (Composed, error) {
owner := map[string]string{}
+15 -4
View File
@@ -38,6 +38,17 @@ func settingsUsed(content string) []string {
return keys
}
// UnsetSettingError is a module whose definition says ${setting:<key>} where nothing sets that key: typed, so
// that whoever reads why a module was left out of a machine can tell a setting nobody gave — the operator's
// to give, named with the command that gives it — from any other reason (novox/hq issue 380). Its words are
// the refusal's, unchanged.
type UnsetSettingError struct {
Module, Setting string
said string
}
func (e *UnsetSettingError) Error() string { return e.said }
// 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 over the module's own
@@ -58,12 +69,12 @@ func settingInto(resource map[string]any, layers []Layer, module string) error {
for _, key := range settingsUsed(content) {
value, set := settingValue(layers, key)
if !set {
return fmt.Errorf(
return &UnsetSettingError{Module: module, Setting: key, said: fmt.Sprintf(
"%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), 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))
module, key, key, module, key, orNoSettings(layers))}
}
content = strings.ReplaceAll(content, "${setting:"+key+"}", plainly(value))
}
@@ -114,11 +125,11 @@ func settingIntoUnit(resource map[string]any, layers []Layer, module string) err
for _, key := range settingsUsed(unit) {
value, set := settingValue(layers, key)
if !set {
return fmt.Errorf(
return &UnsetSettingError{Module: module, Setting: key, said: fmt.Sprintf(
"%s has a service whose unit 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): "+
"`settings set %s <file>` with {%q: …}%s",
module, key, key, module, key, orNoSettings(layers))
module, key, key, module, key, orNoSettings(layers))}
}
v := plainly(value)
if !unitPart.MatchString(v) {
+63 -7
View File
@@ -62,6 +62,9 @@ type DeadLetter struct {
// 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"`
// Original says what became of the ask it was kept from, in its seat's work queue, when it was
// delivered again (novox/hq issue 334).
Original string `json:"original,omitempty"`
// 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"`
@@ -250,9 +253,15 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
// 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.
// filters that subject, so a message is never let go as delivered while nobody receives it.
//
// An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's
// work queue has one worker per subject, so only the worker that gave it up takes it, and DeliverAgain
// removes the original from the queue first, so there are never two (novox/hq issue 334); only while the
// worker that gave it up is on the bus. Any other stream's
// message is refused, with why: an ask to a seat the controller does not ask would need a publish it is
// not granted, in its asker's name (ADR 0264's consequences, ADR 0259 §3), and a stream that is neither
// would reach every consumer of its subject.
func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
switch {
case d.Lost != "":
@@ -260,10 +269,22 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
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 broker.TheControllersAsk(d.Stream, d.Subject):
// The seat's worker takes it, and only while it is on the bus: without one the queue would keep
// the ask for a holder that may never come, and it would be let go as delivered meanwhile.
if _, err := js.ConsumerInfo(d.Stream, d.Consumer); errors.Is(err, nats.ErrConsumerNotFound) {
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so the ask "+
"would only wait in %s for a holder. Nothing was done, and it is still kept: deliver it again once "+
"the seat has a holder, or drop it", d.ID, d.Who, d.Stream)
} else if err != nil {
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
}
return d.Subject, nil
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)
return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+
"delivered again: publishing it again would reach every consumer of its subject, or need a grant the "+
"controller does not hold to speak for its asker. 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) {
@@ -284,7 +305,7 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
}
// 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
// from DEAD_LETTERS; an ask's original is removed from its work queue before. 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)
@@ -295,6 +316,9 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
if err != nil {
return d, "", err
}
if d.Stream != broker.EventsStream {
d.Original = removeOriginal(js, d)
}
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...)
@@ -302,6 +326,10 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
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 {
if d.Original != "" {
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept, so "+
"delivering it again tries once more. Its original: %s", id, to, err, d.Original)
}
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) {
@@ -311,6 +339,34 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
return d, to, nil
}
// removeOriginal takes the ask a dead letter was kept from out of its seat's work queue, and says what
// became of it. Given up on, the original is never acknowledged and would stay beside its copy until the
// stream's age drops it (novox/hq issue 334). It is removed before the copy is published, so a copy that
// cannot be published leaves the kept one to try again, and never two.
//
// **Only when the queue still holds that same message**: its subject and the time it was stored are the
// dead letter's. A seat's stream deleted and made again — the build handover deletes one (builds.go) —
// numbers from one again, and an old dead letter's sequence may then name another, live ask; deleting by
// the number alone would drop that one silently. Anything else is said, never an error: the original is
// gone or is not this one, and the copy is the only one there will be.
func removeOriginal(js nats.JetStreamContext, d DeadLetter) string {
held, err := js.GetMsg(d.Stream, d.Sequence)
switch {
case errors.Is(err, nats.ErrMsgNotFound):
return fmt.Sprintf("%s no longer held message %d, so there was nothing to remove", d.Stream, d.Sequence)
case err != nil:
return fmt.Sprintf("message %d of %s could not be read, so it was left as it is: %v", d.Sequence, d.Stream, err)
case d.Published.IsZero() || held.Subject != d.Subject || !held.Time.Equal(d.Published):
return fmt.Sprintf("message %d of %s is another message now (%s, stored %s), so it was left as it is; "+
"the one given up on is gone", d.Sequence, d.Stream, held.Subject, held.Time.UTC().Format(time.RFC3339))
}
if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
return fmt.Sprintf("message %d of %s could not be removed, so it stays beside its copy until the "+
"stream's age drops it; nothing delivers it again: %v", d.Sequence, d.Stream, err)
}
return fmt.Sprintf("message %d of %s, the one given up on, was removed from the queue", d.Sequence, d.Stream)
}
// DropDeadLetter removes a kept message for good.
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
d, err := DeadLetterNamed(js, id)
+217
View File
@@ -0,0 +1,217 @@
package link
import (
"errors"
"strings"
"testing"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// What a seat's worker gave up on, when the ask is one the controller itself makes (novox/hq issue 334).
//
// The controller already publishes on the accept subjects of the seats it asks (broker's
// seatsTheControllerAsks) and reaches the stream API, so delivering its own ask again needs no grant it
// does not hold: the original is removed from the work queue by its sequence, and the kept copy is
// published on the ask's own subject, where the seat's one worker takes it. An ask to any other seat is
// still refused, and stays kept: delivering it would need a publish the controller is not granted
// (ADR 0264's consequences), in the asker's name (ADR 0259 §3).
// theBuildWorker is the build seat's worker as a holder pulls from it.
func theBuildWorker(t *testing.T, js *broker.JetStream) jetstream.Consumer {
t.Helper()
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
stream := "SEAT_NODE_BUILD_AGENT"
worker, err := api.Consumer(t.Context(), stream, stream+"_worker")
if err != nil {
t.Fatal(err)
}
return worker
}
func TestAnAskTheControllerMadeIsDeliveredAgainToTheSeatsWorker(t *testing.T) {
js := aBusWithTheBuildRole(t)
keeping(t, js)
worker := theBuildWorker(t, js)
subject := BuildWorkOf(TheBuildMachine)
ask := &nats.Msg{Subject: subject, Data: []byte(`{"module":"x"}`), Header: nats.Header{}}
ask.Header.Set(nats.MsgIdHdr, "build-1")
if _, err := js.Context().PublishMsg(ask); err != nil {
t.Fatal(err)
}
d := givenUp(t, js, worker)
if d.Stream != "SEAT_NODE_BUILD_AGENT" || d.Subject != subject || d.Lost != "" {
t.Fatalf("kept as %+v", d)
}
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil {
t.Fatalf("the original is not in the work queue before it is delivered again: %v", err)
}
_, to, err := DeliverAgain(js.Context(), d.ID)
if err != nil {
t.Fatal(err)
}
if to != subject {
t.Fatalf("delivered again on %s, not the ask's own subject %s", to, subject)
}
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) {
t.Fatalf("the original is still in the work queue beside its copy: %v", err)
}
again := next(t, worker)
if again == nil {
t.Fatal("the seat's worker was not handed the ask again")
}
if string(again.Data()) != `{"module":"x"}` || again.Headers().Get(AgainHeader) == "" {
t.Fatalf("handed again as %s %v", again.Data(), again.Headers())
}
_ = again.Ack()
if m := next(t, worker); m != nil {
t.Fatalf("the worker 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)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
t.Fatalf("a second delivery answered %v", err)
}
}
// An ask to a seat the controller does not ask is refused, says why, and nothing is done.
func TestAnAskTheControllerDidNotMakeIsStillOnlyDropped(t *testing.T) {
js := aBus(t)
for _, d := range []DeadLetter{
{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker",
Subject: "mesh.seat.telegram-sender.accept.send"},
// The subject of a seat the controller asks, on a stream that is not that seat's queue.
{ID: 3, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker",
Subject: "mesh.seat.node-build-agent.accept.build"},
} {
if to, err := AgainTo(js.Context(), d); err == nil {
t.Errorf("%s on %s was given %s to be delivered again on", d.Subject, d.Stream, to)
} else if !strings.Contains(err.Error(), "Drop it") {
t.Errorf("refused without saying what to do: %v", err)
}
}
}
// aBuildAskGivenUp publishes one build ask and lets the worker give it up.
func aBuildAskGivenUp(t *testing.T, js *broker.JetStream, worker jetstream.Consumer, body string) DeadLetter {
t.Helper()
if _, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(body)); err != nil {
t.Fatal(err)
}
return givenUp(t, js, worker)
}
// A copy that cannot be published: the original is already out of the queue, the dead letter is kept and
// the answer says both; asked again once it can be published, it is delivered once.
func TestAnAskWhoseCopyIsRefusedStaysKeptAndIsDeliveredOnTheNextTry(t *testing.T) {
js := aBusWithTheBuildRole(t)
keeping(t, js)
worker := theBuildWorker(t, js)
d := aBuildAskGivenUp(t, js, worker, `{"module":"x"}`)
// The queue stops taking the ask's subject, so the copy's publish is refused by the server.
info, err := js.Context().StreamInfo(d.Stream)
if err != nil {
t.Fatal(err)
}
cfg := info.Config
taking := cfg.Subjects
cfg.Subjects = []string{"mesh.seat." + TheBuildMachine + ".accept.nothing"}
if _, err := js.Context().UpdateStream(&cfg); err != nil {
t.Fatal(err)
}
_, _, err = DeliverAgain(js.Context(), d.ID)
if err == nil || !strings.Contains(err.Error(), "still kept") || !strings.Contains(err.Error(), "was removed") {
t.Fatalf("a refused copy answered %v", err)
}
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) {
t.Fatalf("the original is still in the queue: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused copy let the dead letter go: %v", err)
}
cfg.Subjects = taking
if _, err := js.Context().UpdateStream(&cfg); err != nil {
t.Fatal(err)
}
retried, _, err := DeliverAgain(js.Context(), d.ID)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(retried.Original, "no longer held") {
t.Fatalf("the retry said of the original: %q", retried.Original)
}
if again := next(t, worker); again == nil || string(again.Data()) != `{"module":"x"}` {
t.Fatal("the retry did not hand the ask to the worker")
} else {
_ = again.Ack()
}
if m := next(t, worker); m != nil {
t.Fatalf("handed twice: %s", m.Data())
}
}
// A queue made again numbers from one: the old dead letter's sequence then names a live ask, which is left
// alone.
func TestAnAskWhoseSequenceNamesAnotherMessageLeavesThatOneAlone(t *testing.T) {
js := aBusWithTheBuildRole(t)
keeping(t, js)
d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"old"}`)
// The build handover's way: the seat's queue deleted and made again, with its worker.
if err := js.Context().DeleteStream(d.Stream); err != nil {
t.Fatal(err)
}
if err := broker.RaiseSeats(js, []broker.DeclaredSeat{{Name: TheBuildMachine, Accepts: []string{"build"},
Emits: []string{"built"}}}, map[string]broker.Holder{TheBuildMachine: {Node: "anchor", Module: "builder"}}); err != nil {
t.Fatal(err)
}
live, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(`{"module":"live"}`))
if err != nil {
t.Fatal(err)
}
if live.Sequence != d.Sequence {
t.Fatalf("the live ask is message %d, the dead letter names %d: the test does not set up the collision",
live.Sequence, d.Sequence)
}
delivered, _, err := DeliverAgain(js.Context(), d.ID)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(delivered.Original, "another message") {
t.Fatalf("the answer said of the original: %q", delivered.Original)
}
if held, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil || string(held.Data) != `{"module":"live"}` {
t.Fatalf("the live ask at that sequence was touched: %v", err)
}
}
// No worker on the seat's queue: refused, kept, and nothing published.
func TestAnAskIsNotDeliveredAgainWhileTheSeatHasNoWorker(t *testing.T) {
js := aBusWithTheBuildRole(t)
keeping(t, js)
d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"x"}`)
if err := js.Context().DeleteConsumer(d.Stream, d.Consumer); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "wait in") {
t.Fatalf("delivered with no worker: %v", err)
}
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil {
t.Fatalf("a refusal removed the original: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refusal let the dead letter go: %v", err)
}
}