Author SHA1 Message Date
mesh-admin b2808549ee Merge pull request 'Grant the genesis controller its lease and secret-replaced, not mesh.control.> (hq to-be 45 Phase 2)' (#36) from feat/a-core-that-cannot-fail-silently-phase-2-grants into main 2026-10-06 10:30:28 +00:00
jochen bbe74f3fc2 Grant the genesis controller its lease and the secret-replaced event, and not what the machines say (hq to-be 45 Phase 2)
The controller now composes a publish grant for its lease bucket
($KV.mesh-controller_lease.>), which it must write before it acts, and for
the seat event secret-replaced (hq ADR 0228, mesh-controller #82); and it no
longer publishes mesh.control.>: a machine's report has one writer, its
node-engine, and the writers table refuses a second at composition. The
installer's first user list must say what the controller derives, or a new
mesh's first controller is refused its lease and serves without one.
2026-10-06 12:25:12 +02:00
mesh-admin 3e80b7ae32 Merge pull request 'Apply through one queue, and order what is applied and reported (hq to-be 45 Phase 2)' (#35) from feat/a-core-that-cannot-fail-silently-phase-2 into main 2026-10-06 09:55:19 +00:00
jochen 31804bd8c2 Apply through one queue, and order what is applied and reported (hq to-be 45 Phase 2)
A delivery and the five-minute reconcile were two paths that applied, ordered only by a lock, and
each order it allowed was met live (issues 257, 261, 267). Now both only enqueue: one worker takes
the newest declaration held when it starts, applies it once and makes one report, and reports leave
in the order they are made.

A declaration may carry the controller's lease epoch beside its sequence; one older than what this
node applied is refused before anything is touched, counted, logged and reported. A report carries
the declaration's epoch and sequence and the host's own report sequence, kept on disk so it goes on
increasing across restarts and self-updates. Without an epoch, today's behaviour stands.
2026-10-06 11:54:10 +02:00
mesh-admin d7d93f57c5 Merge pull request 'Grant the genesis controller the self-check's ban-list question (hq to-be 45 Phase 1, D8)' (#34) from fix/phase-1-d8-grant-and-d10-version into main 2026-10-06 08:45:16 +00:00
jochen b9774ea426 Grant the genesis controller the self-check's ban-list question (hq to-be 45 Phase 1, D8)
The controller now composes a publish grant for
mesh.seat.node-intrusion-prevention.tool.banned.*, which D8 asks every
machine; the installer's first user list must say what the controller
derives, or a new mesh's first self-check is refused it.
2026-10-06 10:44:39 +02:00
mesh-admin 1545b00a87 Merge pull request 'Say the heartbeat's interval, export the host's validator, grant the genesis controller Phase 1 (hq to-be 45)' (#33) from feat/a-core-that-cannot-fail-silently-phase-1 into main 2026-10-06 08:24:39 +00:00
jochen 6953b5bafd Say the heartbeat's interval, export the host's validator, grant the genesis controller Phase 1 (hq to-be 45)
The controller's watchdog of a machine's heartbeat (S1) is bound to three
of its intervals, and a bound the controller guessed would not move when
the interval does: the heartbeat now carries interval_seconds.

Its self-check (D1) must judge every composed declaration as the host
does, and a second validator written from the host's rules would drift
from them: the host's own parsing is exported, unchanged, as
github.com/novox/mesh-host/validate.

The installer's first user list grants what the controller now composes
for itself: its condition buckets, the condition and doctor-heartbeat
events, the bus's two consumer advisories and $SRV.INFO — or the first
controller would be refused them until the broker's machine is pushed.
2026-10-06 10:18:54 +02:00
mesh-admin 93efe41dc9 Merge pull request 'Grant the genesis controller its buckets and cancelled sets (hq to-be 45 Phase 0, issue 269)' (#32) from feat/a-core-that-cannot-fail-silently-phase-0 into main 2026-10-06 07:13:38 +00:00
jochen 1e463def6b Grant the genesis controller its work queues' cancelled sets (hq issue 269)
The controller now derives the cancelled sets' subjects; the installer's
first user list says the same, or a cancel on a new mesh times out.
2026-10-06 03:01:19 +02:00
jochen a654fa768d Grant the genesis controller its own buckets' subjects (hq to-be 45 Phase 0)
The controller keeps its calls and the hand-act log in two key-value buckets
of its own, and composes a grant to write them. The installer's first user
list must say what the controller derives, or the first controller writes
nothing until the broker's machine is pushed.
2026-10-06 02:59:51 +02:00
19 changed files with 1804 additions and 393 deletions
+77
View File
@@ -0,0 +1,77 @@
package main
import (
"context"
"crypto/ed25519"
"encoding/json"
"errors"
"os"
"path/filepath"
"testing"
"github.com/novox/mesh-host/internal/link"
"github.com/novox/mesh-host/internal/store"
)
// **A declaration from an older lease epoch is refused before anything is applied** (novox/hq to-be
// 45 §6, ADR 0227 rule 2). The controller that sent it lost its lease and goes on sending; this node
// applied the holder's. Refused on the host's own apply path — the one the queue's worker runs — before
// the machine is read or touched, with a report naming what was refused and what is held, and the
// kept declaration left as it was.
func TestADeclarationFromAnOlderEpochIsRefusedBeforeAnythingIsApplied(t *testing.T) {
dir := t.TempDir()
opts := options{state: filepath.Join(dir, "state.json")}
_, private, err := ed25519.GenerateKey(nil)
if err != nil {
t.Fatal(err)
}
target := filepath.Join(dir, "a.conf")
declared := func(epoch, sequence int64) store.Declared {
body, err := json.Marshal(map[string]any{"declaration": 1, "epoch": epoch, "sequence": sequence,
"resources": []any{map[string]any{"id": "a", "type": "file", "path": target, "content": "x\n"}}})
if err != nil {
t.Fatal(err)
}
return store.Declared{Declaration: body, Signature: ed25519.Sign(private, body)}
}
held := declared(57, 3)
if err := store.SaveDeclared(store.DeclaredPath(opts.state), held); err != nil {
t.Fatal(err)
}
stale := declared(41, 12)
report := applyAndKeep(context.Background(), opts, stale.Declaration, &stale, nil, nil)
if report.OlderThan == nil || *report.OlderThan != (link.Order{Epoch: 57, Sequence: 3}) ||
report.Order != (link.Order{Epoch: 41, Sequence: 12}) || report.Declared != digestOf(stale.Declaration) ||
report.Refused == "" {
t.Fatalf("the refusal does not name what was refused and what is held: %+v", report)
}
if _, err := os.Stat(target); !errors.Is(err, os.ErrNotExist) {
t.Fatal("the refused declaration touched the machine")
}
kept, err := store.ReadDeclared(store.DeclaredPath(opts.state))
if err != nil || string(kept.Declaration) != string(held.Declaration) {
t.Fatal("the refused declaration replaced what this node kept")
}
}
// The queue's counts are kept beside the state and read back by the next host; unreadable is an
// error, never zero.
func TestTheNumbersSurviveTheHost(t *testing.T) {
state := filepath.Join(t.TempDir(), "state.json")
if sequence, refused, err := numbersBeside(state).Read(); err != nil || sequence != 0 || refused != 0 {
t.Fatalf("a node that never reported read %d, %d, %v", sequence, refused, err)
}
if err := numbersBeside(state).Save(41, 2); err != nil {
t.Fatal(err)
}
if sequence, refused, err := numbersBeside(state).Read(); err != nil || sequence != 41 || refused != 2 {
t.Fatalf("read back %d, %d, %v", sequence, refused, err)
}
if err := os.WriteFile(store.NumbersPath(state), []byte(`{"report_sequence":`), 0o600); err != nil {
t.Fatal(err)
}
if _, _, err := numbersBeside(state).Read(); err == nil {
t.Fatal("unreadable numbers were read as zero")
}
}
+139 -74
View File
@@ -458,10 +458,10 @@ func short(digest string) string {
// them again — and everything the mesh declared read as no longer declared and removed. Even the
// very declaration the mesh last sent, applied from a file, would plan to remove the foundation.
// `apply FILE` is for a machine the mesh has not spoken to, and is refused saying so.
func refuseStale(known store.State, kept store.Declared, keptErr error, digest string, from provenance, sequence int64) error {
func refuseStale(known store.State, kept store.Declared, keptErr error, digest string, from provenance, order link.Order) error {
switch from {
case fromDeclared:
return refuseOlder(kept, keptErr, sequence)
return refuseOlder(kept, keptErr, order)
case fromBundle:
if known.Genesis == nil || known.Genesis.Digest == digest {
return nil
@@ -558,7 +558,7 @@ func runApply(ctx context.Context, opts options, d *declaration.Declaration, raw
if err := apply.CheckMode(known, d); err != nil {
return err
}
if err := refuseStale(known, kept, keptErr, digest, from, d.Sequence); err != nil {
if err := refuseStale(known, kept, keptErr, digest, from, link.Order{Epoch: d.Epoch, Sequence: d.Sequence}); err != nil {
return err
}
@@ -1057,6 +1057,14 @@ func runLink(ctx context.Context, opts options) error {
aside, standAside := context.WithCancel(ctx)
defer standAside()
stoodAside := false
membership := link.Membership{
Node: mine.Node,
Broker: mine.Membership.Broker,
Fingerprint: mine.Membership.Fingerprint,
Password: mine.Membership.Password,
Transport: mine.Membership.Transport,
Signer: mine.Membership.Signer,
}
applier := func(ctx context.Context, raw, signature []byte) link.Report {
report := applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}, sched, say)
@@ -1084,40 +1092,44 @@ func runLink(ctx context.Context, opts options) error {
}
// Two things at once, and the second is what makes disconnection ordinary. The link brings
// new declarations; this holds the machine in the last one whether the link is up or not. A
// laptop shut for a week comes back and reconciles — it does not come back and ask what it is
// (novox/hq ADR 0004).
// Reports a reconcile has to make unasked — what an adopted node holds changed, or its
// firewall did — go out over the link when it is up (novox/hq ADR 0100).
outbox := make(chan link.Unasked, 1)
// new declarations; the reconcile holds the machine in the last one whether the link is up or
// not. A laptop shut for a week comes back and reconciles — it does not come back and ask what
// it is (novox/hq ADR 0004).
//
// **Both only enqueue** (novox/hq to-be 45 §6). One worker applies, the newest declaration held
// when it starts, once, and makes one report; a reconcile due while a delivery waits is that
// delivery's apply. Reports a reconcile has to make unasked — what an adopted node holds changed,
// its firewall, its outward links — are said by the same worker, in the order they are made
// (novox/hq ADR 0100, issue 267).
watch := &adoptionWatch{}
applier = watch.noting(applier)
go holdTheMachine(ctx, opts, mine, say, sched, func(r link.Report) {
if !watch.differs(r) {
return
}
select {
case <-outbox:
// An older one nobody has published yet; this one says everything it did.
default:
}
queue := &link.Queue{
Membership: membership,
Apply: applier,
Reconcile: reconcileKept(opts, mine.Membership.Signer, sched, say),
News: func(r link.Report) bool { return worthSaying(r) && watch.differs(r) },
// Counted as said only once the broker has taken it: queued and lost — the link down, the
// publish refused — the change would never be said again (novox/hq ADR 0100).
outbox <- link.Unasked{Report: r, Done: func(published bool) {
if published {
watch.said(r)
}
}}
})
Heard: watch.said,
Unsaid: unsaidBeside(opts.state),
Numbers: numbersBeside(opts.state),
Say: say,
Timeout: opts.timeout,
}
worker := make(chan struct{})
go func() {
defer close(worker)
queue.Run(aside)
}()
// **A host starting is a reason to reconcile** (to-be 45 §6: the self-update hand-over is one of
// the four): its scheduled steps are armed from the declaration the node kept by the first apply,
// and a successor that waited five minutes for it left them unarmed for five.
queue.ReconcileDue()
go holdTheMachine(ctx, queue)
held := link.HoldRoused(aside, link.Membership{
Node: mine.Node,
Broker: mine.Membership.Broker,
Fingerprint: mine.Membership.Fingerprint,
Password: mine.Membership.Password,
Transport: mine.Membership.Transport,
Signer: mine.Membership.Signer,
}, applier, say, opts.timeout, rousedBySignal(ctx), outbox, unsaidBeside(opts.state))
held := link.HoldRoused(aside, membership, queue, say, opts.timeout, rousedBySignal(ctx))
// The act in hand finishes and is kept before this process exits: an apply that stood aside with
// no link open has its report only in the unsaid store, and only once the worker has written it.
<-worker
// **Cleanly**, or the launcher counts standing aside as a crash and rolls the new host back
// before it has run once. The context this returns on was cancelled deliberately, so its error
// is not a fault to report.
@@ -1169,6 +1181,23 @@ func (u unsaidFile) Pending() (link.Report, bool, error) {
return r, r.Declared != "", nil
}
// numbersFile keeps the queue's report sequence and its count of declarations refused as older beside
// the node's state (novox/hq to-be 45 §6), where a successor reads them and goes on from there.
type numbersFile struct{ path string }
func numbersBeside(statePath string) numbersFile {
return numbersFile{path: store.NumbersPath(statePath)}
}
func (n numbersFile) Read() (int64, int64, error) {
kept, err := store.ReadNumbers(n.path)
return kept.ReportSequence, kept.RefusedOlder, err
}
func (n numbersFile) Save(reportSequence, refusedOlder int64) error {
return store.SaveNumbers(n.path, store.Numbers{ReportSequence: reportSequence, RefusedOlder: refusedOlder})
}
// adoptionWatch remembers what the node last said about what it holds and its firewall, so a
// reconcile speaks unasked only when that changed.
type adoptionWatch struct {
@@ -1232,15 +1261,6 @@ func (w *adoptionWatch) changed(r link.Report) bool {
return true
}
// noting wraps the applier, so a report the link publishes after a delivery counts as said.
func (w *adoptionWatch) noting(apply link.Applier) link.Applier {
return func(ctx context.Context, raw, signature []byte) link.Report {
r := apply(ctx, raw, signature)
w.changed(r)
return r
}
}
// rousedBySignal is the machine telling this process that its link is probably stale.
//
// **A signal, because nothing may listen on a node** (novox/hq ADR 0004). A socket for this would
@@ -1285,43 +1305,44 @@ func rousedBySignal(ctx context.Context) link.Roused {
// changed it, and then for ever.
const ReconcileEvery = 5 * time.Minute
func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, say link.Announce,
sched *apply.Scheduler, publish func(link.Report)) {
// holdTheMachine asks for a reconcile every ReconcileEvery. **It asks; it does not apply** (novox/hq
// to-be 45 §6): the queue's worker does, when its turn comes, and a reconcile due while a delivery is
// waiting is that delivery's apply.
func holdTheMachine(ctx context.Context, queue *link.Queue) {
ticker := time.NewTicker(ReconcileEvery)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
queue.ReconcileDue()
}
}
}
report, err := reapplyKept(ctx, opts, mine.Membership.Signer, sched, say)
// reconcileKept is the reconcile's act, as the queue's worker runs it: hold the machine to what it
// was last told, read once it is this act's turn, and say what did not go.
//
// A reconcile is otherwise silent. On an adopted node it speaks when what it holds or its firewall
// changed, because that is how a predecessor still writing is caught (novox/hq ADR 0100); and on any
// node when it can say which links face outside (ADR 0140) — a converged node holds nothing and found
// no firewall, so without that it could speak only in reply to a declaration, while the mesh composes
// no declaration for a node that has not said which links face outside. Whether a report is news is
// the queue's News (worthSaying and the watch), so an unchanged answer costs one comparison every
// reconcile and nothing on the bus.
func reconcileKept(opts options, signer ed25519.PublicKey, sched *apply.Scheduler,
say link.Announce) link.Reconcile {
return func(ctx context.Context) (link.Report, bool) {
report, err := reapplyKept(ctx, opts, signer, sched, say)
if errors.Is(err, store.ErrNothingDeclared) {
// Nothing to hold this machine to yet. Ordinary on a node that has enrolled and not
// been assigned anything.
continue
return link.Report{}, false
}
if err != nil {
say("cannot re-apply what this node was told: " + err.Error())
continue
}
// A reconcile is otherwise silent. On an adopted node it speaks when what it holds or
// its firewall changed, because that is how a predecessor still writing is caught
// (novox/hq ADR 0100); publish decides whether anything did.
//
// **And on any node, when it can say which links face outside** (novox/hq ADR 0140). A
// converged node holds nothing and found no firewall, so this gate closed on it and the
// node could speak only in reply to a declaration — while the mesh composes no declaration
// for a node that has not said which links face outside. A machine waiting for a push that
// was waiting for the machine, and measured: three converged machines sat silent while the
// control plane refused to send them a filter.
//
// Offered, not published: whether it is news is still the watch's to decide, so an
// unchanged answer costs one comparison every reconcile and nothing on the bus.
if publish != nil && worthSaying(report) {
publish(report)
return link.Report{}, false
}
switch {
case report.Refused != "":
@@ -1329,6 +1350,7 @@ func holdTheMachine(ctx context.Context, opts options, mine identity.Identity, s
case len(report.Failed) > 0:
say(fmt.Sprintf("holding this machine: %d applied, and %v", len(report.Applied), report.Failed))
}
return report, true
}
}
@@ -1367,13 +1389,13 @@ func announceOr(say link.Announce) func(string) {
// applying serialises applies within this process.
//
// **Two things apply here: the link and the reconcile loop**, and each reads the node's state,
// acts on the machine, and writes the state back. Run at the same time they interleave, and the
// one that saves last writes a state read before the other acted — losing what the first recorded:
// a hold, the firewall found here, a resource just applied. The machine would then be one thing
// and its record another, which is the fault every read-back in this package exists to prevent.
// Across processes — `reconcile` run by hand beside this service — the lock beside the state does
// the same (store.Lock).
// **One thing applies here: the apply queue's worker** (novox/hq to-be 45 §6), for a delivery and for
// the reconcile alike. This lock was the only order between two that applied — the link and the
// reconcile loop — and each order it allowed was met as a fault (issues 257, 261, 267); the queue is
// the order now. It stays as the guard for what may be added beside the worker: two applies
// interleaved lose what one of them recorded — a hold, the firewall found here, a resource just
// applied — and the machine is then one thing and its record another. Across processes — `reconcile`
// run by hand beside this service — the lock beside the state does the same (store.Lock).
var applying sync.Mutex
// applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can
@@ -1418,6 +1440,21 @@ func applyAndKeepHeld(ctx context.Context, opts options, raw []byte, signed *sto
if err != nil {
return link.Report{Refused: err.Error()}
}
order := link.Order{Epoch: declared.Epoch, Sequence: declared.Sequence}
// **Older than what this node applied is refused, before anything is read from the machine**
// (novox/hq to-be 45 §6, ADR 0227 rule 2): a controller that lost its lease and goes on sending.
// Read under the lock, so what it is compared with is what the last apply kept. Only a delivery:
// what a reconcile applies is what was kept.
if signed != nil {
kept, keptErr := store.ReadDeclared(store.DeclaredPath(opts.state))
if held, older := olderThanKept(kept, keptErr, order); older {
return link.Report{Declared: digestOf(raw), Order: order, OlderThan: &held,
Refused: fmt.Sprintf("this declaration was sent under lease epoch %d, sequence %d, and this "+
"node has applied epoch %d, sequence %d — from the controller holding the lease now. "+
"Refused whole; nothing was applied", order.Epoch, order.Sequence, held.Epoch, held.Sequence)}
}
}
known, err := store.Load(opts.state)
if err != nil {
@@ -1490,7 +1527,7 @@ func applyAndKeepHeld(ctx context.Context, opts options, raw []byte, signed *sto
sched.Sync(declared, held)
}
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw), Host: runningVersion(),
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw), Order: order, Host: runningVersion(),
Profile: profileAsReported(profile.Detect(ctx, profile.Default(nil), opts.timeout))}
// Which of this machine's links face outside, for the filter the mesh writes around them
// (novox/hq ADR 0140). Reported whatever the node's mode: a converged node's filter needs it,
@@ -1720,12 +1757,23 @@ func adoptDeliveredMembership(identityPath string, mine *identity.Identity, say
// sent, and one kept with no sequence is one this host received before it understood them; in either
// case there is no order to compare, and refusing on a guess would strand the node the moment the
// controller is older than the host. Equal is the same declaration again, which reconciling is for.
func refuseOlder(kept store.Declared, keptErr error, sequence int64) error {
//
// **And by the lease epoch first** (novox/hq to-be 45 §6): a declaration sent under an older epoch than
// the one kept here is a controller that lost its lease, whatever its sequence — the same rule the
// link's deliveries are held to (olderThanKept). Where both claim an epoch it alone decides, because a
// new lease holder's sequence says nothing against the last one's.
func refuseOlder(kept store.Declared, keptErr error, order link.Order) error {
if held, older := olderThanKept(kept, keptErr, order); older {
return fmt.Errorf("this declaration was sent under lease epoch %d, sequence %d, and the one kept "+
"here under epoch %d, sequence %d: a controller that no longer holds the lease sent it. Refused "+
"whole; nothing was applied", order.Epoch, order.Sequence, held.Epoch, held.Sequence)
}
sequence := order.Sequence
if sequence == 0 || keptErr != nil {
return nil
}
last, err := declaration.ParseTrusted(kept.Declaration)
if err != nil || last.Sequence == 0 {
if err != nil || last.Sequence == 0 || (order.Epoch > 0 && last.Epoch > 0) {
return nil
}
if sequence < last.Sequence {
@@ -1737,6 +1785,23 @@ func refuseOlder(kept store.Declared, keptErr error, sequence int64) error {
return nil
}
// olderThanKept says whether a declaration of this order is older than the one this node kept — the
// last it applied — and the kept one's order (link.Order.Older: only when both claim an epoch). A kept
// declaration that cannot be read claims no order: what it is compared with is unknown, and refusing
// on a guess would hold the machine off everything the mesh sends until a person looked. That kept
// file is read, and refused by name, by the next reconcile.
func olderThanKept(kept store.Declared, keptErr error, order link.Order) (link.Order, bool) {
if keptErr != nil || order.Epoch <= 0 {
return link.Order{}, false
}
last, err := declaration.ParseTrusted(kept.Declaration)
if err != nil {
return link.Order{}, false
}
held := link.Order{Epoch: last.Epoch, Sequence: last.Sequence}
return held, order.Older(held)
}
// profileAsReported is the profile as the mesh reads it — the same bytes enrolment sends, so a
// report's profile and an enrolment's are one shape on the controller's side (novox/hq ADR 0161).
func profileAsReported(detected profile.Profile) map[string]any {
+6 -4
View File
@@ -189,13 +189,15 @@ func TestAReconcileSpeaksOnlyWhenWhatIsHeldChanged(t *testing.T) {
}
}
func TestWhatTheLinkPublishedCountsAsSaid(t *testing.T) {
// What the mesh heard — a delivery's account as much as a reconcile's — counts as said: the queue
// tells the watch every report the broker took (link.Queue.Heard).
func TestWhatTheMeshHeardCountsAsSaid(t *testing.T) {
w := &adoptionWatch{}
report := link.Report{Firewall: "ufw", Held: []link.Held{{ID: "a"}}}
applier := w.noting(func(context.Context, []byte, []byte) link.Report { return report })
applier(context.Background(), nil, nil)
heard := w.said
heard(report)
if w.changed(report) {
t.Error("a reconcile repeated what the link had just published")
t.Error("a reconcile repeated what the mesh had just heard")
}
}
+44 -7
View File
@@ -5,6 +5,7 @@ import (
"strings"
"testing"
"github.com/novox/mesh-host/internal/link"
"github.com/novox/mesh-host/internal/store"
)
@@ -13,9 +14,15 @@ import (
// "older" could not.
func keptWith(t *testing.T, sequence int64) store.Declared {
t.Helper()
return keptAt(t, link.Order{Sequence: sequence})
}
func keptAt(t *testing.T, order link.Order) store.Declared {
t.Helper()
body, err := json.Marshal(map[string]any{
"declaration": 1, "resources": []any{}, "owns_nothing": true, "sequence": sequence,
"declaration": 1, "resources": []any{}, "owns_nothing": true, "sequence": order.Sequence,
"epoch": order.Epoch,
})
if err != nil {
t.Fatal(err)
@@ -24,7 +31,7 @@ func keptWith(t *testing.T, sequence int64) store.Declared {
}
func TestAnOlderDeclarationFromTheMeshIsRefused(t *testing.T) {
err := refuseOlder(keptWith(t, 7), nil, 5)
err := refuseOlder(keptWith(t, 7), nil, link.Order{Sequence: 5})
if err == nil {
t.Fatal("sequence 5 was accepted over a kept 7")
}
@@ -34,11 +41,11 @@ func TestAnOlderDeclarationFromTheMeshIsRefused(t *testing.T) {
}
func TestANewerOrEqualDeclarationIsNot(t *testing.T) {
if err := refuseOlder(keptWith(t, 7), nil, 8); err != nil {
if err := refuseOlder(keptWith(t, 7), nil, link.Order{Sequence: 8}); err != nil {
t.Fatalf("sequence 8 was refused over a kept 7: %v", err)
}
// Equal is the same declaration again, which reconciling is for.
if err := refuseOlder(keptWith(t, 7), nil, 7); err != nil {
if err := refuseOlder(keptWith(t, 7), nil, link.Order{Sequence: 7}); err != nil {
t.Fatalf("the same sequence was refused: %v", err)
}
}
@@ -46,13 +53,43 @@ func TestANewerOrEqualDeclarationIsNot(t *testing.T) {
func TestNoOrderClaimedMeansNoOrderCompared(t *testing.T) {
// An older controller sends none; a host that received before it understood them kept none.
// Refusing on a guess would strand a node the moment the controller is older than the host.
if err := refuseOlder(keptWith(t, 7), nil, 0); err != nil {
if err := refuseOlder(keptWith(t, 7), nil, link.Order{Sequence: 0}); err != nil {
t.Fatalf("a declaration claiming no order was refused: %v", err)
}
if err := refuseOlder(keptWith(t, 0), nil, 3); err != nil {
if err := refuseOlder(keptWith(t, 0), nil, link.Order{Sequence: 3}); err != nil {
t.Fatalf("a declaration was refused against a kept one that claimed no order: %v", err)
}
if err := refuseOlder(store.Declared{}, store.ErrNothingDeclared, 3); err != nil {
if err := refuseOlder(store.Declared{}, store.ErrNothingDeclared, link.Order{Sequence: 3}); err != nil {
t.Fatalf("a first declaration was refused: %v", err)
}
}
// **The lease epoch decides first** (novox/hq to-be 45 §6): a declaration from a controller that lost
// its lease is refused whatever its sequence, and a new lease holder's is taken whatever its sequence.
func TestAnOlderEpochIsRefusedAndANewerTaken(t *testing.T) {
kept := keptAt(t, link.Order{Epoch: 57, Sequence: 3})
err := refuseOlder(kept, nil, link.Order{Epoch: 41, Sequence: 12})
if err == nil || !strings.Contains(err.Error(), "epoch 41") || !strings.Contains(err.Error(), "epoch 57") {
t.Fatalf("an older epoch was not refused by name: %v", err)
}
if err := refuseOlder(kept, nil, link.Order{Epoch: 58, Sequence: 1}); err != nil {
t.Fatalf("a new lease holder's first declaration was refused for its sequence: %v", err)
}
if err := refuseOlder(kept, nil, link.Order{Epoch: 57, Sequence: 2}); err == nil {
t.Fatal("a lower sequence under the same epoch was accepted")
}
// A controller with no lease — older, or rolled back to a build from before it — is today's
// behaviour: no epoch compared.
if err := refuseOlder(kept, nil, link.Order{Sequence: 4}); err != nil {
t.Fatalf("a declaration with no epoch was refused for one: %v", err)
}
if held, older := olderThanKept(kept, nil, link.Order{Epoch: 41, Sequence: 12}); !older ||
held != (link.Order{Epoch: 57, Sequence: 3}) {
t.Fatalf("olderThanKept said %v, holding %+v", older, held)
}
// What is kept cannot be read: no order to compare, so nothing is refused on a guess.
if _, older := olderThanKept(store.Declared{Declaration: []byte("garbled")}, nil,
link.Order{Epoch: 41, Sequence: 1}); older {
t.Fatal("refused against a kept declaration nobody can read")
}
}
+1 -1
View File
@@ -161,7 +161,7 @@
"type": "file",
"path": "/var/lib/mesh-bus-conf/accounts.conf",
"mode": "0600",
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.control.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\"] }\n subscribe: { allow: [\"$JS.API.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
},
{
"id": "broker",
+21 -1
View File
@@ -1297,6 +1297,16 @@ type Declaration struct {
// superseded.
Sequence int64
// Epoch is the controller's lease epoch this declaration was sent under (novox/hq to-be 45 §6):
// the revision at which the sending controller took the lease. A controller that lost its lease
// and goes on sending sends an older epoch than the holder's, and the node-engine refuses what
// is older than what it applied. Zero is a declaration from a controller without a lease — every
// one sent before the lease existed — and carries no claim.
//
// Inside what is signed, beside the sequence, so a message cannot be given a newer epoch than
// the controller gave it.
Epoch int64
// LeftOut names the modules of this machine's set the mesh left out of this declaration,
// because a setting stored for one cannot compose with its definition (novox/hq ADR 0163,
// rule 6). A machine is told everything or nothing about what it IS told; this is what it is
@@ -1462,6 +1472,10 @@ type envelope struct {
// Sequence is optional on the wire, so a controller that does not send one is still
// understood: absent reads as zero, which is "no ordering claimed" rather than "first".
Sequence int64 `json:"sequence,omitempty"`
// Epoch is optional on the wire as the sequence is: absent is a controller without a lease.
// **An older host refuses this key**, decoding strictly; a controller sends it only to a host
// whose reports carry a report sequence, which a host that reads it does.
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"`
}
@@ -1481,8 +1495,14 @@ func parse(raw []byte, allowActions bool) (*Declaration, error) {
}
d := &Declaration{Version: env.Version, For: env.For, Adoption: env.Adoption, Sequence: env.Sequence,
LeftOut: env.LeftOut}
Epoch: env.Epoch, LeftOut: env.LeftOut}
var problems []string
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.
problems = append(problems, fmt.Sprintf("an order below zero (epoch %d, sequence %d) is not "+
"one the mesh assigns", env.Epoch, env.Sequence))
}
if len(env.LeftOut) > 0 && allowActions {
// The bundle is carried with the binary and leaves nothing out: which module a setting
// stopped composing for is the mesh's record (ADR 0163).
+37
View File
@@ -0,0 +1,37 @@
package declaration
import (
"strings"
"testing"
)
// A declaration says where it stands — the controller's lease epoch and its sequence — inside what is
// signed (novox/hq to-be 45 §6). Absent is no claim; below zero is no order the mesh assigns.
func withOrder(order string) string {
return `{"declaration":1,` + order + `"resources":[{"id":"etc","type":"directory","path":"/etc/mesh","mode":"0755"}]}`
}
func TestADeclarationCarriesItsEpochAndSequence(t *testing.T) {
d, err := Parse([]byte(withOrder(`"epoch":57,"sequence":3,`)))
if err != nil {
t.Fatalf("a declaration carrying its order was refused: %v", err)
}
if d.Epoch != 57 || d.Sequence != 3 {
t.Fatalf("read epoch %d, sequence %d", d.Epoch, d.Sequence)
}
// An older controller's: no order, and no claim.
d, err = Parse([]byte(withOrder(``)))
if err != nil || d.Epoch != 0 || d.Sequence != 0 {
t.Fatalf("a declaration with no order read as %+v, %v", d, err)
}
}
func TestAnOrderBelowZeroIsRefused(t *testing.T) {
for _, order := range []string{`"epoch":-1,`, `"sequence":-4,`} {
refusal := refusalFor(t, withOrder(order))
if !strings.Contains(refusal.Error(), "below zero") {
t.Errorf("%s: refused for something else: %v", order, refusal)
}
}
}
+116 -2
View File
@@ -6,6 +6,7 @@ import (
"encoding/json"
"fmt"
"os"
"reflect"
"testing"
"time"
@@ -194,8 +195,8 @@ func TestNatsANodeThatWasAwayGetsOnlyTheNewest(t *testing.T) {
select {
case msg := <-feed:
if declaredIn(msg.Data) != "d3" {
t.Fatalf("the node was given %q rather than the newest", declaredIn(msg.Data))
if want := declaredIn(signedBy(t, private, []byte(`{"declared":"d3"}`))); declaredIn(msg.Data) != want {
t.Fatalf("the node was given %q rather than the newest", msg.Data)
}
case <-time.After(8 * time.Second):
t.Fatal("the node that was away was given nothing")
@@ -226,3 +227,116 @@ func TestNatsAReportIsKeptAndAHeartbeatIsNot(t *testing.T) {
t.Fatalf("%d messages were kept; a report must be and a heartbeat must not", info.State.Msgs)
}
}
// **One apply queue, against a real bus** (novox/hq to-be 45 §6, R2). A declaration is being applied
// when two more are pushed and the five-minute reconcile comes due: the worker applies the one in
// hand, then only the newest — the reconcile is that apply — and the mesh's stream holds one account
// per apply, each naming the order it applied, numbered in the order they were made. Every
// declaration delivered is settled.
func TestNatsTheQueueAppliesTheNewestOnceAndReportsInOrder(t *testing.T) {
conn, js := aBus(t)
const node = "queueing"
theMeshMakes(t, js, node)
public, private, _ := ed25519.GenerateKey(nil)
m := Membership{Node: node, Signer: public}
ctx, stop := context.WithCancel(context.Background())
defer stop()
l := &natsLink{conn: conn, js: js, node: node,
arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)}
feed := make(chan *nats.Msg, drainDepth)
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
go func() {
for msg := range feed {
l.arrived <- natsDeclaration{msg}
}
}()
reports, err := js.SubscribeSync(ReportSubject(node), nats.BindStream("CONTROL"))
if err != nil {
t.Fatal(err)
}
defer func() { _ = reports.Unsubscribe() }()
a := &applies{}
inFirst, release := make(chan struct{}, 1), make(chan struct{})
reconciled := 0
q := &Queue{Membership: m, Timeout: 5 * time.Second, Unsaid: &keptInMemory{},
Apply: func(ctx context.Context, raw, sig []byte) Report {
r := a.apply(ctx, raw, sig)
if r.Sequence == 1 {
inFirst <- struct{}{}
<-release
}
return r
},
Reconcile: func(context.Context) (Report, bool) { reconciled++; return Report{}, false },
}
worker := runWorker(ctx, q)
go func() { _ = serve(ctx, l, m, q, nil, 5*time.Second) }()
push := func(id string, sequence int64) {
inner, _ := json.Marshal(map[string]any{"declaration": 1, "epoch": 7, "sequence": sequence,
"resources": []any{map[string]any{"id": id}}})
if _, err := js.Publish(DeclareSubject(node), signedBy(t, private, inner)); err != nil {
t.Fatal(err)
}
}
push("one", 1)
select {
case <-inFirst:
case <-time.After(8 * time.Second):
t.Fatal("the first declaration never reached the worker")
}
push("two", 2)
push("three", 3)
q.ReconcileDue()
time.Sleep(2 * drainWindow) // the link gathers both and hands them to the queue
close(release)
var heard []Report
for len(accounts(heard)) < 2 {
msg, err := reports.NextMsg(8 * time.Second)
if err != nil {
t.Fatalf("the mesh heard %+v and then nothing: %v", heard, err)
}
var r Report
if err := json.Unmarshal(msg.Data, &r); err != nil {
t.Fatal(err)
}
heard = append(heard, r)
_ = msg.Ack()
}
stop()
<-worker
if got := a.all(); !reflect.DeepEqual(got, []string{"one", "three"}) {
t.Fatalf("applied %v; the one in hand and then only the newest were wanted", got)
}
if reconciled != 0 {
t.Fatal("the reconcile applied beside the delivery it was due with")
}
accounted := accounts(heard)
if accounted[0].Sequence != 1 || accounted[1].Sequence != 3 || accounted[1].Epoch != 7 {
t.Fatalf("the accounts do not name what was applied: %+v", accounted)
}
for i := 1; i < len(heard); i++ {
if heard[i].ReportSequence <= heard[i-1].ReportSequence {
t.Fatalf("the reports left out of the order they were made: %+v", heard)
}
}
deadline := time.Now().Add(5 * time.Second)
for {
info, err := js.ConsumerInfo("NODES", node)
if err == nil && info.NumAckPending == 0 {
break
}
if time.Now().After(deadline) {
t.Fatalf("a delivered declaration was left unsettled: %+v %v", info, err)
}
time.Sleep(20 * time.Millisecond)
}
}
+92
View File
@@ -28,6 +28,10 @@ const (
// look like a stale one.
type Alive struct {
Node string `json:"node"`
// IntervalSeconds is how often this node says it is there (novox/hq to-be 45 §3, S1): the
// controller's watchdog is bound to three of them, so the bound moves with the interval rather
// than with a number the controller guessed. A controller older than that ignores it.
IntervalSeconds int `json:"interval_seconds"`
}
// Signed is a declaration and the signature over it.
@@ -92,6 +96,36 @@ type Report struct {
// cannot answer "which"; the digest is the answer itself.
Declared string `json:"declared,omitempty"`
// Order is where the declaration this report is about stands — its `epoch` and `sequence`, as
// the declaration carried them — so the mesh orders reports by what they are about rather than
// by when they arrived (novox/hq to-be 45 §6). On the wire as two top-level keys beside
// `declared`; absent for a declaration that claimed no order, and for a report about none.
Order
// ReportSequence is this node-engine's own order among everything it reports: one higher for
// every report it makes, kept on disk so it goes on increasing across restarts and self-updates
// (to-be 45 §6). With the declaration's order it lets the mesh refuse an older report about the
// same declaration — a reconcile's account overtaking the apply that followed it (novox/hq issue
// 267) is then refused by number rather than set aside by digest. A report said again because
// it never reached the mesh (issue 264) carries the number it was made with.
//
// Zero claims no order: every report an older host makes, and a one-shot report a command makes
// (a rekey). **Its presence also says this host reads a declaration's `epoch`**, which is how the
// mesh knows it may send one: an older host decodes a declaration strictly and refuses a key it
// does not know.
ReportSequence int64 `json:"report_sequence,omitempty"`
// OlderThan is set on a report refusing a declaration older than one this node has applied: the
// order of the newer declaration it holds (to-be 45 §6, rule 2). The refused declaration is the
// report's own Declared and Order, and Refused says it in words. Applying it would make the
// machine into something a controller that lost its lease said, after one that holds it.
OlderThan *Order `json:"older_than,omitempty"`
// RefusedOlder is how many declarations this node-engine has refused as older, ever — kept on
// disk with the report sequence. On every report, so a count missed with a lost refusal is still
// read from the next one: what the mesh's stale-writer watchdog (to-be 45 §3, S13) counts.
RefusedOlder int64 `json:"refused_older,omitempty"`
// Held is what this 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.
Held []Held `json:"held,omitempty"`
@@ -162,6 +196,64 @@ type Report struct {
Rekey *Rekey `json:"rekey,omitempty"`
}
// Order is where a declaration stands among everything the mesh has sent this node (novox/hq to-be
// 45 §6): the controller's lease **epoch** it was sent under, and its **sequence** — one higher for
// every send to this node (04-ISSUES/107). The declaration carries both inside what is signed, as
// top-level `epoch` and `sequence` beside `declaration`; a report carries back the pair of the
// declaration it is about.
//
// Zero in either is "no order claimed", never "first": every declaration an older controller sent,
// and the bundle genesis applies.
type Order struct {
Epoch int64 `json:"epoch,omitempty"`
Sequence int64 `json:"sequence,omitempty"`
}
// Older says whether a declaration of this order is older than one of order `than`, which this node
// has applied — the one rule by which a node-engine refuses a declaration (to-be 45 §6: "older epoch;
// same epoch, lower sequence").
//
// **Only when both claim an epoch.** A declaration with none is an older controller's — or a
// controller rolled back to a build from before the lease — and one kept with none is what this host
// held before any controller had a lease. Neither has an order to compare, and refusing on a guess
// would strand the machine the moment the controller is older than the host: without an epoch, today's
// behaviour stands. A higher epoch is a new lease holder and is never older, whatever its sequence.
func (o Order) Older(than Order) bool {
if o.Epoch <= 0 || than.Epoch <= 0 {
return false
}
if o.Epoch != than.Epoch {
return o.Epoch < than.Epoch
}
return o.Sequence > 0 && than.Sequence > 0 && o.Sequence < than.Sequence
}
// Supersedes says whether a declaration of order `o`, arriving after one of order `before`, takes its
// place among what is waiting to be applied. By epoch when both claim one and they differ, then by
// sequence when both claim one, and by arrival when they do not — what the drain had to go on before
// declarations said where they stand (04-ISSUES/107).
func (o Order) Supersedes(before Order) bool {
if o.Epoch > 0 && before.Epoch > 0 && o.Epoch != before.Epoch {
return o.Epoch > before.Epoch
}
if o.Sequence > 0 && before.Sequence > 0 {
return o.Sequence >= before.Sequence
}
return true
}
// Words is an order as a person reads it in a log line. Not String: Report embeds Order, and a
// Stringer promoted onto every report would print each one as its order alone.
func (o Order) Words() string {
switch {
case o.Epoch > 0:
return "epoch " + strconv.FormatInt(o.Epoch, 10) + ", sequence " + strconv.FormatInt(o.Sequence, 10)
case o.Sequence > 0:
return "sequence " + strconv.FormatInt(o.Sequence, 10) + ", no epoch"
}
return "no order"
}
// CarriedTunnel is this node's account of the tunnel it took over. State is one of the Carried
// states below; Note is what the host did about it, when it did something.
type CarriedTunnel struct {
+21 -5
View File
@@ -1,19 +1,30 @@
package link
import (
"sync"
"testing"
"time"
)
// said is one declaration as a test hands it over, with no transport under it — which is what the
// seam bought: the drain's reasoning was reachable only through a real broker before.
type said struct {
body []byte
body []byte
mu sync.Mutex
handled bool
}
func (s *said) Body() []byte { return s.body }
func (s *said) Handled() error { s.handled = true; return nil }
func (s *said) Body() []byte { return s.body }
func (s *said) Handled() error {
s.mu.Lock()
defer s.mu.Unlock()
s.handled = true
return nil
}
func (s *said) wasHandled() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.handled
}
func arriving(bodies ...string) chan Declaration {
ch := make(chan Declaration, 8)
@@ -23,6 +34,11 @@ func arriving(bodies ...string) chan Declaration {
return ch
}
// newest is what the queue's worker would take from what the link gathered behind `first`.
func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) {
return pick(gather(arriving, first, window))
}
// A machine asked to be five things becomes the last one: what is already waiting supersedes what
// arrived first, and everything set aside is named so it can be reported.
func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) {
+3 -3
View File
@@ -27,7 +27,7 @@ func TestTheDrainKeepsTheHighestSequenceNotTheLastToArrive(t *testing.T) {
waiting <- sequenced(t, 9)
waiting <- sequenced(t, 4) // arrived last, composed earlier
latest, aside := newest(waiting, sequenced(t, 8), 30*time.Millisecond)
if got := sequenceOf(latest.Body()); got != 9 {
if got := orderOf(latest.Body()).Sequence; got != 9 {
t.Fatalf("the drain kept sequence %d, and 9 was waiting", got)
}
if len(aside) != 2 {
@@ -44,7 +44,7 @@ func TestWithoutSequencesTheLastToArriveStillWins(t *testing.T) {
}
func TestAnUnreadableBodyClaimsNoOrder(t *testing.T) {
if got := sequenceOf([]byte("not json")); got != 0 {
t.Fatalf("garbage claimed sequence %d", got)
if got := orderOf([]byte("not json")); got != (Order{}) {
t.Fatalf("garbage claimed %+v", got)
}
}
-94
View File
@@ -1,94 +0,0 @@
package link
import (
"context"
"testing"
"time"
)
// A reconcile's report about the declaration kept before a delivery is never said after that
// delivery's report (novox/hq issue 267).
// Measured on the home server: the reconcile timer fired three seconds before a declaration
// arrived. The reconcile held the machine to the declaration kept then; the delivery's apply waited
// for it, applied the new one and was reported — and then the reconcile's report, queued meanwhile,
// went out and was stored as the machine's latest account, naming the older declaration. The
// release plan waited on a report it had already been given.
func TestAReconcileReportOlderThanTheApplyIsNotSaidAfterIt(t *testing.T) {
m, key := aMember(t)
l := newQuietLink(&said{body: signedBy(t, key, []byte(`{"declaration":2}`))})
outbox := make(chan Unasked, 1)
ctx, stop := context.WithCancel(context.Background())
defer stop()
settled := make(chan bool, 1)
apply := func(context.Context, []byte, []byte) Report {
// The reconcile ran first, on what was kept then, and queued its report while this waited.
outbox <- Unasked{Report: Report{Declared: "d1", Applied: []string{"a"}, Outward: []string{"eth0"}},
Done: func(published bool) { settled <- published; stop() }}
return Report{Declared: "d2", Applied: []string{"a", "b"}, Outward: []string{"eth0"}}
}
done := make(chan error, 1)
go func() { done <- serve(ctx, l, m, apply, nil, time.Second, outbox, &keptInMemory{}) }()
select {
case published := <-settled:
if published {
t.Fatal("the reconcile's report was counted as said")
}
case <-time.After(5 * time.Second):
t.Fatal("the reconcile's report was never settled")
}
if err := <-done; err != nil {
t.Fatal(err)
}
reports := l.said()
if len(reports) != 1 || reports[0].Declared != "d2" {
t.Fatalf("the mesh heard %+v; it should have heard only the apply of d2", reports)
}
}
// A reconcile's report about the declaration this link applied is news, and is said.
func TestAReconcileReportAboutTheAppliedDeclarationIsSaid(t *testing.T) {
m, key := aMember(t)
l := newQuietLink(&said{body: signedBy(t, key, []byte(`{"declaration":2}`))})
outbox := make(chan Unasked, 1)
ctx, stop := context.WithCancel(context.Background())
defer stop()
settled := make(chan bool, 1)
apply := func(context.Context, []byte, []byte) Report {
outbox <- Unasked{Report: Report{Declared: "d2", Applied: []string{"a", "b"}, Firewall: "nftables"},
Done: func(published bool) { settled <- published; stop() }}
return Report{Declared: "d2", Applied: []string{"a", "b"}}
}
go func() { _ = serve(ctx, l, m, apply, nil, time.Second, outbox, nil) }()
select {
case published := <-settled:
if !published {
t.Fatal("a reconcile's report about the applied declaration was not said")
}
case <-time.After(5 * time.Second):
t.Fatal("the reconcile's report was never settled")
}
if reports := l.said(); len(reports) != 2 || reports[1].Declared != "d2" {
t.Fatalf("the mesh heard %+v", reports)
}
}
func TestOvertaken(t *testing.T) {
for _, c := range []struct {
declared, applied string
want bool
}{
{"d1", "d2", true},
{"d2", "d2", false},
{"", "d2", false}, // names no declaration
{"d1", "", false}, // this link has applied nothing yet
} {
if got := overtaken(Report{Declared: c.declared}, c.applied); got != c.want {
t.Errorf("overtaken(%q, %q) = %v, want %v", c.declared, c.applied, got, c.want)
}
}
}
+450
View File
@@ -0,0 +1,450 @@
package link
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"sync"
"time"
)
// Queue is the node-engine's one apply queue (novox/hq to-be 45 §6, ADR 0227 rule 1).
//
// **One worker applies; everything else asks it to.** A declaration delivered over the link and the
// five-minute reconcile are two reasons to *enqueue* the same act, never two paths that apply. Before
// this they were two, each taking an apply lock, and every order between them was a fault met on the
// live mesh: the reconcile read what was kept before a delivery and applied it over the newer one
// (novox/hq issues 257, 261); the reconcile went first and its report about the older declaration
// reached the mesh after the delivery's (267). Each fix closed one order and left the class open.
//
// The queue holds at most one pending request, coalesced: whatever has been delivered and not yet
// taken, and whether a reconcile is due. When the worker starts it takes the **newest declaration held
// at that moment**, sets the rest aside — each reported as set aside, then settled unapplied — applies
// it once, and makes **one report**, naming the declaration's order and its own report sequence. A
// reconcile due while a delivery waits is that delivery's apply: applying the newest thing the mesh
// said is what reconciling is. Only with nothing delivered does a reconcile apply what was kept, read
// once its turn has come.
//
// **Reports leave in the order they are made**, because the worker makes and says them, one at a time.
// A report made while the link is down waits: an apply's in the unsaid store (issue 264), a
// reconcile's in one slot replaced by anything newer, and both are said, in that order, when the next
// link opens — before anything newly delivered is applied.
type Queue struct {
// Membership is whose declarations these are: their signer, and the node the reports are from.
Membership Membership
// Apply applies a declaration the mesh signed.
Apply Applier
// Reconcile holds the machine to what it was last told, read once it is this act's turn. False
// when there is nothing to hold it to, or the reconcile could not run; it says why itself.
Reconcile Reconcile
// News is whether a reconcile's report is worth saying unasked (novox/hq ADR 0100, ADR 0140). Nil
// is never: a reconcile is otherwise silent.
News func(Report) bool
// Heard is told every account of the machine the mesh has taken — an apply's or a reconcile's,
// not a refusal or a declaration set aside, which say nothing about the machine. Nil is allowed.
Heard func(Report)
// Unsaid keeps an apply's report until the mesh has taken it (issue 264). Nil keeps nothing.
Unsaid Unsaid
// Numbers keeps the report sequence and the refusals counted, across restarts and self-updates.
// Nil counts in memory, which is only right where nothing outlives the process — a test.
Numbers Numbers
Say Announce
Timeout time.Duration
once sync.Once
wake chan struct{}
mu sync.Mutex
idle *sync.Cond
// delivering is a delivered declaration in hand: its report is owed on the link it came over, so
// that link is not let go until it is said (detach). A reconcile in hand holds no link — its report
// waits for the next one if this one goes.
delivering bool
waiting []Declaration // delivered and not yet taken, in the order they arrived
due bool // a reconcile asked for
bus Bus // the link open now; nil while there is none
unasked *Report // a reconcile's report not yet said, made while there was no link
// saying is held while a report is made and said, so the order they leave in is the order they
// are made in — whoever says them, the worker or a link opening.
saying sync.Mutex
counted bool
numbers numbered
}
// Reconcile is the reconcile's act, run by the queue's worker. See Queue.Reconcile.
type Reconcile func(ctx context.Context) (Report, bool)
// Numbers is where the queue keeps what it counts, so a successor goes on from where it stopped.
type Numbers interface {
// Read is what was kept: zero for a node that never reported, an error when it cannot be read —
// never zero for "could not tell".
Read() (reportSequence, refusedOlder int64, err error)
// Save keeps them, on disk when it returns.
Save(reportSequence, refusedOlder int64) error
}
type numbered struct{ sequence, refused int64 }
func (q *Queue) init() {
q.once.Do(func() {
q.wake = make(chan struct{}, 1)
q.idle = sync.NewCond(&q.mu)
if q.Say == nil {
q.Say = func(string) {}
}
})
}
// Deliver enqueues what the link delivered, in the order it arrived. It never waits for an apply.
func (q *Queue) Deliver(arrived ...Declaration) {
q.init()
q.mu.Lock()
q.waiting = append(q.waiting, arrived...)
q.mu.Unlock()
q.signal()
}
// ReconcileDue enqueues a reconcile. One asked for while another is waiting is the same one.
func (q *Queue) ReconcileDue() {
q.init()
q.mu.Lock()
q.due = true
q.mu.Unlock()
q.signal()
}
func (q *Queue) signal() {
select {
case q.wake <- struct{}{}:
default:
// One is already waiting to be read, and the worker takes everything pending when it does.
}
}
// Run is the worker. It returns once the context ends and the act in hand is finished and said:
// asked to stop — the host standing aside for its successor — nothing further is applied.
func (q *Queue) Run(ctx context.Context) {
q.init()
for {
for q.step(ctx) {
}
select {
case <-ctx.Done():
return
case <-q.wake:
}
}
}
// step takes what is pending and does it, once. False when there was nothing to do, or the queue
// was asked to stop.
func (q *Queue) step(ctx context.Context) bool {
q.mu.Lock()
if ctx.Err() != nil || (len(q.waiting) == 0 && !q.due) {
q.mu.Unlock()
return false
}
batch, due := q.waiting, q.due
q.waiting, q.due = nil, false
q.delivering = len(batch) > 0
q.mu.Unlock()
defer func() {
q.mu.Lock()
q.delivering = false
q.idle.Broadcast()
q.mu.Unlock()
}()
if len(batch) == 0 {
q.reconcile(ctx)
return true
}
latest, aside := pick(batch)
for _, old := range aside {
if d := declaredIn(old.Body()); d != "" && d == declaredIn(latest.Body()) {
// The same declaration again — delivered twice while the worker was busy. Nothing is set
// aside: it is about to be applied.
_ = old.Handled()
continue
}
q.Say(fmt.Sprintf("set aside declaration %s (%s): a newer one arrived with it",
short(declaredIn(old.Body())), orderOf(old.Body()).Words()))
q.tell(ctx, Report{Declared: declaredIn(old.Body()), Order: orderOf(old.Body()),
Superseded: declaredIn(latest.Body())}, setAside, old)
}
report := handleBody(ctx, q.Membership, latest.Body(), q.Apply)
switch {
case report.OlderThan != nil:
// **Refused, counted and said** (rule 2): one line here, the count on this and every later
// report, and the report itself naming what was refused and what is held.
total := q.countRefusal()
q.Say(fmt.Sprintf("refused declaration %s (%s): this node applied %s, from a newer lease "+
"holder — %d refused as older so far", short(report.Declared), report.Order.Words(),
report.OlderThan.Words(), total))
case report.Refused != "":
q.Say("refused a declaration: " + report.Refused)
case len(report.Failed) > 0:
q.Say(fmt.Sprintf("applied %d and failed: %v%s",
len(report.Applied), report.Failed, heldNote(report.Held)))
default:
q.Say(fmt.Sprintf("applied %d resource(s)%s", len(report.Applied), heldNote(report.Held)))
}
if due && report.Refused != "" {
// Nothing was applied, so the reconcile this stood in for has not happened. It runs next.
q.mu.Lock()
q.due = true
q.mu.Unlock()
}
// A refusal as older is not kept to be said again: what the mesh waits for from this node is the
// account of the last apply, and that is still what is kept.
kind := applied
if report.OlderThan != nil {
kind = refusedOlder
}
q.tell(ctx, report, kind, latest)
return true
}
func (q *Queue) reconcile(ctx context.Context) {
if q.Reconcile == nil {
return
}
report, ok := q.Reconcile(ctx)
if !ok || q.News == nil || !q.News(report) {
return
}
q.tell(ctx, report, unasked)
}
// What a report is, for what tell does with it besides saying it.
type kindOf int
const (
// applied is an apply's account: kept until said (issue 264), and fresher than any reconcile's
// report still waiting.
applied kindOf = iota
// setAside is a declaration a newer one took the place of, unapplied.
setAside
// refusedOlder is a declaration refused as older than what was applied. Not kept: what the mesh
// waits for from this node is still the account of the last apply.
refusedOlder
// unasked is a reconcile's report. Not said — no link, or the bus refused it — the newest waits
// for the next link, and anything the worker says before then replaces it.
unasked
)
// tell makes a report this node's — its node, its report sequence, the refusals counted — and says it
// on the link open now. Kept first, when it is an apply's, so a report lost between here and the bus
// is said by whichever host links next. Each declaration it settles is settled after it is said,
// either way. True when the broker took it.
func (q *Queue) tell(ctx context.Context, r Report, kind kindOf, settle ...Declaration) bool {
q.saying.Lock()
defer q.saying.Unlock()
r = q.stamp(r)
keep := kind == applied
if keep {
// What this says supersedes any reconcile's report still waiting to be said.
q.unasked = nil
if q.Unsaid != nil {
if err := q.Unsaid.Keep(r); err != nil {
q.Say("cannot keep this apply's report until it is said: " + err.Error())
}
}
}
q.mu.Lock()
bus := q.bus
q.mu.Unlock()
published := false
if bus != nil {
// **Said even when the link is being let go** (issue 264): not cancelled with it, still
// bounded by the timeout, and the link is not closed until this returns (detach).
published = publishReport(context.WithoutCancel(ctx), bus, q.Membership, r, q.Say, q.Timeout)
}
if !published && kind == unasked {
q.unasked = &r
}
if published {
if keep && q.Unsaid != nil {
if err := q.Unsaid.Said(r.Declared); err != nil {
q.Say("said this apply's report, and cannot forget it: " + err.Error())
}
}
if q.Heard != nil && (kind == applied || kind == unasked) && r.Refused == "" {
q.Heard(r)
}
}
// Settled after the report is published. A node that dies between applying and reporting
// leaves the declaration with the mesh and applies it again on return, which is safe because
// applying is reconciliation — it converges rather than repeating.
for _, d := range settle {
_ = d.Handled()
}
return published
}
// stamp gives a report its node, the next report sequence and the refusals counted. Called holding
// saying.
func (q *Queue) stamp(r Report) Report {
q.load()
q.numbers.sequence++
q.keepNumbers()
r.Node = q.Membership.Node
r.ReportSequence = q.numbers.sequence
r.RefusedOlder = q.numbers.refused
return r
}
func (q *Queue) countRefusal() int64 {
q.saying.Lock()
defer q.saying.Unlock()
q.load()
q.numbers.refused++
q.keepNumbers()
return q.numbers.refused
}
// load reads what was counted, once. **Unreadable is said, and the sequence goes on from above
// anything it can have reached** — the time in milliseconds — rather than from one, which would make
// every report this host says next older than what it already said, and refused.
func (q *Queue) load() {
if q.counted {
return
}
q.counted = true
if q.Numbers == nil {
return
}
sequence, refused, err := q.Numbers.Read()
if err != nil {
floor := time.Now().UnixMilli()
q.Say(fmt.Sprintf("cannot read this node's report sequence: %v — going on from %d, above any "+
"it can have reached, and counting refusals again from none", err, floor))
q.numbers = numbered{sequence: floor}
return
}
q.numbers = numbered{sequence: sequence, refused: refused}
}
func (q *Queue) keepNumbers() {
if q.Numbers == nil {
return
}
if err := q.Numbers.Save(q.numbers.sequence, q.numbers.refused); err != nil {
// Said, and not fatal: the report is still worth saying. What is lost is that a successor
// could number one again.
q.Say("cannot keep this node's report sequence: " + err.Error())
}
}
// attach is a link opening. **What the last apply did, if the mesh never heard it**, is said first,
// exactly as it was made; then a reconcile's report that waited; and only then may the worker say
// anything on it — so a report said again can never land after, and read as newer than, a later one.
func (q *Queue) attach(ctx context.Context, bus Bus) {
q.init()
q.saying.Lock()
defer q.saying.Unlock()
reporting := context.WithoutCancel(ctx)
q.load()
said := true
if q.Unsaid != nil {
switch kept, ok, err := q.Unsaid.Pending(); {
case err != nil:
q.Say("cannot read the report this node kept unsaid: " + err.Error())
case ok:
q.Say("saying again what the last apply did: its report never reached the mesh")
if said = publishReport(reporting, bus, q.Membership, kept, q.Say, q.Timeout); said {
if err := q.Unsaid.Said(kept.Declared); err != nil {
q.Say("said the kept report, and cannot forget it: " + err.Error())
}
if q.Heard != nil && kept.Refused == "" {
q.Heard(kept)
}
}
}
}
if said && q.unasked != nil {
if publishReport(reporting, bus, q.Membership, *q.unasked, q.Say, q.Timeout) {
if q.Heard != nil {
q.Heard(*q.unasked)
}
q.unasked = nil
}
}
q.mu.Lock()
q.bus = bus
q.mu.Unlock()
}
// detach is the link closing: it waits for a delivery in hand to be applied and said — an apply that
// stood aside for its successor is reported on this link before it is let go (issue 264) — and then
// nothing more is said on it. A reconcile in hand is not waited for: a link that was dropped or roused
// is opened again at once, and the reconcile's report, if it is news, waits for it.
func (q *Queue) detach(bus Bus) {
q.init()
q.mu.Lock()
for q.delivering {
q.idle.Wait()
}
q.mu.Unlock()
q.saying.Lock()
defer q.saying.Unlock()
q.mu.Lock()
if q.bus == bus {
q.bus = nil
}
q.mu.Unlock()
}
// pick is the newest of what was delivered and the rest, in arrival order, to set aside: by order
// when the declarations claim one, by arrival when they do not (Order.Supersedes).
func pick(batch []Declaration) (Declaration, []Declaration) {
latest := batch[0]
var aside []Declaration
for _, next := range batch[1:] {
if orderOf(next.Body()).Supersedes(orderOf(latest.Body())) {
aside = append(aside, latest)
latest = next
continue
}
aside = append(aside, next)
}
return latest, aside
}
// orderOf is the order a signed declaration claims, or none when it claims none or cannot be read.
// Read from the envelope alone; the signature is verified when the winner is applied, and a forged
// message that lied about its order would only set aside real ones — which are reported as set aside,
// and the next push sends the current one again.
func orderOf(body []byte) Order {
var signed Signed
if err := json.Unmarshal(body, &signed); err != nil {
return Order{}
}
var o Order
if err := json.Unmarshal(signed.Declaration, &o); err != nil || o.Epoch < 0 || o.Sequence < 0 {
return Order{}
}
return o
}
// declaredIn names a signed declaration as the mesh does — the digest of the declaration's bytes, the
// same the apply's report carries — for a report about one that was not applied. Empty if the message
// is not one: a forged or garbled message is refused when its turn comes; here it is only named.
func declaredIn(body []byte) string {
var signed Signed
if err := json.Unmarshal(body, &signed); err != nil || len(signed.Declaration) == 0 {
return ""
}
sum := sha256.Sum256(signed.Declaration)
return hex.EncodeToString(sum[:])
}
+542
View File
@@ -0,0 +1,542 @@
package link
import (
"context"
"crypto/ed25519"
"encoding/json"
"errors"
"reflect"
"strings"
"sync"
"testing"
"time"
)
// The node-engine's one apply queue (novox/hq to-be 45 §6): a delivery and the reconcile are two
// reasons to enqueue the same act, the newest declaration is applied once, and the reports leave in
// the order they are made. These are the races of issues 257, 261, 264 and 267 replayed against it.
func aQueue(m Membership, apply Applier, kept Unsaid) *Queue {
return &Queue{Membership: m, Apply: apply, Unsaid: kept, Timeout: time.Second}
}
// runWorker runs the queue's worker until ctx ends; the channel closes when it has.
func runWorker(ctx context.Context, q *Queue) <-chan struct{} {
done := make(chan struct{})
go func() {
defer close(done)
q.Run(ctx)
}()
return done
}
// ordered is a signed declaration claiming an order, with a resource id so two are distinct bytes.
func ordered(t *testing.T, key ed25519.PrivateKey, id string, epoch, sequence int64) *said {
t.Helper()
inner, err := json.Marshal(map[string]any{"declaration": 1, "resources": []any{map[string]any{"id": id}},
"epoch": epoch, "sequence": sequence})
if err != nil {
t.Fatal(err)
}
return &said{body: signedBy(t, key, inner)}
}
// applies records what the worker applied, by the id of the declaration's first resource, and
// reports it as applied — with the order the declaration carried, as the host's own apply does.
type applies struct {
mu sync.Mutex
seen []string
}
func (a *applies) apply(_ context.Context, raw, _ []byte) Report {
var d struct {
Resources []struct {
ID string `json:"id"`
} `json:"resources"`
Order
}
_ = json.Unmarshal(raw, &d)
id := ""
if len(d.Resources) > 0 {
id = d.Resources[0].ID
}
a.mu.Lock()
a.seen = append(a.seen, id)
a.mu.Unlock()
return Report{Declared: digest(raw), Order: d.Order, Applied: []string{id}}
}
func (a *applies) all() []string {
a.mu.Lock()
defer a.mu.Unlock()
return append([]string(nil), a.seen...)
}
func digest(raw []byte) string {
b, _ := json.Marshal(Signed{Declaration: raw})
return declaredIn(b)
}
// eventually waits for a condition the worker reaches on its own time.
func eventually(t *testing.T, what string, ok func() bool) {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for !ok() {
if time.Now().After(deadline) {
t.Fatal(what)
}
time.Sleep(5 * time.Millisecond)
}
}
// accounts is the reports the mesh heard that are an apply's account, not a set-aside's.
func accounts(reports []Report) []Report {
var out []Report
for _, r := range reports {
if r.Superseded == "" {
out = append(out, r)
}
}
return out
}
// **R2: a reconcile due during a delivery** (issues 257, 261, 267). The reconcile is asked for while a
// delivery waits: the worker applies the delivery once — applying the newest thing the mesh said is
// what reconciling is — and makes one report, naming the newest sequence. The reconcile never applies
// what was kept before.
func TestAReconcileDueDuringADeliveryIsThatDeliverysApply(t *testing.T) {
m, key := aMember(t)
l := newQuietLink()
a := &applies{}
reconciled := 0
q := aQueue(m, a.apply, &keptInMemory{})
q.Reconcile = func(context.Context) (Report, bool) {
reconciled++
return Report{Declared: "the one kept before", Applied: []string{"old"}}, true
}
q.News = func(Report) bool { return true }
q.attach(context.Background(), l)
q.Deliver(ordered(t, key, "new", 7, 12))
q.ReconcileDue()
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 1 })
stop()
<-worker
if got := a.all(); !reflect.DeepEqual(got, []string{"new"}) {
t.Fatalf("applied %v; one apply of the newest was wanted", got)
}
if reconciled != 0 {
t.Fatalf("the reconcile applied what was kept %d time(s) beside the delivery", reconciled)
}
r := l.said()[0]
if r.Sequence != 12 || r.Epoch != 7 || r.ReportSequence != 1 || r.Node != m.Node {
t.Fatalf("the report does not name the order it applied: %+v", r)
}
}
// **The other order** (issue 267): the reconcile is running when a delivery arrives. Its report, about
// the declaration kept then, leaves first; the delivery's leaves after it, with a higher report
// sequence. The mesh's latest word is the newest declaration, whichever arrived first.
func TestAReconcileRunningWhenADeliveryArrivesIsSaidBeforeIt(t *testing.T) {
m, key := aMember(t)
l := newQuietLink()
a := &applies{}
q := aQueue(m, a.apply, &keptInMemory{})
inReconcile, release := make(chan struct{}), make(chan struct{})
q.Reconcile = func(context.Context) (Report, bool) {
close(inReconcile)
<-release
return Report{Declared: "d1", Order: Order{Epoch: 7, Sequence: 11}, Applied: []string{"a"},
Outward: []string{"eth0"}}, true
}
q.News = func(Report) bool { return true }
q.attach(context.Background(), l)
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
q.ReconcileDue()
<-inReconcile
q.Deliver(ordered(t, key, "b", 7, 12))
close(release)
eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 2 })
stop()
<-worker
reports := l.said()
if reports[0].Declared != "d1" || reports[1].Sequence != 12 {
t.Fatalf("the reports left out of the order they were made: %+v", reports)
}
if reports[0].ReportSequence >= reports[1].ReportSequence {
t.Fatalf("the later report carries no higher report sequence: %d then %d",
reports[0].ReportSequence, reports[1].ReportSequence)
}
}
// **Two deliveries in quick succession**: only the newest is applied, and there is one account of an
// apply — the first is reported as set aside, naming the one that took its place.
func TestTwoDeliveriesInQuickSuccessionApplyOnlyTheNewest(t *testing.T) {
m, key := aMember(t)
first, second := ordered(t, key, "first", 7, 12), ordered(t, key, "second", 7, 13)
l := newQuietLink(first, second)
a := &applies{}
q := aQueue(m, a.apply, &keptInMemory{})
ctx, stop := context.WithCancel(context.Background())
defer stop()
worker := runWorker(ctx, q)
go func() { _ = serve(ctx, l, m, q, nil, time.Second) }()
eventually(t, "nothing was reported", func() bool { return len(accounts(l.said())) == 1 })
stop()
<-worker
if got := a.all(); !reflect.DeepEqual(got, []string{"second"}) {
t.Fatalf("applied %v; only the newest was wanted", got)
}
reports := l.said()
if len(reports) != 2 || reports[0].Superseded != declaredIn(second.body) ||
reports[0].Declared != declaredIn(first.body) || reports[0].Sequence != 12 {
t.Fatalf("the first was not reported as set aside for the second: %+v", reports)
}
if !first.wasHandled() || !second.wasHandled() {
t.Fatal("a declaration was left unsettled")
}
}
// Deliveries that arrive while the worker is applying are coalesced: when it is free it takes the
// newest held at that moment — by order, not by arrival — and applies it once.
func TestDeliveriesWhileTheWorkerIsBusyAreOneApplyOfTheNewest(t *testing.T) {
m, key := aMember(t)
l := newQuietLink()
a := &applies{}
inApply, release := make(chan struct{}, 1), make(chan struct{})
q := aQueue(m, func(ctx context.Context, raw, sig []byte) Report {
r := a.apply(ctx, raw, sig)
if r.Sequence == 1 {
inApply <- struct{}{}
<-release
}
return r
}, &keptInMemory{})
q.attach(context.Background(), l)
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
q.Deliver(ordered(t, key, "one", 7, 1))
<-inApply
q.Deliver(ordered(t, key, "three", 7, 3))
q.Deliver(ordered(t, key, "two", 7, 2)) // arrived last, composed earlier
q.ReconcileDue() // and the five-minute timer fired too
close(release)
eventually(t, "the newest was never reported", func() bool { return len(accounts(l.said())) == 2 })
stop()
<-worker
if got := a.all(); !reflect.DeepEqual(got, []string{"one", "three"}) {
t.Fatalf("applied %v; the one in hand and then only the newest were wanted", got)
}
last := accounts(l.said())[1]
if last.Sequence != 3 {
t.Fatalf("the last account names sequence %d, not the newest", last.Sequence)
}
}
// **An older epoch is refused, counted and said** (rule 2): the refusal's report names the refused
// declaration and what this node holds, carries the count, and does not replace the kept account of
// the last apply — which is what the mesh is waiting for from this node.
func TestADeclarationFromAnOlderEpochIsRefusedCountedAndSaid(t *testing.T) {
m, key := aMember(t)
l := newQuietLink()
kept := &keptInMemory{}
lastApply := Report{Declared: "d-applied", Order: Order{Epoch: 57, Sequence: 3}, Applied: []string{"a"}}
_ = kept.Keep(lastApply)
stale := ordered(t, key, "stale", 41, 12)
q := aQueue(m, func(_ context.Context, raw, _ []byte) Report {
// What the host's apply answers for a declaration older than the one it kept.
return Report{Declared: digest(raw), Order: Order{Epoch: 41, Sequence: 12},
OlderThan: &Order{Epoch: 57, Sequence: 3}, Refused: "older than what this node applied"}
}, kept)
var lines saidSoFar
q.Say = lines.say
l.refusing = true // nothing is said on linking: the kept account stays kept
q.attach(context.Background(), l)
l.mu.Lock()
l.refusing = false
l.mu.Unlock()
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
q.Deliver(stale)
eventually(t, "the refusal was never reported", func() bool { return len(l.said()) == 1 })
stop()
<-worker
r := l.said()[0]
if r.OlderThan == nil || *r.OlderThan != (Order{Epoch: 57, Sequence: 3}) || r.Epoch != 41 ||
r.Declared != declaredIn(stale.body) || r.RefusedOlder != 1 {
t.Fatalf("the refusal does not name what was refused, what is held, and the count: %+v", r)
}
if !strings.Contains(lines.all(), "1 refused as older so far") {
t.Fatalf("the refusal was not said in the log:\n%s", lines.all())
}
if still, ok, _ := kept.Pending(); !ok || still.Declared != "d-applied" {
t.Fatalf("the refusal replaced the kept account of the last apply: %+v", still)
}
if !stale.wasHandled() {
t.Fatal("the refused declaration was not settled: the mesh would deliver it again")
}
}
// **Restart with an unsaid report**: the next host says it first, exactly as it was made, and numbers
// what it says next above it — the report sequence goes on across the restart rather than from one.
func TestAfterARestartTheUnsaidReportIsSaidFirstAndTheNumbersGoOn(t *testing.T) {
m, key := aMember(t)
kept := &keptInMemory{}
numbers := &numbersInMemory{}
// The first host applies; its report never reaches the mesh, and it is gone.
first := newQuietLink()
first.refusing = true
q1 := aQueue(m, (&applies{}).apply, kept)
q1.Numbers = numbers
q1.attach(context.Background(), first)
ctx1, stop1 := context.WithCancel(context.Background())
w1 := runWorker(ctx1, q1)
q1.Deliver(ordered(t, key, "a", 7, 12))
eventually(t, "the first host never kept its report", func() bool { _, ok, _ := kept.Pending(); return ok })
stop1()
<-w1
lost, _, _ := kept.Pending()
// The next host links, and is then delivered something new.
next := newQuietLink()
q2 := aQueue(m, (&applies{}).apply, kept)
q2.Numbers = numbers
q2.attach(context.Background(), next)
ctx2, stop2 := context.WithCancel(context.Background())
w2 := runWorker(ctx2, q2)
q2.Deliver(ordered(t, key, "b", 7, 13))
eventually(t, "the next host's apply was never reported", func() bool { return len(next.said()) == 2 })
stop2()
<-w2
reports := next.said()
if !reflect.DeepEqual(reports[0], lost) {
t.Fatalf("the unsaid report was not said first, as it was made: %+v", reports[0])
}
if reports[1].Sequence != 13 || reports[1].ReportSequence <= lost.ReportSequence {
t.Fatalf("the next host numbered from one again: %d after %d",
reports[1].ReportSequence, lost.ReportSequence)
}
}
// A report sequence that cannot be read is said, and numbering goes on above anything it can have
// reached — never from one, which would make every report older than those already said.
func TestUnreadableNumbersAreSaidAndNumberingGoesOnAboveThem(t *testing.T) {
m, _ := aMember(t)
l := newQuietLink()
q := aQueue(m, nil, nil)
q.Numbers = &numbersInMemory{broken: errors.New("unreadable")}
var lines saidSoFar
q.Say = lines.say
q.News = func(Report) bool { return true }
q.Reconcile = func(context.Context) (Report, bool) { return Report{Declared: "d"}, true }
before := time.Now().UnixMilli()
q.attach(context.Background(), l)
q.ReconcileDue()
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
eventually(t, "nothing was reported", func() bool { return len(l.said()) == 1 })
stop()
<-worker
if got := l.said()[0].ReportSequence; got < before {
t.Fatalf("numbering went on from %d, not above anything it can have reached", got)
}
if !strings.Contains(lines.all(), "cannot read this node's report sequence") {
t.Fatalf("the unreadable sequence was not said:\n%s", lines.all())
}
}
// A reconcile's report made with no link waits for the next one; an apply made before then replaces
// it, because the apply's account is fresher.
func TestAReconcileReportMadeOfflineWaitsAndIsReplacedByAnApply(t *testing.T) {
m, key := aMember(t)
q := aQueue(m, (&applies{}).apply, &keptInMemory{})
q.News = func(Report) bool { return true }
q.Reconcile = func(context.Context) (Report, bool) {
return Report{Declared: "d1", Outward: []string{"eth0"}}, true
}
ctx, stop := context.WithCancel(context.Background())
defer stop()
runWorker(ctx, q)
q.ReconcileDue()
eventually(t, "the reconcile's report was not held", func() bool {
q.saying.Lock()
defer q.saying.Unlock()
return q.unasked != nil
})
l := newQuietLink()
q.attach(context.Background(), l)
if reports := l.said(); len(reports) != 1 || reports[0].Declared != "d1" {
t.Fatalf("the reconcile's report was not said when the link opened: %+v", reports)
}
// Offline again: a reconcile's report, then an apply of a delivery still in hand.
q.detach(l)
q.ReconcileDue()
eventually(t, "the second reconcile's report was not held", func() bool {
q.saying.Lock()
defer q.saying.Unlock()
return q.unasked != nil
})
q.Deliver(ordered(t, key, "x", 7, 2))
eventually(t, "the apply did not replace the reconcile's report", func() bool {
q.saying.Lock()
defer q.saying.Unlock()
return q.unasked == nil
})
}
// The contract with the controller, on the wire: the keys a declaration's order is read from, and the
// keys a report carries back. A rename here must break this test before it breaks a machine.
func TestTheOrderOnTheWire(t *testing.T) {
body, err := json.Marshal(Report{Node: "n", Declared: "d", Order: Order{Epoch: 57, Sequence: 3},
ReportSequence: 9, OlderThan: &Order{Epoch: 58, Sequence: 1}, RefusedOlder: 2})
if err != nil {
t.Fatal(err)
}
var wire map[string]any
_ = json.Unmarshal(body, &wire)
want := map[string]any{"node": "n", "declared": "d", "epoch": 57.0, "sequence": 3.0,
"report_sequence": 9.0, "older_than": map[string]any{"epoch": 58.0, "sequence": 1.0},
"refused_older": 2.0}
if !reflect.DeepEqual(wire, want) {
t.Fatalf("a report's order on the wire is\n%s\nnot the contract", body)
}
// An older host's report, and an older controller's declaration, claim no order.
if body, _ := json.Marshal(Report{Node: "n", Declared: "d"}); strings.Contains(string(body), "sequence") ||
strings.Contains(string(body), "epoch") {
t.Fatalf("a report about a declaration with no order claims one: %s", body)
}
_, key := aMember(t)
if got := orderOf(signedBy(t, key, []byte(`{"declaration":1,"epoch":57,"sequence":3}`))); got != (Order{Epoch: 57, Sequence: 3}) {
t.Fatalf("a declaration's order was read as %+v", got)
}
}
func TestOlder(t *testing.T) {
for _, c := range []struct {
in, held Order
older bool
}{
{Order{41, 12}, Order{57, 3}, true}, // an older lease holder, whatever its sequence
{Order{57, 2}, Order{57, 3}, true}, // same holder, lower sequence
{Order{57, 3}, Order{57, 3}, false}, // the same declaration again: reconciling
{Order{58, 1}, Order{57, 3}, false}, // a new lease holder
{Order{0, 1}, Order{57, 3}, false}, // an older controller, or one rolled back: today's behaviour
{Order{41, 1}, Order{0, 9}, false}, // nothing held claims an epoch
{Order{57, 0}, Order{57, 3}, false}, // no sequence to compare
{Order{0, 0}, Order{0, 0}, false}, // neither claims anything
{Order{-1, 0}, Order{57, 3}, false}, // never a claim
{Order{57, 4}, Order{57, 3}, false}, // newer
{Order{56, 99}, Order{57, 1}, true}, // a higher sequence is no excuse for an older epoch
{Order{57, 1}, Order{56, 99}, false}, // nor a lower one a reason to refuse a newer epoch
{Order{57, 3}, Order{57, 0}, false}, // held with no sequence
{Order{100, 0}, Order{57, 0}, false}, // newer epoch, no sequences
{Order{10, 0}, Order{57, 0}, true}, // older epoch, no sequences
{Order{57, 2}, Order{0, 0}, false}, // nothing held at all
{Order{0, 2}, Order{0, 3}, false}, // sequence alone does not refuse on the link (no epoch)
{Order{57, 2}, Order{57, -1}, false}, // a held sequence below zero is not one
{Order{57, -1}, Order{57, 2}, false}, // nor an arriving one
{Order{-5, -5}, Order{-1, -1}, false}, // nothing below zero is an order
} {
if got := c.in.Older(c.held); got != c.older {
t.Errorf("%+v older than %+v = %v, want %v", c.in, c.held, got, c.older)
}
}
}
func TestSupersedes(t *testing.T) {
for _, c := range []struct {
next, before Order
want bool
}{
{Order{58, 1}, Order{57, 9}, true}, // a new lease holder
{Order{57, 9}, Order{58, 1}, false}, // the stale one arriving late
{Order{57, 4}, Order{57, 3}, true},
{Order{57, 2}, Order{57, 3}, false}, // arrived last, composed earlier
{Order{0, 4}, Order{0, 3}, true},
{Order{0, 2}, Order{0, 3}, false},
{Order{}, Order{0, 3}, true}, // no order: by arrival
{Order{0, 3}, Order{}, true},
} {
if got := c.next.Supersedes(c.before); got != c.want {
t.Errorf("%+v supersedes %+v = %v, want %v", c.next, c.before, got, c.want)
}
}
}
// numbersInMemory is Numbers without a disk, surviving a "restart" by being shared.
type numbersInMemory struct {
mu sync.Mutex
sequence, refused int64
broken error
}
func (n *numbersInMemory) Read() (int64, int64, error) {
n.mu.Lock()
defer n.mu.Unlock()
return n.sequence, n.refused, n.broken
}
func (n *numbersInMemory) Save(sequence, refused int64) error {
n.mu.Lock()
defer n.mu.Unlock()
n.sequence, n.refused = sequence, refused
return nil
}
// A link dropped or roused while a reconcile is in hand is let go at once and opened again: the
// reconcile holds no link, and its report, if it is news, waits for the next one. Only a delivery in
// hand keeps its link until its report is said.
func TestALinkIsNotHeldForAReconcileInHand(t *testing.T) {
m, _ := aMember(t)
inReconcile, release := make(chan struct{}), make(chan struct{})
q := aQueue(m, nil, nil)
q.News = func(Report) bool { return true }
q.Reconcile = func(context.Context) (Report, bool) {
close(inReconcile)
<-release
return Report{Declared: "d1", Outward: []string{"eth0"}}, true
}
l := newQuietLink()
q.attach(context.Background(), l)
ctx, stop := context.WithCancel(context.Background())
worker := runWorker(ctx, q)
q.ReconcileDue()
<-inReconcile
detached := make(chan struct{})
go func() { q.detach(l); close(detached) }()
select {
case <-detached:
case <-time.After(2 * time.Second):
t.Fatal("the link was held while a reconcile ran")
}
close(release)
eventually(t, "the reconcile's report was not held for the next link", func() bool {
q.saying.Lock()
defer q.saying.Unlock()
return q.unasked != nil
})
stop()
<-worker
if len(l.said()) != 0 {
t.Fatalf("a report was said on a link let go: %+v", l.said())
}
}
+48 -191
View File
@@ -71,8 +71,8 @@ type Announce func(string)
// Nil is allowed and means nothing ever rouses it, which is every machine that does not suspend.
type Roused <-chan struct{}
func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
return HoldRoused(ctx, m, apply, say, timeout, nil, nil, nil)
func Hold(ctx context.Context, m Membership, queue *Queue, say Announce, timeout time.Duration) error {
return HoldRoused(ctx, m, queue, say, timeout, nil)
}
// Unsaid keeps the report of the last apply until the mesh has taken it, so a report lost between
@@ -94,26 +94,16 @@ type Unsaid interface {
Pending() (Report, bool, error)
}
// Outbox carries reports the node has to say without having been sent anything — what a
// reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up;
// a report made while it is down waits in the channel for the next one. Nil is allowed.
type Outbox <-chan Unasked
// Unasked is one such report, with the way to say whether it reached the mesh. Done is called
// with true only when the broker took it — a node that marked a change said because it queued it
// would never say it again, and the mesh would go on believing nothing changed.
type Unasked struct {
Report Report
Done func(published bool)
}
// HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed
// reports to publish between deliveries.
func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce,
timeout time.Duration, roused Roused, outbox Outbox, unsaid Unsaid) error {
// HoldRoused is Hold, told when the machine has reason to think its link is stale.
//
// **What arrives is enqueued, never applied here** (novox/hq to-be 45 §6). The queue's one worker
// applies, outlives every link, and says its reports on whichever link is open; the caller runs it
// (Queue.Run) for as long as this holds.
func HoldRoused(ctx context.Context, m Membership, queue *Queue, say Announce,
timeout time.Duration, roused Roused) error {
return holdWith(ctx, func(ctx context.Context) error {
return Run(ctx, m, apply, say, timeout, outbox, unsaid)
return Run(ctx, m, queue, say, timeout)
}, say, roused)
}
@@ -206,30 +196,22 @@ func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) err
//
// Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason;
// Hold is what decides whether to open it again.
func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration,
outbox Outbox, unsaid Unsaid) error {
func Run(ctx context.Context, m Membership, queue *Queue, say Announce, timeout time.Duration) error {
link, err := Open(ctx, m, timeout)
if err != nil {
return err
}
defer link.Close()
return serve(ctx, link, m, apply, say, timeout, outbox, unsaid)
return serve(ctx, link, m, queue, say, timeout)
}
// serve is Run on a link already open: separated so what the node says, and when, can be tested
// without a broker.
func serve(ctx context.Context, link Link, m Membership, apply Applier, say Announce,
timeout time.Duration, outbox Outbox, unsaid Unsaid) error {
func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Announce,
timeout time.Duration) error {
if say == nil {
say = func(string) {}
}
// **A report about something done is said even when the link is being let go.** An apply can
// end the link itself — the host stands aside for a successor it delivered (novox/hq ADR 0141)
// — and a report published on the cancelled context was refused before it left: "applied, and
// could not tell the mesh: reporting: context canceled", while the release plan waited for it
// (novox/hq issue 264). Not cancelled with the link, still bounded by the timeout, and sent
// before the deferred close lets go of the connection.
reporting := context.WithoutCancel(ctx)
// Said, because it is the event anybody watching actually wants. Without it a node logs every
// failure and nothing on success, so a log full of "trying again" and then silence reads as
@@ -243,32 +225,18 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
defer beat.Stop()
publishAlive(ctx, link, m, say, timeout)
// **The declaration this link last applied**, so a report made about an older one is not said
// after it (novox/hq issue 267). Empty until something is applied or said again here.
var applied string
// **What the last apply did, if the mesh never heard it** — said before anything newly
// delivered is applied, so it can never land after, and read as newer than, a later report.
if unsaid != nil {
switch kept, ok, err := unsaid.Pending(); {
case err != nil:
say("cannot read the report this node kept unsaid: " + err.Error())
case ok:
say("saying again what the last apply did: its report never reached the mesh")
applied = kept.Declared
if publishReport(reporting, link, m, kept, say, timeout) {
if err := unsaid.Said(kept.Declared); err != nil {
say("said the kept report, and cannot forget it: " + err.Error())
}
}
}
}
// **What the last apply did, if the mesh never heard it**, is said before anything newly
// delivered can be: the queue says it as the link is handed to it (novox/hq issue 264). And the
// link is not let go while the queue's worker is mid-act: an apply that ends the link — the host
// standing aside for a successor it delivered (novox/hq ADR 0141) — is still reported on it,
// before the deferred close.
queue.attach(ctx, link)
defer queue.detach(link)
declarations := link.Declarations()
for {
// Asked to stop — by standing aside, say — nothing further is applied, whichever of the
// cases below the select would otherwise have picked.
// Asked to stop — by standing aside, say — nothing further is taken in.
if ctx.Err() != nil {
return nil
}
@@ -277,30 +245,6 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
return nil
case <-beat.C:
publishAlive(ctx, link, m, say, timeout)
case unasked := <-outbox:
// Said without having been asked: a reconcile found what an adopted node holds, or
// its firewall, changed since it last said.
//
// **Unless a newer declaration was applied while it waited** (novox/hq issue 267). A
// reconcile that held the machine a moment before a delivery arrived made its report
// about the declaration kept then; the delivery's apply waited for it, was reported, and
// only then did this loop get round to the reconcile's — which reached the mesh last,
// named the older declaration, and was stored as the machine's latest account. The plan
// waiting on the machine then waited on a report it had already been given. What the
// reconcile saw, the newer apply has said since, and fresher; the next reconcile says
// anything that is still news.
if overtaken(unasked.Report, applied) {
say(fmt.Sprintf("set aside a reconcile's report of declaration %s: %s was applied and "+
"reported since", short(unasked.Report.Declared), short(applied)))
if unasked.Done != nil {
unasked.Done(false)
}
continue
}
published := publishReport(ctx, link, m, unasked.Report, say, timeout)
if unasked.Done != nil {
unasked.Done(published)
}
case reason := <-link.Lost():
return reason
case declaration, ok := <-declarations:
@@ -314,55 +258,14 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
return errors.New("the mesh stopped sending this node declarations")
}
}
// Whatever else is already waiting supersedes this one. Each set-aside declaration is
// reported as such, then settled unapplied.
declaration, superseded := newest(declarations, declaration, drainWindow)
for _, old := range superseded {
say("set aside a declaration: a newer one arrived with it")
publishReport(reporting, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()),
Superseded: declaredIn(declaration.Body())}, say, timeout)
_ = old.Handled()
}
report := handleBody(ctx, m, declaration.Body(), apply)
if report.Declared != "" {
applied = report.Declared
}
if unsaid != nil {
if err := unsaid.Keep(report); err != nil {
say("cannot keep this apply's report until it is said: " + err.Error())
}
}
switch {
case report.Refused != "":
say("refused a declaration: " + report.Refused)
case len(report.Failed) > 0:
say(fmt.Sprintf("applied %d and failed: %v%s",
len(report.Applied), report.Failed, heldNote(report.Held)))
default:
say(fmt.Sprintf("applied %d resource(s)%s",
len(report.Applied), heldNote(report.Held)))
}
if publishReport(reporting, link, m, report, say, timeout) && unsaid != nil {
if err := unsaid.Said(report.Declared); err != nil {
say("said this apply's report, and cannot forget it: " + err.Error())
}
}
// Settled after the report is published. A node that dies between applying and
// reporting leaves the declaration with the mesh and applies it again on return,
// which is safe because applying is reconciliation — it converges rather than
// repeating.
_ = declaration.Handled()
// Whatever else is already waiting goes with it. **Enqueued, not applied**: the queue
// applies the newest of everything delivered and not yet taken, and reports the rest as
// set aside (novox/hq to-be 45 §6). The link goes on beating meanwhile.
queue.Deliver(gather(declarations, declaration, drainWindow)...)
}
}
}
// overtaken says whether a report made unasked is about a declaration other than the one this link
// has applied since. Unanswerable is not overtaken: a report that names no declaration, or a link
// that has applied nothing yet, has nothing to be older than.
func overtaken(r Report, applied string) bool {
return r.Declared != "" && applied != "" && r.Declared != applied
}
// short is a declaration's digest as a person reads it in a log line.
func short(digest string) string {
if len(digest) > 12 {
@@ -371,91 +274,45 @@ func short(digest string) string {
return digest
}
// drainDepth is how many declarations the host will hold unacknowledged while it looks for a newer
// one; drainWindow is how long it waits for another to follow the one it has. Both small: a push is
// rare and a backlog is the exception this exists for, not the shape of ordinary traffic.
// drainDepth is how many declarations the link holds unread; drainWindow is how long it waits for
// another to follow the one it has. Both small: a push is rare and a backlog is the exception this
// exists for, not the shape of ordinary traffic.
const (
drainDepth = 16
drainWindow = 750 * time.Millisecond
)
// newest takes what is already waiting behind `first` and returns the last of them to apply, and
// the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a
// declaration in flight from the mesh arrives within that; one that does not is the next push.
// gather takes what is already waiting behind `first`, in the order it arrived. It waits `window`
// for a straggler after each arrival and no longer: a declaration in flight from the mesh arrives
// within that; one that does not is the next push. Which of them is applied is the queue's to decide
// (pick), across everything delivered and not yet taken.
//
// **Its job narrows once declarations are state rather than messages, and does not disappear.**
// On the bus being built, a declaration is last-per-subject (novox/hq design 29 §4), so a node
// that was away receives exactly the current one instead of a queue of superseded ones — the
// catch-up half of what this does is then the stream's. And a stream sequence orders them
// definitively, where this window only infers order from arrival time, which is the wire-level
// answer to novox/hq issue 107.
// On the bus a declaration is last-per-subject (novox/hq design 29 §4), so a node that was away
// receives exactly the current one instead of a queue of superseded ones — the catch-up half is the
// stream's. And a declaration's order decides which is newest, where this window only infers it from
// arrival time, which is the wire-level answer to novox/hq issue 107.
//
// What remains is the live case: three pushes in quick succession to a *connected* node are
// three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not
// deleted — and saying which half goes is worth more than a note that it "can probably be
// removed", which is how a load-bearing window gets deleted by somebody in a hurry.
func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) {
latest := first
var superseded []Declaration
// What remains is the live case: three pushes in quick succession to a *connected*, idle node are
// three deliveries, whatever the stream later retains, and gathering them is what makes them one
// apply rather than two. So this is narrowed, not deleted — and saying which half goes is worth more
// than a note that it "can probably be removed", which is how a load-bearing window gets deleted by
// somebody in a hurry.
func gather(arriving <-chan Declaration, first Declaration, window time.Duration) []Declaration {
batch := []Declaration{first}
for {
select {
case next, ok := <-arriving:
if !ok {
return latest, superseded
return batch
}
// **By sequence when both carry one, by arrival when either does not** (novox/hq
// 04-ISSUES/107). Arrival is what this window had to go on, and it is wrong exactly
// when it matters — a backlog drained out of order. A declaration that says where it
// stands is believed over when it turned up; one that does not is the older
// controller's, and arrival is all there is.
if sequenceOf(next.Body()) < sequenceOf(latest.Body()) &&
sequenceOf(next.Body()) > 0 && sequenceOf(latest.Body()) > 0 {
superseded = append(superseded, next)
continue
}
superseded = append(superseded, latest)
latest = next
batch = append(batch, next)
case <-time.After(window):
return latest, superseded
return batch
}
}
}
// sequenceOf is the order a signed declaration claims, or zero when it claims none or cannot be
// read. Read from the envelope alone; the signature is verified later, when the winner is applied,
// and a forged message that lied about its sequence would only set aside real ones — which are
// reported as set aside, and the next push sends the current one again.
func sequenceOf(body []byte) int64 {
var signed Signed
if err := json.Unmarshal(body, &signed); err != nil {
return 0
}
var d struct {
Sequence int64 `json:"sequence"`
}
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
return 0
}
return d.Sequence
}
// declaredIn is the id a signed declaration carries, for a report about one that was not applied.
// Empty if the message is not one — a forged or garbled message is refused by handleBody when its
// turn comes; here it is only named.
func declaredIn(body []byte) string {
var signed Signed
if err := json.Unmarshal(body, &signed); err != nil {
return ""
}
var d struct {
Declared string `json:"declared"`
}
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
return ""
}
return d.Declared
}
// handleBody is the whole of deciding whether to trust a message, separated from the broker so it
// can be tested as the security check it is rather than as message plumbing.
func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report {
@@ -514,7 +371,7 @@ func publishReport(ctx context.Context, bus Bus, m Membership, report Report,
// publishAlive says this node is here, and nothing else.
func publishAlive(ctx context.Context, bus Bus, m Membership, say Announce,
timeout time.Duration) {
body, err := json.Marshal(Alive{Node: m.Node})
body, err := json.Marshal(Alive{Node: m.Node, IntervalSeconds: int(AliveEvery / time.Second)})
if err != nil {
return
}
+38 -11
View File
@@ -61,10 +61,13 @@ func (l *quietLink) said() []Report {
// keptInMemory is Unsaid without a disk.
type keptInMemory struct {
mu sync.Mutex
report *Report
}
func (k *keptInMemory) Keep(r Report) error {
k.mu.Lock()
defer k.mu.Unlock()
if r.Declared == "" {
k.report = nil
return nil
@@ -73,12 +76,16 @@ func (k *keptInMemory) Keep(r Report) error {
return nil
}
func (k *keptInMemory) Said(declared string) error {
k.mu.Lock()
defer k.mu.Unlock()
if k.report != nil && k.report.Declared == declared {
k.report = nil
}
return nil
}
func (k *keptInMemory) Pending() (Report, bool, error) {
k.mu.Lock()
defer k.mu.Unlock()
if k.report == nil {
return Report{}, false, nil
}
@@ -97,21 +104,28 @@ func aMember(t *testing.T) (Membership, ed25519.PrivateKey) {
// Measured on the anchor: a declaration carried a new controller and a new host, the host applied
// it, stood aside for its successor — which cancels the link — and the report of that apply failed
// with "context canceled". The plan waited on it until somebody pushed by hand.
//
// R3 of novox/hq to-be 45 §9, as a unit: the node-engine self-updates during its report, and the
// report still arrives — on the link the apply came over, before it is let go.
func TestAnApplyThatEndsTheLinkIsStillReported(t *testing.T) {
m, key := aMember(t)
declaration := &said{body: signedBy(t, key, []byte(`{"declaration":1}`))}
behind := &said{body: signedBy(t, key, []byte(`{"declaration":2}`))}
l := newQuietLink(declaration)
kept := &keptInMemory{}
ctx, standAside := context.WithCancel(context.Background())
defer standAside()
apply := func(context.Context, []byte, []byte) Report {
applied := 0
q := aQueue(m, func(context.Context, []byte, []byte) Report {
applied++
standAside() // the host delivered its successor, and stands aside for it
return Report{Declared: "d1", Applied: []string{"a"}}
}
}, kept)
worker := runWorker(ctx, q)
done := make(chan error, 1)
go func() { done <- serve(ctx, l, m, apply, nil, time.Second, nil, kept) }()
go func() { done <- serve(ctx, l, m, q, nil, time.Second) }()
select {
case err := <-done:
if err != nil {
@@ -120,14 +134,22 @@ func TestAnApplyThatEndsTheLinkIsStillReported(t *testing.T) {
case <-time.After(5 * time.Second):
t.Fatal("the link did not end after the apply stood aside")
}
// Something delivered after the host stood aside is not applied by it: the successor will be
// sent it again.
q.Deliver(behind)
<-worker
reports := l.said()
if len(reports) != 1 || reports[0].Declared != "d1" || !reflect.DeepEqual(reports[0].Applied, []string{"a"}) {
t.Fatalf("the apply's report did not reach the mesh: %+v", reports)
}
if !declaration.handled {
if !declaration.wasHandled() {
t.Fatal("the declaration was not settled after its report")
}
if applied != 1 || behind.wasHandled() {
t.Fatalf("%d applies, and the declaration behind settled %v: after standing aside nothing "+
"further is applied", applied, behind.wasHandled())
}
if _, ok, _ := kept.Pending(); ok {
t.Fatal("a report the mesh took is still kept to be said again")
}
@@ -144,13 +166,19 @@ func TestAReportLostAfterTheApplyIsSaidOnTheNextLink(t *testing.T) {
first := newQuietLink(&said{body: signedBy(t, key, []byte(`{"declaration":1}`))})
first.refusing = true
ctx, stop := context.WithCancel(context.Background())
apply := func(context.Context, []byte, []byte) Report { stop(); return made }
if err := serve(ctx, first, m, apply, nil, time.Second, nil, kept); err != nil {
q := aQueue(m, func(context.Context, []byte, []byte) Report { stop(); return made }, kept)
worker := runWorker(ctx, q)
if err := serve(ctx, first, m, q, nil, time.Second); err != nil {
t.Fatal(err)
}
<-worker
if len(first.said()) != 0 {
t.Fatal("the broker refused the report and it counted as said")
}
lost, ok, _ := kept.Pending()
if !ok || lost.ReportSequence != 1 {
t.Fatalf("the lost report was not kept as it was made: %+v", lost)
}
// The next host links. It is stopped at once, so all it does is what it does on linking.
next := newQuietLink()
@@ -160,12 +188,11 @@ func TestAReportLostAfterTheApplyIsSaidOnTheNextLink(t *testing.T) {
t.Fatal("nothing was delivered, and something was applied")
return Report{}
}
if err := serve(gone, next, m, never, nil, time.Second, nil, kept); err != nil {
if err := serve(gone, next, m, aQueue(m, never, kept), nil, time.Second); err != nil {
t.Fatal(err)
}
reports := next.said()
made.Node = m.Node
if len(reports) != 1 || !reflect.DeepEqual(reports[0], made) {
if len(reports) != 1 || !reflect.DeepEqual(reports[0], lost) {
t.Fatalf("the lost report was not said again as it was made: %+v", reports)
}
if _, ok, _ := kept.Pending(); ok {
@@ -179,10 +206,10 @@ func TestNothingIsSaidAgainWhenNothingWasLost(t *testing.T) {
l := newQuietLink()
gone, cancel := context.WithCancel(context.Background())
cancel()
if err := serve(gone, l, m, nil, nil, time.Second, nil, &keptInMemory{}); err != nil {
if err := serve(gone, l, m, aQueue(m, nil, &keptInMemory{}), nil, time.Second); err != nil {
t.Fatal(err)
}
if err := serve(gone, l, m, nil, nil, time.Second, nil, nil); err != nil {
if err := serve(gone, l, m, aQueue(m, nil, nil), nil, time.Second); err != nil {
t.Fatal(err)
}
if reports := l.said(); len(reports) != 0 {
+94
View File
@@ -0,0 +1,94 @@
package store
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
)
// What the node-engine counts about what it says, kept so the counting outlives the process.
//
// **The report sequence must increase across restarts and self-updates** (novox/hq to-be 45 §6). The
// mesh keeps, per machine, the highest report it accepted and refuses an older one; a host that began
// again from one after every restart — and a self-update is a restart — would have every report it
// made refused until it had counted past where its predecessor stopped. So the number is kept beside
// the state, where the successor reads it.
//
// And the count of declarations refused as older than what this node applied (rule 2), which the
// mesh's stale-writer watchdog reads from every report.
// NumbersName is where they live, beside the state.
const NumbersName = "numbers.json"
// NumbersPath is where the numbers live, given where the state lives.
func NumbersPath(statePath string) string {
return filepath.Join(filepath.Dir(statePath), NumbersName)
}
// Numbers is what is kept.
type Numbers struct {
// ReportSequence is the number of the last report this node made.
ReportSequence int64 `json:"report_sequence"`
// RefusedOlder is how many declarations it has refused as older, ever.
RefusedOlder int64 `json:"refused_older"`
}
// ReadNumbers is what was kept, or zero when nothing was — a node that has never reported.
//
// **Unreadable is an error, never zero** (novox/hq ADR 0227, rule 4): read as zero, every report the
// node makes next would be older than the ones it already made, and refused.
func ReadNumbers(path string) (Numbers, error) {
raw, err := os.ReadFile(path)
if errors.Is(err, os.ErrNotExist) {
return Numbers{}, nil
}
if err != nil {
return Numbers{}, err
}
var n Numbers
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(&n); err != nil {
return Numbers{}, fmt.Errorf("the numbers kept at %s are unreadable: %w", path, err)
}
if n.ReportSequence < 0 || n.RefusedOlder < 0 {
return Numbers{}, fmt.Errorf("the numbers kept at %s are below zero, which nothing counting writes", path)
}
return n, nil
}
// SaveNumbers keeps them, replacing what was kept, and is on disk when it returns: a number used and
// not kept is one the next host would use again.
func SaveNumbers(path string, n Numbers) error {
raw, err := json.Marshal(n)
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(path), ".numbers-*")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
if err := tmp.Chmod(0o600); err != nil {
tmp.Close()
return err
}
if _, err := tmp.Write(raw); err != nil {
tmp.Close()
return err
}
if err := tmp.Sync(); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
return os.Rename(tmp.Name(), path)
}
+38
View File
@@ -0,0 +1,38 @@
// Package validate is the node-engine's own judgement of a declaration, for whoever must know before
// a declaration is sent whether a machine would take it (novox/hq to-be 45 §4, D1; §9, the merge gate).
//
// **One validator.** The controller composed declarations the host then refused whole — a manifest
// the catalogue check passed (novox/hq issue 236), a consumer's identity too long for the machine's
// names (issue 263) — because the only validator was the one the host runs as it applies. A second,
// written in the controller from the host's rules, would drift from them the first time either
// changed. So the host's own parsing is exported here, unchanged: what the controller's self-check and
// the merge gate run is what every machine runs.
//
// It judges a declaration as it arrives over the link: actions are refused, as the host refuses them
// from the link (novox/hq ADR 0005). Nothing here touches a machine.
package validate
import (
"errors"
"github.com/novox/mesh-host/internal/declaration"
)
// Declaration is every problem the node-engine would refuse a declaration for, each in its own words;
// nil when it would take it. The body is the declaration as the controller composes it — the bytes a
// signature is made over — not the signed envelope.
func Declaration(raw []byte) []string {
_, err := declaration.Parse(raw)
if err == nil {
return nil
}
var refused *declaration.RefusalError
if errors.As(err, &refused) {
return refused.Problems
}
return []string{err.Error()}
}
// Version is the declaration vocabulary this validator speaks: a declaration of another version is
// refused whole by Declaration, as by the host.
const Version = declaration.Version
+37
View File
@@ -0,0 +1,37 @@
package validate
import (
"strings"
"testing"
)
// **The validator a controller imports is the host's**: a declaration the host takes passes, and one
// it refuses is refused with the host's own words, every problem at once.
func TestTheValidatorIsTheHosts(t *testing.T) {
good := `{"declaration": 1, "resources": [{"id": "d", "type": "directory", "path": "/var/lib/x", "mode": "0755"}]}`
if problems := Declaration([]byte(good)); problems != nil {
t.Fatalf("a declaration the host takes was refused: %v", problems)
}
// And one carrying the lease epoch and sequence the controller sends it under (novox/hq to-be 45
// §6): a controller validating with this package may send them.
ordered := `{"declaration": 1, "epoch": 57, "sequence": 3, "resources": [{"id": "d", "type": "directory", "path": "/var/lib/x", "mode": "0755"}]}`
if problems := Declaration([]byte(ordered)); problems != nil {
t.Fatalf("a declaration carrying its order was refused: %v", problems)
}
for _, bad := range []string{
`{"declaration": 1, "resources": []}`,
`{"declaration": 99, "resources": [{"id": "d", "type": "directory", "path": "/x"}]}`,
`{"declaration": 1, "resources": [{"id": "d", "type": "nothing-the-host-knows"}]}`,
`not json`,
} {
problems := Declaration([]byte(bad))
if len(problems) == 0 {
t.Errorf("a declaration the host refuses was passed: %s", bad)
}
for _, p := range problems {
if strings.TrimSpace(p) == "" {
t.Errorf("a problem said nothing, for %s", bad)
}
}
}
}