Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5dd0e3c056 | ||
|
|
313efa826c | ||
|
|
72fde0ee1a | ||
|
|
7bd04342b6 |
@@ -514,6 +514,10 @@ func composedAndValidated(ctx context.Context, open *stores, node string, gens m
|
||||
if declared.Epoch, err = open.inventory.SentEpoch(ctx, record.ID); err != nil {
|
||||
return sendable{}, nil, err
|
||||
}
|
||||
// And the generation it was last sent, as the would-send is (novox/hq issue 234).
|
||||
if declared.Generation, err = open.inventory.SentGeneration(ctx, record.ID); err != nil {
|
||||
return sendable{}, nil, err
|
||||
}
|
||||
body, err := declared.Body()
|
||||
if err != nil {
|
||||
return sendable{}, nil, err
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The assignment generation a declaration was composed from, and the record of every send (novox/hq
|
||||
// issue 234).
|
||||
//
|
||||
// On 2026-10-04 a declaration newer in sequence than every other named four fewer modules than the
|
||||
// assignments held, and a machine applied it and undeclared all four. The sequence orders arrival and
|
||||
// cannot tell a later send that carries an older view of the assignments. So a declaration now carries
|
||||
// the generation of the assignments it was composed from — a counter the store raises in the same
|
||||
// transaction as every assignment change — and a machine refuses one older than it applied, unless it is
|
||||
// a gate's put-back. And every send is written down with who sent it, from which generation and naming
|
||||
// which modules: the controller logged no send then, and which process sent the stale declaration could
|
||||
// not be read back from anything the mesh kept.
|
||||
|
||||
// kindOlderGeneration is S20's kind: a machine refused a send for the generation it was composed from.
|
||||
const kindOlderGeneration = "older-generation"
|
||||
|
||||
// generationRefusalsSaid is how long a refused send is said after its refusal: an hour, as an advisory is.
|
||||
const generationRefusalsSaid = advisoryQuiet
|
||||
|
||||
// stamp gives a composed declaration the order allotted to it before it was composed: its sequence, and
|
||||
// the epoch and generation when the machine reads them — and keeps the generation and acting epoch for the
|
||||
// record of the send whether or not the machine is sent them.
|
||||
//
|
||||
// **No put-back mark is sent.** A gate puts a machine back by composing it again (sendRollout), so its
|
||||
// declaration is allotted here like any other and carries the generation as it stands — never older than
|
||||
// what the machine applied. The node-engine reads a `put_back` key (mesh-host#82); this controller never
|
||||
// needs to send one.
|
||||
func (o order) stamp(d *sendable) {
|
||||
d.Sequence, d.Epoch = o.sequence, o.epoch
|
||||
d.composedFrom, d.actingEpoch = o.generation, o.acting
|
||||
d.Generation = 0
|
||||
if o.readsGeneration && o.generation > 0 {
|
||||
d.Generation = o.generation
|
||||
}
|
||||
}
|
||||
|
||||
// sentRecord is what the record of a send keeps beside its digest.
|
||||
type sentRecord struct {
|
||||
sequence int64
|
||||
epoch uint64
|
||||
generation int64
|
||||
// toldGeneration is the generation on the wire, which the would-send is stamped with: zero for a machine
|
||||
// whose node-engine has not said it reads one.
|
||||
toldGeneration int64
|
||||
sender string
|
||||
modules []string
|
||||
}
|
||||
|
||||
// sentRecordOf is a composed send's record: its sender is the caller this process acts for.
|
||||
func sentRecordOf(ctx context.Context, d sendable) sentRecord {
|
||||
return sentRecord{sequence: d.Sequence, epoch: d.actingEpoch, generation: d.composedFrom,
|
||||
toldGeneration: d.Generation, sender: senderOf(callerOf(ctx)), modules: d.modules}
|
||||
}
|
||||
|
||||
// callerOf is who asked for what this process does: the seat call's caller when ctx belongs to one, the
|
||||
// caller the controller ran this command for, or the account at the shell.
|
||||
func callerOf(ctx context.Context) string {
|
||||
if c := link.CallerIn(ctx); c != "" {
|
||||
return c
|
||||
}
|
||||
return link.Caller()
|
||||
}
|
||||
|
||||
// senderOf is a send's sender in words: who asked, and the process and build that composed it — the
|
||||
// answer issue 234 could not find, because two controllers and a one-shot push were all sending then.
|
||||
func senderOf(caller string) string {
|
||||
host, _ := os.Hostname()
|
||||
return fmt.Sprintf("%s (pid %d on %s, build %s)", caller, os.Getpid(), host, version)
|
||||
}
|
||||
|
||||
// namedModules are the modules a declaration names: its machine's set, less what was left out of it.
|
||||
func namedModules(plan catalogue.Resolution, leftOut map[string]string) []string {
|
||||
out := make([]string, 0, len(plan.Modules))
|
||||
for _, m := range plan.Modules {
|
||||
if _, left := leftOut[m.Module]; !left {
|
||||
out = append(out, m.Module)
|
||||
}
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
|
||||
// heardGenerationRefusal keeps a machine's refusal of a send for its generation, so S20 names its sender
|
||||
// and says it (novox/hq issue 234): on the send it refused or — when the mesh has no record of sending that
|
||||
// sequence — as a refusal of its own, naming the sender as unknown. Either way S20 is raised: a refusal the
|
||||
// mesh cannot attribute is the louder fact, not a quieter one.
|
||||
//
|
||||
// **And raises the mesh's counter past what the machine applied, when the refused send was composed from
|
||||
// the counter as it stands** (counterBehind): then the sender's view was not stale — the counter is behind
|
||||
// the machine, which is a store put back from a backup — and every send after would be refused for ever. A
|
||||
// stale sender's generation is below the counter, and the counter is left alone for it.
|
||||
//
|
||||
// Nothing is sent from here. This runs in the controller's receive loop, and a push from it would hold that
|
||||
// loop, and the machine's hold, for as long as the push takes — the deaf controller of issues 184 and 185.
|
||||
// S20 names `push <node>` for that case instead.
|
||||
func heardGenerationRefusal(ctx context.Context, inv *inventory.Inventory, report link.Report) {
|
||||
r := report.OlderGeneration
|
||||
if r == nil || inv == nil || report.Node == "" {
|
||||
return
|
||||
}
|
||||
node, err := inv.NodeByName(ctx, report.Node)
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: %s refused a send for its generation, and the machine cannot be "+
|
||||
"read, so it is not raised: %v\n", report.Node, err)
|
||||
return
|
||||
}
|
||||
var raised int64
|
||||
if now, err := inv.AssignmentGeneration(ctx); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: %s refused a send for its generation, and the mesh's own cannot be "+
|
||||
"read: %v\n", report.Node, err)
|
||||
} else if counterBehind(*r, now) {
|
||||
if raised, err = inv.RaiseAssignmentGeneration(ctx, r.Applied); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: the assignment generation (%d) is behind what %s applied (%d), "+
|
||||
"and could not be raised: %v\n", now, report.Node, r.Applied, err)
|
||||
raised = 0
|
||||
} else {
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: the assignment generation was %d, behind what %s applied (%d) — "+
|
||||
"a store put back from a backup — and is raised to %d; `push %s` sends it what the mesh holds now\n",
|
||||
now, report.Node, r.Applied, raised, report.Node)
|
||||
}
|
||||
}
|
||||
send, found, err := inv.RefusedSend(ctx, node.ID, report.Sequence, r.Applied, raised)
|
||||
if err == nil && !found {
|
||||
send, err = inv.RecordUnrecordedRefusal(ctx, inventory.Send{Node: node.ID, Sequence: report.Sequence,
|
||||
Epoch: report.Epoch, Generation: r.Generation, Digest: report.Declared, RefusedApplied: r.Applied,
|
||||
CounterRaisedTo: raised})
|
||||
}
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: %s refused send %d for its generation, and the refusal could not be "+
|
||||
"kept, so it is not raised: %v\n", report.Node, report.Sequence, err)
|
||||
return
|
||||
}
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: %s refused send %d from %s: composed from assignment generation %d, "+
|
||||
"and it applied %d\n", report.Node, report.Sequence, send.Sender, r.Generation, r.Applied)
|
||||
}
|
||||
|
||||
// counterBehind says a refusal shows the mesh's counter behind the machine rather than a stale sender: the
|
||||
// refused send carried the counter as it stands now, and the machine applied more than that.
|
||||
func counterBehind(r link.GenerationRefusal, now int64) bool {
|
||||
return r.Generation >= now && r.Applied > r.Generation
|
||||
}
|
||||
|
||||
// watchGenerationRefusals is S20: every send a machine refused within the hour for the generation it was
|
||||
// composed from, naming who sent it (novox/hq issue 234). One per machine, the newest refusal.
|
||||
func watchGenerationRefusals(f *signalFacts) []conditions.Observation {
|
||||
seen := map[string]bool{}
|
||||
var out []conditions.Observation
|
||||
for _, s := range f.refusedSends {
|
||||
if s.RefusedAt == nil || f.now.Sub(*s.RefusedAt) > generationRefusalsSaid || seen[s.NodeName] {
|
||||
continue
|
||||
}
|
||||
seen[s.NodeName] = true
|
||||
o := conditions.Observation{Scope: conditions.ScopeMachine, ID: s.NodeName, Kind: kindOlderGeneration,
|
||||
Machine: s.NodeName, Severity: conditions.Warning,
|
||||
Summary: fmt.Sprintf("%s refused sequence %d from %s: it was composed from assignment generation %d, and "+
|
||||
"%s applied generation %d — a sender composing from a view of the assignments the mesh has moved "+
|
||||
"past; it named %s", s.NodeName, s.Sequence, s.Sender, s.Generation, s.NodeName, s.RefusedApplied,
|
||||
modulesWords(s.Modules)),
|
||||
Said: fmt.Sprintf("refused at %s", s.RefusedAt.UTC().Format(time.RFC3339)),
|
||||
Headline: fmt.Sprintf("%s refused an out-of-date update", s.NodeName),
|
||||
Explanation: fmt.Sprintf("Something sent %s an update made from an older list of what runs there. %s "+
|
||||
"refused it, so nothing was removed. Nothing for you to do unless it repeats.", s.NodeName, s.NodeName),
|
||||
Resolved: fmt.Sprintf("%s has had no out-of-date update for an hour", s.NodeName)}
|
||||
if s.CounterRaisedTo > 0 {
|
||||
// Not a stale sender: the mesh's counter was behind the machine (a store put back from a backup) and
|
||||
// was raised; nothing has sent the machine its declaration since, so a person is asked to.
|
||||
o.Summary = fmt.Sprintf("%s refused sequence %d from %s: it was composed from assignment generation %d, "+
|
||||
"the mesh's own, and %s applied generation %d — the mesh's counter was behind the machine (a store "+
|
||||
"put back from a backup?) and is raised to %d; `push %s` sends it what the mesh holds now",
|
||||
s.NodeName, s.Sequence, s.Sender, s.Generation, s.NodeName, s.RefusedApplied, s.CounterRaisedTo,
|
||||
s.NodeName)
|
||||
o.Explanation = fmt.Sprintf("Needs you: push %s. The mesh's record was older than %s, so %s refused its "+
|
||||
"update and kept what it had. The record is repaired; a push sends the update again.",
|
||||
s.NodeName, s.NodeName, s.NodeName)
|
||||
}
|
||||
out = append(out, o)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// modulesWords is a send's modules as a sentence says them.
|
||||
func modulesWords(modules []string) string {
|
||||
if len(modules) == 0 {
|
||||
return "no module the mesh recorded"
|
||||
}
|
||||
return strings.Join(modules, ", ")
|
||||
}
|
||||
|
||||
// writeLastSend says what a machine was last told, by whom and from which generation (novox/hq issue 234):
|
||||
// nothing when the mesh has not recorded a send to it.
|
||||
func writeLastSend(ctx context.Context, w io.Writer, inv *inventory.Inventory, node string) error {
|
||||
s, found, err := inv.LastSend(ctx, node)
|
||||
if err != nil || !found {
|
||||
return err
|
||||
}
|
||||
generation := "no generation recorded"
|
||||
if s.Generation > 0 {
|
||||
generation = fmt.Sprintf("assignment generation %d", s.Generation)
|
||||
}
|
||||
fmt.Fprintf(w, "%s was last sent sequence %d at %s by %s, composed from %s, naming %s\n", node, s.Sequence, s.SentAt.Local().Format("2006-01-02 15:04:05"), s.Sender, generation, modulesWords(s.Modules))
|
||||
if s.RefusedAt != nil {
|
||||
fmt.Fprintf(w, " and refused it at %s: it had applied generation %d\n",
|
||||
s.RefusedAt.Local().Format("2006-01-02 15:04:05"), s.RefusedApplied)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The assignment generation a declaration was composed from (novox/hq issue 234).
|
||||
|
||||
func bodyKeys(t *testing.T, s sendable) map[string]any {
|
||||
t.Helper()
|
||||
raw, err := s.Body()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var keys map[string]any
|
||||
if err := json.Unmarshal(raw, &keys); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return keys
|
||||
}
|
||||
|
||||
// **The generation goes on the wire only to a machine whose node-engine said it reads one**: an older
|
||||
// node-engine decodes strictly and refuses an unknown key, whole. No put-back mark is ever sent: a gate
|
||||
// composes its put-back afresh, with the generation as it stands.
|
||||
func TestTheGenerationIsSentOnlyToAMachineThatReadsOne(t *testing.T) {
|
||||
resources := []map[string]any{{"id": "a", "type": "file", "path": "/etc/a", "content": "x\n"}}
|
||||
reads := order{sequence: 12, epoch: 57, generation: 40, readsGeneration: true, acting: 57}
|
||||
var told sendable
|
||||
told.Resources = resources
|
||||
reads.stamp(&told)
|
||||
keys := bodyKeys(t, told)
|
||||
if keys["generation"] != float64(40) || keys["sequence"] != float64(12) {
|
||||
t.Fatalf("a machine that reads a generation was sent %v", keys)
|
||||
}
|
||||
if _, there := keys["put_back"]; there {
|
||||
t.Fatalf("a put-back mark was sent: %v", keys)
|
||||
}
|
||||
if told.composedFrom != 40 || told.actingEpoch != 57 {
|
||||
t.Errorf("the send does not keep what it was composed from for its record: %+v", told)
|
||||
}
|
||||
|
||||
older := order{sequence: 12, epoch: 57, generation: 40, readsGeneration: false, acting: 57}
|
||||
var untold sendable
|
||||
untold.Resources = resources
|
||||
older.stamp(&untold)
|
||||
if _, there := bodyKeys(t, untold)["generation"]; there {
|
||||
t.Fatal("a machine whose node-engine never said it reads a generation was sent one")
|
||||
}
|
||||
if untold.composedFrom != 40 {
|
||||
t.Errorf("the generation it was composed from is recorded whether or not it was sent: %+v", untold)
|
||||
}
|
||||
}
|
||||
|
||||
// **The sender is named**: the caller of the seat call when there is one, else the shell's account, and the
|
||||
// process and build that composed it.
|
||||
func TestASendNamesItsSender(t *testing.T) {
|
||||
said := senderOf("g14.node-tools, through the mesh-controller seat")
|
||||
if !strings.Contains(said, "g14.node-tools") || !strings.Contains(said, "pid ") || !strings.Contains(said, "build ") {
|
||||
t.Fatalf("the sender reads %q", said)
|
||||
}
|
||||
}
|
||||
|
||||
// refusedSend is a send a machine refused for its generation at a moment.
|
||||
func refusedSend(at time.Time) inventory.Send {
|
||||
return inventory.Send{NodeName: "anchor", Sequence: 12, Epoch: 57, Generation: 38, RefusedApplied: 40,
|
||||
Sender: "a one-shot push by jochen at a shell on anchor (pid 4242 on anchor, build 2026.10.11)",
|
||||
Modules: []string{"docker"}, SentAt: at.Add(-time.Second), RefusedAt: &at, Recorded: true}
|
||||
}
|
||||
|
||||
// **A refused send is raised naming its sender** (novox/hq issue 234): who sent it, from which generation,
|
||||
// against which the machine applied, and its sequence — the facts that took a morning to look for.
|
||||
func TestARefusedSendIsRaisedNamingItsSender(t *testing.T) {
|
||||
now := time.Date(2026, 10, 11, 12, 0, 0, 0, time.UTC)
|
||||
f := calm(now)
|
||||
f.refusedSends = []inventory.Send{refusedSend(now.Add(-time.Minute))}
|
||||
got := watchGenerationRefusals(f)
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("%+v", got)
|
||||
}
|
||||
o := got[0]
|
||||
if o.Key() != "machine.anchor.older-generation" || o.Machine != "anchor" {
|
||||
t.Errorf("raised as %s about %q", o.Key(), o.Machine)
|
||||
}
|
||||
for _, want := range []string{"a one-shot push by jochen", "generation 38", "generation 40", "sequence 12"} {
|
||||
if !strings.Contains(o.Summary, want) {
|
||||
t.Errorf("the summary does not say %q: %s", want, o.Summary)
|
||||
}
|
||||
}
|
||||
if o.Headline == "" || o.Explanation == "" || o.Resolved == "" {
|
||||
t.Errorf("the condition is not worded for the operator: %+v", o)
|
||||
}
|
||||
}
|
||||
|
||||
// **A refusal of a send the mesh has no record of is raised too, its sender said to be unknown** — never
|
||||
// only a line on stderr.
|
||||
func TestARefusalOfAnUnrecordedSendIsRaisedSayingTheSenderIsUnknown(t *testing.T) {
|
||||
now := time.Date(2026, 10, 11, 12, 0, 0, 0, time.UTC)
|
||||
f := calm(now)
|
||||
s := refusedSend(now.Add(-time.Minute))
|
||||
s.Recorded, s.Sender, s.Modules = false, "a sender the mesh has no record of (no send of this sequence was recorded)", nil
|
||||
f.refusedSends = []inventory.Send{s}
|
||||
got := watchGenerationRefusals(f)
|
||||
if len(got) != 1 || !strings.Contains(got[0].Summary, "no record of") {
|
||||
t.Fatalf("an unattributed refusal was not raised as one: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// **When the refusal showed the mesh's counter behind the machine, the condition asks for a push** — it was
|
||||
// raised, and nothing has sent the machine its declaration since — and does not say there is nothing to do.
|
||||
func TestACounterBehindTheMachineAsksForAPush(t *testing.T) {
|
||||
now := time.Date(2026, 10, 11, 12, 0, 0, 0, time.UTC)
|
||||
f := calm(now)
|
||||
s := refusedSend(now.Add(-time.Minute))
|
||||
s.Generation, s.RefusedApplied, s.CounterRaisedTo = 12, 40, 41
|
||||
f.refusedSends = []inventory.Send{s}
|
||||
got := watchGenerationRefusals(f)
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("%+v", got)
|
||||
}
|
||||
if !strings.Contains(got[0].Summary, "`push anchor`") || !strings.Contains(got[0].Explanation, "push anchor") ||
|
||||
strings.Contains(got[0].Explanation, "Nothing for you to do") {
|
||||
t.Errorf("a counter behind the machine does not ask for a push: %s / %s", got[0].Summary, got[0].Explanation)
|
||||
}
|
||||
}
|
||||
|
||||
// **Only a refusal of the counter as it stands is the counter behind**: a stale sender's generation is below
|
||||
// it, and the counter is left alone for that.
|
||||
func TestOnlyARefusalOfTheCounterAsItStandsRaisesIt(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
refused, applied, now int64
|
||||
behind bool
|
||||
}{
|
||||
{refused: 12, applied: 40, now: 12, behind: true}, // the store was put back: the mesh says 12, the machine had 40
|
||||
{refused: 38, applied: 40, now: 41, behind: false}, // a stale sender: the counter is already past
|
||||
{refused: 38, applied: 40, now: 40, behind: false}, // a stale sender: the counter is where the machine is
|
||||
} {
|
||||
r := link.GenerationRefusal{Generation: c.refused, Applied: c.applied}
|
||||
if got := counterBehind(r, c.now); got != c.behind {
|
||||
t.Errorf("refused %d, applied %d, counter %d: behind %v, want %v", c.refused, c.applied, c.now, got, c.behind)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -128,7 +128,7 @@ func (r *recordedDelivery) grant(context.Context, []readyNode) error { return ni
|
||||
|
||||
func (r *recordedDelivery) declare(ctx context.Context, s readyNode, body []byte) (string, error) {
|
||||
r.declared = append(r.declared, s.node)
|
||||
return recordSent(ctx, r.inv, s.node, body, s.declared.Builds, s.declared.Epoch)
|
||||
return recordSent(ctx, r.inv, s.node, body, s.declared.Builds, s.declared.Epoch, sentRecordOf(ctx, s.declared))
|
||||
}
|
||||
|
||||
// aResolver is a module built from a repository, at a commit, with something on the machine that
|
||||
|
||||
@@ -77,7 +77,7 @@ func TestASendIsRecordedEvenWhenTheSenderIsBeingCancelled(t *testing.T) {
|
||||
}
|
||||
cancel() // the sender is going away: its context is cancelled between the send and the record
|
||||
body := []byte(`{"declaration":1,"resources":[]}`)
|
||||
digest, err := recordSent(ctx, inv, "anchor", body, nil, 0)
|
||||
digest, err := recordSent(ctx, inv, "anchor", body, nil, 0, sentRecord{sender: "a test"})
|
||||
if err != nil {
|
||||
// NodeByName on the cancelled context may itself refuse; the record must still be possible
|
||||
// through the detached context, so look the node up again on a live one.
|
||||
|
||||
@@ -420,7 +420,8 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
||||
out := sendable{Resources: composed.Resources, Adoption: adoption,
|
||||
Received: composed.Received, Mesh: with.Mesh, BusUsers: with.BusUsers,
|
||||
LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut, withheld: with.Withheld,
|
||||
unbound: with.Unbound, foreseen: composed.Foreseen, unplaced: composed.Unplaced}
|
||||
unbound: with.Unbound, foreseen: composed.Foreseen, unplaced: composed.Unplaced,
|
||||
modules: namedModules(plan, composed.LeftOut)}
|
||||
// And which build of each module it carries, for the send to record (novox/hq issue 259, ADR
|
||||
// 0221). Read only on the send path: a question about what would be sent records nothing.
|
||||
if choosing == Allocating {
|
||||
@@ -1290,6 +1291,11 @@ func planCommand(ctx context.Context, args []string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// What the machine was last told, by whom and from which assignment generation (novox/hq issue 234),
|
||||
// beside what it would be told now. A record that cannot be read is said, not taken for none.
|
||||
if err := writeLastSend(ctx, os.Stdout, open.inventory, args[0]); err != nil {
|
||||
fmt.Printf("what %s was last sent could not be read: %v\n", args[0], err)
|
||||
}
|
||||
// Which modules a push would leave out, and why — said before the plan, since the plan is of
|
||||
// what the machine would be told (novox/hq ADR 0163, rule 6). Judged, never composed: `plan`
|
||||
// without --json allocates nothing.
|
||||
|
||||
@@ -356,10 +356,16 @@ func declare(ctx context.Context, args []string) error {
|
||||
// the mesh did not compose it, so a push that does not name this machine treats it as held.
|
||||
// The epoch it carried, if a person wrote one in, is what the machine heard.
|
||||
var carried struct {
|
||||
Epoch uint64 `json:"epoch"`
|
||||
Epoch uint64 `json:"epoch"`
|
||||
Sequence int64 `json:"sequence"`
|
||||
Generation int64 `json:"generation"`
|
||||
}
|
||||
_ = json.Unmarshal(raw, &carried)
|
||||
if _, err := recordSent(ctx, inv, node, raw, nil, carried.Epoch); err != nil {
|
||||
// And recorded as every send is (novox/hq issue 234): by hand, from what it carried; the modules it
|
||||
// named are not known, since the mesh did not compose it.
|
||||
if _, err := recordSent(ctx, inv, node, raw, nil, carried.Epoch, sentRecord{sequence: carried.Sequence,
|
||||
epoch: carried.Epoch, generation: carried.Generation, toldGeneration: carried.Generation,
|
||||
sender: senderOf(callerOf(ctx)) + ", a declaration sent by hand"}); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||
@@ -798,7 +804,7 @@ func composeEach(names []string, allot func(node string) (order, error),
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
}
|
||||
declared.Sequence, declared.Epoch = numbered.sequence, numbered.epoch
|
||||
numbered.stamp(&declared)
|
||||
if len(declared.Resources) == 0 {
|
||||
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to
|
||||
// nothing may have HELD something before — the broker opening a placement gave it,
|
||||
@@ -1016,7 +1022,8 @@ func (b overTheBus) declare(ctx context.Context, s readyNode, body []byte) (stri
|
||||
}
|
||||
// After it is away, not before. A digest recorded for something that failed to send would make
|
||||
// the machine look current for a declaration it never received.
|
||||
digest, err := recordSent(ctx, b.open.inventory, s.node, body, s.declared.Builds, s.declared.Epoch)
|
||||
digest, err := recordSent(ctx, b.open.inventory, s.node, body, s.declared.Builds, s.declared.Epoch,
|
||||
sentRecordOf(ctx, s.declared))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -1228,7 +1235,7 @@ func sendToEach(ctx context.Context, open *stores, names []string) ([]string, er
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
}
|
||||
declared.Sequence, declared.Epoch = numbered.sequence, numbered.epoch
|
||||
numbered.stamp(&declared)
|
||||
reportLeftOut(name, declared)
|
||||
sending = append(sending, readyNode{name, declared})
|
||||
}
|
||||
@@ -1393,6 +1400,11 @@ func wouldSendFrom(ctx context.Context, open *stores,
|
||||
if declared.Epoch, err = open.inventory.SentEpoch(ctx, n.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// And the generation it was last sent, for the same reason (novox/hq issue 234): an assignment
|
||||
// elsewhere in the mesh is not a change of this machine.
|
||||
if declared.Generation, err = open.inventory.SentGeneration(ctx, n.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
body, err := declared.Body()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -1531,6 +1543,12 @@ var epochForActs = func(ctx context.Context) (uint64, error) { return theLease.e
|
||||
type order struct {
|
||||
sequence int64
|
||||
epoch uint64
|
||||
// generation is the assignment generation read before the composition, and readsGeneration whether the
|
||||
// machine's node-engine said it reads one (novox/hq issue 234); acting is the lease epoch the sender acts
|
||||
// under, whether or not the machine is sent it — both kept for the record of the send.
|
||||
generation int64
|
||||
readsGeneration bool
|
||||
acting uint64
|
||||
}
|
||||
|
||||
// allot takes the next sequence for a machine — the number its next declaration carries — under the
|
||||
@@ -1544,6 +1562,7 @@ func allot(ctx context.Context, inv *inventory.Inventory, node string) (order, e
|
||||
if err != nil {
|
||||
return order{}, err
|
||||
}
|
||||
acting := epoch
|
||||
if epoch > 0 {
|
||||
reads, err := inv.ReadsEpoch(ctx, record.ID)
|
||||
if err != nil {
|
||||
@@ -1553,11 +1572,24 @@ func allot(ctx context.Context, inv *inventory.Inventory, node string) (order, e
|
||||
epoch = 0
|
||||
}
|
||||
}
|
||||
// **The generation is read here, before the composition reads a single assignment** (novox/hq issue
|
||||
// 234). A generation only grows, so one read before is never newer than the view composed after it:
|
||||
// the declaration may claim a generation older than its content, never newer — and a claim newer than
|
||||
// the content is exactly the stale send the machine must be able to refuse.
|
||||
generation, err := inv.AssignmentGeneration(ctx)
|
||||
if err != nil {
|
||||
return order{}, err
|
||||
}
|
||||
readsGeneration, err := inv.ReadsGeneration(ctx, record.ID)
|
||||
if err != nil {
|
||||
return order{}, err
|
||||
}
|
||||
seq, err := inv.NextSequence(ctx, record.ID)
|
||||
if err != nil {
|
||||
return order{}, err
|
||||
}
|
||||
return order{sequence: seq, epoch: epoch}, nil
|
||||
return order{sequence: seq, epoch: epoch, generation: generation, readsGeneration: readsGeneration,
|
||||
acting: acting}, nil
|
||||
}
|
||||
|
||||
// recordSent writes down what a machine was just sent, and returns the digest.
|
||||
@@ -1572,7 +1604,7 @@ func allot(ctx context.Context, inv *inventory.Inventory, node string) (order, e
|
||||
// And the build of each module it carried (novox/hq issue 259, ADR 0221), nil when that is not known:
|
||||
// what tells a machine held back by a policy or a plan from one a push left behind.
|
||||
func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte,
|
||||
builds map[string]string, epoch uint64) (string, error) {
|
||||
builds map[string]string, epoch uint64, sent sentRecord) (string, error) {
|
||||
kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
||||
defer cancel()
|
||||
record, err := inv.NodeByName(kept, node)
|
||||
@@ -1580,7 +1612,21 @@ func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body
|
||||
return "", err
|
||||
}
|
||||
digest := digestOf(body)
|
||||
if err := inv.RecordSentUnder(kept, record.ID, digest, builds, epoch); err != nil {
|
||||
// The digest and the generation it carried are written in one transaction (novox/hq issue 234), so the
|
||||
// would-send is never stamped with a generation the machine was not sent; and the send itself, who sent
|
||||
// it and from which generation, beside them. A record of the send that cannot be written does not make a
|
||||
// send that is away look failed: it is said, loudly, and the machine's own record stands.
|
||||
err = inv.RecordSentWith(kept, record.ID, digest, builds, epoch, sent.toldGeneration, inventory.Send{
|
||||
Sequence: sent.sequence, Epoch: int64(sent.epoch), Sender: sent.sender, Generation: sent.generation,
|
||||
Digest: digest, Modules: sent.modules})
|
||||
var unrecorded *inventory.SendNotRecordedError
|
||||
switch {
|
||||
case errors.As(err, &unrecorded):
|
||||
fmt.Fprintf(os.Stderr, "mesh-controller: SEND NOT RECORDED: %s was sent declaration %s (sequence %d), and who "+
|
||||
"sent it and from which generation could not be written down, so a refusal of it will name no sender "+
|
||||
"(novox/hq issue 234): %v\n", node, short(digest), sent.sequence, unrecorded.Err)
|
||||
fmt.Printf("%s: sent, and the record of who sent it could NOT be written: %v\n", node, unrecorded.Err)
|
||||
case err != nil:
|
||||
return "", err
|
||||
}
|
||||
// And what it was, summarised, so the next push can be compared with it (novox/hq ADR 0217).
|
||||
|
||||
@@ -27,6 +27,19 @@ type sendable struct {
|
||||
// said it reads one is sent none, because an older node-engine refuses a key it does not know, whole
|
||||
// (link/order.go, the contract).
|
||||
Epoch uint64
|
||||
// Generation is the assignment generation it was composed from (novox/hq issue 234): a machine that
|
||||
// applied a later one refuses it, so a send composed from a view of the assignments the mesh has moved
|
||||
// past cannot undeclare what is still assigned. Zero is not sent at all — every machine whose
|
||||
// node-engine has not said it reads one is sent none, for the reason the epoch is not (an older
|
||||
// node-engine refuses a key it does not know, whole).
|
||||
Generation int64
|
||||
// composedFrom is the generation read before this declaration was composed, and actingEpoch the
|
||||
// lease epoch its sender acted under — whether or not the machine is sent either — for the record of
|
||||
// the send; never on the wire.
|
||||
composedFrom int64
|
||||
actingEpoch uint64
|
||||
// modules are the modules this declaration names, for the record of the send; never on the wire.
|
||||
modules []string
|
||||
// Adoption is nil for a converged node, and then the body is byte for byte what it was before
|
||||
// adoption existed: an older host parses the envelope strictly and would refuse the key.
|
||||
Adoption *adoptionEnvelope
|
||||
@@ -94,6 +107,9 @@ func (s sendable) Body() ([]byte, error) {
|
||||
if len(s.LeftOut) > 0 {
|
||||
envelope["left_out"] = s.LeftOut
|
||||
}
|
||||
if s.Generation > 0 {
|
||||
envelope["generation"] = s.Generation
|
||||
}
|
||||
// An empty declaration is deliberate here — the node owns nothing the mesh put there
|
||||
// (novox/hq issue 127) — and the host refuses an empty body unless it is told the emptiness
|
||||
// is meant, so a truncated or mis-composed body is never mistaken for "own nothing".
|
||||
|
||||
@@ -226,6 +226,20 @@ var signalsTable = []signalRow{
|
||||
newest: func(f *signalFacts) time.Time {
|
||||
return newestOf(f.batches, func(b batchFacts) time.Time { return b.closed })
|
||||
}},
|
||||
{Row: "S20", Signal: "a send refused for the assignment generation it was composed from",
|
||||
Emitter: "node-engine", Trigger: "each refusal (novox/hq issue 234)",
|
||||
Bound: "none: raised at the first refusal, naming its sender, the generation it came from and the one the " +
|
||||
"machine applied; cleared an hour after the last",
|
||||
Kind: kindOlderGeneration, Severity: conditions.Warning, Phase: 2,
|
||||
needs: func(f *signalFacts) error { return f.refusedSendsErr }, watch: watchGenerationRefusals,
|
||||
newest: func(f *signalFacts) time.Time {
|
||||
return newestOf(f.refusedSends, func(s inventory.Send) time.Time {
|
||||
if s.RefusedAt == nil {
|
||||
return time.Time{}
|
||||
}
|
||||
return *s.RefusedAt
|
||||
})
|
||||
}},
|
||||
{Row: "S17", Signal: "a send held for the bus's planned step is told to a person", Emitter: "controller's plan",
|
||||
Trigger: "each send refused because it would replace the bus outside its step (novox/hq issue 336)",
|
||||
Bound: "none: raised at the first refusal, for the operator, naming what waits, the bus build from and to, " +
|
||||
|
||||
@@ -178,6 +178,14 @@ var suppressions = map[string]suppression{
|
||||
commit: "c0ffee001122", modules: []string{"app"}, since: f.now.Add(-time.Second)}}}
|
||||
},
|
||||
},
|
||||
// A send refused by its machine for the generation it was composed from (novox/hq issue 234): said at
|
||||
// once, naming its sender, for an hour after. Inside: the last such refusal was more than an hour ago.
|
||||
"S20": {
|
||||
inside: func(f *signalFacts) {
|
||||
f.refusedSends = []inventory.Send{refusedSend(f.now.Add(-61 * time.Minute))}
|
||||
},
|
||||
past: func(f *signalFacts) { f.refusedSends = []inventory.Send{refusedSend(f.now.Add(-time.Second))} },
|
||||
},
|
||||
// Twice by hand within a fortnight is a healer wanted; once, or the first of two a day too old, is not.
|
||||
"S15": {
|
||||
inside: func(f *signalFacts) {
|
||||
|
||||
@@ -196,6 +196,11 @@ func (l nudgingListener) Heard(ctx context.Context, report link.Report) (bool, e
|
||||
if report.Ordered() {
|
||||
link.StaleRefusals.Lifetime(report.Node, report.RefusedOlder, now)
|
||||
}
|
||||
// A send refused for the generation it was composed from is kept on the send it refused, so S20 names
|
||||
// its sender (novox/hq issue 234).
|
||||
if report.OlderGeneration != nil {
|
||||
heardGenerationRefusal(ctx, l.Enrolment.Inventory, report)
|
||||
}
|
||||
// What the machine's witnesses put back and stand by (novox/hq ADR 0236): read by the gate and its
|
||||
// probe. Only from an account of the machine — not a word that a declaration was set aside, nor a rekey.
|
||||
if report.Superseded == "" && report.Rekey == nil && report.Node != "" {
|
||||
|
||||
@@ -55,6 +55,10 @@ type signalFacts struct {
|
||||
waits []waitFacts
|
||||
// batches are the batches not yet cut (S18, S19, novox/hq ADR 0276).
|
||||
batches []batchFacts
|
||||
// refusedSends are the sends a machine refused within the hour for the generation they were composed
|
||||
// from, each naming its sender (S20, novox/hq issue 234).
|
||||
refusedSends []inventory.Send
|
||||
refusedSendsErr error
|
||||
|
||||
loop loopFacts
|
||||
loopErr error
|
||||
@@ -364,6 +368,7 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
|
||||
}
|
||||
}
|
||||
f.handActs, f.handActsErr = w.gatherHandActs(ctx, now)
|
||||
f.refusedSends, f.refusedSendsErr = inv.RefusedSendsSince(ctx, now.Add(-generationRefusalsSaid))
|
||||
f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last()
|
||||
return f
|
||||
}
|
||||
|
||||
@@ -3,7 +3,7 @@ module github.com/novox/mesh-controller
|
||||
go 1.26.0
|
||||
|
||||
require (
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.14-0.20261010191001-33917f91ac3c
|
||||
github.com/jackc/pgx/v5 v5.10.0
|
||||
github.com/nats-io/nats-server/v2 v2.11.17
|
||||
github.com/nats-io/nats.go v1.54.0
|
||||
@@ -35,4 +35,4 @@ require (
|
||||
// committed. Every build (the build agent's `go build`, the Dockerfile) compiles from vendor/ and
|
||||
// fetches nothing; go refuses to build when vendor/ and this file disagree, so a pin moved without
|
||||
// `go mod vendor` fails loudly, at once, everywhere.
|
||||
replace github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
replace github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261011092230-20d5af9b2d5f
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812 h1:pzVzwF5VMWaTECxu8+Pd1dNoOHNEm7upC5wPadQTkBw=
|
||||
git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812/go.mod h1:K3/xEzVgmrNKLMV2vv4M80MwmPnQNXqvQ4C5Jj0fJT4=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970 h1:9tFDQsgmI+4X7/BpZGXIr+HemPKE7YddYGqWV0lINAI=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
|
||||
git.novox.be/novox/mesh-host v0.0.0-20261011092230-20d5af9b2d5f h1:8N5OW2mdTNIck2pe4EciTYX5NsrPsHrTLENGNIWYNTU=
|
||||
git.novox.be/novox/mesh-host v0.0.0-20261011092230-20d5af9b2d5f/go.mod h1:tcTK4LMs1d6JpZUwy3hSHMFG/MIuDr/XFs7/nbeaW/c=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.14-0.20261010191001-33917f91ac3c h1:l5onJwoIeH8yE/PjVlEeoPyXBsHp2NQiWLllz2fwX7s=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.14-0.20261010191001-33917f91ac3c/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
|
||||
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op h1:Z/MZK75wC/NSrkgqeNIa7jexam9uWzhLmFTSCPI/kn0=
|
||||
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
-- A declaration says the assignment generation it was composed from, and every send is recorded
|
||||
-- (novox/hq issue 234).
|
||||
--
|
||||
-- On 2026-10-04 a declaration newer in sequence than every other named four fewer modules than the
|
||||
-- assignments held, and the machine undeclared all four. The sequence orders arrival; it cannot tell a
|
||||
-- later send that carries an older view of the assignments. And nothing the mesh kept said which process
|
||||
-- sent it, from what view: the node row holds only the digest of the last send.
|
||||
|
||||
-- The assignment generation: one counter for the mesh, raised in the same transaction as every change to
|
||||
-- what is assigned where. A trigger, not a line in each writer: an assignment is also taken by a node's
|
||||
-- removal (on delete cascade) and renamed with its module (on update cascade), which no writer in Go
|
||||
-- sees, and a rule kept only by the writers that remember it is the rule this issue found broken. Seeded
|
||||
-- at 1, so every declaration composed after this claims a generation and none claims zero ("none").
|
||||
create table assignment_generation (
|
||||
one boolean primary key default true check (one),
|
||||
generation bigint not null
|
||||
);
|
||||
insert into assignment_generation (one, generation) values (true, 1);
|
||||
|
||||
create function raise_assignment_generation() returns trigger language plpgsql as $$
|
||||
begin
|
||||
update assignment_generation set generation = generation + 1;
|
||||
return null;
|
||||
end
|
||||
$$;
|
||||
|
||||
-- Per row, so a statement that changes nothing raises nothing: an assignment repeated (`on conflict do
|
||||
-- nothing` inserts no row) and an update that leaves a row as it was (the WHEN below; a plain UPDATE fires a
|
||||
-- row trigger whether or not it changed anything). A statement that changes several rows raises once per
|
||||
-- row, which only ever moves it forward. Two triggers, because only an update has both OLD and NEW.
|
||||
create trigger assignment_generation_raised
|
||||
after insert or delete on assignment
|
||||
for each row execute function raise_assignment_generation();
|
||||
create trigger assignment_generation_raised_by_update
|
||||
after update on assignment
|
||||
for each row when (old is distinct from new) execute function raise_assignment_generation();
|
||||
|
||||
-- The generation a machine was last sent, written in the same transaction as the digest of that send, so
|
||||
-- what the mesh WOULD send is stamped the same way and reads as byte for byte what it DID send when
|
||||
-- nothing else changed (as sent_epoch, migration 0068). Null for a send without one.
|
||||
alter table node add column sent_generation bigint;
|
||||
-- Whether the machine's node-engine said it reads a generation (its reports' `reads_generation`): until
|
||||
-- it has, it is sent none, because an older node-engine refuses a key it does not know, whole.
|
||||
alter table node add column reads_generation boolean not null default false;
|
||||
|
||||
-- Every send: the machine, its sequence, who sent it — the lease epoch the sender acted under and the
|
||||
-- caller, process and build — the generation it was composed from, its digest and the modules it named.
|
||||
-- Written after the declaration is away, with the node row's record of it. A refusal by the machine of a
|
||||
-- send for its generation is kept on the send it refused, so the condition it raises names the sender
|
||||
-- from here; a refusal of a sequence the mesh has no send for is kept as a row of its own, `recorded`
|
||||
-- false, whose sender is said to be unknown. `counter_raised_to` is the generation the mesh's counter was
|
||||
-- raised to when the refusal showed the counter behind the machine (a store put back from a backup). The
|
||||
-- newest 200 per machine are kept.
|
||||
create table declaration_send (
|
||||
id bigserial primary key,
|
||||
node uuid not null references node (id) on delete cascade,
|
||||
sequence bigint,
|
||||
epoch bigint,
|
||||
sender text not null,
|
||||
generation bigint,
|
||||
digest text not null,
|
||||
modules text[] not null default '{}',
|
||||
sent_at timestamptz not null default now(),
|
||||
refused_at timestamptz,
|
||||
-- refused_applied is the generation the machine said it had applied when it refused this send.
|
||||
refused_applied bigint,
|
||||
counter_raised_to bigint,
|
||||
recorded boolean not null default true
|
||||
);
|
||||
create index declaration_send_node on declaration_send (node, id desc);
|
||||
create index declaration_send_sequence on declaration_send (node, sequence);
|
||||
create index declaration_send_refused on declaration_send (refused_at) where refused_at is not null;
|
||||
@@ -0,0 +1,287 @@
|
||||
package inventory
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
)
|
||||
|
||||
// The assignment generation and the record of every send (novox/hq issue 234, migration 0094).
|
||||
//
|
||||
// A declaration carries the generation of the assignments it was composed from, and a machine refuses one
|
||||
// composed from an older generation than it applied — the sequence orders arrival and cannot tell a later
|
||||
// send carrying an older view. And every send is written down with who sent it, from which generation and
|
||||
// naming which modules, so what a machine was last told, by whom, is read back rather than guessed at.
|
||||
|
||||
// sendsKept is how many sends are kept per machine: the newest. Enough to read back a morning of pushes
|
||||
// on a busy machine; the record is for finding who sent what, not an archive.
|
||||
const sendsKept = 200
|
||||
|
||||
// AssignmentGeneration is the mesh's assignment generation now: raised by the store itself in the same
|
||||
// transaction as every change to what is assigned where (a trigger on the assignment table).
|
||||
func (i *Inventory) AssignmentGeneration(ctx context.Context) (int64, error) {
|
||||
var g int64
|
||||
if err := i.store.Pool().QueryRow(ctx, `select generation from assignment_generation`).Scan(&g); err != nil {
|
||||
return 0, fmt.Errorf("reading the assignment generation: %w", err)
|
||||
}
|
||||
return g, nil
|
||||
}
|
||||
|
||||
// RaiseAssignmentGeneration raises the generation past one a machine applied, and answers it: one more
|
||||
// than that, or what it already was when it is past it. It never lowers it.
|
||||
//
|
||||
// For one case only: a machine refused a send composed from the generation the mesh holds now, so the
|
||||
// machine applied a higher one than the mesh's counter — a store put back from a backup. Every send after
|
||||
// would be refused for ever; raised past it, the next send carries a generation the machine takes, composed
|
||||
// from the assignments the store holds now. A stale sender never meets this: its generation is below the
|
||||
// counter, and the counter is not touched for it.
|
||||
func (i *Inventory) RaiseAssignmentGeneration(ctx context.Context, past int64) (int64, error) {
|
||||
var g int64
|
||||
err := i.store.Pool().QueryRow(ctx,
|
||||
`update assignment_generation set generation = greatest(generation, $1 + 1) returning generation`, past).Scan(&g)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("raising the assignment generation past %d: %w", past, err)
|
||||
}
|
||||
return g, nil
|
||||
}
|
||||
|
||||
// SentGeneration is the generation a machine was last sent, by its id; zero for one sent without. What the
|
||||
// mesh WOULD send is stamped with it, so it reads as byte for byte what it DID send when nothing else changed.
|
||||
func (i *Inventory) SentGeneration(ctx context.Context, id string) (int64, error) {
|
||||
var g *int64
|
||||
if err := i.store.Pool().QueryRow(ctx, `select sent_generation from node where id = $1`, id).Scan(&g); err != nil {
|
||||
return 0, fmt.Errorf("reading the generation %s was last sent: %w", id, err)
|
||||
}
|
||||
if g == nil {
|
||||
return 0, nil
|
||||
}
|
||||
return *g, nil
|
||||
}
|
||||
|
||||
// ReadsGeneration says a machine's node-engine said it reads a generation in a declaration, by its id.
|
||||
func (i *Inventory) ReadsGeneration(ctx context.Context, id string) (bool, error) {
|
||||
var reads bool
|
||||
if err := i.store.Pool().QueryRow(ctx, `select reads_generation from node where id = $1`, id).Scan(&reads); err != nil {
|
||||
return false, fmt.Errorf("reading whether %s reads a generation: %w", id, err)
|
||||
}
|
||||
return reads, nil
|
||||
}
|
||||
|
||||
// RecordReadsGeneration keeps what a machine's latest report said of reading a generation.
|
||||
func (i *Inventory) RecordReadsGeneration(ctx context.Context, id string, reads bool) error {
|
||||
_, err := i.store.Pool().Exec(ctx, `update node set reads_generation = $2 where id = $1`, id, reads)
|
||||
return err
|
||||
}
|
||||
|
||||
// Send is one declaration sent to a machine, as the mesh records it.
|
||||
type Send struct {
|
||||
// Node is the machine's id; NodeName its name, filled where it is read back.
|
||||
Node string
|
||||
NodeName string
|
||||
Sequence int64
|
||||
// Epoch is the lease epoch its sender acted under, zero for one that acted under none.
|
||||
Epoch int64
|
||||
// Sender is who sent it, in words: the caller, and the process and build that composed it.
|
||||
Sender string
|
||||
// Generation is the assignment generation it was composed from, zero when not known.
|
||||
Generation int64
|
||||
Digest string
|
||||
// Modules are the modules it named.
|
||||
Modules []string
|
||||
SentAt time.Time
|
||||
// RefusedAt is when the machine refused it for its generation, nil when it did not; RefusedApplied the
|
||||
// generation the machine said it had applied then; CounterRaisedTo what the mesh's counter was raised to
|
||||
// because the refusal showed it behind the machine, zero when it was not.
|
||||
RefusedAt *time.Time
|
||||
RefusedApplied int64
|
||||
CounterRaisedTo int64
|
||||
// Recorded is false for a refusal of a sequence the mesh has no send for: its sender is unknown.
|
||||
Recorded bool
|
||||
}
|
||||
|
||||
// unknownSender is the sender of a refused sequence the mesh has no record of sending.
|
||||
const unknownSender = "a sender the mesh has no record of (no send of this sequence was recorded: sent by " +
|
||||
"hand, before sends were recorded, or by a controller whose record was not written)"
|
||||
|
||||
// SendNotRecordedError is a send whose machine's record was written — its digest and the generation it
|
||||
// carried — and whose record of who sent it was not. The send is away and the machine's record stands; the
|
||||
// caller says this loudly and does not take the send for failed.
|
||||
type SendNotRecordedError struct{ Err error }
|
||||
|
||||
func (e *SendNotRecordedError) Error() string {
|
||||
return "the record of the send was not written: " + e.Err.Error()
|
||||
}
|
||||
func (e *SendNotRecordedError) Unwrap() error { return e.Err }
|
||||
|
||||
// RecordSentWith writes down what a machine was just sent — its digest, the builds it carried, the epoch and
|
||||
// the generation on the wire — and the send itself, who sent it and from which generation (novox/hq issue
|
||||
// 234), in one transaction. The digest and the generation are written together, so the would-send is never
|
||||
// stamped with a generation the machine was not sent. The send's own record is written under a savepoint: if
|
||||
// it cannot be, the machine's record is still committed and a *SendNotRecordedError says so. The oldest
|
||||
// sends beyond the newest 200 for the machine are let go.
|
||||
func (i *Inventory) RecordSentWith(ctx context.Context, node, digest string, builds map[string]string, epoch uint64,
|
||||
generation int64, s Send) error {
|
||||
var sentEpoch *int64
|
||||
if epoch > 0 {
|
||||
e := int64(epoch)
|
||||
sentEpoch = &e
|
||||
}
|
||||
var carried *string
|
||||
if builds != nil {
|
||||
raw, err := json.Marshal(builds)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
text := string(raw)
|
||||
carried = &text
|
||||
}
|
||||
tx, err := i.store.Pool().Begin(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
if _, err := tx.Exec(ctx,
|
||||
`update node set sent = $2, sent_at = now(), sent_builds = $3::jsonb, sent_epoch = $4, sent_generation = $5
|
||||
where id = $1`, node, digest, carried, sentEpoch, nullIfZero(generation)); err != nil {
|
||||
return err
|
||||
}
|
||||
recordErr := recordSend(ctx, tx, node, s)
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
if recordErr != nil {
|
||||
return &SendNotRecordedError{Err: recordErr}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// recordSend writes one send under a savepoint of tx, which it rolls back when the send cannot be written.
|
||||
func recordSend(ctx context.Context, tx pgx.Tx, node string, s Send) error {
|
||||
sp, err := tx.Begin(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = sp.Rollback(ctx) }()
|
||||
modules := s.Modules
|
||||
if modules == nil {
|
||||
modules = []string{}
|
||||
}
|
||||
if _, err := sp.Exec(ctx,
|
||||
`insert into declaration_send (node, sequence, epoch, sender, generation, digest, modules)
|
||||
values ($1, $2, $3, $4, $5, $6, $7)`,
|
||||
node, nullIfZero(s.Sequence), nullIfZero(s.Epoch), s.Sender, nullIfZero(s.Generation), s.Digest,
|
||||
modules); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := pruneSends(ctx, sp, node); err != nil {
|
||||
return err
|
||||
}
|
||||
return sp.Commit(ctx)
|
||||
}
|
||||
|
||||
// pruneSends lets go of a machine's sends beyond the newest 200.
|
||||
func pruneSends(ctx context.Context, q queries, node string) error {
|
||||
_, err := q.Exec(ctx,
|
||||
`delete from declaration_send where node = $1 and id not in
|
||||
(select id from declaration_send where node = $1 order by id desc limit $2)`, node, sendsKept)
|
||||
return err
|
||||
}
|
||||
|
||||
// sendColumns are a send's columns as scanSends reads them.
|
||||
const sendColumns = `s.node, n.name, coalesce(s.sequence, 0), coalesce(s.epoch, 0), s.sender, coalesce(s.generation, 0),
|
||||
s.digest, s.modules, s.sent_at, s.refused_at, coalesce(s.refused_applied, 0), coalesce(s.counter_raised_to, 0),
|
||||
s.recorded`
|
||||
|
||||
func scanSends(rows pgx.Rows) ([]Send, error) {
|
||||
defer rows.Close()
|
||||
var out []Send
|
||||
for rows.Next() {
|
||||
var s Send
|
||||
if err := rows.Scan(&s.Node, &s.NodeName, &s.Sequence, &s.Epoch, &s.Sender, &s.Generation, &s.Digest,
|
||||
&s.Modules, &s.SentAt, &s.RefusedAt, &s.RefusedApplied, &s.CounterRaisedTo, &s.Recorded); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, s)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// LastSend is the last declaration a machine was sent, by its name; false when none was recorded. A refusal
|
||||
// of a sequence the mesh has no send for is not a send, and is not it.
|
||||
func (i *Inventory) LastSend(ctx context.Context, name string) (Send, bool, error) {
|
||||
rows, err := i.store.Pool().Query(ctx, `select `+sendColumns+`
|
||||
from declaration_send s join node n on n.id = s.node
|
||||
where n.name = $1 and s.recorded order by s.id desc limit 1`, name)
|
||||
if err != nil {
|
||||
return Send{}, false, err
|
||||
}
|
||||
sends, err := scanSends(rows)
|
||||
if err != nil || len(sends) == 0 {
|
||||
return Send{}, false, err
|
||||
}
|
||||
return sends[0], true, nil
|
||||
}
|
||||
|
||||
// RefusedSend keeps a machine's refusal of the send of a sequence for the generation it came from, and
|
||||
// answers that send, sender and all; false when no send of that sequence was recorded (RecordUnrecordedRefusal
|
||||
// keeps that one). The newest of that sequence when there are several: a declaration a person sent by hand may
|
||||
// repeat one. raisedTo is what the mesh's counter was raised to for it, zero for none.
|
||||
func (i *Inventory) RefusedSend(ctx context.Context, node string, sequence, applied, raisedTo int64) (Send, bool, error) {
|
||||
rows, err := i.store.Pool().Query(ctx, `with refused as (
|
||||
update declaration_send set refused_at = now(), refused_applied = $3, counter_raised_to = $4
|
||||
where id = (select id from declaration_send where node = $1 and sequence = $2 and recorded
|
||||
order by id desc limit 1)
|
||||
returning *)
|
||||
select `+sendColumns+` from refused s join node n on n.id = s.node`, node, sequence, applied,
|
||||
nullIfZero(raisedTo))
|
||||
if err != nil {
|
||||
return Send{}, false, err
|
||||
}
|
||||
sends, err := scanSends(rows)
|
||||
if err != nil || len(sends) == 0 {
|
||||
return Send{}, false, err
|
||||
}
|
||||
return sends[0], true, nil
|
||||
}
|
||||
|
||||
// RecordUnrecordedRefusal keeps a machine's refusal of a sequence the mesh has no send for, as a row of its
|
||||
// own whose sender is unknown, and answers it: S20 raises it like any other, because a refusal nobody can
|
||||
// attribute is no quieter for that. s carries the machine, the sequence, epoch, generation and digest the
|
||||
// refused declaration carried, and the refusal's applied generation and counter raise.
|
||||
func (i *Inventory) RecordUnrecordedRefusal(ctx context.Context, s Send) (Send, error) {
|
||||
rows, err := i.store.Pool().Query(ctx, `with refused as (
|
||||
insert into declaration_send (node, sequence, epoch, sender, generation, digest, refused_at, refused_applied,
|
||||
counter_raised_to, recorded)
|
||||
values ($1, $2, $3, $4, $5, $6, now(), $7, $8, false) returning *)
|
||||
select `+sendColumns+` from refused s join node n on n.id = s.node`,
|
||||
s.Node, nullIfZero(s.Sequence), nullIfZero(s.Epoch), unknownSender, nullIfZero(s.Generation), s.Digest,
|
||||
s.RefusedApplied, nullIfZero(s.CounterRaisedTo))
|
||||
if err != nil {
|
||||
return Send{}, err
|
||||
}
|
||||
sends, err := scanSends(rows)
|
||||
if err != nil {
|
||||
return Send{}, err
|
||||
}
|
||||
if len(sends) == 0 {
|
||||
return Send{}, fmt.Errorf("the refusal of sequence %d was not kept", s.Sequence)
|
||||
}
|
||||
if err := pruneSends(ctx, i.store.Pool(), s.Node); err != nil {
|
||||
return Send{}, err
|
||||
}
|
||||
return sends[0], nil
|
||||
}
|
||||
|
||||
// RefusedSendsSince is every send refused for its generation since a moment, newest first.
|
||||
func (i *Inventory) RefusedSendsSince(ctx context.Context, since time.Time) ([]Send, error) {
|
||||
rows, err := i.store.Pool().Query(ctx, `select `+sendColumns+`
|
||||
from declaration_send s join node n on n.id = s.node
|
||||
where s.refused_at >= $1 order by s.refused_at desc`, since)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return scanSends(rows)
|
||||
}
|
||||
@@ -0,0 +1,207 @@
|
||||
package inventory
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The assignment generation and the record of every send (novox/hq issue 234). On 2026-10-04 a declaration
|
||||
// newer in sequence than every other named four fewer modules than the assignments held, a machine applied
|
||||
// it, and nothing the mesh kept could say who sent it or from what view.
|
||||
|
||||
func generationNow(t *testing.T, inv *Inventory) int64 {
|
||||
t.Helper()
|
||||
g, err := inv.AssignmentGeneration(t.Context())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return g
|
||||
}
|
||||
|
||||
// **Every change to what is assigned where raises the generation, in its own transaction** — an
|
||||
// assignment, an unassignment, and one taken by the machine's removal, which no writer in Go makes — and
|
||||
// a repeated assignment, which changes nothing, does not.
|
||||
func TestEveryAssignmentChangeRaisesTheGeneration(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t, "postgres", "web")
|
||||
ctx := t.Context()
|
||||
start := generationNow(t, inv)
|
||||
if start < 1 {
|
||||
t.Fatalf("the generation starts at %d; zero is \"none claimed\" on the wire", start)
|
||||
}
|
||||
if _, err := inv.Assign(ctx, node, "postgres"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assigned := generationNow(t, inv)
|
||||
if assigned <= start {
|
||||
t.Fatalf("an assignment left the generation at %d", assigned)
|
||||
}
|
||||
if _, err := inv.Assign(ctx, node, "postgres"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if again := generationNow(t, inv); again != assigned {
|
||||
t.Errorf("an assignment repeated, which changed nothing, moved the generation from %d to %d", assigned, again)
|
||||
}
|
||||
if err := inv.Unassign(ctx, node, "postgres"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
unassigned := generationNow(t, inv)
|
||||
if unassigned <= assigned {
|
||||
t.Fatalf("an unassignment left the generation at %d", unassigned)
|
||||
}
|
||||
if _, err := inv.Assign(ctx, node, "web"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before := generationNow(t, inv)
|
||||
if _, err := inv.RemoveNodeForTest(ctx, node); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if removed := generationNow(t, inv); removed <= before {
|
||||
t.Errorf("a machine's removal took its assignments and left the generation at %d", removed)
|
||||
}
|
||||
}
|
||||
|
||||
// **The generation never goes down, and is raised past what a machine applied when the mesh's own counter
|
||||
// is behind it** — a store put back from a backup — so the next send is not refused for ever.
|
||||
func TestTheGenerationIsRaisedPastWhatAMachineApplied(t *testing.T) {
|
||||
inv, _ := aNodeWithModules(t)
|
||||
now := generationNow(t, inv)
|
||||
if raised, err := inv.RaiseAssignmentGeneration(t.Context(), now+40); err != nil || raised != now+41 {
|
||||
t.Fatalf("raised past %d to %d (%v)", now+40, raised, err)
|
||||
}
|
||||
if raised, err := inv.RaiseAssignmentGeneration(t.Context(), 2); err != nil || raised != now+41 {
|
||||
t.Fatalf("a raise below the counter moved it to %d (%v)", raised, err)
|
||||
}
|
||||
}
|
||||
|
||||
// **Every send is recorded: its sequence, its sender, the generation it came from and the modules it
|
||||
// named** — and the machine's last send is what `plan` reads, the would-send is stamped with the generation
|
||||
// written with the digest, and a refusal is kept on the send it refused, naming its sender.
|
||||
func TestASendIsRecordedWithItsSenderGenerationAndModules(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t)
|
||||
ctx := t.Context()
|
||||
record, err := inv.NodeByName(ctx, node)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
first := Send{Sequence: 11, Epoch: 57, Sender: "the controller's daemon (pid 7 on anchor)",
|
||||
Generation: 40, Digest: "d11", Modules: []string{"docker", "pacman", "sudo"}}
|
||||
if err := inv.RecordSentWith(ctx, record.ID, "d11", nil, 57, 40, first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stale := Send{Sequence: 12, Epoch: 57, Sender: "a one-shot push by jochen at a shell on anchor",
|
||||
Generation: 38, Digest: "d12", Modules: []string{"docker"}}
|
||||
if err := inv.RecordSentWith(ctx, record.ID, "d12", nil, 57, 38, stale); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
last, found, err := inv.LastSend(ctx, node)
|
||||
if err != nil || !found {
|
||||
t.Fatalf("the last send is not kept: %v", err)
|
||||
}
|
||||
if last.Sequence != 12 || last.Sender != stale.Sender || last.Generation != 38 || !last.Recorded ||
|
||||
!slices.Equal(last.Modules, []string{"docker"}) || last.SentAt.IsZero() {
|
||||
t.Fatalf("the last send reads %+v", last)
|
||||
}
|
||||
if generation, err := inv.SentGeneration(ctx, record.ID); err != nil || generation != 38 {
|
||||
t.Fatalf("what the machine was last sent of it reads %d (%v)", generation, err)
|
||||
}
|
||||
|
||||
refused, found, err := inv.RefusedSend(ctx, record.ID, 12, 40, 0)
|
||||
if err != nil || !found || refused.Sender != stale.Sender || refused.RefusedApplied != 40 ||
|
||||
refused.CounterRaisedTo != 0 {
|
||||
t.Fatalf("the refusal is not kept on the send it refused: %+v, %v, %v", refused, found, err)
|
||||
}
|
||||
since, err := inv.RefusedSendsSince(ctx, time.Now().Add(-time.Hour))
|
||||
if err != nil || len(since) != 1 || since[0].NodeName != node || since[0].Sequence != 12 {
|
||||
t.Fatalf("the refused sends within the hour read %+v (%v)", since, err)
|
||||
}
|
||||
if _, found, err := inv.RefusedSend(ctx, record.ID, 99, 40, 0); err != nil || found {
|
||||
t.Errorf("a refusal of a send never recorded was found: %v, %v", found, err)
|
||||
}
|
||||
}
|
||||
|
||||
// **A refusal of a sequence the mesh has no send for is kept too, its sender unknown**, so S20 raises it;
|
||||
// it is not the machine's last send, and it says what the counter was raised to.
|
||||
func TestARefusalOfASendNeverRecordedIsKeptWithAnUnknownSender(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t)
|
||||
ctx := t.Context()
|
||||
record, err := inv.NodeByName(ctx, node)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
kept, err := inv.RecordUnrecordedRefusal(ctx, Send{Node: record.ID, Sequence: 99, Epoch: 57, Generation: 41,
|
||||
Digest: "d99", RefusedApplied: 44, CounterRaisedTo: 45})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if kept.Recorded || kept.Sender != unknownSender || kept.RefusedAt == nil || kept.CounterRaisedTo != 45 ||
|
||||
kept.NodeName != node {
|
||||
t.Fatalf("the refusal reads %+v", kept)
|
||||
}
|
||||
since, err := inv.RefusedSendsSince(ctx, time.Now().Add(-time.Hour))
|
||||
if err != nil || len(since) != 1 || since[0].Sequence != 99 {
|
||||
t.Fatalf("the refused sends within the hour read %+v (%v)", since, err)
|
||||
}
|
||||
if _, found, err := inv.LastSend(ctx, node); err != nil || found {
|
||||
t.Errorf("a refusal of a send never recorded reads as the machine's last send: %v, %v", found, err)
|
||||
}
|
||||
}
|
||||
|
||||
// **The machine's record and the send's are written together, and a send record that cannot be written
|
||||
// leaves the machine's record standing and says so**: the digest and generation are the machine's, and a
|
||||
// send that is away must not read as failed, nor the machine as behind.
|
||||
func TestASendRecordThatCannotBeWrittenLeavesTheMachinesRecord(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t)
|
||||
ctx := t.Context()
|
||||
record, err := inv.NodeByName(ctx, node)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// A NUL byte is text PostgreSQL refuses: the send's own row cannot be written.
|
||||
broken := Send{Sequence: 13, Sender: "a test", Generation: 42, Digest: "d\x0013"}
|
||||
err = inv.RecordSentWith(ctx, record.ID, "d13", nil, 0, 42, broken)
|
||||
var unrecorded *SendNotRecordedError
|
||||
if !errors.As(err, &unrecorded) {
|
||||
t.Fatalf("a send record that could not be written read as %v", err)
|
||||
}
|
||||
if generation, err := inv.SentGeneration(ctx, record.ID); err != nil || generation != 42 {
|
||||
t.Fatalf("the machine's generation was not written with its digest: %d (%v)", generation, err)
|
||||
}
|
||||
if sent, err := inv.Outstanding(ctx, node); err != nil || sent != "d13" {
|
||||
t.Fatalf("the machine's digest was not written: %q (%v)", sent, err)
|
||||
}
|
||||
}
|
||||
|
||||
// **An update that leaves an assignment as it was raises nothing** (the trigger's WHEN).
|
||||
func TestAnUpdateThatChangesNoAssignmentRaisesNothing(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t, "web")
|
||||
ctx := t.Context()
|
||||
if _, err := inv.Assign(ctx, node, "web"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before := generationNow(t, inv)
|
||||
if _, err := inv.store.Pool().Exec(ctx, `update assignment set module = module`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if after := generationNow(t, inv); after != before {
|
||||
t.Errorf("an update that changed nothing moved the generation from %d to %d", before, after)
|
||||
}
|
||||
}
|
||||
|
||||
// **A machine that reads a generation is recorded from its reports**, and one rolled back says so no longer.
|
||||
func TestWhetherAMachineReadsAGenerationIsItsLatestWord(t *testing.T) {
|
||||
inv, node := aNodeWithModules(t)
|
||||
record, err := inv.NodeByName(t.Context(), node)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, reads := range []bool{false, true, false} {
|
||||
if err := inv.RecordReadsGeneration(t.Context(), record.ID, reads); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got, err := inv.ReadsGeneration(t.Context(), record.ID); err != nil || got != reads {
|
||||
t.Fatalf("recorded %v, read %v (%v)", reads, got, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -447,6 +447,13 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
|
||||
if err := e.Inventory.RecordReadsEpoch(ctx, node.ID, report.ReadsEpoch()); err != nil {
|
||||
return false, err
|
||||
}
|
||||
// And whether it reads a generation (novox/hq issue 234), the same way: only from a node-engine that
|
||||
// orders its reports, so a one-shot report (a rekey) does not say an engine stopped reading one.
|
||||
if report.Ordered() {
|
||||
if err := e.Inventory.RecordReadsGeneration(ctx, node.ID, report.ReadsGeneration); err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
if report.Ordered() {
|
||||
account := AccountOf(report)
|
||||
news, err = e.Inventory.RecordOrderedDoing(ctx, node.ID, doing,
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// A report's word on the assignment generation (novox/hq issue 234), as mesh-host's report writes it:
|
||||
// `reads_generation` on every report of a node-engine that reads one, and `older_generation` on a refusal
|
||||
// of a declaration composed from an older generation than the machine applied.
|
||||
func TestTheGenerationOnTheWire(t *testing.T) {
|
||||
raw := []byte(`{"node":"anchor","refused":"older generation","declared":"d12","epoch":57,"sequence":12,` +
|
||||
`"report_sequence":8,"reads_generation":true,"older_generation":{"generation":38,"applied":40}}`)
|
||||
var r Report
|
||||
if err := json.Unmarshal(raw, &r); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !r.ReadsGeneration || r.OlderGeneration == nil || *r.OlderGeneration != (GenerationRefusal{38, 40}) ||
|
||||
r.Sequence != 12 {
|
||||
t.Fatalf("a node-engine's generation refusal reads as %+v", r)
|
||||
}
|
||||
// Not a refusal by the lease's epoch: S13 counts those, and this one is raised naming its sender.
|
||||
if r.StaleRefusalOf() {
|
||||
t.Errorf("a generation refusal reads as a stale writer's")
|
||||
}
|
||||
var older Report
|
||||
if err := json.Unmarshal([]byte(`{"node":"anchor","applied":["a"],"report_sequence":3}`), &older); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if older.ReadsGeneration || older.OlderGeneration != nil {
|
||||
t.Fatalf("a node-engine that never said it reads a generation reads as one that does: %+v", older)
|
||||
}
|
||||
}
|
||||
@@ -197,6 +197,15 @@ type Report struct {
|
||||
// RefusedOlder is how many declarations the node-engine has refused as older, ever, on every report:
|
||||
// a refusal whose own report was lost is still counted from the next.
|
||||
RefusedOlder int64 `json:"refused_older,omitempty"`
|
||||
// ReadsGeneration says the node-engine reads a declaration's assignment generation and its put-back
|
||||
// mark (novox/hq issue 234), on every report it makes: the mesh sends them only to a machine that said
|
||||
// so, because an older node-engine refuses a key it does not know, whole. Mirrors mesh-host
|
||||
// internal/link/messages.go.
|
||||
ReadsGeneration bool `json:"reads_generation,omitempty"`
|
||||
// OlderGeneration is set on a report refusing a declaration composed from an older assignment
|
||||
// generation than the machine applied (novox/hq issue 234). The refused declaration is the report's
|
||||
// own Declared and Order, and its sender is read from the record of sends by that sequence.
|
||||
OlderGeneration *GenerationRefusal `json:"older_generation,omitempty"`
|
||||
|
||||
// Held is what an adopted node found and is keeping as it was until its module is taken
|
||||
// (novox/hq ADR 0100). Without it an adopted node reads as converged.
|
||||
@@ -605,3 +614,10 @@ type HealthSaid struct {
|
||||
Node string `json:"node"`
|
||||
Health Health `json:"health"`
|
||||
}
|
||||
|
||||
// GenerationRefusal is a declaration refused for the assignment generation it was composed from: that
|
||||
// generation, and the highest the machine applied (novox/hq issue 234; mesh-host's GenerationRefusal).
|
||||
type GenerationRefusal struct {
|
||||
Generation int64 `json:"generation"`
|
||||
Applied int64 `json:"applied"`
|
||||
}
|
||||
|
||||
Vendored
+1
-1
@@ -9,7 +9,7 @@ a second place to clean besides the catalogue itself.
|
||||
|
||||
mesh-catalog b9de001833b1b61b297182b8e3bcdae1cbdeace9 modules/*/module.json, modules/nats/Dockerfile
|
||||
mesh-host bd5cc6980419c1bd4824be6d58dfebd2381a9de3 examples/foundation-first-node-nats.lock
|
||||
mesh-sdk f047d0a4a9702f5d7f3d7eadd3e62496311acc00 conformance/events/module-event.json
|
||||
mesh-sdk 33917f91ac3cc213b879b32f3c70e2aa70395f98 conformance/events/module-event.json
|
||||
|
||||
To move them, from this repository's root, with the two repositories checked out beside it:
|
||||
|
||||
|
||||
+16
-1
@@ -154,7 +154,9 @@ func canonical(b *strings.Builder, v string) {
|
||||
//
|
||||
// It is over the ask's named fields in a fixed order, each written canonically, and the expiry as UTC
|
||||
// RFC 3339 to the nanosecond — never over a language's encoding of the struct, so a field added to Ask
|
||||
// later changes no digest until it is added here, on purpose.
|
||||
// later changes no digest until it is added here, on purpose. Whole and Details (novox/hq issue 383) follow
|
||||
// the options only when the ask gives either: an ask without them digests as it did before they existed, so
|
||||
// a router and an asker of different builds still agree on every such ask.
|
||||
func (a Ask) Digest() string {
|
||||
var b strings.Builder
|
||||
b.WriteString("novox.ask.v1\n")
|
||||
@@ -168,6 +170,10 @@ func (a Ask) Digest() string {
|
||||
canonical(&b, v)
|
||||
}
|
||||
}
|
||||
if a.Whole != "" || a.Details != "" {
|
||||
canonical(&b, a.Whole)
|
||||
canonical(&b, a.Details)
|
||||
}
|
||||
sum := sha256.Sum256([]byte(b.String()))
|
||||
return "sha256:" + hex.EncodeToString(sum[:])
|
||||
}
|
||||
@@ -190,6 +196,15 @@ type Ask struct {
|
||||
// About is what the ask is about (a condition's key): a newer ask about it replaces the older.
|
||||
About string `json:"about,omitempty"`
|
||||
Urgent bool `json:"urgent,omitempty"`
|
||||
// Whole is the explanation with its exact values whole, shown in place of Explanation on a channel kind
|
||||
// that proves who answers (verified-sender) and carries the ask's answers (novox/hq issue 383): what the
|
||||
// person approves — a mount point, a share, a private address — must be readable where they approve it.
|
||||
// The asker withholds a value shaped like a secret in it, and the router refuses the ask when one is left;
|
||||
// Explanation stays under the whole content rule everywhere else. Empty, Explanation is shown everywhere.
|
||||
Whole string `json:"whole,omitempty"`
|
||||
// Details is what the Details answer on the ask's message shows, line by line under the content rule: a
|
||||
// fingerprint, how to read the proposal whole at the terminal. Empty, an ask's own message offers no Details.
|
||||
Details string `json:"details,omitempty"`
|
||||
}
|
||||
|
||||
// The bounds of an ask (novox/hq ADR 0234 §8, ADR 0259 §4).
|
||||
|
||||
+29
-1
@@ -1378,6 +1378,20 @@ type Declaration struct {
|
||||
// what it wrote for it — its resources are absent from the declaration, and absence would
|
||||
// otherwise read as removal.
|
||||
LeftOut []string
|
||||
|
||||
// Generation is the assignment generation this declaration was composed from: the controller's
|
||||
// counter, raised in the same transaction as every change to what is assigned where (novox/hq
|
||||
// issue 234). The sequence orders arrival and cannot tell a later send that carries an older view
|
||||
// of the assignments; this can. A node-engine refuses a declaration composed from a generation
|
||||
// older than the highest it applied — unless it is a put-back — because applying it would
|
||||
// undeclare what the mesh still assigns. Zero is a declaration from a controller that claims none,
|
||||
// and is applied as before.
|
||||
Generation int64
|
||||
|
||||
// PutBack says this declaration is a gate putting a machine back (novox/hq issue 234): it carries
|
||||
// the generation of what it puts back, which may be older than what the machine applied, and is
|
||||
// not refused for it.
|
||||
PutBack bool
|
||||
}
|
||||
|
||||
// LeftOutModuleOf says which left-out module a recorded resource belongs to, if any: its id is the
|
||||
@@ -1542,6 +1556,11 @@ type envelope struct {
|
||||
Epoch int64 `json:"epoch,omitempty"`
|
||||
// LeftOut is optional on the wire too, and absent when nothing was left out (ADR 0163).
|
||||
LeftOut []string `json:"left_out,omitempty"`
|
||||
// Generation and PutBack are optional on the wire too (novox/hq issue 234): absent is a controller
|
||||
// that claims no generation. **An older host refuses these keys**, decoding strictly; a controller
|
||||
// sends them only to a host whose reports say `reads_generation`.
|
||||
Generation int64 `json:"generation,omitempty"`
|
||||
PutBack bool `json:"put_back,omitempty"`
|
||||
}
|
||||
|
||||
func parse(raw []byte, allowActions bool) (*Declaration, error) {
|
||||
@@ -1559,8 +1578,17 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
|
||||
}
|
||||
|
||||
d := &Declaration{Version: env.Version, For: env.For, Adoption: env.Adoption, Sequence: env.Sequence,
|
||||
Epoch: env.Epoch, LeftOut: env.LeftOut}
|
||||
Epoch: env.Epoch, LeftOut: env.LeftOut, Generation: env.Generation, PutBack: env.PutBack}
|
||||
var problems []string
|
||||
if env.Generation < 0 {
|
||||
// Below zero would read as "none claimed" and pass every refusal of an older generation.
|
||||
problems = append(problems, fmt.Sprintf("an assignment generation below zero (%d) is not one the "+
|
||||
"mesh assigns", env.Generation))
|
||||
}
|
||||
if env.PutBack && allowActions {
|
||||
// A put-back is a gate's, sent by the mesh; a carried bundle puts nothing back.
|
||||
problems = append(problems, "a carried bundle says it is a put-back, and only the mesh's gate can say that")
|
||||
}
|
||||
if env.Sequence < 0 || env.Epoch < 0 {
|
||||
// Below zero is no order any controller assigns, and read as "none claimed" it would let the
|
||||
// declaration past every refusal of what is older.
|
||||
|
||||
Vendored
+3
-3
@@ -1,4 +1,4 @@
|
||||
# git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970
|
||||
# git.novox.be/novox/mesh-sdk/go v0.1.14-0.20261010191001-33917f91ac3c
|
||||
## explicit; go 1.22
|
||||
git.novox.be/novox/mesh-sdk/go/asks
|
||||
# github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op
|
||||
@@ -78,7 +78,7 @@ github.com/nats-io/nkeys
|
||||
# github.com/nats-io/nuid v1.0.1
|
||||
## explicit
|
||||
github.com/nats-io/nuid
|
||||
# github.com/novox/mesh-host v0.0.0 => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
# github.com/novox/mesh-host v0.0.0 => git.novox.be/novox/mesh-host v0.0.0-20261011092230-20d5af9b2d5f
|
||||
## explicit; go 1.26.0
|
||||
github.com/novox/mesh-host/internal/declaration
|
||||
github.com/novox/mesh-host/rootsearch
|
||||
@@ -135,4 +135,4 @@ golang.org/x/text/width
|
||||
# golang.org/x/time v0.15.0
|
||||
## explicit; go 1.25.0
|
||||
golang.org/x/time/rate
|
||||
# github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
# github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261011092230-20d5af9b2d5f
|
||||
|
||||
Reference in New Issue
Block a user