Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a330c564ba | ||
|
|
0a2e58c070 | ||
|
|
ef551fdfb6 | ||
|
|
f510b46319 |
@@ -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
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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 == "" {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
@@ -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 }
|
||||
|
||||
|
||||
@@ -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{}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user