Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fe2e5e91b5 | ||
|
|
51db0b5273 | ||
|
|
5c4fa43f8b | ||
|
|
c494ed03a7 | ||
|
|
02c61b983c | ||
|
|
83ce19b9c0 | ||
|
|
7758301444 | ||
|
|
094d3d5bc6 | ||
|
|
ebca7816bf | ||
|
|
9bed7d6398 |
@@ -514,10 +514,6 @@ 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
|
||||
|
||||
@@ -1,222 +0,0 @@
|
||||
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
|
||||
}
|
||||
@@ -1,148 +0,0 @@
|
||||
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, sentRecordOf(ctx, s.declared))
|
||||
return recordSent(ctx, r.inv, s.node, body, s.declared.Builds, s.declared.Epoch)
|
||||
}
|
||||
|
||||
// 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, sentRecord{sender: "a test"})
|
||||
digest, err := recordSent(ctx, inv, "anchor", body, nil, 0)
|
||||
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,8 +420,7 @@ 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,
|
||||
modules: namedModules(plan, composed.LeftOut)}
|
||||
unbound: with.Unbound, foreseen: composed.Foreseen, unplaced: composed.Unplaced}
|
||||
// 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 {
|
||||
@@ -1291,11 +1290,6 @@ 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.
|
||||
|
||||
@@ -288,13 +288,24 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) {
|
||||
"those walks answer them; a newer merge, or `rebuild <module>`, builds again", p.ID)
|
||||
}
|
||||
}
|
||||
// Settled from the build records before anything is judged (novox/hq issue 457): a build of the tier
|
||||
// that failed after the plan did is as failed as the one that failed it. What toRetry says is the
|
||||
// failed set every step below works from.
|
||||
var failed []string
|
||||
if p.Tier < len(p.Tiers) {
|
||||
recorded, byID, err := recordsOfAsked(ctx, inv, &p, p.Tiers[p.Tier])
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
failed = toRetry(&p, recorded, byID)
|
||||
}
|
||||
if err := retryRefusal(p, plans); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if again := unjudgedAtGate(p); len(failedIn(p)) == 0 && len(again) > 0 {
|
||||
if again := unjudgedAtGate(p); len(failed) == 0 && len(again) > 0 {
|
||||
return retryTierWhole(ctx, open, &p, again)
|
||||
}
|
||||
if len(failedIn(p)) == 0 {
|
||||
if len(failed) == 0 {
|
||||
return retryRollouts(ctx, open, &p)
|
||||
}
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
@@ -305,7 +316,6 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) {
|
||||
for _, e := range entries {
|
||||
byName[e.Manifest.Module] = e
|
||||
}
|
||||
failed := failedIn(p)
|
||||
var asked []string
|
||||
for _, m := range failed {
|
||||
askModule(ctx, &p, m, byName)
|
||||
@@ -525,3 +535,16 @@ func retryTierWhole(ctx context.Context, open *stores, p *inventory.Plan, again
|
||||
func sendAgain(s *inventory.PlanModule) {
|
||||
s.First, s.FirstAt, s.Gate, s.GatedBy, s.Previous, s.Why = nil, nil, nil, "", "", ""
|
||||
}
|
||||
|
||||
// toRetry is the modules a retry asks again: every one of the tier that failed, the plan's state
|
||||
// settled from the build records first (novox/hq issue 457). A plan fails on the first failure in its
|
||||
// tier, and an outcome arriving after that finds no open plan to answer — it is kept only in the build
|
||||
// records. Read from the plan alone, a retry asked only the build that failed first, and the records
|
||||
// then failed the plan again on the next: each failed build of a tier took a retry of its own. What
|
||||
// still runs is left asked, and its outcome is the plan's once the retry sets it building.
|
||||
func toRetry(p *inventory.Plan, recorded map[string][]inventory.Build, byID map[string]inventory.Build) []string {
|
||||
if p.Tier < len(p.Tiers) {
|
||||
settleFromRecords(p, p.Tiers[p.Tier], recorded, byID)
|
||||
}
|
||||
return failedIn(*p)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
)
|
||||
|
||||
// A retry asks every failed build of the tier, not one (novox/hq issue 457). Seen 2026-10-11: gitea
|
||||
// and plex both failed in tier 0 while the registry was held still; gitea's failure failed the plan,
|
||||
// and plex's, arriving after, found no open plan and was kept only in the build records. The first
|
||||
// retry asked gitea alone, the records then failed the plan again on plex, and a second retry asked
|
||||
// plex.
|
||||
func TestARetryAsksEveryFailedBuildOfTheTier(t *testing.T) {
|
||||
asked := time.Date(2026, 10, 11, 1, 29, 0, 0, time.UTC)
|
||||
failedAt := asked.Add(2 * time.Minute)
|
||||
p := inventory.Plan{ID: "plan-457", State: inventory.PlanFailed, Tier: 0,
|
||||
Tiers: [][]string{{"gitea", "plex"}}, Note: "gitea failed to build in tier 0",
|
||||
Modules: map[string]*inventory.PlanModule{
|
||||
"gitea": {State: "failed", AskedAt: &asked, Build: "build-gitea", Why: "cannot reach the registry"},
|
||||
"plex": {State: "asked", AskedAt: &asked, Build: "build-plex"},
|
||||
}}
|
||||
byID := map[string]inventory.Build{
|
||||
"build-plex": {ID: "build-plex", Module: "plex", At: failedAt, Failed: "cannot reach the registry"},
|
||||
}
|
||||
|
||||
if got := toRetry(&p, nil, byID); !reflect.DeepEqual(got, []string{"gitea", "plex"}) {
|
||||
t.Fatalf("a retry of a tier where gitea and plex failed asks %v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A build of the tier still running when the plan failed is not asked again: its outcome is the plan's
|
||||
// once the retry sets it building, and one that is recorded built is taken as built.
|
||||
func TestARetryLeavesABuildThatRunsOrWorked(t *testing.T) {
|
||||
asked := time.Date(2026, 10, 11, 1, 29, 0, 0, time.UTC)
|
||||
builtAt := asked.Add(3 * time.Minute)
|
||||
p := inventory.Plan{ID: "plan-457", State: inventory.PlanFailed, Tier: 0,
|
||||
Tiers: [][]string{{"a", "b", "c"}},
|
||||
Modules: map[string]*inventory.PlanModule{
|
||||
"a": {State: "failed", AskedAt: &asked, Build: "build-a", Why: "broken"},
|
||||
"b": {State: "asked", AskedAt: &asked, Build: "build-b"},
|
||||
"c": {State: "asked", AskedAt: &asked, Build: "build-c"},
|
||||
}}
|
||||
byID := map[string]inventory.Build{"build-c": {ID: "build-c", Module: "c", At: builtAt, Commit: "c0ffee"}}
|
||||
|
||||
if got := toRetry(&p, nil, byID); !reflect.DeepEqual(got, []string{"a"}) {
|
||||
t.Fatalf("asks %v, want a alone", got)
|
||||
}
|
||||
if s := p.Modules["c"]; s.State != "built" || s.Commit != "c0ffee" {
|
||||
t.Errorf("c, recorded built, is %+v", s)
|
||||
}
|
||||
if s := p.Modules["b"]; s.State != "asked" {
|
||||
t.Errorf("b, still building, is %+v", s)
|
||||
}
|
||||
}
|
||||
@@ -356,16 +356,10 @@ 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"`
|
||||
Sequence int64 `json:"sequence"`
|
||||
Generation int64 `json:"generation"`
|
||||
Epoch uint64 `json:"epoch"`
|
||||
}
|
||||
_ = json.Unmarshal(raw, &carried)
|
||||
// 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 {
|
||||
if _, err := recordSent(ctx, inv, node, raw, nil, carried.Epoch); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||
@@ -804,7 +798,7 @@ func composeEach(names []string, allot func(node string) (order, error),
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
}
|
||||
numbered.stamp(&declared)
|
||||
declared.Sequence, declared.Epoch = numbered.sequence, numbered.epoch
|
||||
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,
|
||||
@@ -1022,8 +1016,7 @@ 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,
|
||||
sentRecordOf(ctx, s.declared))
|
||||
digest, err := recordSent(ctx, b.open.inventory, s.node, body, s.declared.Builds, s.declared.Epoch)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -1235,7 +1228,7 @@ func sendToEach(ctx context.Context, open *stores, names []string) ([]string, er
|
||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||
continue
|
||||
}
|
||||
numbered.stamp(&declared)
|
||||
declared.Sequence, declared.Epoch = numbered.sequence, numbered.epoch
|
||||
reportLeftOut(name, declared)
|
||||
sending = append(sending, readyNode{name, declared})
|
||||
}
|
||||
@@ -1400,11 +1393,6 @@ 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
|
||||
@@ -1543,12 +1531,6 @@ 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
|
||||
@@ -1562,7 +1544,6 @@ 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 {
|
||||
@@ -1572,24 +1553,11 @@ 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, generation: generation, readsGeneration: readsGeneration,
|
||||
acting: acting}, nil
|
||||
return order{sequence: seq, epoch: epoch}, nil
|
||||
}
|
||||
|
||||
// recordSent writes down what a machine was just sent, and returns the digest.
|
||||
@@ -1604,7 +1572,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, sent sentRecord) (string, error) {
|
||||
builds map[string]string, epoch uint64) (string, error) {
|
||||
kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
||||
defer cancel()
|
||||
record, err := inv.NodeByName(kept, node)
|
||||
@@ -1612,21 +1580,7 @@ func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body
|
||||
return "", err
|
||||
}
|
||||
digest := digestOf(body)
|
||||
// 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:
|
||||
if err := inv.RecordSentUnder(kept, record.ID, digest, builds, epoch); err != nil {
|
||||
return "", err
|
||||
}
|
||||
// And what it was, summarised, so the next push can be compared with it (novox/hq ADR 0217).
|
||||
|
||||
@@ -649,26 +649,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
// the controller in its first tier: the build that produced the new one is recorded, and the
|
||||
// plan never hears it. The record is the fact; a build recorded after the ask is that tier's
|
||||
// outcome, whoever was listening.
|
||||
recorded := map[string][]inventory.Build{}
|
||||
byID := map[string]inventory.Build{}
|
||||
for _, m := range tier {
|
||||
if s := p.Modules[m]; s != nil && s.State == "asked" {
|
||||
builds, err := inv.Builds(ctx, m, 5)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
recorded[m] = builds
|
||||
// Its own ask's record, by id — found even when the outcome named no module (ADR 0219).
|
||||
if s.Build != "" {
|
||||
b, found, err := inv.BuildByID(ctx, s.Build)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if found {
|
||||
byID[s.Build] = b
|
||||
}
|
||||
}
|
||||
}
|
||||
recorded, byID, err := recordsOfAsked(ctx, inv, p, tier)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if settleFromRecords(p, tier, recorded, byID) {
|
||||
return true, nil
|
||||
@@ -1776,6 +1759,34 @@ func splitList(s string) []string {
|
||||
return out
|
||||
}
|
||||
|
||||
// recordsOfAsked is what settleFromRecords reads: for every module of the tier still `asked`, its last
|
||||
// builds and the record of its own ask, by id.
|
||||
func recordsOfAsked(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan, tier []string) (
|
||||
map[string][]inventory.Build, map[string]inventory.Build, error) {
|
||||
recorded := map[string][]inventory.Build{}
|
||||
byID := map[string]inventory.Build{}
|
||||
for _, m := range tier {
|
||||
if s := p.Modules[m]; s != nil && s.State == "asked" {
|
||||
builds, err := inv.Builds(ctx, m, 5)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
recorded[m] = builds
|
||||
// Its own ask's record, by id — found even when the outcome named no module (ADR 0219).
|
||||
if s.Build != "" {
|
||||
b, found, err := inv.BuildByID(ctx, s.Build)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
if found {
|
||||
byID[s.Build] = b
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return recorded, byID, nil
|
||||
}
|
||||
|
||||
// settleFromRecords marks every module of the tier still `asked` built — or failed — from a build
|
||||
// recorded after it was asked, and says whether it changed anything (novox/hq 04-ISSUES/214).
|
||||
// Newest first, as Builds answers: the first record after the ask is the outcome of that ask.
|
||||
|
||||
@@ -27,19 +27,6 @@ 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
|
||||
@@ -107,9 +94,6 @@ 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,20 +226,6 @@ 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,14 +178,6 @@ 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,11 +196,6 @@ 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,10 +55,6 @@ 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
|
||||
@@ -368,7 +364,6 @@ 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.14-0.20261010191001-33917f91ac3c
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970
|
||||
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-20261011092230-20d5af9b2d5f
|
||||
replace github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
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=
|
||||
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=
|
||||
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,117 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// A tool call names its caller (novox/hq issue 365, by ADR 0259 §3's precedent for asks).
|
||||
//
|
||||
// mesh.mod.<module>.call.<tool>[.<node>].<caller>
|
||||
// mesh.seat.<seat>.call.<verb>[.<node>].<caller>
|
||||
//
|
||||
// The last token is the bus user that published the call, and each user is granted these subjects with its own
|
||||
// name there and no other — the caller is a fact the server enforces, and the tool runtime hands it to the
|
||||
// module from the subject a call arrived on (mesh-tools node-tools internal/bus caller.go holds the same shape).
|
||||
//
|
||||
// **A kind of its own, `call`, not the `tool` subject with a token appended.** Every grant to call a tool is a
|
||||
// wildcard over the `tool` kind — `.tool.<t>.*` for the instance on any machine, `mesh.mod.*.tool.>` for every
|
||||
// tool — and either would match a `tool` subject with any caller appended, so a caller could name another.
|
||||
// Under `call` nothing is granted but the caller-named subjects. No stream's filter covers it: a tool call is
|
||||
// never persisted (design 25 §3).
|
||||
//
|
||||
// **Derived from the `tool` grants, in one place** (callerNamed): wherever a principal may call or answer a
|
||||
// tool, it may call it in its own name, or answer it naming any caller — so no kind of principal is left
|
||||
// unable to call in its own name, and a new grant to call is one in both shapes by construction.
|
||||
//
|
||||
// The `tool` grants stay beside these for one release, so every caller and every runtime moves without a gap:
|
||||
// their retirement is hq issue 464.
|
||||
const CallKind = "call"
|
||||
|
||||
var userName = regexp.MustCompile(`^[A-Za-z0-9_-]+(\.[A-Za-z0-9_-]+)*$`)
|
||||
|
||||
// CallerToken is a bus user's name as the last token of a call it publishes: each dot written `~`, which no part
|
||||
// of a user's name may hold (safeSubject), so the token names that user and no other. "" for a name that is not
|
||||
// a user's.
|
||||
func CallerToken(user string) string {
|
||||
if !userName.MatchString(user) {
|
||||
return ""
|
||||
}
|
||||
return strings.ReplaceAll(user, ".", "~")
|
||||
}
|
||||
|
||||
// toolGrant splits a grant on the `tool` kind — `mesh.mod.<m>.tool.<rest>` or `mesh.seat.<s>.tool.<rest>`, any
|
||||
// part possibly a wildcard — at the kind; ok is false for any other subject.
|
||||
func toolGrant(subject string) (head, rest string, ok bool) {
|
||||
parts := strings.SplitN(subject, ".", 5)
|
||||
if len(parts) != 5 || parts[0] != "mesh" || (parts[1] != "mod" && parts[1] != "seat") || parts[3] != "tool" ||
|
||||
parts[2] == "" || parts[4] == "" {
|
||||
return "", "", false
|
||||
}
|
||||
return parts[0] + "." + parts[1] + "." + parts[2], parts[4], true
|
||||
}
|
||||
|
||||
// CalledSubject is where a call to a tool subject goes naming its caller; "" when the subject is not a tool's
|
||||
// or the user is not a bus user.
|
||||
func CalledSubject(subject, user string) string {
|
||||
head, rest, ok := toolGrant(subject)
|
||||
token := CallerToken(user)
|
||||
if !ok || token == "" || strings.ContainsAny(rest, "*>") {
|
||||
return ""
|
||||
}
|
||||
return head + "." + CallKind + "." + rest + "." + token
|
||||
}
|
||||
|
||||
// CalledPattern is what a holder answering a tool subject also subscribes, to hear the calls that name their
|
||||
// caller: the same address under the `call` kind, any caller last. "" for a subject that is not a tool's.
|
||||
func CalledPattern(subject string) string {
|
||||
head, rest, ok := toolGrant(subject)
|
||||
if !ok || strings.ContainsAny(rest, "*>") {
|
||||
return ""
|
||||
}
|
||||
return head + "." + CallKind + "." + rest + ".*"
|
||||
}
|
||||
|
||||
// calledPublish is a grant to call on the `tool` kind, as the same grant in the caller's own name: its last
|
||||
// token the caller's, and nothing that reaches past it. A `>` is every tail a call carries — the tool alone or
|
||||
// the tool and the machine — so it becomes both, each ending in the caller.
|
||||
func calledPublish(grant, token string) []string {
|
||||
head, rest, ok := toolGrant(grant)
|
||||
if !ok || token == "" {
|
||||
return nil
|
||||
}
|
||||
base := head + "." + CallKind + "."
|
||||
switch {
|
||||
case rest == ">":
|
||||
return []string{base + "*." + token, base + "*.*." + token}
|
||||
case strings.HasSuffix(rest, ".>"):
|
||||
return []string{base + strings.TrimSuffix(rest, ">") + "*." + token}
|
||||
}
|
||||
return []string{base + rest + "." + token}
|
||||
}
|
||||
|
||||
// calledSubscribe is a grant to answer on the `tool` kind, as the same grant for the calls naming any caller.
|
||||
func calledSubscribe(grant string) []string {
|
||||
head, rest, ok := toolGrant(grant)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
base := head + "." + CallKind + "."
|
||||
if strings.HasSuffix(rest, ">") {
|
||||
return []string{base + rest}
|
||||
}
|
||||
return []string{base + rest + ".*"}
|
||||
}
|
||||
|
||||
// callerNamed adds, beside a principal's grants on the `tool` kind, the same grants under `call`: to publish in
|
||||
// its own name, to subscribe naming anybody.
|
||||
func callerNamed(user string, pub, sub []string) ([]string, []string) {
|
||||
token := CallerToken(user)
|
||||
for _, g := range pub {
|
||||
pub = append(pub, calledPublish(g, token)...)
|
||||
}
|
||||
for _, g := range sub {
|
||||
sub = append(sub, calledSubscribe(g)...)
|
||||
}
|
||||
return unique(pub), unique(sub)
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats-server/v2/server"
|
||||
"github.com/nats-io/nats.go"
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
)
|
||||
|
||||
// **On a real server, as composed** (novox/hq issue 365): a machine's runtime calls in its own name, and the
|
||||
// server refuses it a call naming another — the grant, not the runtime, is what makes the caller a fact.
|
||||
func TestAServerComposedFromTheGrantsRefusesACallNamingAnother(t *testing.T) {
|
||||
hash, err := bcrypt.GenerateFromPassword([]byte("pw"), bcrypt.MinCost)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ledger := Seat{Name: "issue-tracker", Scope: "mesh", Serves: []string{"open"}}
|
||||
accounts, err := ComposeAccounts([]Principal{
|
||||
{Kind: KindNodeTools, Node: "novox", Module: RuntimeModule, PasswordHash: string(hash),
|
||||
Carries: []Declared{{Module: "mesh-issues", Holds: []Seat{ledger}}}},
|
||||
{Kind: KindNodeTools, Node: "shanks", Module: RuntimeModule, PasswordHash: string(hash)},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
conf := filepath.Join(t.TempDir(), "bus.conf")
|
||||
if err := os.WriteFile(conf, []byte("listen: 127.0.0.1:-1\njetstream { store_dir: "+
|
||||
`"`+t.TempDir()+`"`+" }\n"+accounts), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
opts, err := server.ProcessConfigFile(conf)
|
||||
if err != nil {
|
||||
t.Fatalf("the composed accounts do not parse: %v", err)
|
||||
}
|
||||
opts.NoLog, opts.NoSigs = true, true
|
||||
s, err := server.NewServer(opts)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
go s.Start()
|
||||
if !s.ReadyForConnections(10 * time.Second) {
|
||||
t.Fatal("the bus did not come up")
|
||||
}
|
||||
defer s.Shutdown()
|
||||
|
||||
holder, err := nats.Connect(s.ClientURL(), nats.UserInfo("novox.node-tools", "pw"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer holder.Close()
|
||||
heard := make(chan string, 4)
|
||||
sub, err := holder.Subscribe("mesh.seat.issue-tracker.call.open.*", func(m *nats.Msg) {
|
||||
heard <- m.Subject
|
||||
_ = m.Respond([]byte(`{"result":{}}`))
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = holder.Flush()
|
||||
if !sub.IsValid() {
|
||||
t.Fatal("the holder may not hear the caller-named calls to its seat")
|
||||
}
|
||||
|
||||
refusals := make(chan error, 4)
|
||||
caller, err := nats.Connect(s.ClientURL(), nats.UserInfo("shanks.node-tools", "pw"),
|
||||
nats.CustomInboxPrefix("_INBOX.shanks.node-tools"),
|
||||
nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { refusals <- err }))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer caller.Close()
|
||||
if _, err := caller.Request("mesh.seat.issue-tracker.call.open.shanks~node-tools", []byte(`{}`), 3*time.Second); err != nil {
|
||||
t.Fatalf("a call in the caller's own name was not answered: %v", err)
|
||||
}
|
||||
if got := <-heard; got != "mesh.seat.issue-tracker.call.open.shanks~node-tools" {
|
||||
t.Fatalf("heard %s", got)
|
||||
}
|
||||
if err := caller.Publish("mesh.seat.issue-tracker.call.open.novox~node-tools", []byte(`{}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_ = caller.Flush()
|
||||
select {
|
||||
case err := <-refusals:
|
||||
if !errors.Is(err, nats.ErrPermissionViolation) {
|
||||
t.Errorf("the server said %v, want a permissions violation", err)
|
||||
}
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Error("the server did not refuse a call naming another caller")
|
||||
}
|
||||
select {
|
||||
case got := <-heard:
|
||||
t.Errorf("a call naming another reached the holder: %s", got)
|
||||
case <-time.After(200 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// **A tool call names its caller, and the bus lets each user name itself alone** (novox/hq issue 365, by ADR
|
||||
// 0259 §3's precedent for asks): wherever a principal may call a tool, it may call it on the caller-named
|
||||
// subject — kind `call`, its own bus user last, dots written `~` — and on no subject naming anybody else.
|
||||
func TestEachCredentialMayCallOnlyInItsOwnName(t *testing.T) {
|
||||
ledger := Seat{Name: "node-desk", Scope: "node", Serves: []string{"who"}}
|
||||
tracker := Seat{Name: "issue-tracker", Scope: "mesh", Serves: []string{"open"}}
|
||||
for _, c := range []struct {
|
||||
p Principal
|
||||
may []string
|
||||
mayNot []string
|
||||
subject []string // what it subscribes, to answer calls naming any caller
|
||||
}{
|
||||
{p: Principal{Kind: KindPerson, Module: "jochen", Invokes: []string{"*"}},
|
||||
may: []string{"mesh.mod.ledger.call.who.person~jochen", "mesh.mod.ledger.call.who.anchor.person~jochen",
|
||||
"mesh.seat.issue-tracker.call.open.person~jochen", "mesh.seat.node-desk.call.who.anchor.person~jochen"},
|
||||
mayNot: []string{"mesh.mod.ledger.call.who.shanks~node-tools", "mesh.mod.ledger.call.who.anchor.controller",
|
||||
"mesh.seat.issue-tracker.call.open.person~somebody", "mesh.mod.ledger.call.who.person"}},
|
||||
{p: Principal{Kind: KindNodeTools, Node: "shanks", Module: RuntimeModule,
|
||||
Carries: []Declared{{Module: "ledger", Holds: []Seat{ledger, tracker}}}},
|
||||
may: []string{"mesh.mod.ledger.call.who.shanks~node-tools", "mesh.seat.issue-tracker.call.open.shanks~node-tools"},
|
||||
mayNot: []string{"mesh.mod.ledger.call.who.novox~node-tools", "mesh.seat.issue-tracker.call.open.controller"},
|
||||
subject: []string{"mesh.mod.ledger.call.who.novox~node-tools", "mesh.seat.node-desk.call.who.shanks.person~jochen", "mesh.seat.issue-tracker.call.open.controller"}},
|
||||
{p: Principal{Kind: KindModule, Node: "two", Module: "shop", Invokes: []string{"ledger.who", "seat:issue-tracker.open"}},
|
||||
may: []string{"mesh.mod.ledger.call.who.two~shop", "mesh.mod.ledger.call.who.anchor.two~shop",
|
||||
"mesh.seat.issue-tracker.call.open.two~shop"},
|
||||
mayNot: []string{"mesh.mod.ledger.call.other.two~shop", "mesh.mod.ledger.call.who.one~shop",
|
||||
"mesh.seat.issue-tracker.call.open.one~telegram"}},
|
||||
{p: Principal{Kind: KindNode, Node: "one", Checks: []string{"ledger.health"}},
|
||||
may: []string{"mesh.mod.ledger.call.health.one.node~one"},
|
||||
mayNot: []string{"mesh.mod.ledger.call.health.two.node~one", "mesh.mod.ledger.call.health.one.node~two"}},
|
||||
{p: Principal{Kind: KindModule, Node: "one", Module: "ledger", Holds: []Seat{ledger, tracker}},
|
||||
subject: []string{"mesh.mod.ledger.call.who.person~jochen", "mesh.mod.ledger.call.who.one.controller",
|
||||
"mesh.seat.node-desk.call.who.one.two~shop", "mesh.seat.issue-tracker.call.open.two~shop"}},
|
||||
} {
|
||||
perms, err := PermissionsFor(c.p)
|
||||
if err != nil {
|
||||
t.Fatalf("%s: %v", c.p.Username(), err)
|
||||
}
|
||||
for _, s := range c.may {
|
||||
if !MayPublish(perms, s) {
|
||||
t.Errorf("%s may not call %s, in its own name", c.p.Username(), s)
|
||||
}
|
||||
}
|
||||
for _, s := range c.mayNot {
|
||||
if MayPublish(perms, s) {
|
||||
t.Errorf("%s may call %s, naming somebody else", c.p.Username(), s)
|
||||
}
|
||||
}
|
||||
for _, s := range c.subject {
|
||||
if !MaySubscribe(perms, s) {
|
||||
t.Errorf("%s does not hear %s, a call to what it serves", c.p.Username(), s)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The subjects that name no caller stay granted beside the caller-named ones for one release, so a caller
|
||||
// moves over without a gap (their retirement: hq issue 464).
|
||||
func TestTheSubjectsThatNameNoCallerStayGrantedForOneRelease(t *testing.T) {
|
||||
perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "two", Module: "shop", Invokes: []string{"ledger.who"}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, s := range []string{"mesh.mod.ledger.tool.who", "mesh.mod.ledger.tool.who.anchor"} {
|
||||
if !MayPublish(perms, s) {
|
||||
t.Errorf("the old subject %s is no longer granted", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The desk's hidden prompt is the controller's alone on the caller-named subjects too (the review of
|
||||
// 2026-10-09, M4): a grant of every tool does not reach it there either.
|
||||
func TestTheDesksPromptIsTheControllersAloneUnderCallToo(t *testing.T) {
|
||||
for _, p := range []Principal{{Kind: KindPerson, Module: "jochen", Invokes: []string{"*"}},
|
||||
{Kind: KindNodeTools, Node: "shanks", Module: RuntimeModule}} {
|
||||
perms, err := PermissionsFor(p)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
token := CallerToken(p.Username())
|
||||
for _, s := range []string{"mesh.seat.node-launcher.call.secret.shanks." + token,
|
||||
"mesh.mod.shell.call.node-launcher.secret.shanks." + token, "mesh.mod.shell.call.node-launcher.secret." + token} {
|
||||
if MayPublish(perms, s) {
|
||||
t.Errorf("%s may publish %s, the desk's hidden prompt", p.Username(), s)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The controller's grants are the installer's first user list too, so they move with hq issue 464: this release
|
||||
// it calls and serves on the subjects that name no caller alone, and nobody may call in its name.
|
||||
func TestTheControllerKeepsTheSubjectsThatNameNoCallerThisRelease(t *testing.T) {
|
||||
perms, err := PermissionsFor(Principal{Kind: KindController})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, s := range append(perms.Publish, perms.Subscribe...) {
|
||||
if strings.Contains(s, "."+CallKind+".") {
|
||||
t.Errorf("the controller is granted %s, which the installer's first user list does not carry", s)
|
||||
}
|
||||
}
|
||||
for _, p := range []Principal{{Kind: KindPerson, Module: "jochen", Invokes: []string{"*"}},
|
||||
{Kind: KindNodeTools, Node: "shanks", Module: RuntimeModule}} {
|
||||
perms, err := PermissionsFor(p)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if MayPublish(perms, "mesh.mod.ledger.call.who.controller") || MayPublish(perms, "mesh.seat.mesh-controller.call.status.controller") {
|
||||
t.Errorf("%s may call in the controller's name", p.Username())
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -42,8 +42,9 @@ func TestInvokingGrantsNothingButTheCall(t *testing.T) {
|
||||
if strings.Contains(p, ".event.") {
|
||||
t.Errorf("a module that only invokes may publish %q, an event it never declared", p)
|
||||
}
|
||||
// A role's tools are tools (ADR 0132); a role's work queue and events are not.
|
||||
if strings.HasPrefix(p, "mesh.seat.") && !strings.Contains(p, ".tool.") {
|
||||
// A role's tools are tools (ADR 0132), named by their caller or not (novox/hq issue 365); a role's work
|
||||
// queue and events are not.
|
||||
if strings.HasPrefix(p, "mesh.seat.") && !strings.Contains(p, ".tool.") && !strings.Contains(p, ".call.") {
|
||||
t.Errorf("a module that only invokes may publish %q, a seat it neither holds nor uses", p)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -212,6 +212,10 @@ func ControllerOnly() []string {
|
||||
for _, base := range []string{"mesh.seat." + v.Seat + ".tool." + v.Verb, "mesh.mod.*.tool." + v.Seat + "." + v.Verb} {
|
||||
out = append(out, base, base+".*")
|
||||
}
|
||||
// And where a call names its caller (novox/hq issue 365): any machine, any caller.
|
||||
for _, base := range []string{"mesh.seat." + v.Seat + "." + CallKind + "." + v.Verb, "mesh.mod.*." + CallKind + "." + v.Seat + "." + v.Verb} {
|
||||
out = append(out, base+".>")
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -875,6 +879,19 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
pub = append(pub, "$JS.ACK."+consumerStream(p)+"."+consumerDurable(p)+".>")
|
||||
}
|
||||
|
||||
// **And every tool grant again, naming its caller** (novox/hq issue 365): a call in this principal's own
|
||||
// name, and an answer to a call naming anybody. The `tool` grants above stay for one release beside these,
|
||||
// so callers and runtimes move without a gap; their retirement is hq issue 464.
|
||||
//
|
||||
// **Not the controller's, this release.** Its grants are also the installer's first user list
|
||||
// (TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose), judged against the node-engine the mesh
|
||||
// runs, so a change to them waits on a node-engine delivery. It calls and serves on the subjects that name
|
||||
// no caller meanwhile — a caller-named call to its seat reaches nobody and is asked again on those at once —
|
||||
// and moves with issue 464.
|
||||
if p.Kind != KindController {
|
||||
pub, sub = callerNamed(p.Username(), pub, sub)
|
||||
}
|
||||
|
||||
sort.Strings(pub)
|
||||
sort.Strings(sub)
|
||||
// One writer per piece of state (novox/hq to-be 45 §1): a grant that would make a second is
|
||||
|
||||
@@ -240,9 +240,9 @@ func TestAPersonReachesNothingButTools(t *testing.T) {
|
||||
perms, _ := PermissionsFor(Principal{Kind: KindPerson, Module: "jo",
|
||||
Invokes: []string{"*"}, PasswordHash: "x"})
|
||||
for _, p := range perms.Publish {
|
||||
// A tool call, or asking what answers (novox/hq ADR 0197) — a question every service
|
||||
// answers about itself, which claims nothing and controls nothing.
|
||||
if !strings.Contains(p, ".tool.") && !strings.HasPrefix(p, "$SRV.") {
|
||||
// A tool call — named by its caller or not (novox/hq issue 365) — or asking what answers (novox/hq
|
||||
// ADR 0197), a question every service answers about itself, which claims nothing and controls nothing.
|
||||
if !strings.Contains(p, ".tool.") && !strings.Contains(p, ".call.") && !strings.HasPrefix(p, "$SRV.") {
|
||||
t.Errorf("a person may publish %q, which is not a tool call", p)
|
||||
}
|
||||
}
|
||||
|
||||
+3
-3
@@ -43,17 +43,17 @@ accounts {
|
||||
} }
|
||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "$SRV.STATS", "$SRV.STATS.telegram", "$SRV.STATS.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "$SRV.STATS", "$SRV.STATS.telegram", "$SRV.STATS.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.call.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.audit", "$SRV.INFO.audit.>", "$SRV.PING", "$SRV.PING.audit", "$SRV.PING.audit.>", "$SRV.STATS", "$SRV.STATS.audit", "$SRV.STATS.audit.>", "_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.audit", "$SRV.INFO.audit.>", "$SRV.PING", "$SRV.PING.audit", "$SRV.PING.audit.>", "$SRV.STATS", "$SRV.STATS.audit", "$SRV.STATS.audit.>", "_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.call.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
|
||||
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.call.>", "mesh.mod.shop.tool.>"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
]
|
||||
|
||||
+107
-6
@@ -7,8 +7,10 @@ import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/fs"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path"
|
||||
@@ -159,17 +161,18 @@ func build(ctx context.Context, run Runner, publish Publisher,
|
||||
}
|
||||
// The credential is a file git reads, never an argument: a URL carrying a password in argv
|
||||
// would be readable by anything that can list processes for as long as a clone runs.
|
||||
credentials := ""
|
||||
if forge.URL != "" {
|
||||
credentials = filepath.Join(workspace, "git-credentials")
|
||||
if err := os.WriteFile(credentials, []byte(forge.URL+"\n"), 0o600); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
credentials, forget, err := storeCredential(workspace, forge)
|
||||
if err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
defer forget()
|
||||
tree := filepath.Join(workspace, "source")
|
||||
if err := os.RemoveAll(tree); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
// Removed when the build ends, whatever happens: the tree holds the .npmrc a build may be handed, in the
|
||||
// workspace a later check's container mounts as its HOME (novox/hq issue 462).
|
||||
defer os.RemoveAll(tree)
|
||||
// A fresh clone every time rather than a fetch into a tree that is already there. A build
|
||||
// that reuses a working tree can succeed because of something a previous build left behind,
|
||||
// and that is a build nobody can reproduce.
|
||||
@@ -177,6 +180,9 @@ func build(ctx context.Context, run Runner, publish Publisher,
|
||||
say("clone", "FAILED: %v", err)
|
||||
return Result{}, fmt.Errorf("cannot clone %s: %w", repository, err)
|
||||
}
|
||||
if err := recordedWithoutUserinfo(ctx, run, tree, repository); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
say("clone", "done")
|
||||
if ref != "" {
|
||||
if _, err := run(ctx, tree, "git", "checkout", "--quiet", ref); err != nil {
|
||||
@@ -248,6 +254,8 @@ func build(ctx context.Context, run Runner, publish Publisher,
|
||||
if err := os.WriteFile(npmrcPath, []byte(content), 0o600); err != nil {
|
||||
return Result{}, fmt.Errorf("cannot write the package-registry credential for the build: %w", err)
|
||||
}
|
||||
// Only for as long as the build: a later check's container mounts this workspace (novox/hq issue 462).
|
||||
defer os.Remove(npmrcPath)
|
||||
say("packages", "resolving %s from the mesh's package registry", npmrc.Scope)
|
||||
src.notPinned("it resolves packages from the mesh's registry at build time")
|
||||
}
|
||||
@@ -441,6 +449,9 @@ func contextFrom(ctx context.Context, run Runner, workspace, artifact, credentia
|
||||
if _, err := run(ctx, workspace, "git", cloneWith(credentials, "clone", "--quiet", url, dir)...); err != nil {
|
||||
return "", fmt.Errorf("cannot clone %s: %w", url, err)
|
||||
}
|
||||
if err := recordedWithoutUserinfo(ctx, run, dir, url); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if from.Ref != "" {
|
||||
if _, err := run(ctx, dir, "git", "checkout", "--quiet", from.Ref); err != nil {
|
||||
return "", fmt.Errorf("%s has no %s: %w", from.Repository, from.Ref, err)
|
||||
@@ -467,6 +478,96 @@ func contextURL(from catalogue.ArtifactContext, seats map[string]string) (string
|
||||
return strings.TrimRight(base, "/") + "/" + strings.TrimSuffix(strings.Trim(from.Repository, "/"), ".git") + ".git", nil
|
||||
}
|
||||
|
||||
// storeCredential writes the forge credential where git's credential store reads it, and the function that
|
||||
// removes it again; "" and nothing to remove when the builder holds none.
|
||||
//
|
||||
// **Never inside the workspace** (novox/hq issue 462): a merge check's toolchain container mounts the
|
||||
// workspace as its HOME, so a credential kept there — even one a build left behind — is readable by any pull
|
||||
// request's merge-check.sh, and what it prints is kept on the bus. So it lives in a directory of its own
|
||||
// outside the workspace, made private, and a credential an older builder left in the workspace is removed.
|
||||
func storeCredential(workspace string, forge GitCredential) (string, func(), error) {
|
||||
if err := os.Remove(filepath.Join(workspace, "git-credentials")); err != nil && !errors.Is(err, os.ErrNotExist) {
|
||||
return "", nil, fmt.Errorf("a credential left in the workspace cannot be removed: %w", err)
|
||||
}
|
||||
if forge.URL == "" {
|
||||
return "", func() {}, nil
|
||||
}
|
||||
dir, err := os.MkdirTemp("", "mesh-forge-credential-")
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
forget := func() { os.RemoveAll(dir) }
|
||||
if withinDir(workspace, dir) {
|
||||
forget()
|
||||
return "", nil, fmt.Errorf("the temporary directory %s is inside the workspace %s, which a check's "+
|
||||
"container mounts: the forge credential is not written where it could read it", dir, workspace)
|
||||
}
|
||||
path := filepath.Join(dir, "git-credentials")
|
||||
if err := os.WriteFile(path, []byte(forge.URL+"\n"), 0o600); err != nil {
|
||||
forget()
|
||||
return "", nil, err
|
||||
}
|
||||
return path, forget, nil
|
||||
}
|
||||
|
||||
// recordedWithoutUserinfo makes a clone record the URL it came from without userinfo: git keeps it as given in
|
||||
// the clone's .git/config, inside the workspace a check's container mounts (novox/hq issue 462).
|
||||
func recordedWithoutUserinfo(ctx context.Context, run Runner, clone, repository string) error {
|
||||
bare, carried := withoutUserinfo(repository)
|
||||
if !carried {
|
||||
return nil
|
||||
}
|
||||
if _, err := run(ctx, clone, "git", "remote", "set-url", "origin", bare); err != nil {
|
||||
return fmt.Errorf("the clone of %s keeps its credential: %w", bare, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// forgetLeftCredentials removes what an earlier build or an older builder left in the workspace that holds
|
||||
// a credential: the forge's git-credentials, and every .npmrc a build wrote into a tree it cloned. Run
|
||||
// before a check, whose container mounts the workspace as its HOME (novox/hq issue 462). The Go caches are
|
||||
// not walked: no build writes a credential there, and they hold more files than everything else.
|
||||
func forgetLeftCredentials(workspace string) error {
|
||||
if err := os.Remove(filepath.Join(workspace, "git-credentials")); err != nil && !errors.Is(err, os.ErrNotExist) {
|
||||
return err
|
||||
}
|
||||
return filepath.WalkDir(workspace, func(p string, d fs.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
if d.IsDir() && p != workspace && (d.Name() == "go-cache" || d.Name() == "go-modules") &&
|
||||
filepath.Dir(p) == workspace {
|
||||
return filepath.SkipDir
|
||||
}
|
||||
if !d.IsDir() && d.Name() == ".npmrc" {
|
||||
if err := os.Remove(p); err != nil && !errors.Is(err, os.ErrNotExist) {
|
||||
return fmt.Errorf("a package-registry credential left at %s cannot be removed: %w", p, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
// withinDir is whether path is dir or under it, both made absolute and resolved.
|
||||
func withinDir(dir, path string) bool {
|
||||
d, err1 := filepath.Abs(dir)
|
||||
p, err2 := filepath.Abs(path)
|
||||
if err1 != nil || err2 != nil {
|
||||
return true
|
||||
}
|
||||
if r, err := filepath.EvalSymlinks(d); err == nil {
|
||||
d = r
|
||||
}
|
||||
if r, err := filepath.EvalSymlinks(p); err == nil {
|
||||
p = r
|
||||
}
|
||||
rel, err := filepath.Rel(d, p)
|
||||
return err != nil || rel == "." || filepath.IsLocal(rel)
|
||||
}
|
||||
|
||||
// cloneWith is a git invocation that may offer a stored credential.
|
||||
//
|
||||
// The first `-c credential.helper=` clears every helper the environment might carry, so exactly
|
||||
|
||||
@@ -438,30 +438,34 @@ func TestABuildOffersTheForgesCredentialThroughGitsOwnStore(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stored := filepath.Join(workspace, "git-credentials")
|
||||
clone := r.ran[0]
|
||||
if !strings.Contains(clone, "credential.helper=store --file="+stored) {
|
||||
t.Fatalf("the clone does not name the credential store: %s", clone)
|
||||
}
|
||||
stored := storeNamedIn(t, clone)
|
||||
for _, line := range r.ran {
|
||||
if strings.Contains(line, "sw0rdfi5h") {
|
||||
t.Fatalf("the secret is in a command line, readable by anything that can list processes: %s", line)
|
||||
}
|
||||
}
|
||||
raw, err := os.ReadFile(stored)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
// Outside the workspace a check's container mounts, and gone once the build is (novox/hq issue 462); what
|
||||
// it held while git read it is checked with the check's clones (TestACheckContainerSeesNoForgeCredential).
|
||||
if withinDir(workspace, stored) {
|
||||
t.Fatalf("the credential store %s is inside the workspace %s", stored, workspace)
|
||||
}
|
||||
if strings.TrimSpace(string(raw)) != "http://mesh_novox_builder:sw0rdfi5h@forge.invalid:20000" {
|
||||
t.Fatalf("the store does not hold the credential as given: %q", raw)
|
||||
if _, err := os.Stat(stored); !os.IsNotExist(err) {
|
||||
t.Fatalf("the credential store outlives the build: %v", err)
|
||||
}
|
||||
info, err := os.Stat(stored)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
if _, err := os.Stat(filepath.Join(workspace, "git-credentials")); !os.IsNotExist(err) {
|
||||
t.Fatal("a credential file is left in the workspace")
|
||||
}
|
||||
if info.Mode().Perm() != 0o600 {
|
||||
t.Fatalf("the credential file is readable beyond its owner: %v", info.Mode())
|
||||
}
|
||||
|
||||
// storeNamedIn is the credential store a git command line offers.
|
||||
func storeNamedIn(t *testing.T, line string) string {
|
||||
t.Helper()
|
||||
_, after, ok := strings.Cut(line, "credential.helper=store --file=")
|
||||
if !ok {
|
||||
t.Fatalf("the clone does not name the credential store: %s", line)
|
||||
}
|
||||
return strings.Fields(after)[0]
|
||||
}
|
||||
|
||||
// Without a credential, a clone is exactly the invocation it always was, and no credential file
|
||||
@@ -506,7 +510,7 @@ func TestAContextCloneCarriesTheSameCredentialStore(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stored := filepath.Join(workspace, "git-credentials")
|
||||
stored := storeNamedIn(t, r.ran[0])
|
||||
var contextClone string
|
||||
for _, line := range r.ran {
|
||||
if strings.Contains(line, "clone") && strings.Contains(line, "source.git") {
|
||||
|
||||
+32
-12
@@ -229,7 +229,13 @@ func ScriptToolchain(script []byte) string {
|
||||
// Check runs one merge check. An error is that it could not run; the verdict is then "error".
|
||||
func Check(ctx context.Context, run Runner, spec CheckSpec, workspace, registry string, forge GitCredential,
|
||||
log Log) (CheckVerdict, error) {
|
||||
say := logging(log)
|
||||
// Everything a check says goes to the build's log, which the bus keeps: said redacted (novox/hq issue 462).
|
||||
redact := redactorFor(forge)
|
||||
say := logging(func(step, message string) {
|
||||
if log != nil {
|
||||
log(step, redact.redact(message))
|
||||
}
|
||||
})
|
||||
began := time.Now()
|
||||
ctx, stop := context.WithTimeout(ctx, CheckTimeout)
|
||||
defer stop()
|
||||
@@ -242,20 +248,26 @@ func Check(ctx context.Context, run Runner, spec CheckSpec, workspace, registry
|
||||
return CheckVerdict{}, err
|
||||
}
|
||||
defer os.RemoveAll(root)
|
||||
credentials := ""
|
||||
if forge.URL != "" {
|
||||
credentials = filepath.Join(workspace, "git-credentials")
|
||||
if err := os.WriteFile(credentials, []byte(forge.URL+"\n"), 0o600); err != nil {
|
||||
return CheckVerdict{}, err
|
||||
}
|
||||
// The forge credential lives outside the workspace the check's containers mount, and only while the
|
||||
// check clones (novox/hq issue 462).
|
||||
credentials, forget, err := storeCredential(workspace, forge)
|
||||
if err != nil {
|
||||
return CheckVerdict{}, err
|
||||
}
|
||||
defer forget()
|
||||
if err := forgetLeftCredentials(workspace); err != nil {
|
||||
return CheckVerdict{}, err
|
||||
}
|
||||
clone := func(repository, ref, dir string) error {
|
||||
if _, err := run(ctx, root, "git", cloneWith(credentials, "clone", "--quiet", repository, dir)...); err != nil {
|
||||
return fmt.Errorf("cannot clone %s: %w", repository, err)
|
||||
return fmt.Errorf("cannot clone %s: %s", redact.redact(repository), redact.redact(err.Error()))
|
||||
}
|
||||
if err := recordedWithoutUserinfo(ctx, run, filepath.Join(root, dir), repository); err != nil {
|
||||
return err
|
||||
}
|
||||
if ref != "" {
|
||||
if _, err := run(ctx, filepath.Join(root, dir), "git", "checkout", "--quiet", ref); err != nil {
|
||||
return fmt.Errorf("%s has no %s: %w", repository, ref, err)
|
||||
return fmt.Errorf("%s has no %s: %w", redact.redact(repository), ref, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
@@ -323,6 +335,9 @@ func Check(ctx context.Context, run Runner, spec CheckSpec, workspace, registry
|
||||
}
|
||||
}
|
||||
|
||||
// Every clone is made: the credential is gone before anything of the check runs (novox/hq issue 462).
|
||||
forget()
|
||||
|
||||
// The facts, and the versions they say the mesh runs.
|
||||
if registry == "" {
|
||||
return CheckVerdict{}, errors.New("no artifact store to read the facts snapshot from")
|
||||
@@ -479,8 +494,11 @@ func Check(ctx context.Context, run Runner, spec CheckSpec, workspace, registry
|
||||
}
|
||||
}
|
||||
|
||||
// What the check printed travels in the verdict to the forge and the controller: redacted as the log is.
|
||||
v.Gate.Summary, v.Repo.Summary, v.Repo.Failed = redact.redact(v.Gate.Summary), redact.redact(v.Repo.Summary),
|
||||
redact.redact(v.Repo.Failed)
|
||||
v.Verdict, v.Summary = v.Gate.Verdict, v.Gate.Summary
|
||||
v.Report, v.Took = withWhatFailed(out.String(), v.Repo), time.Since(began)
|
||||
v.Report, v.Took = redact.redact(withWhatFailed(out.String(), v.Repo)), time.Since(began)
|
||||
say("check", "gate %s — %s; repository %s — %s (%s)", strings.ToUpper(v.Gate.Verdict), v.Gate.Summary,
|
||||
strings.ToUpper(v.Repo.Verdict), v.Repo.Summary, v.Took.Round(time.Second))
|
||||
return v, nil
|
||||
@@ -606,7 +624,9 @@ func (l *toTheLog) line(line string) {
|
||||
line = line[:logLineBytes] + fmt.Sprintf(" … (%d bytes more)", len(line)-logLineBytes)
|
||||
}
|
||||
l.said++
|
||||
l.say("output", "%s", line)
|
||||
// Redacted by its shape before the bus keeps it (novox/hq issue 462); Check's own say adds the secrets
|
||||
// the builder knows.
|
||||
l.say("output", "%s", redactor{}.redact(line))
|
||||
}
|
||||
|
||||
// close says the last line, and, when lines were left out, how many and what failed.
|
||||
@@ -624,7 +644,7 @@ func (l *toTheLog) close(failed string) {
|
||||
}
|
||||
l.say("output", "--- what failed, picked from the whole of its output")
|
||||
for _, line := range strings.Split(failed, "\n") {
|
||||
l.say("output", "%s", line)
|
||||
l.say("output", "%s", redactor{}.redact(line))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,251 @@
|
||||
package builder
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io/fs"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/facts"
|
||||
)
|
||||
|
||||
// **A check's container never sees the forge credential** (novox/hq issue 462): the toolchain container
|
||||
// mounts the workspace as HOME, so a credential kept there — or a clone's .git/config carrying one — is
|
||||
// readable by any pull request's merge-check.sh, and printed, kept on the bus for days.
|
||||
|
||||
const (
|
||||
forgeSecret = "sw0rdfi5h-forge"
|
||||
forgeURL = "http://mesh_novox_builder:" + forgeSecret + "@forge.invalid:20000"
|
||||
besideSecret = "b3side-t0ken"
|
||||
npmSecret = "npm-s3cret-t0ken"
|
||||
)
|
||||
|
||||
// aFactsRegistry is an artifact store holding the facts snapshot, and nothing else.
|
||||
func aFactsRegistry(t *testing.T) string {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(facts.Facts{Format: facts.Format, Taken: time.Now().UTC(),
|
||||
Versions: facts.Versions{Bus: "2.11.17", Store: "17.11"}, Machines: []facts.Machine{{Name: "abcdef", Length: 6}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sum := sha256.Sum256(body)
|
||||
digest := "sha256:" + hex.EncodeToString(sum[:])
|
||||
manifest, _ := json.Marshal(map[string]any{"schemaVersion": 2, "layers": []map[string]any{
|
||||
{"mediaType": facts.MediaType, "digest": digest, "size": len(body)}}})
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch {
|
||||
case strings.Contains(r.URL.Path, "/manifests/"):
|
||||
w.Write(manifest)
|
||||
case strings.HasSuffix(r.URL.Path, "/blobs/"+digest):
|
||||
w.Write(body)
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
return strings.TrimPrefix(srv.URL, "http://")
|
||||
}
|
||||
|
||||
// bareURL is a URL with its userinfo left out.
|
||||
func bareURL(t *testing.T, raw string) string {
|
||||
u, err := url.Parse(raw)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
u.User = nil
|
||||
return u.String()
|
||||
}
|
||||
|
||||
func TestACheckContainerSeesNoForgeCredential(t *testing.T) {
|
||||
repo, head := aCheckedRepository(t, map[string]string{CheckScript: "echo checked\n"})
|
||||
besideRepo, besideHead := aCheckedRepository(t, map[string]string{"README": "beside"})
|
||||
// A clone source that carries userinfo, as a forge's clone URL may: git records it as given in the
|
||||
// clone's .git/config, which the container reads.
|
||||
besideURL := "file://beside-user:" + besideSecret + "@" + besideRepo
|
||||
workspace := t.TempDir()
|
||||
|
||||
var stores []string
|
||||
reached := false
|
||||
var leaks []string
|
||||
run := func(ctx context.Context, dir, name string, args ...string) (string, error) {
|
||||
switch name {
|
||||
case "git":
|
||||
for _, a := range args {
|
||||
if f, ok := strings.CutPrefix(a, "credential.helper=store --file="); ok {
|
||||
stores = append(stores, f)
|
||||
raw, err := os.ReadFile(f)
|
||||
if err != nil || strings.TrimSpace(string(raw)) != forgeURL {
|
||||
t.Errorf("git is offered a store that does not hold the credential as given: %q, %v", raw, err)
|
||||
}
|
||||
if info, err := os.Stat(f); err == nil && info.Mode().Perm() != 0o600 {
|
||||
t.Errorf("the credential store is readable beyond its owner: %v", info.Mode())
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(args) >= 2 && args[len(args)-3] == "--quiet" && hasString(args, "clone") {
|
||||
source := args[len(args)-2]
|
||||
if u, err := url.Parse(source); err == nil && u.User != nil {
|
||||
// Git cannot reach a file:// URL with userinfo; clone it without, then record it as git
|
||||
// would have: as given.
|
||||
clone := append(append([]string{}, args[:len(args)-2]...), bareURL(t, source), args[len(args)-1])
|
||||
if out, err := Command(ctx, dir, "git", clone...); err != nil {
|
||||
return out, err
|
||||
}
|
||||
return Command(ctx, filepath.Join(dir, args[len(args)-1]), "git", "remote", "set-url", "origin", source)
|
||||
}
|
||||
}
|
||||
return Command(ctx, dir, name, args...)
|
||||
case "docker":
|
||||
if !reached {
|
||||
reached = true
|
||||
// The first container: everything the workspace holds is what the toolchain container sees.
|
||||
filepath.WalkDir(workspace, func(path string, d fs.DirEntry, err error) error {
|
||||
if err != nil || d.IsDir() {
|
||||
return nil
|
||||
}
|
||||
raw, _ := os.ReadFile(path)
|
||||
s := string(raw)
|
||||
if strings.Contains(s, forgeSecret) || strings.Contains(s, besideSecret) ||
|
||||
strings.Contains(s, npmSecret) || d.Name() == "git-credentials" || d.Name() == ".npmrc" ||
|
||||
strings.Contains(s, "credential.helper") {
|
||||
leaks = append(leaks, path)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
for _, f := range stores {
|
||||
if _, err := os.Stat(f); !errors.Is(err, os.ErrNotExist) {
|
||||
leaks = append(leaks, f+" (still there when the first container runs)")
|
||||
}
|
||||
if rel, err := filepath.Rel(workspace, f); err == nil && !strings.HasPrefix(rel, "..") {
|
||||
leaks = append(leaks, f+" (inside the workspace the container mounts)")
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(args) > 0 && args[0] == "ps" {
|
||||
return "", nil
|
||||
}
|
||||
return "", errors.New("no container runtime in this test")
|
||||
}
|
||||
return "", errors.New("unexpected command " + name)
|
||||
}
|
||||
// A credential an older builder left in the workspace is removed too: the forge's, and the .npmrc a
|
||||
// build wrote into the tree it cloned.
|
||||
if err := os.WriteFile(filepath.Join(workspace, "git-credentials"), []byte(forgeURL+"\n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, dir := range []string{"source/x", "context-server"} {
|
||||
if err := os.MkdirAll(filepath.Join(workspace, dir), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(workspace, dir, ".npmrc"),
|
||||
[]byte("//forge.invalid/api/packages/novox/npm/:_authToken="+npmSecret+"\n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
_, err := Check(t.Context(), run, CheckSpec{ID: "check-462", Repository: repo, Ref: head, Owner: "novox",
|
||||
Repo: "mesh-controller", Number: 1, Toolchain: "golang", Beside: map[string]Beside{
|
||||
"mesh-catalog": {Repository: besideURL, Ref: besideHead}}}, workspace, aFactsRegistry(t),
|
||||
GitCredential{URL: forgeURL}, nil)
|
||||
if err == nil {
|
||||
t.Fatal("the check ran past its first container in a test with none")
|
||||
}
|
||||
if !reached {
|
||||
t.Fatalf("the check never reached its first container: %v", err)
|
||||
}
|
||||
if len(stores) == 0 {
|
||||
t.Fatal("no clone was offered the forge credential")
|
||||
}
|
||||
if len(leaks) > 0 {
|
||||
t.Fatalf("the check's container sees the credential:\n%s", strings.Join(leaks, "\n"))
|
||||
}
|
||||
}
|
||||
|
||||
// Every line a repository's own check prints is published to the build's log redacted.
|
||||
func TestACheckLinePublishedToTheLogIsRedacted(t *testing.T) {
|
||||
var said []string
|
||||
say := func(step, format string, args ...any) {
|
||||
if step == "output" && len(args) > 0 {
|
||||
said = append(said, args[0].(string))
|
||||
}
|
||||
}
|
||||
var out tail
|
||||
layer := ownCheck(t.Context(), CheckSpec{Toolchain: "golang"}, []ScriptPart{{Toolchain: "go", Script: CheckScript}},
|
||||
t.TempDir(), &out, func() bool { return false }, func(string, string) *exec.Cmd {
|
||||
return exec.CommandContext(t.Context(), "sh", "-c", "echo cloning http://mesh_builder:t0ps3cret-forge@forge.invalid/novox/x.git; "+
|
||||
"echo token ghp_abcdefghijklmnopqrstuvwxyz0123456789")
|
||||
}, say)
|
||||
if layer == nil || layer.Verdict != "pass" {
|
||||
t.Fatalf("the check answered %+v\n%s", layer, out.String())
|
||||
}
|
||||
joined := strings.Join(said, "\n")
|
||||
if strings.Contains(joined, "t0ps3cret-forge") || strings.Contains(joined, "ghp_abcdef") {
|
||||
t.Fatalf("a secret the check printed is published to the build's log:\n%s", joined)
|
||||
}
|
||||
if !strings.Contains(joined, "http://mesh_builder:[redacted: a password in a URI]@forge.invalid/novox/x.git") {
|
||||
t.Fatalf("the line is not said with what was there named:\n%s", joined)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheRedactorHidesTheForgeCredentialAndShapes(t *testing.T) {
|
||||
r := redactorFor(GitCredential{URL: forgeURL})
|
||||
for in, want := range map[string]string{
|
||||
"the secret alone: " + forgeSecret: "the secret alone: [redacted: the forge credential]",
|
||||
"go test ./... ok": "go test ./... ok",
|
||||
"--password hunter22 and done": "--password [redacted: the word after --password] and done",
|
||||
"commit 3b6b54a0c1d2e3f4a5b6c7d8e9f0": "commit 3b6b54a0c1d2e3f4a5b6c7d8e9f0",
|
||||
} {
|
||||
if got := r.redact(in); got != want {
|
||||
t.Errorf("%q redacted as %q, want %q", in, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A build's clone of a URL carrying userinfo records it without, so no build leaves a credential in
|
||||
// workspace/source/.git/config while it runs either (novox/hq issue 462); the tree itself is gone when the
|
||||
// build ends.
|
||||
func TestABuildsCloneRecordsNoUserinfo(t *testing.T) {
|
||||
repo, _ := aCheckedRepository(t, map[string]string{ManifestName: `{"module":"plain","version":"1"}`})
|
||||
source := "file://build-user:" + besideSecret + "@" + repo
|
||||
workspace := t.TempDir()
|
||||
tree := filepath.Join(workspace, "source")
|
||||
var configs []string
|
||||
run := func(ctx context.Context, dir, name string, args ...string) (string, error) {
|
||||
if name == "git" && hasString(args, "clone") && args[len(args)-2] == source {
|
||||
clone := append(append([]string{}, args[:len(args)-2]...), bareURL(t, source), args[len(args)-1])
|
||||
if out, err := Command(ctx, dir, "git", clone...); err != nil {
|
||||
return out, err
|
||||
}
|
||||
return Command(ctx, args[len(args)-1], "git", "remote", "set-url", "origin", source)
|
||||
}
|
||||
if name == "git" && len(args) > 0 && args[0] == "rev-parse" {
|
||||
raw, _ := os.ReadFile(filepath.Join(tree, ".git", "config"))
|
||||
configs = append(configs, string(raw))
|
||||
}
|
||||
return Command(ctx, dir, name, args...)
|
||||
}
|
||||
if _, err := Build(t.Context(), run, &recorded{}, source, "", "", workspace, nil, Npmrc{}, GitCredential{}, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(configs) == 0 {
|
||||
t.Fatal("the build never read its clone")
|
||||
}
|
||||
for _, c := range configs {
|
||||
if strings.Contains(c, besideSecret) || strings.Contains(c, "build-user") {
|
||||
t.Fatalf("the build's clone records the credential it was cloned with:\n%s", c)
|
||||
}
|
||||
}
|
||||
if _, err := os.Stat(tree); !os.IsNotExist(err) {
|
||||
t.Fatalf("the build's tree outlives the build: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -429,6 +429,11 @@ func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, dig
|
||||
if response.ContentLength > 0 {
|
||||
put.ContentLength = response.ContentLength
|
||||
}
|
||||
// **Not waited for if refused** (novox/hq issue 457): the body streams from upstream and cannot be
|
||||
// read twice, so a registry held still between the POST above and this PUT fails the copy with
|
||||
// "its body cannot be read twice" rather than waiting. Accepted: the POST a moment before already
|
||||
// waited the registry out, so the window is the length of one upstream fetch, and the build fails
|
||||
// loudly, to be asked again, rather than buffering every base blob in memory.
|
||||
done, err := r.client().Do(put)
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot upload blob %s: %w", digest, err)
|
||||
|
||||
@@ -70,7 +70,16 @@ func TestNpmrcDisabledUntilThereIsARegistry(t *testing.T) {
|
||||
func TestAnImageBuildGetsTheCredentialInTheContextAndHostNetwork(t *testing.T) {
|
||||
r, workspace := aRepository(t, withBoth, map[string]string{"Dockerfile": "FROM scratch\nCOPY .npmrc ./", "files/x": "y"})
|
||||
n := Npmrc{Scope: "@novox", Registry: "https://forge.invalid/api/packages/novox/npm/", Token: "t"}
|
||||
if _, err := Build(context.Background(), r.run, r,
|
||||
npmrc := filepath.Join(workspace, "source", ".npmrc")
|
||||
inContext := false
|
||||
run := func(ctx context.Context, dir, name string, args ...string) (string, error) {
|
||||
if name == "docker" && len(args) > 0 && args[0] == "build" {
|
||||
_, err := os.Stat(npmrc)
|
||||
inContext = err == nil
|
||||
}
|
||||
return r.run(ctx, dir, name, args...)
|
||||
}
|
||||
if _, err := Build(context.Background(), run, r,
|
||||
"https://forge.invalid/meshboard.git", "", "", workspace, nil, n, GitCredential{}, nil); err != nil {
|
||||
t.Fatalf("the build failed: %v", err)
|
||||
}
|
||||
@@ -90,11 +99,14 @@ func TestAnImageBuildGetsTheCredentialInTheContextAndHostNetwork(t *testing.T) {
|
||||
if !strings.Contains(build, "--network host") {
|
||||
t.Fatalf("the build was not given the host network to reach the registry: %s", build)
|
||||
}
|
||||
// The .npmrc is written into the build context (the source tree), where a Dockerfile COPYs it.
|
||||
tree := filepath.Join(workspace, "source")
|
||||
npmrc := filepath.Join(tree, ".npmrc")
|
||||
if _, err := os.Stat(npmrc); err != nil {
|
||||
t.Fatalf("the credential was not written into the build context: %v", err)
|
||||
// The .npmrc is written into the build context (the source tree), where a Dockerfile COPYs it, and
|
||||
// removed with the tree when the build ends: a later check's container mounts the workspace (novox/hq
|
||||
// issue 462).
|
||||
if !inContext {
|
||||
t.Fatal("the credential was not in the build context when the image was built")
|
||||
}
|
||||
if _, err := os.Stat(npmrc); !os.IsNotExist(err) {
|
||||
t.Fatalf("the credential outlives the build in the workspace: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
package builder
|
||||
|
||||
// **Every line of a check's output is redacted before it is kept** (novox/hq issue 462).
|
||||
//
|
||||
// A repository's own check prints into the build's log, which the bus keeps for days and anyone who may read
|
||||
// its events reads, and into the verdict, which the forge shows on the pull request. Software prints what it
|
||||
// was given — a URL carrying a password, a token in a flag — and a pull request may print on purpose.
|
||||
// So a line is said only after every secret the builder knows (the forge credential's password) and every
|
||||
// value whose shape says it is one is replaced by a mark naming what was there, as the journal verb does.
|
||||
//
|
||||
// Copied from the journal tool's redactor (mesh-catalog, modules/systemd/cmd/systemd-tools/secrets.go,
|
||||
// itself a copy of the docker module's), narrowed to a line's shapes, with the token shapes a check's output
|
||||
// may carry added. A third copy: sharing them through mesh-sdk is novox/hq issue 471.
|
||||
|
||||
import (
|
||||
"net/url"
|
||||
"regexp"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// secretName is a variable name that says its value is a secret.
|
||||
var secretName = regexp.MustCompile(`(?i)(pass(word|wd|phrase)?|secret|token|api_?key|private_?key|access_?key|credential|auth)`)
|
||||
|
||||
// notAValue is a name that says its value is where a secret is, not the secret: a file or a path.
|
||||
var notAValue = regexp.MustCompile(`(?i)(_FILE|FILE|_PATH|_DIR)$`)
|
||||
|
||||
// uriPassword is a URI carrying a password in its userinfo: scheme://user:password@.
|
||||
var uriPassword = regexp.MustCompile(`[A-Za-z][A-Za-z0-9+.-]*://[^\s/:@'"]*:([^\s/@'"]+)@`)
|
||||
|
||||
// tokenShaped are tokens recognised by their own prefix, whatever surrounds them: a forge's or a host's
|
||||
// access token, a JSON web token, a NATS seed.
|
||||
var tokenShaped = []struct {
|
||||
name string
|
||||
re *regexp.Regexp
|
||||
}{
|
||||
{"an access token", regexp.MustCompile(`\b(gh[pousr]_[A-Za-z0-9]{20,}|github_pat_[A-Za-z0-9_]{20,}|glpat-[A-Za-z0-9_-]{20,}|xox[abpr]-[A-Za-z0-9-]{10,}|sk-ant-[A-Za-z0-9_-]{20,})`)},
|
||||
{"a JSON web token", regexp.MustCompile(`\beyJ[A-Za-z0-9_-]{8,}\.eyJ[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]+`)},
|
||||
{"a NATS seed", regexp.MustCompile(`\bS[ACNOU][A-Z2-7]{56}\b`)},
|
||||
}
|
||||
|
||||
// masked is a password a program already hid: ***, xxx, <redacted>, [REDACTED].
|
||||
var masked = regexp.MustCompile(`^(\*+|x+|X+|<[^>]*>|\[[^\]]*\]|%2A+)$`)
|
||||
|
||||
// ordinary is a value under a secret's name that is not one: a path, an address, a number, a switch.
|
||||
var ordinary = regexp.MustCompile(`^(/.*|[A-Za-z][A-Za-z0-9+.-]*://.*|[0-9.]+[a-z]?|(?i:true|false|yes|no|on|off|none|null))$`)
|
||||
|
||||
// leastSecret is the shortest value compared as a secret: a shorter one matches ordinary words.
|
||||
const leastSecret = 6
|
||||
|
||||
// passwordFlags take a secret as their next word, or after `=`, whatever the program.
|
||||
var passwordFlags = map[string]bool{
|
||||
"-P": true, "--password": true, "--pass": true, "--passwd": true, "--secret": true, "--secret-key": true,
|
||||
"--token": true, "--api-key": true, "--apikey": true, "--auth": true,
|
||||
}
|
||||
|
||||
// knownSecret is one value the builder holds, by the name it is said under.
|
||||
type knownSecret struct {
|
||||
Name string
|
||||
Value string
|
||||
}
|
||||
|
||||
// redactor hides the secrets it knows and those a line's shapes say are secrets.
|
||||
type redactor struct{ known []knownSecret }
|
||||
|
||||
// redactorFor knows the forge credential's password, and its user's name with it, in every form git or a
|
||||
// program may print them.
|
||||
func redactorFor(forge GitCredential) redactor {
|
||||
var r redactor
|
||||
if forge.URL == "" {
|
||||
return r
|
||||
}
|
||||
for _, m := range uriPassword.FindAllStringSubmatch(forge.URL, -1) {
|
||||
r.add("the forge credential", m[1])
|
||||
if dec, err := url.PathUnescape(m[1]); err == nil && dec != m[1] {
|
||||
r.add("the forge credential", dec)
|
||||
}
|
||||
}
|
||||
return r
|
||||
}
|
||||
|
||||
func (r *redactor) add(name, value string) {
|
||||
if len(value) < leastSecret || masked.MatchString(value) {
|
||||
return
|
||||
}
|
||||
for _, k := range r.known {
|
||||
if k.Value == value {
|
||||
return
|
||||
}
|
||||
}
|
||||
r.known = append(r.known, knownSecret{name, value})
|
||||
}
|
||||
|
||||
// redact is a text with every known secret, every value its shape says is one, and every password inside a
|
||||
// URI replaced by a mark naming what was there. Line by line: a shape is judged within its line.
|
||||
func (r redactor) redact(text string) string {
|
||||
if !strings.ContainsAny(text, "\n") {
|
||||
return r.line(text)
|
||||
}
|
||||
lines := strings.Split(text, "\n")
|
||||
for i, l := range lines {
|
||||
lines[i] = r.line(l)
|
||||
}
|
||||
return strings.Join(lines, "\n")
|
||||
}
|
||||
|
||||
func (r redactor) line(line string) string {
|
||||
replace := func(s knownSecret) {
|
||||
for _, f := range forms(s.Value) {
|
||||
line = strings.ReplaceAll(line, f, "[redacted: "+s.Name+"]")
|
||||
}
|
||||
}
|
||||
for _, s := range r.known {
|
||||
replace(s)
|
||||
}
|
||||
for _, s := range shaped(line) {
|
||||
replace(s)
|
||||
}
|
||||
line = uriPassword.ReplaceAllStringFunc(line, func(m string) string {
|
||||
sub := uriPassword.FindStringSubmatch(m)
|
||||
if masked.MatchString(sub[1]) || strings.HasPrefix(sub[1], "[redacted") {
|
||||
return m
|
||||
}
|
||||
return strings.TrimSuffix(m, sub[1]+"@") + "[redacted: a password in a URI]@"
|
||||
})
|
||||
for _, t := range tokenShaped {
|
||||
line = t.re.ReplaceAllString(line, "[redacted: "+t.name+"]")
|
||||
}
|
||||
return line
|
||||
}
|
||||
|
||||
// forms are the ways a value may appear printed: as given, and URL-encoded.
|
||||
func forms(value string) []string {
|
||||
out := []string{value}
|
||||
for _, f := range []string{url.QueryEscape(value), url.PathEscape(value)} {
|
||||
if f != value && !hasString(out, f) {
|
||||
out = append(out, f)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func hasString(list []string, s string) bool {
|
||||
for _, x := range list {
|
||||
if x == s {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// shaped are the values a line carries by their shape: the word after a password flag, or the value of one
|
||||
// given with `=`, and a NAME=value whose name says secret.
|
||||
func shaped(line string) []knownSecret {
|
||||
var out []knownSecret
|
||||
add := func(name, value string) {
|
||||
value = strings.Trim(value, `"',;`)
|
||||
if len(value) < leastSecret || masked.MatchString(value) || ordinary.MatchString(value) ||
|
||||
strings.HasPrefix(value, "[redacted") {
|
||||
return
|
||||
}
|
||||
out = append(out, knownSecret{name, value})
|
||||
}
|
||||
words := strings.Fields(line)
|
||||
for i, w := range words {
|
||||
if flag, value, ok := strings.Cut(w, "="); ok && strings.HasPrefix(flag, "-") {
|
||||
if passwordFlags[flag] {
|
||||
add("the value of "+flag, value)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if name, value, ok := strings.Cut(w, "="); ok && name != "" && secretName.MatchString(name) &&
|
||||
!notAValue.MatchString(name) && !strings.ContainsAny(name, "/:") {
|
||||
add("the value of "+name, value)
|
||||
continue
|
||||
}
|
||||
if i+1 < len(words) && passwordFlags[w] {
|
||||
add("the word after "+w, words[i+1])
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// withoutUserinfo is a URL with its userinfo left out, and whether it carried any.
|
||||
func withoutUserinfo(raw string) (string, bool) {
|
||||
u, err := url.Parse(raw)
|
||||
if err != nil || u.User == nil {
|
||||
return raw, false
|
||||
}
|
||||
u.User = nil
|
||||
return u.String(), true
|
||||
}
|
||||
@@ -52,7 +52,11 @@ func (r Registry) PublishImage(ctx context.Context, localTag, repository string)
|
||||
if _, err := r.Run(ctx, "", "docker", "tag", localTag, remote); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if _, err := r.Run(ctx, "", "docker", "push", remote); err != nil {
|
||||
// The push waits for a registry held still, as every request to it does (novox/hq issue 457).
|
||||
if err := waitForRegistry(ctx, r.Address, "docker push "+remote, func() error {
|
||||
_, err := r.Run(ctx, "", "docker", "push", remote)
|
||||
return err
|
||||
}); err != nil {
|
||||
return "", err
|
||||
}
|
||||
out, err := r.Run(ctx, "", "docker", "inspect", "--format", "{{index .RepoDigests 0}}", remote)
|
||||
@@ -169,11 +173,13 @@ func (r Registry) has(ctx context.Context, url string, accept ...string) (bool,
|
||||
return response.StatusCode == http.StatusOK, nil
|
||||
}
|
||||
|
||||
// client is the client for the registry and for upstream, whose requests to the registry wait out a
|
||||
// registry held still (novox/hq issue 457).
|
||||
func (r Registry) client() *http.Client {
|
||||
if r.HTTP != nil {
|
||||
return r.HTTP
|
||||
return waiting(r.HTTP, r.Address)
|
||||
}
|
||||
return http.DefaultClient
|
||||
return waiting(http.DefaultClient, r.Address)
|
||||
}
|
||||
|
||||
// separator is whether the upload location already carries a query.
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
package builder
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A build waits out a registry that refuses connections, for a bounded time (novox/hq issue 457).
|
||||
//
|
||||
// **Why waiting, and why here.** The store's nightly collection holds the registry still for its run —
|
||||
// about a minute and a half, measured on 2026-10-11 — and a build that reached the registry in that
|
||||
// window failed on "connection refused", and its whole delivery plan with it: a plan failed for a
|
||||
// pause the mesh itself scheduled. The other design weighed was the collection telling the controller
|
||||
// it holds the registry, and the controller holding build asks while it runs. Waiting here is smaller
|
||||
// and covers more: it is local to the one place that talks to the registry, needs no new message
|
||||
// between modules, and also carries a build over any other short outage — a registry restarted by its
|
||||
// own update, say. A refusal is the one error waited for: nothing was sent, so trying again cannot
|
||||
// do anything twice, and it is what a registry that is stopped answers.
|
||||
//
|
||||
// **Bounded, and loud past the bound.** registryWait is longer than the collection holds the registry
|
||||
// (five minutes against about one and a half), so the pause the mesh schedules is always waited out,
|
||||
// and a registry that is really down still fails the build — saying how long it was refused — rather
|
||||
// than holding a build machine for ever. Every wait is said in the build's log, with why, and so is
|
||||
// the registry answering again.
|
||||
|
||||
var (
|
||||
// registryWait is how long a build waits for a registry that refuses, per call that found it so.
|
||||
registryWait = 5 * time.Minute
|
||||
// registryFirstPause is the first pause between tries; each pause doubles, up to registryMostPause.
|
||||
registryFirstPause = time.Second
|
||||
)
|
||||
|
||||
// registryMostPause is the longest pause between two tries: short enough that a build goes on within
|
||||
// seconds of the registry answering again.
|
||||
const registryMostPause = 10 * time.Second
|
||||
|
||||
// refused is whether an error is a connection refused: from a dial here, or as a command such as docker
|
||||
// said it in its output.
|
||||
func refused(err error) bool {
|
||||
return err != nil && (errors.Is(err, syscall.ECONNREFUSED) || strings.Contains(err.Error(), "connection refused"))
|
||||
}
|
||||
|
||||
// waitForRegistry runs try, and while it fails because the registry at address refuses connections,
|
||||
// tries again with a growing pause until registryWait has passed. what names the call, for the log.
|
||||
func waitForRegistry(ctx context.Context, address, what string, try func() error) error {
|
||||
err := try()
|
||||
if !refused(err) {
|
||||
return err
|
||||
}
|
||||
started := time.Now()
|
||||
pause := registryFirstPause
|
||||
tell("registry", "%s: refused; the build waits for the registry at %s, for up to %s — it is held still while "+
|
||||
"the store's nightly collection runs, about a minute and a half (novox/hq issue 457)",
|
||||
what, address, registryWait)
|
||||
for {
|
||||
left := registryWait - time.Since(started)
|
||||
if left <= 0 {
|
||||
tell("registry", "%s: the registry at %s still refuses after %s; the build fails", what, address,
|
||||
time.Since(started).Round(time.Second))
|
||||
return fmt.Errorf("the registry at %s refused every connection for %s, longer than its nightly "+
|
||||
"collection holds it still, so it is down, not paused: %w",
|
||||
address, time.Since(started).Round(time.Millisecond), err)
|
||||
}
|
||||
wait := min(pause, left)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return fmt.Errorf("stopped while waiting for the registry at %s: %w (last: %v)", address, ctx.Err(), err)
|
||||
case <-time.After(wait):
|
||||
}
|
||||
pause = min(pause*2, registryMostPause)
|
||||
if err = try(); !refused(err) {
|
||||
// Said as it is: the registry answering is only the build going on when the call worked.
|
||||
if err == nil {
|
||||
tell("registry", "%s: the registry at %s answers again after %s; the build goes on", what, address,
|
||||
time.Since(started).Round(time.Millisecond))
|
||||
} else {
|
||||
tell("registry", "%s: the registry at %s no longer refuses after %s, and answered with: %v", what,
|
||||
address, time.Since(started).Round(time.Millisecond), err)
|
||||
}
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// waitingTransport waits for the registry on every request to it, and on none to anywhere else: the
|
||||
// same client copies from upstream registries, whose refusals are theirs to answer.
|
||||
type waitingTransport struct {
|
||||
base http.RoundTripper
|
||||
address string
|
||||
}
|
||||
|
||||
func (t waitingTransport) RoundTrip(request *http.Request) (*http.Response, error) {
|
||||
if request.URL.Host != t.address {
|
||||
return t.base.RoundTrip(request)
|
||||
}
|
||||
var response *http.Response
|
||||
tries := 0
|
||||
err := waitForRegistry(request.Context(), t.address, request.Method+" "+request.URL.Path, func() error {
|
||||
attempt := request
|
||||
if tries > 0 && request.Body != nil && request.Body != http.NoBody {
|
||||
// A body is sent again only when it can be read again; one that cannot is not retried.
|
||||
if request.GetBody == nil {
|
||||
return fmt.Errorf("%s %s cannot be sent again: its body cannot be read twice", request.Method, request.URL)
|
||||
}
|
||||
body, err := request.GetBody()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
attempt = request.Clone(request.Context())
|
||||
attempt.Body = body
|
||||
}
|
||||
tries++
|
||||
var err error
|
||||
response, err = t.base.RoundTrip(attempt)
|
||||
return err
|
||||
})
|
||||
return response, err
|
||||
}
|
||||
|
||||
// waiting is a client like c whose requests to the registry wait for it.
|
||||
func waiting(c *http.Client, address string) *http.Client {
|
||||
if _, already := c.Transport.(waitingTransport); already {
|
||||
return c
|
||||
}
|
||||
copied := *c
|
||||
base := c.Transport
|
||||
if base == nil {
|
||||
base = http.DefaultTransport
|
||||
}
|
||||
copied.Transport = waitingTransport{base: base, address: address}
|
||||
return &copied
|
||||
}
|
||||
@@ -0,0 +1,154 @@
|
||||
package builder
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"net"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A registry held still for a while (novox/hq issue 457): the store's nightly collection stops it for
|
||||
// about a minute and a half, and a build in that window was failed, and its whole plan with it, for a
|
||||
// pause the mesh itself scheduled. A build waits it out — for a bound longer than the collection
|
||||
// holds it — says so in its log, and still fails, loudly, past the bound.
|
||||
|
||||
// heldStill is a registry address that refuses every connection until it starts answering after
|
||||
// pause, or never when pause is negative.
|
||||
func heldStill(t *testing.T, f *fakeRegistry, pause time.Duration) string {
|
||||
t.Helper()
|
||||
handler := f.serve(t).Config.Handler
|
||||
reserved, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
address := reserved.Addr().String()
|
||||
reserved.Close() // refused from here on: nothing listens
|
||||
if pause < 0 {
|
||||
return address
|
||||
}
|
||||
server := &http.Server{Handler: handler}
|
||||
go func() {
|
||||
time.Sleep(pause)
|
||||
l, err := net.Listen("tcp", address)
|
||||
if err != nil {
|
||||
t.Errorf("cannot answer at %s again: %v", address, err)
|
||||
return
|
||||
}
|
||||
_ = server.Serve(l)
|
||||
}()
|
||||
t.Cleanup(func() { _ = server.Close() })
|
||||
return address
|
||||
}
|
||||
|
||||
// saying collects what a build says, as the build machine's per-build Said does.
|
||||
func saying(t *testing.T) func() []string {
|
||||
t.Helper()
|
||||
var mu sync.Mutex
|
||||
var lines []string
|
||||
was := Said
|
||||
Said = func(step, message string) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
lines = append(lines, step+": "+message)
|
||||
}
|
||||
t.Cleanup(func() { Said = was })
|
||||
return func() []string {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return append([]string(nil), lines...)
|
||||
}
|
||||
}
|
||||
|
||||
// waitingFor shortens the bound and the first pause, so a test waits for milliseconds.
|
||||
func waitingFor(t *testing.T, bound time.Duration) {
|
||||
t.Helper()
|
||||
wasBound, wasFirst := registryWait, registryFirstPause
|
||||
registryWait, registryFirstPause = bound, 20*time.Millisecond
|
||||
t.Cleanup(func() { registryWait, registryFirstPause = wasBound, wasFirst })
|
||||
}
|
||||
|
||||
func TestABuildWaitsForARegistryHeldStillAndGoesOn(t *testing.T) {
|
||||
waitingFor(t, 5*time.Second)
|
||||
said := saying(t)
|
||||
f := &fakeRegistry{}
|
||||
r := Registry{Address: heldStill(t, f, 400*time.Millisecond)}
|
||||
body := []byte("a theme")
|
||||
sum := sha256.Sum256(body)
|
||||
digest := "sha256:" + hex.EncodeToString(sum[:])
|
||||
|
||||
if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil {
|
||||
t.Fatalf("a registry refusing for 400ms failed the build: %v", err)
|
||||
}
|
||||
if string(f.blobs[digest]) != "a theme" {
|
||||
t.Fatalf("the registry holds %q", f.blobs[digest])
|
||||
}
|
||||
log := strings.Join(said(), "\n")
|
||||
if !strings.Contains(log, "waits for the registry at "+r.Address) || !strings.Contains(log, "nightly collection") {
|
||||
t.Errorf("the build's log does not say it waited for the registry, and why:\n%s", log)
|
||||
}
|
||||
if !strings.Contains(log, "answers again") {
|
||||
t.Errorf("the build's log does not say the registry came back:\n%s", log)
|
||||
}
|
||||
}
|
||||
|
||||
func TestABuildFailsLoudlyOnARegistryRefusingPastTheBound(t *testing.T) {
|
||||
waitingFor(t, 300*time.Millisecond)
|
||||
said := saying(t)
|
||||
f := &fakeRegistry{}
|
||||
r := Registry{Address: heldStill(t, f, -1)}
|
||||
body := []byte("a theme")
|
||||
sum := sha256.Sum256(body)
|
||||
digest := "sha256:" + hex.EncodeToString(sum[:])
|
||||
|
||||
started := time.Now()
|
||||
_, err := r.PublishArchive(context.Background(), "shell/config", body, digest)
|
||||
if err == nil {
|
||||
t.Fatal("a registry that never answered published the archive")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "refused every connection for") || !strings.Contains(err.Error(), "connection refused") {
|
||||
t.Errorf("the failure does not say it waited and was refused throughout: %v", err)
|
||||
}
|
||||
if waited := time.Since(started); waited < 300*time.Millisecond {
|
||||
t.Errorf("failed after %s, before the bound", waited)
|
||||
}
|
||||
if log := strings.Join(said(), "\n"); !strings.Contains(log, "waits for the registry") {
|
||||
t.Errorf("the build's log does not say it waited:\n%s", log)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnImagePushWaitsForARegistryHeldStill(t *testing.T) {
|
||||
waitingFor(t, 5*time.Second)
|
||||
said := saying(t)
|
||||
pushes := 0
|
||||
run := func(_ context.Context, _ string, name string, args ...string) (string, error) {
|
||||
switch args[0] {
|
||||
case "push":
|
||||
pushes++
|
||||
if pushes < 3 {
|
||||
return "", errorString("docker push: dial tcp 127.0.0.1:5000: connect: connection refused")
|
||||
}
|
||||
case "inspect":
|
||||
return "127.0.0.1:5000/m/server@sha256:abc\n", nil
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
r := Registry{Address: "127.0.0.1:5000", Run: run}
|
||||
if _, err := r.PublishImage(context.Background(), "local", "m/server"); err != nil {
|
||||
t.Fatalf("a push refused twice failed the build: %v", err)
|
||||
}
|
||||
if pushes != 3 {
|
||||
t.Errorf("pushed %d times, want 3", pushes)
|
||||
}
|
||||
if log := strings.Join(said(), "\n"); !strings.Contains(log, "waits for the registry") {
|
||||
t.Errorf("the build's log does not say it waited:\n%s", log)
|
||||
}
|
||||
}
|
||||
|
||||
type errorString string
|
||||
|
||||
func (e errorString) Error() string { return string(e) }
|
||||
@@ -1,72 +0,0 @@
|
||||
-- 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;
|
||||
@@ -44,15 +44,17 @@ func TestAPersonMayCallToolsAndNothingElse(t *testing.T) {
|
||||
}
|
||||
// The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance
|
||||
// answers, and to the instance on one machine. Nothing else.
|
||||
// And asking what answers (novox/hq ADR 0197), which claims nothing and calls nothing.
|
||||
// And asking what answers (novox/hq ADR 0197), which claims nothing and calls nothing. Each way again
|
||||
// naming ada as the caller and nobody else (novox/hq issue 365).
|
||||
var tools []string
|
||||
for _, s := range perms.Publish {
|
||||
if !strings.HasPrefix(s, "$SRV.") {
|
||||
tools = append(tools, s)
|
||||
}
|
||||
}
|
||||
if len(tools) != 2 || tools[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" ||
|
||||
tools[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" {
|
||||
want := []string{"mesh.mod.mesh-catalog.call.catalog_tools.*.person~ada", "mesh.mod.mesh-catalog.call.catalog_tools.person~ada",
|
||||
"mesh.mod.mesh-catalog.tool.catalog_tools", "mesh.mod.mesh-catalog.tool.catalog_tools.*"}
|
||||
if strings.Join(tools, " ") != strings.Join(want, " ") {
|
||||
t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish)
|
||||
}
|
||||
for _, s := range perms.Publish {
|
||||
|
||||
@@ -1,287 +0,0 @@
|
||||
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)
|
||||
}
|
||||
@@ -1,207 +0,0 @@
|
||||
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,13 +447,6 @@ 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,
|
||||
|
||||
@@ -1,33 +0,0 @@
|
||||
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,15 +197,6 @@ 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.
|
||||
@@ -614,10 +605,3 @@ 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 33917f91ac3cc213b879b32f3c70e2aa70395f98 conformance/events/module-event.json
|
||||
mesh-sdk f047d0a4a9702f5d7f3d7eadd3e62496311acc00 conformance/events/module-event.json
|
||||
|
||||
To move them, from this repository's root, with the two repositories checked out beside it:
|
||||
|
||||
|
||||
+1
-16
@@ -154,9 +154,7 @@ 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. 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.
|
||||
// later changes no digest until it is added here, on purpose.
|
||||
func (a Ask) Digest() string {
|
||||
var b strings.Builder
|
||||
b.WriteString("novox.ask.v1\n")
|
||||
@@ -170,10 +168,6 @@ 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[:])
|
||||
}
|
||||
@@ -196,15 +190,6 @@ 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).
|
||||
|
||||
+1
-29
@@ -1378,20 +1378,6 @@ 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
|
||||
@@ -1556,11 +1542,6 @@ 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) {
|
||||
@@ -1578,17 +1559,8 @@ 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, Generation: env.Generation, PutBack: env.PutBack}
|
||||
Epoch: env.Epoch, LeftOut: env.LeftOut}
|
||||
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.14-0.20261010191001-33917f91ac3c
|
||||
# git.novox.be/novox/mesh-sdk/go v0.1.11-0.20261009143344-f047d0a4a970
|
||||
## 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-20261011092230-20d5af9b2d5f
|
||||
# github.com/novox/mesh-host v0.0.0 => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
## 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-20261011092230-20d5af9b2d5f
|
||||
# github.com/novox/mesh-host => git.novox.be/novox/mesh-host v0.0.0-20261009231844-b8c854611812
|
||||
|
||||
Reference in New Issue
Block a user