diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index c2e4194..7675985 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -338,6 +338,12 @@ func pushCommand(ctx context.Context, args []string) error { sentDigest := map[string]string{} defer release() for _, s := range sending { + // Numbered under the hold, one higher than the last, before the body exists — the number is + // inside the signed bytes, so a replayed older declaration cannot borrow a newer one's + // (novox/hq 04-ISSUES/107). + if err := number(ctx, inv, &s); err != nil { + return err + } body, err := s.declared.Body() if err != nil { return err @@ -564,6 +570,9 @@ func sendRound(ctx context.Context, open *stores, names []string, return compose(held, node) }) for _, s := range sending { + if err := number(ctx, open.inventory, &s); err != nil { + return refused, err + } body, err := s.declared.Body() if err != nil { return refused, err @@ -645,6 +654,9 @@ func sendTo(ctx context.Context, open *stores, names []string) error { defer server.Close() for _, s := range sending { + if err := number(ctx, inv, &s); err != nil { + return err + } body, err := s.declared.Body() if err != nil { return err @@ -694,6 +706,12 @@ func wouldSend(ctx context.Context, open *stores, if err != nil { continue } + // Composed with the number the machine was LAST sent, so this is byte for byte what it was + // sent when nothing else changed. A fresh number here would make every machine read as + // behind for ever (novox/hq 04-ISSUES/107). + if declared.Sequence, err = open.inventory.Sequence(ctx, n.ID); err != nil { + return nil, err + } body, err := declared.Body() if err != nil { return nil, err @@ -827,3 +845,17 @@ func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]brok } return out, nil } + +// number gives one send the next sequence for its node (novox/hq 04-ISSUES/107). +func number(ctx context.Context, inv *inventory.Inventory, s *readyNode) error { + record, err := inv.NodeByName(ctx, s.node) + if err != nil { + return err + } + seq, err := inv.NextSequence(ctx, record.ID) + if err != nil { + return err + } + s.declared.Sequence = seq + return nil +} diff --git a/cmd/mesh-controller/sendable.go b/cmd/mesh-controller/sendable.go index fdebf7e..f6a69d5 100644 --- a/cmd/mesh-controller/sendable.go +++ b/cmd/mesh-controller/sendable.go @@ -18,6 +18,10 @@ import ( // make a machine look out of date for ever, or send something `plan` never showed. type sendable struct { Resources []map[string]any + // Sequence orders this send against every other to the same node: one higher each time, taken + // under the node's hold just before the body is made (novox/hq 04-ISSUES/107). Zero is not sent + // at all, which a host reads as "no order claimed" — the shape of every declaration before this. + Sequence int64 // Adoption is nil for a converged node, and then the body is byte for byte what it was before // adoption existed: an older host parses the envelope strictly and would refuse the key. Adoption *adoptionEnvelope @@ -38,6 +42,9 @@ func (s sendable) Body() ([]byte, error) { if s.Adoption != nil { envelope["adoption"] = s.Adoption } + if s.Sequence > 0 { + envelope["sequence"] = s.Sequence + } // An empty declaration is deliberate here — the node owns nothing the mesh put there // (novox/hq issue 127) — and the host refuses an empty body unless it is told the emptiness // is meant, so a truncated or mis-composed body is never mistaken for "own nothing". diff --git a/cmd/mesh-controller/sequence_test.go b/cmd/mesh-controller/sequence_test.go new file mode 100644 index 0000000..5d2dbd6 --- /dev/null +++ b/cmd/mesh-controller/sequence_test.go @@ -0,0 +1,77 @@ +package main + +import ( + "encoding/json" + "testing" +) + +// A declaration's only identity was the digest of its bytes; the controller already held a per-node +// lock and recorded each send, so the order existed and was thrown away at the wire (novox/hq +// 04-ISSUES/107). + +func TestASendCarriesItsNumberInsideTheSignedBytes(t *testing.T) { + body, err := sendable{Resources: []map[string]any{{"id": "x", "type": "file"}}, Sequence: 7}.Body() + if err != nil { + t.Fatal(err) + } + var env map[string]any + if err := json.Unmarshal(body, &env); err != nil { + t.Fatal(err) + } + if got, _ := env["sequence"].(float64); got != 7 { + t.Fatalf("the body carries sequence %v, wanted 7", env["sequence"]) + } +} + +func TestAnUnnumberedSendIsByteForByteWhatItWasBefore(t *testing.T) { + // Zero is not sent at all. A host reads absence as "no order claimed" — the shape of every + // declaration before this — so an older host, or the read-only comparison against a machine + // sent nothing since sends were numbered, sees exactly the bytes it always saw. + body, err := sendable{Resources: []map[string]any{{"id": "x", "type": "file"}}}.Body() + if err != nil { + t.Fatal(err) + } + var env map[string]any + if err := json.Unmarshal(body, &env); err != nil { + t.Fatal(err) + } + if _, present := env["sequence"]; present { + t.Fatalf("a send numbered zero put a sequence on the wire: %s", body) + } +} + +func TestEachSendToANodeIsOneHigherAndReadable(t *testing.T) { + open := aMesh(t) + record, err := open.inventory.NodeByName(t.Context(), "anchor") + if err != nil { + t.Fatal(err) + } + // Sent nothing since numbering existed: what it would be sent is composed with zero, which is + // not on the wire, which is what it was actually sent. + if n, err := open.inventory.Sequence(t.Context(), record.ID); err != nil || n != 0 { + t.Fatalf("a fresh node reads sequence %d, %v", n, err) + } + first, err := open.inventory.NextSequence(t.Context(), record.ID) + if err != nil { + t.Fatal(err) + } + second, err := open.inventory.NextSequence(t.Context(), record.ID) + if err != nil { + t.Fatal(err) + } + if first != 1 || second != 2 { + t.Fatalf("two sends were numbered %d and %d", first, second) + } + // And the read path sees the last one taken, so the comparison composes what was sent. + if n, err := open.inventory.Sequence(t.Context(), record.ID); err != nil || n != 2 { + t.Fatalf("after two sends the node reads sequence %d, %v", n, err) + } + // Another node counts on its own. + other, err := open.inventory.NodeByName(t.Context(), "laptop") + if err != nil { + t.Fatal(err) + } + if n, err := open.inventory.NextSequence(t.Context(), other.ID); err != nil || n != 1 { + t.Fatalf("a second node's first send was numbered %d, %v", n, err) + } +} diff --git a/internal/inventory/migrations/0046-each-send-is-numbered.sql b/internal/inventory/migrations/0046-each-send-is-numbered.sql new file mode 100644 index 0000000..6102d3d --- /dev/null +++ b/internal/inventory/migrations/0046-each-send-is-numbered.sql @@ -0,0 +1,11 @@ +-- The order of what a machine was sent, so a host can tell an older declaration from a newer. +-- +-- novox/hq 04-ISSUES/107. A declaration's only identity was the digest of its bytes. The host could +-- say "this is not the last one" and could not say "this is older", so a backlog drained out of +-- order applied a declaration the mesh had already superseded. The controller already holds a +-- per-node lock while it composes and records each send — the order existed and was thrown away at +-- the wire. +-- +-- One counter per node, taken under that hold, one higher per send. Null is a machine sent nothing +-- since this existed, which is not the same as a machine sent nothing. +alter table node add column sequence bigint; diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index 30c1423..fef70aa 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -1005,3 +1005,34 @@ func (i *Inventory) RecordHostVersion(ctx context.Context, id, version string) e `update node set host_version = $2, last_seen = now() where id = $1`, id, version) return err } + +// NextSequence takes the next number for a declaration to this node, one higher than the last it +// was sent (novox/hq 04-ISSUES/107). +// +// **One statement, so two composers cannot take the same number.** The caller holds the node while it +// composes and sends, so in practice there is one; the increment is atomic anyway, because a rule +// that is true only while a lock is held somewhere else is a rule nobody can see from here. +func (i *Inventory) NextSequence(ctx context.Context, id string) (int64, error) { + var n int64 + err := i.store.Pool().QueryRow(ctx, + `update node set sequence = coalesce(sequence, 0) + 1 where id = $1 returning sequence`, id).Scan(&n) + if err != nil { + return 0, fmt.Errorf("taking the next sequence for %s: %w", id, err) + } + return n, nil +} + +// Sequence is the number of the last declaration this node was sent, and zero for one sent nothing +// since sends were numbered. Read, not taken: what the mesh WOULD send is composed with this, so it +// is byte for byte what it DID send when nothing else changed — a comparison that took a fresh number +// would read every machine as behind for ever (novox/hq 04-ISSUES/107). +func (i *Inventory) Sequence(ctx context.Context, id string) (int64, error) { + var n *int64 + if err := i.store.Pool().QueryRow(ctx, `select sequence from node where id = $1`, id).Scan(&n); err != nil { + return 0, fmt.Errorf("reading the sequence of %s: %w", id, err) + } + if n == nil { + return 0, nil + } + return *n, nil +}