Merge pull request 'Each send is numbered, inside the signed bytes' (#160) from feat/107-each-send-is-numbered into main
This commit was merged in pull request #160.
This commit is contained in:
@@ -338,6 +338,12 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
sentDigest := map[string]string{}
|
sentDigest := map[string]string{}
|
||||||
defer release()
|
defer release()
|
||||||
for _, s := range sending {
|
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()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -564,6 +570,9 @@ func sendRound(ctx context.Context, open *stores, names []string,
|
|||||||
return compose(held, node)
|
return compose(held, node)
|
||||||
})
|
})
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
|
if err := number(ctx, open.inventory, &s); err != nil {
|
||||||
|
return refused, err
|
||||||
|
}
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return refused, err
|
return refused, err
|
||||||
@@ -645,6 +654,9 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
defer server.Close()
|
defer server.Close()
|
||||||
|
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
|
if err := number(ctx, inv, &s); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -694,6 +706,12 @@ func wouldSend(ctx context.Context, open *stores,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
continue
|
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()
|
body, err := declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -827,3 +845,17 @@ func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]brok
|
|||||||
}
|
}
|
||||||
return out, nil
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -18,6 +18,10 @@ import (
|
|||||||
// make a machine look out of date for ever, or send something `plan` never showed.
|
// make a machine look out of date for ever, or send something `plan` never showed.
|
||||||
type sendable struct {
|
type sendable struct {
|
||||||
Resources []map[string]any
|
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 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 existed: an older host parses the envelope strictly and would refuse the key.
|
||||||
Adoption *adoptionEnvelope
|
Adoption *adoptionEnvelope
|
||||||
@@ -38,6 +42,9 @@ func (s sendable) Body() ([]byte, error) {
|
|||||||
if s.Adoption != nil {
|
if s.Adoption != nil {
|
||||||
envelope["adoption"] = s.Adoption
|
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
|
// 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
|
// (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".
|
// is meant, so a truncated or mis-composed body is never mistaken for "own nothing".
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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;
|
||||||
@@ -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)
|
`update node set host_version = $2, last_seen = now() where id = $1`, id, version)
|
||||||
return err
|
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
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user