Compare commits
11
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b2808549ee | ||
|
|
bbe74f3fc2 | ||
|
|
3e80b7ae32 | ||
|
|
31804bd8c2 | ||
|
|
d7d93f57c5 | ||
|
|
b9774ea426 | ||
|
|
1545b00a87 | ||
|
|
6953b5bafd | ||
|
|
93efe41dc9 | ||
|
|
1e463def6b | ||
|
|
a654fa768d |
@@ -0,0 +1,77 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-host/internal/link"
|
||||
"github.com/novox/mesh-host/internal/store"
|
||||
)
|
||||
|
||||
// **A declaration from an older lease epoch is refused before anything is applied** (novox/hq to-be
|
||||
// 45 §6, ADR 0227 rule 2). The controller that sent it lost its lease and goes on sending; this node
|
||||
// applied the holder's. Refused on the host's own apply path — the one the queue's worker runs — before
|
||||
// the machine is read or touched, with a report naming what was refused and what is held, and the
|
||||
// kept declaration left as it was.
|
||||
func TestADeclarationFromAnOlderEpochIsRefusedBeforeAnythingIsApplied(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
opts := options{state: filepath.Join(dir, "state.json")}
|
||||
_, private, err := ed25519.GenerateKey(nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
target := filepath.Join(dir, "a.conf")
|
||||
declared := func(epoch, sequence int64) store.Declared {
|
||||
body, err := json.Marshal(map[string]any{"declaration": 1, "epoch": epoch, "sequence": sequence,
|
||||
"resources": []any{map[string]any{"id": "a", "type": "file", "path": target, "content": "x\n"}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return store.Declared{Declaration: body, Signature: ed25519.Sign(private, body)}
|
||||
}
|
||||
held := declared(57, 3)
|
||||
if err := store.SaveDeclared(store.DeclaredPath(opts.state), held); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
stale := declared(41, 12)
|
||||
report := applyAndKeep(context.Background(), opts, stale.Declaration, &stale, nil, nil)
|
||||
if report.OlderThan == nil || *report.OlderThan != (link.Order{Epoch: 57, Sequence: 3}) ||
|
||||
report.Order != (link.Order{Epoch: 41, Sequence: 12}) || report.Declared != digestOf(stale.Declaration) ||
|
||||
report.Refused == "" {
|
||||
t.Fatalf("the refusal does not name what was refused and what is held: %+v", report)
|
||||
}
|
||||
if _, err := os.Stat(target); !errors.Is(err, os.ErrNotExist) {
|
||||
t.Fatal("the refused declaration touched the machine")
|
||||
}
|
||||
kept, err := store.ReadDeclared(store.DeclaredPath(opts.state))
|
||||
if err != nil || string(kept.Declaration) != string(held.Declaration) {
|
||||
t.Fatal("the refused declaration replaced what this node kept")
|
||||
}
|
||||
}
|
||||
|
||||
// The queue's counts are kept beside the state and read back by the next host; unreadable is an
|
||||
// error, never zero.
|
||||
func TestTheNumbersSurviveTheHost(t *testing.T) {
|
||||
state := filepath.Join(t.TempDir(), "state.json")
|
||||
if sequence, refused, err := numbersBeside(state).Read(); err != nil || sequence != 0 || refused != 0 {
|
||||
t.Fatalf("a node that never reported read %d, %d, %v", sequence, refused, err)
|
||||
}
|
||||
if err := numbersBeside(state).Save(41, 2); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if sequence, refused, err := numbersBeside(state).Read(); err != nil || sequence != 41 || refused != 2 {
|
||||
t.Fatalf("read back %d, %d, %v", sequence, refused, err)
|
||||
}
|
||||
if err := os.WriteFile(store.NumbersPath(state), []byte(`{"report_sequence":`), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := numbersBeside(state).Read(); err == nil {
|
||||
t.Fatal("unreadable numbers were read as zero")
|
||||
}
|
||||
}
|
||||
+139
-74
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -161,7 +161,7 @@
|
||||
"type": "file",
|
||||
"path": "/var/lib/mesh-bus-conf/accounts.conf",
|
||||
"mode": "0600",
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.control.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\"] }\n subscribe: { allow: [\"$JS.API.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
},
|
||||
{
|
||||
"id": "broker",
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,10 @@ const (
|
||||
// look like a stale one.
|
||||
type Alive struct {
|
||||
Node string `json:"node"`
|
||||
// IntervalSeconds is how often this node says it is there (novox/hq to-be 45 §3, S1): the
|
||||
// controller's watchdog is bound to three of them, so the bound moves with the interval rather
|
||||
// than with a number the controller guessed. A controller older than that ignores it.
|
||||
IntervalSeconds int `json:"interval_seconds"`
|
||||
}
|
||||
|
||||
// Signed is a declaration and the signature over it.
|
||||
@@ -92,6 +96,36 @@ type Report struct {
|
||||
// cannot answer "which"; the digest is the answer itself.
|
||||
Declared string `json:"declared,omitempty"`
|
||||
|
||||
// Order is where the declaration this report is about stands — its `epoch` and `sequence`, as
|
||||
// the declaration carried them — so the mesh orders reports by what they are about rather than
|
||||
// by when they arrived (novox/hq to-be 45 §6). On the wire as two top-level keys beside
|
||||
// `declared`; absent for a declaration that claimed no order, and for a report about none.
|
||||
Order
|
||||
|
||||
// ReportSequence is this node-engine's own order among everything it reports: one higher for
|
||||
// every report it makes, kept on disk so it goes on increasing across restarts and self-updates
|
||||
// (to-be 45 §6). With the declaration's order it lets the mesh refuse an older report about the
|
||||
// same declaration — a reconcile's account overtaking the apply that followed it (novox/hq issue
|
||||
// 267) is then refused by number rather than set aside by digest. A report said again because
|
||||
// it never reached the mesh (issue 264) carries the number it was made with.
|
||||
//
|
||||
// Zero claims no order: every report an older host makes, and a one-shot report a command makes
|
||||
// (a rekey). **Its presence also says this host reads a declaration's `epoch`**, which is how the
|
||||
// mesh knows it may send one: an older host decodes a declaration strictly and refuses a key it
|
||||
// does not know.
|
||||
ReportSequence int64 `json:"report_sequence,omitempty"`
|
||||
|
||||
// OlderThan is set on a report refusing a declaration older than one this node has applied: the
|
||||
// order of the newer declaration it holds (to-be 45 §6, rule 2). The refused declaration is the
|
||||
// report's own Declared and Order, and Refused says it in words. Applying it would make the
|
||||
// machine into something a controller that lost its lease said, after one that holds it.
|
||||
OlderThan *Order `json:"older_than,omitempty"`
|
||||
|
||||
// RefusedOlder is how many declarations this node-engine has refused as older, ever — kept on
|
||||
// disk with the report sequence. On every report, so a count missed with a lost refusal is still
|
||||
// read from the next one: what the mesh's stale-writer watchdog (to-be 45 §3, S13) counts.
|
||||
RefusedOlder int64 `json:"refused_older,omitempty"`
|
||||
|
||||
// Held is what this adopted node found and is keeping as it was until its module is taken
|
||||
// (novox/hq ADR 0100). Without it an adopted node reads as converged.
|
||||
Held []Held `json:"held,omitempty"`
|
||||
@@ -162,6 +196,64 @@ type Report struct {
|
||||
Rekey *Rekey `json:"rekey,omitempty"`
|
||||
}
|
||||
|
||||
// Order is where a declaration stands among everything the mesh has sent this node (novox/hq to-be
|
||||
// 45 §6): the controller's lease **epoch** it was sent under, and its **sequence** — one higher for
|
||||
// every send to this node (04-ISSUES/107). The declaration carries both inside what is signed, as
|
||||
// top-level `epoch` and `sequence` beside `declaration`; a report carries back the pair of the
|
||||
// declaration it is about.
|
||||
//
|
||||
// Zero in either is "no order claimed", never "first": every declaration an older controller sent,
|
||||
// and the bundle genesis applies.
|
||||
type Order struct {
|
||||
Epoch int64 `json:"epoch,omitempty"`
|
||||
Sequence int64 `json:"sequence,omitempty"`
|
||||
}
|
||||
|
||||
// Older says whether a declaration of this order is older than one of order `than`, which this node
|
||||
// has applied — the one rule by which a node-engine refuses a declaration (to-be 45 §6: "older epoch;
|
||||
// same epoch, lower sequence").
|
||||
//
|
||||
// **Only when both claim an epoch.** A declaration with none is an older controller's — or a
|
||||
// controller rolled back to a build from before the lease — and one kept with none is what this host
|
||||
// held before any controller had a lease. Neither has an order to compare, and refusing on a guess
|
||||
// would strand the machine the moment the controller is older than the host: without an epoch, today's
|
||||
// behaviour stands. A higher epoch is a new lease holder and is never older, whatever its sequence.
|
||||
func (o Order) Older(than Order) bool {
|
||||
if o.Epoch <= 0 || than.Epoch <= 0 {
|
||||
return false
|
||||
}
|
||||
if o.Epoch != than.Epoch {
|
||||
return o.Epoch < than.Epoch
|
||||
}
|
||||
return o.Sequence > 0 && than.Sequence > 0 && o.Sequence < than.Sequence
|
||||
}
|
||||
|
||||
// Supersedes says whether a declaration of order `o`, arriving after one of order `before`, takes its
|
||||
// place among what is waiting to be applied. By epoch when both claim one and they differ, then by
|
||||
// sequence when both claim one, and by arrival when they do not — what the drain had to go on before
|
||||
// declarations said where they stand (04-ISSUES/107).
|
||||
func (o Order) Supersedes(before Order) bool {
|
||||
if o.Epoch > 0 && before.Epoch > 0 && o.Epoch != before.Epoch {
|
||||
return o.Epoch > before.Epoch
|
||||
}
|
||||
if o.Sequence > 0 && before.Sequence > 0 {
|
||||
return o.Sequence >= before.Sequence
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// Words is an order as a person reads it in a log line. Not String: Report embeds Order, and a
|
||||
// Stringer promoted onto every report would print each one as its order alone.
|
||||
func (o Order) Words() string {
|
||||
switch {
|
||||
case o.Epoch > 0:
|
||||
return "epoch " + strconv.FormatInt(o.Epoch, 10) + ", sequence " + strconv.FormatInt(o.Sequence, 10)
|
||||
case o.Sequence > 0:
|
||||
return "sequence " + strconv.FormatInt(o.Sequence, 10) + ", no epoch"
|
||||
}
|
||||
return "no order"
|
||||
}
|
||||
|
||||
// CarriedTunnel is this node's account of the tunnel it took over. State is one of the Carried
|
||||
// states below; Note is what the host did about it, when it did something.
|
||||
type CarriedTunnel struct {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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[:])
|
||||
}
|
||||
@@ -0,0 +1,542 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The node-engine's one apply queue (novox/hq to-be 45 §6): a delivery and the reconcile are two
|
||||
// reasons to enqueue the same act, the newest declaration is applied once, and the reports leave in
|
||||
// the order they are made. These are the races of issues 257, 261, 264 and 267 replayed against it.
|
||||
|
||||
func aQueue(m Membership, apply Applier, kept Unsaid) *Queue {
|
||||
return &Queue{Membership: m, Apply: apply, Unsaid: kept, Timeout: time.Second}
|
||||
}
|
||||
|
||||
// runWorker runs the queue's worker until ctx ends; the channel closes when it has.
|
||||
func runWorker(ctx context.Context, q *Queue) <-chan struct{} {
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
q.Run(ctx)
|
||||
}()
|
||||
return done
|
||||
}
|
||||
|
||||
// ordered is a signed declaration claiming an order, with a resource id so two are distinct bytes.
|
||||
func ordered(t *testing.T, key ed25519.PrivateKey, id string, epoch, sequence int64) *said {
|
||||
t.Helper()
|
||||
inner, err := json.Marshal(map[string]any{"declaration": 1, "resources": []any{map[string]any{"id": id}},
|
||||
"epoch": epoch, "sequence": sequence})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return &said{body: signedBy(t, key, inner)}
|
||||
}
|
||||
|
||||
// applies records what the worker applied, by the id of the declaration's first resource, and
|
||||
// reports it as applied — with the order the declaration carried, as the host's own apply does.
|
||||
type applies struct {
|
||||
mu sync.Mutex
|
||||
seen []string
|
||||
}
|
||||
|
||||
func (a *applies) apply(_ context.Context, raw, _ []byte) Report {
|
||||
var d struct {
|
||||
Resources []struct {
|
||||
ID string `json:"id"`
|
||||
} `json:"resources"`
|
||||
Order
|
||||
}
|
||||
_ = json.Unmarshal(raw, &d)
|
||||
id := ""
|
||||
if len(d.Resources) > 0 {
|
||||
id = d.Resources[0].ID
|
||||
}
|
||||
a.mu.Lock()
|
||||
a.seen = append(a.seen, id)
|
||||
a.mu.Unlock()
|
||||
return Report{Declared: digest(raw), Order: d.Order, Applied: []string{id}}
|
||||
}
|
||||
|
||||
func (a *applies) all() []string {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
return append([]string(nil), a.seen...)
|
||||
}
|
||||
|
||||
func digest(raw []byte) string {
|
||||
b, _ := json.Marshal(Signed{Declaration: raw})
|
||||
return declaredIn(b)
|
||||
}
|
||||
|
||||
// eventually waits for a condition the worker reaches on its own time.
|
||||
func eventually(t *testing.T, what string, ok func() bool) {
|
||||
t.Helper()
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
for !ok() {
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatal(what)
|
||||
}
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
// accounts is the reports the mesh heard that are an apply's account, not a set-aside's.
|
||||
func accounts(reports []Report) []Report {
|
||||
var out []Report
|
||||
for _, r := range reports {
|
||||
if r.Superseded == "" {
|
||||
out = append(out, r)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// **R2: a reconcile due during a delivery** (issues 257, 261, 267). The reconcile is asked for while a
|
||||
// delivery waits: the worker applies the delivery once — applying the newest thing the mesh said is
|
||||
// what reconciling is — and makes one report, naming the newest sequence. The reconcile never applies
|
||||
// what was kept before.
|
||||
func TestAReconcileDueDuringADeliveryIsThatDeliverysApply(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
l := newQuietLink()
|
||||
a := &applies{}
|
||||
reconciled := 0
|
||||
q := aQueue(m, a.apply, &keptInMemory{})
|
||||
q.Reconcile = func(context.Context) (Report, bool) {
|
||||
reconciled++
|
||||
return Report{Declared: "the one kept before", Applied: []string{"old"}}, true
|
||||
}
|
||||
q.News = func(Report) bool { return true }
|
||||
q.attach(context.Background(), l)
|
||||
|
||||
q.Deliver(ordered(t, key, "new", 7, 12))
|
||||
q.ReconcileDue()
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 1 })
|
||||
stop()
|
||||
<-worker
|
||||
|
||||
if got := a.all(); !reflect.DeepEqual(got, []string{"new"}) {
|
||||
t.Fatalf("applied %v; one apply of the newest was wanted", got)
|
||||
}
|
||||
if reconciled != 0 {
|
||||
t.Fatalf("the reconcile applied what was kept %d time(s) beside the delivery", reconciled)
|
||||
}
|
||||
r := l.said()[0]
|
||||
if r.Sequence != 12 || r.Epoch != 7 || r.ReportSequence != 1 || r.Node != m.Node {
|
||||
t.Fatalf("the report does not name the order it applied: %+v", r)
|
||||
}
|
||||
}
|
||||
|
||||
// **The other order** (issue 267): the reconcile is running when a delivery arrives. Its report, about
|
||||
// the declaration kept then, leaves first; the delivery's leaves after it, with a higher report
|
||||
// sequence. The mesh's latest word is the newest declaration, whichever arrived first.
|
||||
func TestAReconcileRunningWhenADeliveryArrivesIsSaidBeforeIt(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
l := newQuietLink()
|
||||
a := &applies{}
|
||||
q := aQueue(m, a.apply, &keptInMemory{})
|
||||
inReconcile, release := make(chan struct{}), make(chan struct{})
|
||||
q.Reconcile = func(context.Context) (Report, bool) {
|
||||
close(inReconcile)
|
||||
<-release
|
||||
return Report{Declared: "d1", Order: Order{Epoch: 7, Sequence: 11}, Applied: []string{"a"},
|
||||
Outward: []string{"eth0"}}, true
|
||||
}
|
||||
q.News = func(Report) bool { return true }
|
||||
q.attach(context.Background(), l)
|
||||
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
q.ReconcileDue()
|
||||
<-inReconcile
|
||||
q.Deliver(ordered(t, key, "b", 7, 12))
|
||||
close(release)
|
||||
eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 2 })
|
||||
stop()
|
||||
<-worker
|
||||
|
||||
reports := l.said()
|
||||
if reports[0].Declared != "d1" || reports[1].Sequence != 12 {
|
||||
t.Fatalf("the reports left out of the order they were made: %+v", reports)
|
||||
}
|
||||
if reports[0].ReportSequence >= reports[1].ReportSequence {
|
||||
t.Fatalf("the later report carries no higher report sequence: %d then %d",
|
||||
reports[0].ReportSequence, reports[1].ReportSequence)
|
||||
}
|
||||
}
|
||||
|
||||
// **Two deliveries in quick succession**: only the newest is applied, and there is one account of an
|
||||
// apply — the first is reported as set aside, naming the one that took its place.
|
||||
func TestTwoDeliveriesInQuickSuccessionApplyOnlyTheNewest(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
first, second := ordered(t, key, "first", 7, 12), ordered(t, key, "second", 7, 13)
|
||||
l := newQuietLink(first, second)
|
||||
a := &applies{}
|
||||
q := aQueue(m, a.apply, &keptInMemory{})
|
||||
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
defer stop()
|
||||
worker := runWorker(ctx, q)
|
||||
go func() { _ = serve(ctx, l, m, q, nil, time.Second) }()
|
||||
eventually(t, "nothing was reported", func() bool { return len(accounts(l.said())) == 1 })
|
||||
stop()
|
||||
<-worker
|
||||
|
||||
if got := a.all(); !reflect.DeepEqual(got, []string{"second"}) {
|
||||
t.Fatalf("applied %v; only the newest was wanted", got)
|
||||
}
|
||||
reports := l.said()
|
||||
if len(reports) != 2 || reports[0].Superseded != declaredIn(second.body) ||
|
||||
reports[0].Declared != declaredIn(first.body) || reports[0].Sequence != 12 {
|
||||
t.Fatalf("the first was not reported as set aside for the second: %+v", reports)
|
||||
}
|
||||
if !first.wasHandled() || !second.wasHandled() {
|
||||
t.Fatal("a declaration was left unsettled")
|
||||
}
|
||||
}
|
||||
|
||||
// Deliveries that arrive while the worker is applying are coalesced: when it is free it takes the
|
||||
// newest held at that moment — by order, not by arrival — and applies it once.
|
||||
func TestDeliveriesWhileTheWorkerIsBusyAreOneApplyOfTheNewest(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
l := newQuietLink()
|
||||
a := &applies{}
|
||||
inApply, release := make(chan struct{}, 1), make(chan struct{})
|
||||
q := aQueue(m, func(ctx context.Context, raw, sig []byte) Report {
|
||||
r := a.apply(ctx, raw, sig)
|
||||
if r.Sequence == 1 {
|
||||
inApply <- struct{}{}
|
||||
<-release
|
||||
}
|
||||
return r
|
||||
}, &keptInMemory{})
|
||||
q.attach(context.Background(), l)
|
||||
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
q.Deliver(ordered(t, key, "one", 7, 1))
|
||||
<-inApply
|
||||
q.Deliver(ordered(t, key, "three", 7, 3))
|
||||
q.Deliver(ordered(t, key, "two", 7, 2)) // arrived last, composed earlier
|
||||
q.ReconcileDue() // and the five-minute timer fired too
|
||||
close(release)
|
||||
eventually(t, "the newest was never reported", func() bool { return len(accounts(l.said())) == 2 })
|
||||
stop()
|
||||
<-worker
|
||||
|
||||
if got := a.all(); !reflect.DeepEqual(got, []string{"one", "three"}) {
|
||||
t.Fatalf("applied %v; the one in hand and then only the newest were wanted", got)
|
||||
}
|
||||
last := accounts(l.said())[1]
|
||||
if last.Sequence != 3 {
|
||||
t.Fatalf("the last account names sequence %d, not the newest", last.Sequence)
|
||||
}
|
||||
}
|
||||
|
||||
// **An older epoch is refused, counted and said** (rule 2): the refusal's report names the refused
|
||||
// declaration and what this node holds, carries the count, and does not replace the kept account of
|
||||
// the last apply — which is what the mesh is waiting for from this node.
|
||||
func TestADeclarationFromAnOlderEpochIsRefusedCountedAndSaid(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
l := newQuietLink()
|
||||
kept := &keptInMemory{}
|
||||
lastApply := Report{Declared: "d-applied", Order: Order{Epoch: 57, Sequence: 3}, Applied: []string{"a"}}
|
||||
_ = kept.Keep(lastApply)
|
||||
stale := ordered(t, key, "stale", 41, 12)
|
||||
q := aQueue(m, func(_ context.Context, raw, _ []byte) Report {
|
||||
// What the host's apply answers for a declaration older than the one it kept.
|
||||
return Report{Declared: digest(raw), Order: Order{Epoch: 41, Sequence: 12},
|
||||
OlderThan: &Order{Epoch: 57, Sequence: 3}, Refused: "older than what this node applied"}
|
||||
}, kept)
|
||||
var lines saidSoFar
|
||||
q.Say = lines.say
|
||||
l.refusing = true // nothing is said on linking: the kept account stays kept
|
||||
q.attach(context.Background(), l)
|
||||
l.mu.Lock()
|
||||
l.refusing = false
|
||||
l.mu.Unlock()
|
||||
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
q.Deliver(stale)
|
||||
eventually(t, "the refusal was never reported", func() bool { return len(l.said()) == 1 })
|
||||
stop()
|
||||
<-worker
|
||||
|
||||
r := l.said()[0]
|
||||
if r.OlderThan == nil || *r.OlderThan != (Order{Epoch: 57, Sequence: 3}) || r.Epoch != 41 ||
|
||||
r.Declared != declaredIn(stale.body) || r.RefusedOlder != 1 {
|
||||
t.Fatalf("the refusal does not name what was refused, what is held, and the count: %+v", r)
|
||||
}
|
||||
if !strings.Contains(lines.all(), "1 refused as older so far") {
|
||||
t.Fatalf("the refusal was not said in the log:\n%s", lines.all())
|
||||
}
|
||||
if still, ok, _ := kept.Pending(); !ok || still.Declared != "d-applied" {
|
||||
t.Fatalf("the refusal replaced the kept account of the last apply: %+v", still)
|
||||
}
|
||||
if !stale.wasHandled() {
|
||||
t.Fatal("the refused declaration was not settled: the mesh would deliver it again")
|
||||
}
|
||||
}
|
||||
|
||||
// **Restart with an unsaid report**: the next host says it first, exactly as it was made, and numbers
|
||||
// what it says next above it — the report sequence goes on across the restart rather than from one.
|
||||
func TestAfterARestartTheUnsaidReportIsSaidFirstAndTheNumbersGoOn(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
kept := &keptInMemory{}
|
||||
numbers := &numbersInMemory{}
|
||||
|
||||
// The first host applies; its report never reaches the mesh, and it is gone.
|
||||
first := newQuietLink()
|
||||
first.refusing = true
|
||||
q1 := aQueue(m, (&applies{}).apply, kept)
|
||||
q1.Numbers = numbers
|
||||
q1.attach(context.Background(), first)
|
||||
ctx1, stop1 := context.WithCancel(context.Background())
|
||||
w1 := runWorker(ctx1, q1)
|
||||
q1.Deliver(ordered(t, key, "a", 7, 12))
|
||||
eventually(t, "the first host never kept its report", func() bool { _, ok, _ := kept.Pending(); return ok })
|
||||
stop1()
|
||||
<-w1
|
||||
lost, _, _ := kept.Pending()
|
||||
|
||||
// The next host links, and is then delivered something new.
|
||||
next := newQuietLink()
|
||||
q2 := aQueue(m, (&applies{}).apply, kept)
|
||||
q2.Numbers = numbers
|
||||
q2.attach(context.Background(), next)
|
||||
ctx2, stop2 := context.WithCancel(context.Background())
|
||||
w2 := runWorker(ctx2, q2)
|
||||
q2.Deliver(ordered(t, key, "b", 7, 13))
|
||||
eventually(t, "the next host's apply was never reported", func() bool { return len(next.said()) == 2 })
|
||||
stop2()
|
||||
<-w2
|
||||
|
||||
reports := next.said()
|
||||
if !reflect.DeepEqual(reports[0], lost) {
|
||||
t.Fatalf("the unsaid report was not said first, as it was made: %+v", reports[0])
|
||||
}
|
||||
if reports[1].Sequence != 13 || reports[1].ReportSequence <= lost.ReportSequence {
|
||||
t.Fatalf("the next host numbered from one again: %d after %d",
|
||||
reports[1].ReportSequence, lost.ReportSequence)
|
||||
}
|
||||
}
|
||||
|
||||
// A report sequence that cannot be read is said, and numbering goes on above anything it can have
|
||||
// reached — never from one, which would make every report older than those already said.
|
||||
func TestUnreadableNumbersAreSaidAndNumberingGoesOnAboveThem(t *testing.T) {
|
||||
m, _ := aMember(t)
|
||||
l := newQuietLink()
|
||||
q := aQueue(m, nil, nil)
|
||||
q.Numbers = &numbersInMemory{broken: errors.New("unreadable")}
|
||||
var lines saidSoFar
|
||||
q.Say = lines.say
|
||||
q.News = func(Report) bool { return true }
|
||||
q.Reconcile = func(context.Context) (Report, bool) { return Report{Declared: "d"}, true }
|
||||
before := time.Now().UnixMilli()
|
||||
q.attach(context.Background(), l)
|
||||
|
||||
q.ReconcileDue()
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
eventually(t, "nothing was reported", func() bool { return len(l.said()) == 1 })
|
||||
stop()
|
||||
<-worker
|
||||
if got := l.said()[0].ReportSequence; got < before {
|
||||
t.Fatalf("numbering went on from %d, not above anything it can have reached", got)
|
||||
}
|
||||
if !strings.Contains(lines.all(), "cannot read this node's report sequence") {
|
||||
t.Fatalf("the unreadable sequence was not said:\n%s", lines.all())
|
||||
}
|
||||
}
|
||||
|
||||
// A reconcile's report made with no link waits for the next one; an apply made before then replaces
|
||||
// it, because the apply's account is fresher.
|
||||
func TestAReconcileReportMadeOfflineWaitsAndIsReplacedByAnApply(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
q := aQueue(m, (&applies{}).apply, &keptInMemory{})
|
||||
q.News = func(Report) bool { return true }
|
||||
q.Reconcile = func(context.Context) (Report, bool) {
|
||||
return Report{Declared: "d1", Outward: []string{"eth0"}}, true
|
||||
}
|
||||
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
defer stop()
|
||||
runWorker(ctx, q)
|
||||
q.ReconcileDue()
|
||||
eventually(t, "the reconcile's report was not held", func() bool {
|
||||
q.saying.Lock()
|
||||
defer q.saying.Unlock()
|
||||
return q.unasked != nil
|
||||
})
|
||||
l := newQuietLink()
|
||||
q.attach(context.Background(), l)
|
||||
if reports := l.said(); len(reports) != 1 || reports[0].Declared != "d1" {
|
||||
t.Fatalf("the reconcile's report was not said when the link opened: %+v", reports)
|
||||
}
|
||||
|
||||
// Offline again: a reconcile's report, then an apply of a delivery still in hand.
|
||||
q.detach(l)
|
||||
q.ReconcileDue()
|
||||
eventually(t, "the second reconcile's report was not held", func() bool {
|
||||
q.saying.Lock()
|
||||
defer q.saying.Unlock()
|
||||
return q.unasked != nil
|
||||
})
|
||||
q.Deliver(ordered(t, key, "x", 7, 2))
|
||||
eventually(t, "the apply did not replace the reconcile's report", func() bool {
|
||||
q.saying.Lock()
|
||||
defer q.saying.Unlock()
|
||||
return q.unasked == nil
|
||||
})
|
||||
}
|
||||
|
||||
// The contract with the controller, on the wire: the keys a declaration's order is read from, and the
|
||||
// keys a report carries back. A rename here must break this test before it breaks a machine.
|
||||
func TestTheOrderOnTheWire(t *testing.T) {
|
||||
body, err := json.Marshal(Report{Node: "n", Declared: "d", Order: Order{Epoch: 57, Sequence: 3},
|
||||
ReportSequence: 9, OlderThan: &Order{Epoch: 58, Sequence: 1}, RefusedOlder: 2})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var wire map[string]any
|
||||
_ = json.Unmarshal(body, &wire)
|
||||
want := map[string]any{"node": "n", "declared": "d", "epoch": 57.0, "sequence": 3.0,
|
||||
"report_sequence": 9.0, "older_than": map[string]any{"epoch": 58.0, "sequence": 1.0},
|
||||
"refused_older": 2.0}
|
||||
if !reflect.DeepEqual(wire, want) {
|
||||
t.Fatalf("a report's order on the wire is\n%s\nnot the contract", body)
|
||||
}
|
||||
|
||||
// An older host's report, and an older controller's declaration, claim no order.
|
||||
if body, _ := json.Marshal(Report{Node: "n", Declared: "d"}); strings.Contains(string(body), "sequence") ||
|
||||
strings.Contains(string(body), "epoch") {
|
||||
t.Fatalf("a report about a declaration with no order claims one: %s", body)
|
||||
}
|
||||
_, key := aMember(t)
|
||||
if got := orderOf(signedBy(t, key, []byte(`{"declaration":1,"epoch":57,"sequence":3}`))); got != (Order{Epoch: 57, Sequence: 3}) {
|
||||
t.Fatalf("a declaration's order was read as %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOlder(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
in, held Order
|
||||
older bool
|
||||
}{
|
||||
{Order{41, 12}, Order{57, 3}, true}, // an older lease holder, whatever its sequence
|
||||
{Order{57, 2}, Order{57, 3}, true}, // same holder, lower sequence
|
||||
{Order{57, 3}, Order{57, 3}, false}, // the same declaration again: reconciling
|
||||
{Order{58, 1}, Order{57, 3}, false}, // a new lease holder
|
||||
{Order{0, 1}, Order{57, 3}, false}, // an older controller, or one rolled back: today's behaviour
|
||||
{Order{41, 1}, Order{0, 9}, false}, // nothing held claims an epoch
|
||||
{Order{57, 0}, Order{57, 3}, false}, // no sequence to compare
|
||||
{Order{0, 0}, Order{0, 0}, false}, // neither claims anything
|
||||
{Order{-1, 0}, Order{57, 3}, false}, // never a claim
|
||||
{Order{57, 4}, Order{57, 3}, false}, // newer
|
||||
{Order{56, 99}, Order{57, 1}, true}, // a higher sequence is no excuse for an older epoch
|
||||
{Order{57, 1}, Order{56, 99}, false}, // nor a lower one a reason to refuse a newer epoch
|
||||
{Order{57, 3}, Order{57, 0}, false}, // held with no sequence
|
||||
{Order{100, 0}, Order{57, 0}, false}, // newer epoch, no sequences
|
||||
{Order{10, 0}, Order{57, 0}, true}, // older epoch, no sequences
|
||||
{Order{57, 2}, Order{0, 0}, false}, // nothing held at all
|
||||
{Order{0, 2}, Order{0, 3}, false}, // sequence alone does not refuse on the link (no epoch)
|
||||
{Order{57, 2}, Order{57, -1}, false}, // a held sequence below zero is not one
|
||||
{Order{57, -1}, Order{57, 2}, false}, // nor an arriving one
|
||||
{Order{-5, -5}, Order{-1, -1}, false}, // nothing below zero is an order
|
||||
} {
|
||||
if got := c.in.Older(c.held); got != c.older {
|
||||
t.Errorf("%+v older than %+v = %v, want %v", c.in, c.held, got, c.older)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSupersedes(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
next, before Order
|
||||
want bool
|
||||
}{
|
||||
{Order{58, 1}, Order{57, 9}, true}, // a new lease holder
|
||||
{Order{57, 9}, Order{58, 1}, false}, // the stale one arriving late
|
||||
{Order{57, 4}, Order{57, 3}, true},
|
||||
{Order{57, 2}, Order{57, 3}, false}, // arrived last, composed earlier
|
||||
{Order{0, 4}, Order{0, 3}, true},
|
||||
{Order{0, 2}, Order{0, 3}, false},
|
||||
{Order{}, Order{0, 3}, true}, // no order: by arrival
|
||||
{Order{0, 3}, Order{}, true},
|
||||
} {
|
||||
if got := c.next.Supersedes(c.before); got != c.want {
|
||||
t.Errorf("%+v supersedes %+v = %v, want %v", c.next, c.before, got, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// numbersInMemory is Numbers without a disk, surviving a "restart" by being shared.
|
||||
type numbersInMemory struct {
|
||||
mu sync.Mutex
|
||||
sequence, refused int64
|
||||
broken error
|
||||
}
|
||||
|
||||
func (n *numbersInMemory) Read() (int64, int64, error) {
|
||||
n.mu.Lock()
|
||||
defer n.mu.Unlock()
|
||||
return n.sequence, n.refused, n.broken
|
||||
}
|
||||
|
||||
func (n *numbersInMemory) Save(sequence, refused int64) error {
|
||||
n.mu.Lock()
|
||||
defer n.mu.Unlock()
|
||||
n.sequence, n.refused = sequence, refused
|
||||
return nil
|
||||
}
|
||||
|
||||
// A link dropped or roused while a reconcile is in hand is let go at once and opened again: the
|
||||
// reconcile holds no link, and its report, if it is news, waits for the next one. Only a delivery in
|
||||
// hand keeps its link until its report is said.
|
||||
func TestALinkIsNotHeldForAReconcileInHand(t *testing.T) {
|
||||
m, _ := aMember(t)
|
||||
inReconcile, release := make(chan struct{}), make(chan struct{})
|
||||
q := aQueue(m, nil, nil)
|
||||
q.News = func(Report) bool { return true }
|
||||
q.Reconcile = func(context.Context) (Report, bool) {
|
||||
close(inReconcile)
|
||||
<-release
|
||||
return Report{Declared: "d1", Outward: []string{"eth0"}}, true
|
||||
}
|
||||
l := newQuietLink()
|
||||
q.attach(context.Background(), l)
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
q.ReconcileDue()
|
||||
<-inReconcile
|
||||
|
||||
detached := make(chan struct{})
|
||||
go func() { q.detach(l); close(detached) }()
|
||||
select {
|
||||
case <-detached:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("the link was held while a reconcile ran")
|
||||
}
|
||||
close(release)
|
||||
eventually(t, "the reconcile's report was not held for the next link", func() bool {
|
||||
q.saying.Lock()
|
||||
defer q.saying.Unlock()
|
||||
return q.unasked != nil
|
||||
})
|
||||
stop()
|
||||
<-worker
|
||||
if len(l.said()) != 0 {
|
||||
t.Fatalf("a report was said on a link let go: %+v", l.said())
|
||||
}
|
||||
}
|
||||
+48
-191
@@ -71,8 +71,8 @@ type Announce func(string)
|
||||
// Nil is allowed and means nothing ever rouses it, which is every machine that does not suspend.
|
||||
type Roused <-chan struct{}
|
||||
|
||||
func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration) error {
|
||||
return HoldRoused(ctx, m, apply, say, timeout, nil, nil, nil)
|
||||
func Hold(ctx context.Context, m Membership, queue *Queue, say Announce, timeout time.Duration) error {
|
||||
return HoldRoused(ctx, m, queue, say, timeout, nil)
|
||||
}
|
||||
|
||||
// Unsaid keeps the report of the last apply until the mesh has taken it, so a report lost between
|
||||
@@ -94,26 +94,16 @@ type Unsaid interface {
|
||||
Pending() (Report, bool, error)
|
||||
}
|
||||
|
||||
// Outbox carries reports the node has to say without having been sent anything — what a
|
||||
// reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up;
|
||||
// a report made while it is down waits in the channel for the next one. Nil is allowed.
|
||||
type Outbox <-chan Unasked
|
||||
|
||||
// Unasked is one such report, with the way to say whether it reached the mesh. Done is called
|
||||
// with true only when the broker took it — a node that marked a change said because it queued it
|
||||
// would never say it again, and the mesh would go on believing nothing changed.
|
||||
type Unasked struct {
|
||||
Report Report
|
||||
Done func(published bool)
|
||||
}
|
||||
|
||||
// HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed
|
||||
// reports to publish between deliveries.
|
||||
func HoldRoused(ctx context.Context, m Membership, apply Applier, say Announce,
|
||||
timeout time.Duration, roused Roused, outbox Outbox, unsaid Unsaid) error {
|
||||
// HoldRoused is Hold, told when the machine has reason to think its link is stale.
|
||||
//
|
||||
// **What arrives is enqueued, never applied here** (novox/hq to-be 45 §6). The queue's one worker
|
||||
// applies, outlives every link, and says its reports on whichever link is open; the caller runs it
|
||||
// (Queue.Run) for as long as this holds.
|
||||
func HoldRoused(ctx context.Context, m Membership, queue *Queue, say Announce,
|
||||
timeout time.Duration, roused Roused) error {
|
||||
|
||||
return holdWith(ctx, func(ctx context.Context) error {
|
||||
return Run(ctx, m, apply, say, timeout, outbox, unsaid)
|
||||
return Run(ctx, m, queue, say, timeout)
|
||||
}, say, roused)
|
||||
}
|
||||
|
||||
@@ -206,30 +196,22 @@ func holdWith(ctx context.Context, run attempt, say Announce, roused Roused) err
|
||||
//
|
||||
// Outbound only, and nothing listens on this machine. Returns when the link ends, for any reason;
|
||||
// Hold is what decides whether to open it again.
|
||||
func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout time.Duration,
|
||||
outbox Outbox, unsaid Unsaid) error {
|
||||
func Run(ctx context.Context, m Membership, queue *Queue, say Announce, timeout time.Duration) error {
|
||||
link, err := Open(ctx, m, timeout)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer link.Close()
|
||||
return serve(ctx, link, m, apply, say, timeout, outbox, unsaid)
|
||||
return serve(ctx, link, m, queue, say, timeout)
|
||||
}
|
||||
|
||||
// serve is Run on a link already open: separated so what the node says, and when, can be tested
|
||||
// without a broker.
|
||||
func serve(ctx context.Context, link Link, m Membership, apply Applier, say Announce,
|
||||
timeout time.Duration, outbox Outbox, unsaid Unsaid) error {
|
||||
func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Announce,
|
||||
timeout time.Duration) error {
|
||||
if say == nil {
|
||||
say = func(string) {}
|
||||
}
|
||||
// **A report about something done is said even when the link is being let go.** An apply can
|
||||
// end the link itself — the host stands aside for a successor it delivered (novox/hq ADR 0141)
|
||||
// — and a report published on the cancelled context was refused before it left: "applied, and
|
||||
// could not tell the mesh: reporting: context canceled", while the release plan waited for it
|
||||
// (novox/hq issue 264). Not cancelled with the link, still bounded by the timeout, and sent
|
||||
// before the deferred close lets go of the connection.
|
||||
reporting := context.WithoutCancel(ctx)
|
||||
|
||||
// Said, because it is the event anybody watching actually wants. Without it a node logs every
|
||||
// failure and nothing on success, so a log full of "trying again" and then silence reads as
|
||||
@@ -243,32 +225,18 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
|
||||
defer beat.Stop()
|
||||
publishAlive(ctx, link, m, say, timeout)
|
||||
|
||||
// **The declaration this link last applied**, so a report made about an older one is not said
|
||||
// after it (novox/hq issue 267). Empty until something is applied or said again here.
|
||||
var applied string
|
||||
|
||||
// **What the last apply did, if the mesh never heard it** — said before anything newly
|
||||
// delivered is applied, so it can never land after, and read as newer than, a later report.
|
||||
if unsaid != nil {
|
||||
switch kept, ok, err := unsaid.Pending(); {
|
||||
case err != nil:
|
||||
say("cannot read the report this node kept unsaid: " + err.Error())
|
||||
case ok:
|
||||
say("saying again what the last apply did: its report never reached the mesh")
|
||||
applied = kept.Declared
|
||||
if publishReport(reporting, link, m, kept, say, timeout) {
|
||||
if err := unsaid.Said(kept.Declared); err != nil {
|
||||
say("said the kept report, and cannot forget it: " + err.Error())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// **What the last apply did, if the mesh never heard it**, is said before anything newly
|
||||
// delivered can be: the queue says it as the link is handed to it (novox/hq issue 264). And the
|
||||
// link is not let go while the queue's worker is mid-act: an apply that ends the link — the host
|
||||
// standing aside for a successor it delivered (novox/hq ADR 0141) — is still reported on it,
|
||||
// before the deferred close.
|
||||
queue.attach(ctx, link)
|
||||
defer queue.detach(link)
|
||||
|
||||
declarations := link.Declarations()
|
||||
|
||||
for {
|
||||
// Asked to stop — by standing aside, say — nothing further is applied, whichever of the
|
||||
// cases below the select would otherwise have picked.
|
||||
// Asked to stop — by standing aside, say — nothing further is taken in.
|
||||
if ctx.Err() != nil {
|
||||
return nil
|
||||
}
|
||||
@@ -277,30 +245,6 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
|
||||
return nil
|
||||
case <-beat.C:
|
||||
publishAlive(ctx, link, m, say, timeout)
|
||||
case unasked := <-outbox:
|
||||
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||
// its firewall, changed since it last said.
|
||||
//
|
||||
// **Unless a newer declaration was applied while it waited** (novox/hq issue 267). A
|
||||
// reconcile that held the machine a moment before a delivery arrived made its report
|
||||
// about the declaration kept then; the delivery's apply waited for it, was reported, and
|
||||
// only then did this loop get round to the reconcile's — which reached the mesh last,
|
||||
// named the older declaration, and was stored as the machine's latest account. The plan
|
||||
// waiting on the machine then waited on a report it had already been given. What the
|
||||
// reconcile saw, the newer apply has said since, and fresher; the next reconcile says
|
||||
// anything that is still news.
|
||||
if overtaken(unasked.Report, applied) {
|
||||
say(fmt.Sprintf("set aside a reconcile's report of declaration %s: %s was applied and "+
|
||||
"reported since", short(unasked.Report.Declared), short(applied)))
|
||||
if unasked.Done != nil {
|
||||
unasked.Done(false)
|
||||
}
|
||||
continue
|
||||
}
|
||||
published := publishReport(ctx, link, m, unasked.Report, say, timeout)
|
||||
if unasked.Done != nil {
|
||||
unasked.Done(published)
|
||||
}
|
||||
case reason := <-link.Lost():
|
||||
return reason
|
||||
case declaration, ok := <-declarations:
|
||||
@@ -314,55 +258,14 @@ func serve(ctx context.Context, link Link, m Membership, apply Applier, say Anno
|
||||
return errors.New("the mesh stopped sending this node declarations")
|
||||
}
|
||||
}
|
||||
// Whatever else is already waiting supersedes this one. Each set-aside declaration is
|
||||
// reported as such, then settled unapplied.
|
||||
declaration, superseded := newest(declarations, declaration, drainWindow)
|
||||
for _, old := range superseded {
|
||||
say("set aside a declaration: a newer one arrived with it")
|
||||
publishReport(reporting, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()),
|
||||
Superseded: declaredIn(declaration.Body())}, say, timeout)
|
||||
_ = old.Handled()
|
||||
}
|
||||
report := handleBody(ctx, m, declaration.Body(), apply)
|
||||
if report.Declared != "" {
|
||||
applied = report.Declared
|
||||
}
|
||||
if unsaid != nil {
|
||||
if err := unsaid.Keep(report); err != nil {
|
||||
say("cannot keep this apply's report until it is said: " + err.Error())
|
||||
}
|
||||
}
|
||||
switch {
|
||||
case report.Refused != "":
|
||||
say("refused a declaration: " + report.Refused)
|
||||
case len(report.Failed) > 0:
|
||||
say(fmt.Sprintf("applied %d and failed: %v%s",
|
||||
len(report.Applied), report.Failed, heldNote(report.Held)))
|
||||
default:
|
||||
say(fmt.Sprintf("applied %d resource(s)%s",
|
||||
len(report.Applied), heldNote(report.Held)))
|
||||
}
|
||||
if publishReport(reporting, link, m, report, say, timeout) && unsaid != nil {
|
||||
if err := unsaid.Said(report.Declared); err != nil {
|
||||
say("said this apply's report, and cannot forget it: " + err.Error())
|
||||
}
|
||||
}
|
||||
// Settled after the report is published. A node that dies between applying and
|
||||
// reporting leaves the declaration with the mesh and applies it again on return,
|
||||
// which is safe because applying is reconciliation — it converges rather than
|
||||
// repeating.
|
||||
_ = declaration.Handled()
|
||||
// Whatever else is already waiting goes with it. **Enqueued, not applied**: the queue
|
||||
// applies the newest of everything delivered and not yet taken, and reports the rest as
|
||||
// set aside (novox/hq to-be 45 §6). The link goes on beating meanwhile.
|
||||
queue.Deliver(gather(declarations, declaration, drainWindow)...)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// overtaken says whether a report made unasked is about a declaration other than the one this link
|
||||
// has applied since. Unanswerable is not overtaken: a report that names no declaration, or a link
|
||||
// that has applied nothing yet, has nothing to be older than.
|
||||
func overtaken(r Report, applied string) bool {
|
||||
return r.Declared != "" && applied != "" && r.Declared != applied
|
||||
}
|
||||
|
||||
// short is a declaration's digest as a person reads it in a log line.
|
||||
func short(digest string) string {
|
||||
if len(digest) > 12 {
|
||||
@@ -371,91 +274,45 @@ func short(digest string) string {
|
||||
return digest
|
||||
}
|
||||
|
||||
// drainDepth is how many declarations the host will hold unacknowledged while it looks for a newer
|
||||
// one; drainWindow is how long it waits for another to follow the one it has. Both small: a push is
|
||||
// rare and a backlog is the exception this exists for, not the shape of ordinary traffic.
|
||||
// drainDepth is how many declarations the link holds unread; drainWindow is how long it waits for
|
||||
// another to follow the one it has. Both small: a push is rare and a backlog is the exception this
|
||||
// exists for, not the shape of ordinary traffic.
|
||||
const (
|
||||
drainDepth = 16
|
||||
drainWindow = 750 * time.Millisecond
|
||||
)
|
||||
|
||||
// newest takes what is already waiting behind `first` and returns the last of them to apply, and
|
||||
// the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a
|
||||
// declaration in flight from the mesh arrives within that; one that does not is the next push.
|
||||
// gather takes what is already waiting behind `first`, in the order it arrived. It waits `window`
|
||||
// for a straggler after each arrival and no longer: a declaration in flight from the mesh arrives
|
||||
// within that; one that does not is the next push. Which of them is applied is the queue's to decide
|
||||
// (pick), across everything delivered and not yet taken.
|
||||
//
|
||||
// **Its job narrows once declarations are state rather than messages, and does not disappear.**
|
||||
// On the bus being built, a declaration is last-per-subject (novox/hq design 29 §4), so a node
|
||||
// that was away receives exactly the current one instead of a queue of superseded ones — the
|
||||
// catch-up half of what this does is then the stream's. And a stream sequence orders them
|
||||
// definitively, where this window only infers order from arrival time, which is the wire-level
|
||||
// answer to novox/hq issue 107.
|
||||
// On the bus a declaration is last-per-subject (novox/hq design 29 §4), so a node that was away
|
||||
// receives exactly the current one instead of a queue of superseded ones — the catch-up half is the
|
||||
// stream's. And a declaration's order decides which is newest, where this window only infers it from
|
||||
// arrival time, which is the wire-level answer to novox/hq issue 107.
|
||||
//
|
||||
// What remains is the live case: three pushes in quick succession to a *connected* node are
|
||||
// three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not
|
||||
// deleted — and saying which half goes is worth more than a note that it "can probably be
|
||||
// removed", which is how a load-bearing window gets deleted by somebody in a hurry.
|
||||
func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) {
|
||||
latest := first
|
||||
var superseded []Declaration
|
||||
// What remains is the live case: three pushes in quick succession to a *connected*, idle node are
|
||||
// three deliveries, whatever the stream later retains, and gathering them is what makes them one
|
||||
// apply rather than two. So this is narrowed, not deleted — and saying which half goes is worth more
|
||||
// than a note that it "can probably be removed", which is how a load-bearing window gets deleted by
|
||||
// somebody in a hurry.
|
||||
func gather(arriving <-chan Declaration, first Declaration, window time.Duration) []Declaration {
|
||||
batch := []Declaration{first}
|
||||
for {
|
||||
select {
|
||||
case next, ok := <-arriving:
|
||||
if !ok {
|
||||
return latest, superseded
|
||||
return batch
|
||||
}
|
||||
// **By sequence when both carry one, by arrival when either does not** (novox/hq
|
||||
// 04-ISSUES/107). Arrival is what this window had to go on, and it is wrong exactly
|
||||
// when it matters — a backlog drained out of order. A declaration that says where it
|
||||
// stands is believed over when it turned up; one that does not is the older
|
||||
// controller's, and arrival is all there is.
|
||||
if sequenceOf(next.Body()) < sequenceOf(latest.Body()) &&
|
||||
sequenceOf(next.Body()) > 0 && sequenceOf(latest.Body()) > 0 {
|
||||
superseded = append(superseded, next)
|
||||
continue
|
||||
}
|
||||
superseded = append(superseded, latest)
|
||||
latest = next
|
||||
batch = append(batch, next)
|
||||
case <-time.After(window):
|
||||
return latest, superseded
|
||||
return batch
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// sequenceOf is the order a signed declaration claims, or zero when it claims none or cannot be
|
||||
// read. Read from the envelope alone; the signature is verified later, when the winner is applied,
|
||||
// and a forged message that lied about its sequence would only set aside real ones — which are
|
||||
// reported as set aside, and the next push sends the current one again.
|
||||
func sequenceOf(body []byte) int64 {
|
||||
var signed Signed
|
||||
if err := json.Unmarshal(body, &signed); err != nil {
|
||||
return 0
|
||||
}
|
||||
var d struct {
|
||||
Sequence int64 `json:"sequence"`
|
||||
}
|
||||
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
|
||||
return 0
|
||||
}
|
||||
return d.Sequence
|
||||
}
|
||||
|
||||
// declaredIn is the id a signed declaration carries, for a report about one that was not applied.
|
||||
// Empty if the message is not one — a forged or garbled message is refused by handleBody when its
|
||||
// turn comes; here it is only named.
|
||||
func declaredIn(body []byte) string {
|
||||
var signed Signed
|
||||
if err := json.Unmarshal(body, &signed); err != nil {
|
||||
return ""
|
||||
}
|
||||
var d struct {
|
||||
Declared string `json:"declared"`
|
||||
}
|
||||
if err := json.Unmarshal(signed.Declaration, &d); err != nil {
|
||||
return ""
|
||||
}
|
||||
return d.Declared
|
||||
}
|
||||
|
||||
// handleBody is the whole of deciding whether to trust a message, separated from the broker so it
|
||||
// can be tested as the security check it is rather than as message plumbing.
|
||||
func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report {
|
||||
@@ -514,7 +371,7 @@ func publishReport(ctx context.Context, bus Bus, m Membership, report Report,
|
||||
// publishAlive says this node is here, and nothing else.
|
||||
func publishAlive(ctx context.Context, bus Bus, m Membership, say Announce,
|
||||
timeout time.Duration) {
|
||||
body, err := json.Marshal(Alive{Node: m.Node})
|
||||
body, err := json.Marshal(Alive{Node: m.Node, IntervalSeconds: int(AliveEvery / time.Second)})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
// Package validate is the node-engine's own judgement of a declaration, for whoever must know before
|
||||
// a declaration is sent whether a machine would take it (novox/hq to-be 45 §4, D1; §9, the merge gate).
|
||||
//
|
||||
// **One validator.** The controller composed declarations the host then refused whole — a manifest
|
||||
// the catalogue check passed (novox/hq issue 236), a consumer's identity too long for the machine's
|
||||
// names (issue 263) — because the only validator was the one the host runs as it applies. A second,
|
||||
// written in the controller from the host's rules, would drift from them the first time either
|
||||
// changed. So the host's own parsing is exported here, unchanged: what the controller's self-check and
|
||||
// the merge gate run is what every machine runs.
|
||||
//
|
||||
// It judges a declaration as it arrives over the link: actions are refused, as the host refuses them
|
||||
// from the link (novox/hq ADR 0005). Nothing here touches a machine.
|
||||
package validate
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"github.com/novox/mesh-host/internal/declaration"
|
||||
)
|
||||
|
||||
// Declaration is every problem the node-engine would refuse a declaration for, each in its own words;
|
||||
// nil when it would take it. The body is the declaration as the controller composes it — the bytes a
|
||||
// signature is made over — not the signed envelope.
|
||||
func Declaration(raw []byte) []string {
|
||||
_, err := declaration.Parse(raw)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
var refused *declaration.RefusalError
|
||||
if errors.As(err, &refused) {
|
||||
return refused.Problems
|
||||
}
|
||||
return []string{err.Error()}
|
||||
}
|
||||
|
||||
// Version is the declaration vocabulary this validator speaks: a declaration of another version is
|
||||
// refused whole by Declaration, as by the host.
|
||||
const Version = declaration.Version
|
||||
@@ -0,0 +1,37 @@
|
||||
package validate
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// **The validator a controller imports is the host's**: a declaration the host takes passes, and one
|
||||
// it refuses is refused with the host's own words, every problem at once.
|
||||
func TestTheValidatorIsTheHosts(t *testing.T) {
|
||||
good := `{"declaration": 1, "resources": [{"id": "d", "type": "directory", "path": "/var/lib/x", "mode": "0755"}]}`
|
||||
if problems := Declaration([]byte(good)); problems != nil {
|
||||
t.Fatalf("a declaration the host takes was refused: %v", problems)
|
||||
}
|
||||
// And one carrying the lease epoch and sequence the controller sends it under (novox/hq to-be 45
|
||||
// §6): a controller validating with this package may send them.
|
||||
ordered := `{"declaration": 1, "epoch": 57, "sequence": 3, "resources": [{"id": "d", "type": "directory", "path": "/var/lib/x", "mode": "0755"}]}`
|
||||
if problems := Declaration([]byte(ordered)); problems != nil {
|
||||
t.Fatalf("a declaration carrying its order was refused: %v", problems)
|
||||
}
|
||||
for _, bad := range []string{
|
||||
`{"declaration": 1, "resources": []}`,
|
||||
`{"declaration": 99, "resources": [{"id": "d", "type": "directory", "path": "/x"}]}`,
|
||||
`{"declaration": 1, "resources": [{"id": "d", "type": "nothing-the-host-knows"}]}`,
|
||||
`not json`,
|
||||
} {
|
||||
problems := Declaration([]byte(bad))
|
||||
if len(problems) == 0 {
|
||||
t.Errorf("a declaration the host refuses was passed: %s", bad)
|
||||
}
|
||||
for _, p := range problems {
|
||||
if strings.TrimSpace(p) == "" {
|
||||
t.Errorf("a problem said nothing, for %s", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user