diff --git a/cmd/mesh-host/epoch_test.go b/cmd/mesh-host/epoch_test.go new file mode 100644 index 0000000..56bb0bc --- /dev/null +++ b/cmd/mesh-host/epoch_test.go @@ -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") + } +} diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 51c572d..c01d355 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -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 { diff --git a/cmd/mesh-host/main_test.go b/cmd/mesh-host/main_test.go index 86bc6e5..b860dc6 100644 --- a/cmd/mesh-host/main_test.go +++ b/cmd/mesh-host/main_test.go @@ -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") } } diff --git a/cmd/mesh-host/order_test.go b/cmd/mesh-host/order_test.go index 0380aa4..1cf9d90 100644 --- a/cmd/mesh-host/order_test.go +++ b/cmd/mesh-host/order_test.go @@ -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") + } +} diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index f64c02d..dc908cb 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -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). diff --git a/internal/declaration/order_test.go b/internal/declaration/order_test.go new file mode 100644 index 0000000..cdb4ad3 --- /dev/null +++ b/internal/declaration/order_test.go @@ -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) + } + } +} diff --git a/internal/link/hearing_nats_test.go b/internal/link/hearing_nats_test.go index 9a33e38..0226bc7 100644 --- a/internal/link/hearing_nats_test.go +++ b/internal/link/hearing_nats_test.go @@ -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) + } +} diff --git a/internal/link/messages.go b/internal/link/messages.go index a427eb7..ef3fff2 100644 --- a/internal/link/messages.go +++ b/internal/link/messages.go @@ -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 { diff --git a/internal/link/newest_test.go b/internal/link/newest_test.go index 5eab4a0..5eb12a5 100644 --- a/internal/link/newest_test.go +++ b/internal/link/newest_test.go @@ -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) { diff --git a/internal/link/order_test.go b/internal/link/order_test.go index 64df68a..529d06c 100644 --- a/internal/link/order_test.go +++ b/internal/link/order_test.go @@ -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) } } diff --git a/internal/link/overtaken_test.go b/internal/link/overtaken_test.go deleted file mode 100644 index 5087b27..0000000 --- a/internal/link/overtaken_test.go +++ /dev/null @@ -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) - } - } -} diff --git a/internal/link/queue.go b/internal/link/queue.go new file mode 100644 index 0000000..51c0bfc --- /dev/null +++ b/internal/link/queue.go @@ -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[:]) +} diff --git a/internal/link/queue_test.go b/internal/link/queue_test.go new file mode 100644 index 0000000..dd75b69 --- /dev/null +++ b/internal/link/queue_test.go @@ -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()) + } +} diff --git a/internal/link/run.go b/internal/link/run.go index 4e42fc0..dfaeedb 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -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 { diff --git a/internal/link/unsaid_test.go b/internal/link/unsaid_test.go index d42da39..85e3361 100644 --- a/internal/link/unsaid_test.go +++ b/internal/link/unsaid_test.go @@ -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 { diff --git a/internal/store/numbers.go b/internal/store/numbers.go new file mode 100644 index 0000000..61fe263 --- /dev/null +++ b/internal/store/numbers.go @@ -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) +} diff --git a/validate/validate_test.go b/validate/validate_test.go index cd3bfa5..dba68e4 100644 --- a/validate/validate_test.go +++ b/validate/validate_test.go @@ -12,6 +12,12 @@ func TestTheValidatorIsTheHosts(t *testing.T) { 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"}]}`,