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.
This commit is contained in:
jochen
2026-10-06 11:54:10 +02:00
parent d7d93f57c5
commit 31804bd8c2
17 changed files with 1729 additions and 391 deletions
+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)
}
}
+88
View File
@@ -96,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"`
@@ -166,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())
}
}
+47 -190
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 {
+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)
}