From 1a41b88ed3b145a650d3812aeebf9a931849a813 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 14:06:28 +0200 Subject: [PATCH 1/5] Nothing the control queue carries is lost while the store restarts: an enrolment claims its token and spends it last, and is asked to try again; build results, upgrades and catch-ups are handed back, bounded (novox/hq issue 083) --- internal/inventory/inventory_test.go | 64 ++++++++ ...-a-token-is-claimed-before-it-is-spent.sql | 11 ++ internal/inventory/nodes.go | 54 +++++++ internal/link/enrol_claim_test.go | 44 ++++++ internal/link/enrolment.go | 51 ++++-- internal/link/protocol.go | 5 + internal/link/serve.go | 147 +++++++++++++----- internal/link/store_window_test.go | 116 ++++++++++++++ 8 files changed, 433 insertions(+), 59 deletions(-) create mode 100644 internal/inventory/migrations/0028-a-token-is-claimed-before-it-is-spent.sql create mode 100644 internal/link/enrol_claim_test.go create mode 100644 internal/link/store_window_test.go diff --git a/internal/inventory/inventory_test.go b/internal/inventory/inventory_test.go index 0ace1cc..86aea3e 100644 --- a/internal/inventory/inventory_test.go +++ b/internal/inventory/inventory_test.go @@ -383,3 +383,67 @@ func TestBeingHeardFromDoesNotChangeWhatANodeOwns(t *testing.T) { t.Errorf("after a bare word that the node is here, the mesh believes it owns %v", owned) } } + +// **A token is claimed, then spent** (novox/hq issue 083). A presenter may claim it again — an +// enrolment interrupted by a restarting store asks again with the same key — while another is held +// off until the lease lapses; and it is spent only by the one holding the claim. +func TestATokenIsClaimedByOnePresenterAndSpentOnlyByIt(t *testing.T) { + inv := fresh(t) + ctx := t.Context() + if _, err := inv.AddNode(ctx, "laptop"); err != nil { + t.Fatal(err) + } + issued, err := inv.IssueToken(ctx, "laptop", time.Hour) + if err != nil { + t.Fatal(err) + } + node, err := inv.Claim(ctx, issued.Secret, "key-a") + if err != nil || node.Name != "laptop" { + t.Fatalf("a fresh token was not claimed for its node: %v %q", err, node.Name) + } + if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { + t.Fatalf("the presenter holding the claim could not claim again after an interruption: %v", err) + } + if _, err := inv.Claim(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenInUse) { + t.Fatalf("a second presenter was not held off while the claim is live: %v", err) + } + if err := inv.Spend(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("a presenter not holding the claim spent the token: %v", err) + } + if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil { + t.Fatalf("the presenter holding the claim could not spend it: %v", err) + } + for _, by := range []string{"key-a", "key-b"} { + if _, err := inv.Claim(ctx, issued.Secret, by); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("a spent token was claimed again by %s: %v", by, err) + } + } +} + +// A claim lapses: a host that gave up and was started over, with keys of its own, is not held off +// for longer than the lease. +func TestAClaimThatLapsedCanBeTakenByAnotherPresenter(t *testing.T) { + inv := fresh(t) + ctx := t.Context() + if _, err := inv.AddNode(ctx, "laptop"); err != nil { + t.Fatal(err) + } + issued, err := inv.IssueToken(ctx, "laptop", time.Hour) + if err != nil { + t.Fatal(err) + } + if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { + t.Fatal(err) + } + if _, err := inv.store.Pool().Exec(ctx, + `update enrolment_token set claimed_until = now() - interval '1 second' where secret = $1`, + hashSecret(issued.Secret)); err != nil { + t.Fatal(err) + } + if _, err := inv.Claim(ctx, issued.Secret, "key-b"); err != nil { + t.Fatalf("a lapsed claim held off a new presenter: %v", err) + } + if err := inv.Spend(ctx, issued.Secret, "key-a"); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("the presenter whose claim lapsed could still spend the token: %v", err) + } +} diff --git a/internal/inventory/migrations/0028-a-token-is-claimed-before-it-is-spent.sql b/internal/inventory/migrations/0028-a-token-is-claimed-before-it-is-spent.sql new file mode 100644 index 0000000..a1cc42c --- /dev/null +++ b/internal/inventory/migrations/0028-a-token-is-claimed-before-it-is-spent.sql @@ -0,0 +1,11 @@ +-- A token is claimed by the enrolment presenting it, and spent only when that enrolment has +-- written everything it needs (novox/hq 04-ISSUES/083). +-- +-- Spending came first and the node's keys after, in another database: a store that went away +-- between the two left a spent token and a node with no key, and the host — which makes new keys +-- on every attempt — could not try again. The claim holds the token for one presenter for a short +-- lease, so an attempt that failed part-way can be made again by the same presenter, and a second +-- presenter cannot interleave with the first. + +alter table enrolment_token add column claimed_by text; +alter table enrolment_token add column claimed_until timestamptz; diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index 8cb1ccb..fc4e0b9 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -212,6 +212,60 @@ var ErrTokenRefused = errors.New("that token cannot be used") // The update is the check: one statement that both finds a live token and marks it used, so two // simultaneous redemptions of one secret cannot both succeed. Reading first and writing second // would leave exactly that gap. +// ClaimLease is how long a claimed token is held for the one presenter that claimed it. Long +// enough for an enrolment to be tried again through a store restart; short enough that a host +// which gave up and was started over, with keys of its own, is not kept waiting long. +const ClaimLease = 2 * time.Minute + +// ErrTokenInUse is a token another presenter holds a claim on right now. Not a refusal: the claim +// lapses, and asking again after it is the answer. +var ErrTokenInUse = errors.New("the token is being used by another enrolment") + +// Claim takes a token for one presenter — `by`, which names the key presenting it — for the length +// of a lease, and says which node it enrols. The same presenter may claim it again, as may anyone +// once the lease has lapsed; nothing is spent until Spend (novox/hq 04-ISSUES/083). +func (i *Inventory) Claim(ctx context.Context, secret, by string) (Node, error) { + var id string + err := i.store.Pool().QueryRow(ctx, + `update enrolment_token set claimed_by = $2, claimed_until = now() + $3::interval + where secret = $1 and redeemed is null and expires > now() + and (claimed_by is null or claimed_by = $2 or claimed_until < now()) + returning node`, hashSecret(secret), by, ClaimLease.String()).Scan(&id) + if errors.Is(err, pgx.ErrNoRows) { + // Unusable, or held by someone else — told apart, because the second passes. + var held bool + probe := i.store.Pool().QueryRow(ctx, + `select true from enrolment_token + where secret = $1 and redeemed is null and expires > now()`, hashSecret(secret)).Scan(&held) + if probe == nil && held { + return Node{}, ErrTokenInUse + } + return Node{}, ErrTokenRefused + } + if err != nil { + return Node{}, err + } + var n Node + err = i.store.Pool().QueryRow(ctx, + `select id, name, created from node where id = $1`, id).Scan(&n.ID, &n.Name, &n.Created) + return n, err +} + +// Spend makes a claimed token used, only for the presenter holding the claim. The last write of an +// enrolment, so a token is spent exactly when the node it enrolled is complete. +func (i *Inventory) Spend(ctx context.Context, secret, by string) error { + tag, err := i.store.Pool().Exec(ctx, + `update enrolment_token set redeemed = now() + where secret = $1 and redeemed is null and claimed_by = $2`, hashSecret(secret), by) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + return ErrTokenRefused + } + return nil +} + func (i *Inventory) Redeem(ctx context.Context, secret string) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, diff --git a/internal/link/enrol_claim_test.go b/internal/link/enrol_claim_test.go new file mode 100644 index 0000000..bf373f3 --- /dev/null +++ b/internal/link/enrol_claim_test.go @@ -0,0 +1,44 @@ +package link_test + +import ( + "crypto/ed25519" + "crypto/rand" + "errors" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A token another enrolment holds for the moment is "not now", not a refusal: the node asks again +// with the same request, and nothing is spent (novox/hq issue 083). +func TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain(t *testing.T) { + inv := inventory.ForTest(t) + ctx := t.Context() + if _, err := inv.AddNode(ctx, "laptop"); err != nil { + t.Fatal(err) + } + issued, err := inv.IssueToken(ctx, "laptop", time.Hour) + if err != nil { + t.Fatal(err) + } + if _, err := inv.Claim(ctx, issued.Secret, "another enrolment's key"); err != nil { + t.Fatal(err) + } + public, _, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: issued.Secret, PublicKey: public}) + if !errors.Is(err, link.ErrTryAgain) { + t.Fatalf("an enrolment met by a held token was not asked to try again: %v", err) + } + // And a token that cannot be used at all is still refused outright. + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: "not-a-token", PublicKey: public}) + if err == nil || errors.Is(err, link.ErrTryAgain) { + t.Fatalf("a token that cannot be used was not refused outright: %v", err) + } +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index b7975bb..c730404 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -4,7 +4,9 @@ import ( "context" "crypto/ed25519" "crypto/rand" + "crypto/sha256" "encoding/base64" + "encoding/hex" "errors" "fmt" "log" @@ -34,26 +36,35 @@ type Enrolment struct { // single statement that both finds and marks it, so two machines racing on one secret produce one // winner. Only then is a key recorded — because recording a key for a node whose token turned out // to be spent would leave the mesh believing a machine that never had the right to join. -func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, error) { +func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply EnrolReply, err error) { secret, public, profile := request.Secret, ed25519.PublicKey(request.PublicKey), request.Profile + // A store that could not be asked right now, or a token another presenter holds for the + // moment, is "not now": the node asks again with the same request (novox/hq issue 083). + defer func() { + if inventory.Unreachable(err) || errors.Is(err, inventory.ErrTokenInUse) { + err = fmt.Errorf("%w: %w", ErrTryAgain, err) + } + }() + if len(public) != ed25519.PublicKeySize { return EnrolReply{}, fmt.Errorf("a node presented a %d-byte key, and an identity is %d", len(public), ed25519.PublicKeySize) } - node, err := e.Inventory.Redeem(ctx, secret) + // Claimed, not spent: the token is held for this presenter while the node is written, and + // spent only as the last write. The node's identity lives in another database than the + // token, so the two cannot be one transaction; a failure between them used to leave a spent + // token and a node with no key, which the host — making new keys on every attempt — could not + // recover from. Every write below overwrites, so an attempt made again is safe. + by := claimant(public) + node, err := e.Inventory.Claim(ctx, secret, by) if err != nil { return EnrolReply{}, err } - // From here the token is gone whatever happens next, so anything that fails leaves a node - // record with no live key — which is visible and fixable with a new token, where a spent - // token believed to be unspent is neither. if _, err := e.Identity.RecordNodeKey(ctx, node.ID, public); err != nil { - return EnrolReply{}, fmt.Errorf( - "the token was spent and the key could not be recorded, so %s has no identity and "+ - "needs a new token: %w", node.Name, err) + return EnrolReply{}, fmt.Errorf("%s's key could not be recorded: %w", node.Name, err) } key, err := e.Identity.Active(ctx) @@ -61,7 +72,7 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, return EnrolReply{}, err } - reply := EnrolReply{ + reply = EnrolReply{ Accepted: true, Node: node.Name, Queue: QueueFor(node.Name), @@ -80,8 +91,7 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, } if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil { return EnrolReply{}, fmt.Errorf( - "the token was spent and %s's broker password could not be replaced: %w", - node.Name, err) + "%s's broker password could not be replaced: %w", node.Name, err) } reply.Password = password } @@ -98,24 +108,27 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, if request.ServingKey != "" { if err := e.Identity.RecordServingKey(ctx, node.ID, request.ServingKey); err != nil { return EnrolReply{}, fmt.Errorf( - "the token was spent and %s's serving key could not be recorded: %w", - node.Name, err) + "%s's serving key could not be recorded: %w", node.Name, err) } } if request.SealingKey != "" { if err := e.Inventory.RecordSealingKey(ctx, node.ID, request.SealingKey); err != nil { return EnrolReply{}, fmt.Errorf( - "the token was spent and %s's sealing key could not be recorded: %w", - node.Name, err) + "%s's sealing key could not be recorded: %w", node.Name, err) } } if request.OverlayKey != "" { if err := e.Inventory.RecordOverlayKey(ctx, node.ID, request.OverlayKey); err != nil { return EnrolReply{}, fmt.Errorf( - "the token was spent and %s's overlay key could not be recorded: %w", node.Name, err) + "%s's overlay key could not be recorded: %w", node.Name, err) } } + // Spent last, so a token is used exactly when the node it enrolled is complete. + if err := e.Inventory.Spend(ctx, secret, by); err != nil { + return EnrolReply{}, fmt.Errorf("%s was written and its token could not be spent: %w", node.Name, err) + } + if profile != nil { // Not fatal if it fails. The profile is what the control plane needs in order to decide // what this machine should run, and it is reported again on every connection — so losing @@ -125,6 +138,12 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, return reply, nil } +// claimant names the key presenting a token, so a claim can be held for it alone. +func claimant(public ed25519.PublicKey) string { + sum := sha256.Sum256(public) + return hex.EncodeToString(sum[:]) +} + // freshPassword is the node's own broker credential from enrolment onward. func freshPassword() (string, error) { raw := make([]byte, 32) diff --git a/internal/link/protocol.go b/internal/link/protocol.go index efac2ca..25e24fa 100644 --- a/internal/link/protocol.go +++ b/internal/link/protocol.go @@ -109,6 +109,11 @@ type EnrolReply struct { // Accepted says whether the node is now known. Accepted bool `json:"accepted"` + // TryAgain says the mesh cannot answer right now — its store is restarting, or the token is + // held for a moment by another enrolment — and the node should ask again with the same + // request. Nothing was spent (novox/hq issue 083). + TryAgain bool `json:"try_again,omitempty"` + // Node is the name the mesh has for this machine, which settles any disagreement: the token // was issued for a node record, and that record's name wins over what the machine called // itself. diff --git a/internal/link/serve.go b/internal/link/serve.go index a3c7f9c..3c55b8b 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -2,10 +2,13 @@ package link import ( "context" + "crypto/sha256" + "encoding/hex" "encoding/json" "errors" "fmt" "github.com/novox/mesh-controller/internal/envfile" + "github.com/novox/mesh-controller/internal/inventory" "log" "os" "time" @@ -91,7 +94,8 @@ type Upgrader interface { // Upgraded is told which module moved and between which commits. An error is logged and the // message is not requeued: an upgrade the control plane could not act on is not one it will // act on by being handed the same message again, and a poison message on a durable queue - // would stop every upgrade behind it. + // would stop every upgrade behind it — except the store unreachable for the moment, which is + // asked again for a bounded time (novox/hq issue 083). Upgraded(ctx context.Context, u Upgraded) error } @@ -332,27 +336,13 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { } if s.listener != nil { if err := s.listener.Heard(context.Background(), report); err != nil { - if errors.Is(err, ErrTryAgain) && s.keepTrying(report) { - // Kept, not acknowledged. The node reports an apply once, and a report lost here - // is a node the mesh never hears from again: the store restarting under the - // adoption that node just applied lost exactly that (novox/hq issue 082). - again := s.again - if again == 0 { - again = TryAgainAfter - } - s.log.Printf("could not record %s's report yet, and will again in %s: %v", report.Node, again, err) - select { - case <-ctx.Done(): - case <-time.After(again): - } - _ = delivery.Nack(false, true) + // Kept, not acknowledged, while the store cannot take it: the node reports an apply + // once, and a report lost here is a node the mesh never hears from again — the store + // restarting under the adoption that node just applied lost exactly that (issue 082). + what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) + if s.tryLater(ctx, delivery, what, err) { return } - if errors.Is(err, ErrTryAgain) { - s.log.Printf("LOST %s's report of declaration %s: the store has not come back in %s, "+ - "and the control queue cannot wait longer — the node is current but the mesh will "+ - "read it as unanswered until its next push: %v", report.Node, report.Declared, s.giveUpAfter(), err) - } // Said rather than swallowed. A report the mesh heard and failed to write down is a // node whose recovery copy is silently older than it looks. s.log.Printf("could not record %s's report: %v", report.Node, err) @@ -367,26 +357,58 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { default: s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied)) } - s.stopTrying(report) + s.settled(delivery) _ = delivery.Ack(false) } -// keepTrying says whether a report the store could not take is still within the time it may hold -// the queue, starting that clock on its first failure. -func (s *Server) keepTrying(r Report) bool { +// tryLater hands a message the store could not take right now back to the broker, to be asked +// again after a pause, and says whether it did (novox/hq issues 082, 083). +// +// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error +// the inventory reads as an outage. Anything else is an answer, and is left to the caller to +// settle. One message holds its queue at most giveUpAfter: the consumer takes one message at a +// time, so a message tried again holds everything behind it, and a store that has not come back +// in that long is not restarting. Past it, the message is let go with a line saying it was lost. +func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, what string, err error) bool { + if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) { + return false + } + key := waitingKey(delivery) if s.waiting == nil { s.waiting = map[string]time.Time{} } - key := r.Node + " " + r.Declared first, seen := s.waiting[key] if !seen { - s.waiting[key] = time.Now() - return true + first = time.Now() + s.waiting[key] = first } - return time.Since(first) < s.giveUpAfter() + if time.Since(first) >= s.giveUpAfter() { + delete(s.waiting, key) + s.log.Printf("LOST %s: the store has not come back in %s, and the queue cannot wait longer: %v", + what, s.giveUpAfter(), err) + return false + } + again := s.again + if again == 0 { + again = TryAgainAfter + } + s.log.Printf("could not keep %s yet, and will again in %s: %v", what, again, err) + select { + case <-ctx.Done(): + case <-time.After(again): + } + _ = delivery.Nack(false, true) + return true } -func (s *Server) stopTrying(r Report) { delete(s.waiting, r.Node+" "+r.Declared) } +// settled forgets a message's time spent waiting, once it has been handled either way. +func (s *Server) settled(delivery amqp.Delivery) { delete(s.waiting, waitingKey(delivery)) } + +// waitingKey is a message by its content: the same message handed back is the same key. +func waitingKey(delivery amqp.Delivery) string { + sum := sha256.Sum256(append([]byte(delivery.RoutingKey+"\x00"), delivery.Body...)) + return hex.EncodeToString(sum[:]) +} func (s *Server) giveUpAfter() time.Duration { if s.giveUp == 0 { @@ -403,11 +425,18 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { s.log.Printf("an enrolment request could not be read: %v", err) } else { accepted, err := s.enroller.Enrol(ctx, request) - if err != nil { + switch { + case errors.Is(err, ErrTryAgain): + // Not a refusal: nothing was spent, and the same request asked again will be + // answered. Replied at once rather than held, so the node — which is waiting on + // this answer — decides when to ask, and the queue behind it moves (issue 083). + reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"} + s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err) + case err != nil: // Logged in full here, where an operator can see it; sent back as one refusal, so // that somebody guessing learns nothing from which reason came back. s.log.Printf("refusing enrolment for %q: %v", request.Node, err) - } else { + default: reply = accepted s.log.Printf("enrolled %s", accepted.Node) } @@ -465,10 +494,17 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { return } if err := s.recorder.Built(ctx, result); err != nil { + // Kept while the store cannot take it: a build result lost here is never announced, and + // recording one twice is harmless — the build is kept by its id (issue 083). + if s.tryLater(ctx, delivery, fmt.Sprintf("a build result from %s", result.On), err) { + return + } s.log.Printf("cannot keep a build result from %s: %v", result.On, err) + s.settled(delivery) _ = delivery.Reject(false) return } + s.settled(delivery) switch { case result.Failed != "": s.log.Printf("%s could not build %s", result.On, result.Repository) @@ -478,26 +514,30 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { _ = delivery.Ack(false) } -// upgraded hands one announcement to whatever is following them. -// -// **Acknowledged whatever happens.** A failure here is the control plane being unable to act on an -// upgrade — a machine that cannot be resolved, a broker that will not take a declaration — and -// none of those get better by being handed the same message again. Requeuing would put a poison -// message at the head of a durable queue and stop every upgrade behind it, which turns one module -// nobody can push into a mesh that stops following its own catalogue. // catchingUp answers a catalogue that has just started and may have missed builds. // -// Acknowledged before the work, deliberately: a replay that fails is not one that succeeds by -// being handed the same request again, and the catalogue asks every time it starts. Requeueing a -// poison request would stop every later catch-up behind it. +// Acknowledged after the work. A replay that fails for a reason other than the store is not one +// that succeeds by being handed the same request again, so that is acknowledged and said; but a +// store that could not be read right now is asked again after a pause, bounded, rather than the +// request lost until the catalogue next restarts (issue 083). func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { - defer func() { _ = delivery.Ack(false) }() + requeued := false + defer func() { + if !requeued { + s.settled(delivery) + _ = delivery.Ack(false) + } + }() if s.replayer == nil { s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") return } announcements, err := s.replayer.Announceable(ctx) if err != nil { + if s.tryLater(ctx, delivery, "a catalogue's request to catch up", err) { + requeued = true + return + } s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err) return } @@ -516,8 +556,22 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent) } +// upgraded hands one announcement to whatever is following them. +// +// **Acknowledged whatever happens, but one thing.** A failure here is usually the control plane +// being unable to act on an upgrade — a machine that cannot be resolved, a broker that will not +// take a declaration — and none of those get better by being handed the same message again. +// Requeuing those would put a poison message at the head of a durable queue and stop every +// upgrade behind it. The one exception is the store unreachable for the moment, which does get +// better: that is asked again after a pause, for a bounded time (novox/hq issue 083). func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { - defer func() { _ = delivery.Ack(false) }() + requeued := false + defer func() { + if !requeued { + s.settled(delivery) + _ = delivery.Ack(false) + } + }() var u Upgraded if err := json.Unmarshal(delivery.Body, &u); err != nil { s.log.Printf("an upgrade announcement could not be read: %v", err) @@ -528,6 +582,13 @@ func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { return } if err := s.upgrader.Upgraded(ctx, u); err != nil { + // A store that could not be read right now is asked again, bounded: acting on an upgrade + // twice pushes the same declarations twice, which converges (issue 083). Any other + // failure is acknowledged, as before — requeued, it would stop every upgrade behind it. + if s.tryLater(ctx, delivery, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err) { + requeued = true + return + } s.log.Printf("%s moved to %s and the mesh could not act on it: %v", u.Module, short(u.Commit), err) } diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go new file mode 100644 index 0000000..d894335 --- /dev/null +++ b/internal/link/store_window_test.go @@ -0,0 +1,116 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "io" + "log" + "testing" + + "github.com/jackc/pgx/v5/pgconn" + amqp "github.com/rabbitmq/amqp091-go" +) + +// What a store restarting under an adoption answers with (novox/hq issues 082, 083). +var restarting = &pgconn.PgError{Code: "57P03", Message: "the database system is starting up"} + +type recordsWith struct{ err error } + +func (r recordsWith) Built(context.Context, BuildResult) error { return r.err } + +type upgradesWith struct{ err error } + +func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err } + +type replaysWith struct{ err error } + +func (r replaysWith) Announceable(context.Context) ([]Announcement, error) { return nil, r.err } + +type settledAs struct{ acked, nacked, requeued, rejected bool } + +func (a *settledAs) Ack(uint64, bool) error { a.acked = true; return nil } +func (a *settledAs) Nack(_ uint64, _ bool, requeue bool) error { + a.nacked, a.requeued = true, requeue + return nil +} +func (a *settledAs) Reject(uint64, bool) error { a.rejected = true; return nil } + +func a(t *testing.T, to *settledAs, key string, v any) amqp.Delivery { + t.Helper() + body, err := json.Marshal(v) + if err != nil { + t.Fatal(err) + } + return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body} +} + +func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0), again: 1} } + +// A build result the store could not take right now is handed back; one it refused is rejected, +// as before; one it kept is acknowledged. +func TestABuildResultWaitsOutARestartingStore(t *testing.T) { + built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"} + for _, c := range []struct { + what string + err error + want func(*settledAs) bool + }{ + {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked && !s.rejected }}, + {"refused", errors.New("no such module"), func(s *settledAs) bool { return s.rejected && !s.nacked }}, + {"kept", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + } { + s := quietServer() + s.recorder = recordsWith{err: c.err} + to := &settledAs{} + s.handleBuilt(context.Background(), a(t, to, KeyBuilt, built)) + if !c.want(to) { + t.Errorf("%s: a build result was settled as %+v", c.what, *to) + } + } +} + +// An upgrade announcement arriving while the store restarts is asked again; any other failure is +// acknowledged, so it cannot stop every upgrade behind it. +func TestAnUpgradeWaitsOutARestartingStoreAndNothingElse(t *testing.T) { + moved := Upgraded{Module: "gitea", Commit: "abcdef0123"} + for _, c := range []struct { + what string + err error + want func(*settledAs) bool + }{ + {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked }}, + {"cannot act", errors.New("anchor cannot be resolved"), func(s *settledAs) bool { return s.acked && !s.nacked }}, + {"acted", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + } { + s := quietServer() + s.upgrader = upgradesWith{err: c.err} + to := &settledAs{} + s.upgraded(context.Background(), a(t, to, "upgraded", moved)) + if !c.want(to) { + t.Errorf("%s: an upgrade was settled as %+v", c.what, *to) + } + } +} + +// A catalogue's request to catch up is acknowledged after the work, and asked again while the +// store cannot be read — not lost until the catalogue next restarts. +func TestACatchUpWaitsOutARestartingStore(t *testing.T) { + for _, c := range []struct { + what string + err error + want func(*settledAs) bool + }{ + {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked }}, + {"unreadable", errors.New("a build row is malformed"), func(s *settledAs) bool { return s.acked && !s.nacked }}, + {"nothing to replay", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, + } { + s := quietServer() + s.replayer = replaysWith{err: c.err} + to := &settledAs{} + s.catchingUp(context.Background(), a(t, to, "catch-up", map[string]string{})) + if !c.want(to) { + t.Errorf("%s: a catch-up request was settled as %+v", c.what, *to) + } + } +} From a3b7e830c8a41722e8a852309ef08b5277d90051 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 14:23:03 +0200 Subject: [PATCH 2/5] Review of 083: a message the store cannot take is held and retried on a ticker, not slept on, so enrolments are answered meanwhile; a newer one per subject supersedes; the password is replaced after the spend; the same presenter may finish after a lost answer; upgrades retry only on the store --- cmd/mesh-controller/upgrades.go | 15 +- internal/inventory/inventory_test.go | 16 +- internal/inventory/nodes.go | 43 ++++-- internal/link/enrolment.go | 49 +++--- internal/link/report_retry_test.go | 108 ++++++------- internal/link/serve.go | 219 +++++++++++++++++---------- internal/link/store_window_test.go | 14 +- 7 files changed, 280 insertions(+), 184 deletions(-) diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 63b5769..3476864 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -28,13 +28,15 @@ type following struct{ open *stores } func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { inv := f.open.inventory + // The store read first, and an outage there said as one, so the announcement is held and asked + // again (novox/hq issue 083). Only here: a push that fails further down is not asked again. decision, err := inv.UpgradeOf(ctx, u.Module) if err != nil { - return err + return notNow(err) } on, err := inv.Running(ctx, u.Module) if err != nil { - return err + return notNow(err) } if len(on) == 0 { fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit)) @@ -194,3 +196,12 @@ func (f following) Announceable(ctx context.Context) ([]link.Announcement, error } return out, nil } + +// notNow marks a store that could not be read right now, so the announcement is held rather than +// lost; anything else is returned as it was. +func notNow(err error) error { + if inventory.Unreachable(err) { + return fmt.Errorf("%w: %w", link.ErrTryAgain, err) + } + return err +} diff --git a/internal/inventory/inventory_test.go b/internal/inventory/inventory_test.go index 86aea3e..7ae1532 100644 --- a/internal/inventory/inventory_test.go +++ b/internal/inventory/inventory_test.go @@ -413,10 +413,18 @@ func TestATokenIsClaimedByOnePresenterAndSpentOnlyByIt(t *testing.T) { if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil { t.Fatalf("the presenter holding the claim could not spend it: %v", err) } - for _, by := range []string{"key-a", "key-b"} { - if _, err := inv.Claim(ctx, issued.Secret, by); !errors.Is(err, ErrTokenRefused) { - t.Fatalf("a spent token was claimed again by %s: %v", by, err) - } + // Spent: nobody else may claim it, ever. + if _, err := inv.Claim(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("a spent token was claimed by another presenter: %v", err) + } + // But the presenter that spent it may, and spend it again: its spend reached the store and the + // answer did not reach the node, which asked again — refusing it would lock out a machine the + // mesh holds as enrolled. + if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { + t.Fatalf("the presenter whose answer was lost after its spend was refused: %v", err) + } + if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil { + t.Fatalf("spending again by the same presenter failed: %v", err) } } diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index fc4e0b9..b3792c6 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -203,15 +203,6 @@ func (i *Inventory) IssueToken(ctx context.Context, nodeName string, validFor ti // token that had expired. var ErrTokenRefused = errors.New("that token cannot be used") -// Redeem spends a token and reports which node it was for. -// -// It does not issue an identity. What a node presents afterwards to prove it is that node is not -// decided anywhere (novox/hq ADR 0004 names the property, not the mechanism), and guessing at it -// in a migration is the most expensive guess available here. -// -// The update is the check: one statement that both finds a live token and marks it used, so two -// simultaneous redemptions of one secret cannot both succeed. Reading first and writing second -// would leave exactly that gap. // ClaimLease is how long a claimed token is held for the one presenter that claimed it. Long // enough for an enrolment to be tried again through a store restart; short enough that a host // which gave up and was started over, with keys of its own, is not kept waiting long. @@ -224,12 +215,17 @@ var ErrTokenInUse = errors.New("the token is being used by another enrolment") // Claim takes a token for one presenter — `by`, which names the key presenting it — for the length // of a lease, and says which node it enrols. The same presenter may claim it again, as may anyone // once the lease has lapsed; nothing is spent until Spend (novox/hq 04-ISSUES/083). +// +// A token this same presenter already spent is claimed again too: its spend reached the store and +// the answer did not reach the node, which asked again. Refusing it then would lock out a machine +// the mesh holds as enrolled — with the key it is still presenting. func (i *Inventory) Claim(ctx context.Context, secret, by string) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, `update enrolment_token set claimed_by = $2, claimed_until = now() + $3::interval - where secret = $1 and redeemed is null and expires > now() - and (claimed_by is null or claimed_by = $2 or claimed_until < now()) + where secret = $1 and expires > now() + and ((redeemed is null and (claimed_by is null or claimed_by = $2 or claimed_until < now())) + or (redeemed is not null and claimed_by = $2)) returning node`, hashSecret(secret), by, ClaimLease.String()).Scan(&id) if errors.Is(err, pgx.ErrNoRows) { // Unusable, or held by someone else — told apart, because the second passes. @@ -237,8 +233,12 @@ func (i *Inventory) Claim(ctx context.Context, secret, by string) (Node, error) probe := i.store.Pool().QueryRow(ctx, `select true from enrolment_token where secret = $1 and redeemed is null and expires > now()`, hashSecret(secret)).Scan(&held) - if probe == nil && held { + switch { + case probe == nil && held: return Node{}, ErrTokenInUse + case probe != nil && !errors.Is(probe, pgx.ErrNoRows): + // The store went away between the two questions: "not now", not a refusal. + return Node{}, probe } return Node{}, ErrTokenRefused } @@ -251,12 +251,13 @@ func (i *Inventory) Claim(ctx context.Context, secret, by string) (Node, error) return n, err } -// Spend makes a claimed token used, only for the presenter holding the claim. The last write of an -// enrolment, so a token is spent exactly when the node it enrolled is complete. +// Spend makes a claimed token used, only for the presenter holding the claim. The last write to the +// store in an enrolment, so a token is spent exactly when the node it enrolled is complete. Spent +// again by the same presenter is not an error: an answer lost after the first spend. func (i *Inventory) Spend(ctx context.Context, secret, by string) error { tag, err := i.store.Pool().Exec(ctx, - `update enrolment_token set redeemed = now() - where secret = $1 and redeemed is null and claimed_by = $2`, hashSecret(secret), by) + `update enrolment_token set redeemed = coalesce(redeemed, now()) + where secret = $1 and claimed_by = $2`, hashSecret(secret), by) if err != nil { return err } @@ -266,6 +267,16 @@ func (i *Inventory) Spend(ctx context.Context, secret, by string) error { return nil } +// Redeem spends a token in one step and reports which node it was for. Enrolment claims and then +// spends (Claim, Spend); this is the one-step form, kept for what spends a token outright. +// +// It does not issue an identity. What a node presents afterwards to prove it is that node is not +// decided anywhere (novox/hq ADR 0004 names the property, not the mechanism), and guessing at it +// in a migration is the most expensive guess available here. +// +// The update is the check: one statement that both finds a live token and marks it used, so two +// simultaneous redemptions of one secret cannot both succeed. Reading first and writing second +// would leave exactly that gap. func (i *Inventory) Redeem(ctx context.Context, secret string) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index c730404..b17ddf2 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -30,12 +30,15 @@ type Enrolment struct { Broker broker.Broker } -// Enrol spends the token and records what the node presented. +// Enrol records what the node presented and spends the token. // -// Order matters and it is the order things become irreversible. The token is spent first, in a -// single statement that both finds and marks it, so two machines racing on one secret produce one -// winner. Only then is a key recorded — because recording a key for a node whose token turned out -// to be spent would leave the mesh believing a machine that never had the right to join. +// Order matters and it is the order things become irreversible (novox/hq issue 083). The token is +// claimed first, in a single statement that both finds it and holds it for this presenter's key, so +// two machines racing on one secret produce one holder. Then everything the node presented is +// written — each write an overwrite, so an attempt interrupted by the store going away can be made +// again by the same presenter. Then the token is spent. Last, the token's secret stops being the +// node's broker password: done after the spend, because a password replaced by an attempt that +// then failed would be one nobody holds, and the node could not even log in to ask again. func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply EnrolReply, err error) { secret, public, profile := request.Secret, ed25519.PublicKey(request.PublicKey), request.Profile @@ -81,21 +84,6 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol Signer: key.Public, } - // The token's secret was the broker password up to this moment, which is what let this - // connection exist at all. It is replaced now, so the one-time thing stays one-time and the - // credential the node keeps for years is not the one that was pasted into a terminal. - if e.Management != nil { - password, err := freshPassword() - if err != nil { - return EnrolReply{}, err - } - if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil { - return EnrolReply{}, fmt.Errorf( - "%s's broker password could not be replaced: %w", node.Name, err) - } - reply.Password = password - } - // Recorded before the profile because the overlay is the first declaration this node will // receive, and without this key the mesh cannot compose one. A node enrolled with no overlay // key is a node the graph skips — an ordinary in-between state, and one worth leaving as @@ -124,11 +112,30 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol } } - // Spent last, so a token is used exactly when the node it enrolled is complete. + // Spent once the node is complete in the store. if err := e.Inventory.Spend(ctx, secret, by); err != nil { return EnrolReply{}, fmt.Errorf("%s was written and its token could not be spent: %w", node.Name, err) } + // The token's secret was the broker password up to this moment, which is what let this + // connection exist at all. It is replaced now, so the one-time thing stays one-time and the + // credential the node keeps for years is not the one that was pasted into a terminal. After + // the spend and not before: a replaced password on an attempt that failed would be held by + // nobody. If the broker will not take it now, the enrolment still stands — the node keeps + // the token's secret as its password, which it is told, and which is said here. + if e.Management != nil { + password, err := freshPassword() + if err != nil { + return EnrolReply{}, err + } + if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil { + log.Printf("%s is enrolled and its broker password could not be replaced, so it keeps "+ + "the token's secret as its password: %v", node.Name, err) + } else { + reply.Password = password + } + } + if profile != nil { // Not fatal if it fails. The profile is what the control plane needs in order to decide // what this machine should run, and it is reported again on every connection — so losing diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go index 7e2f61e..8a175d7 100644 --- a/internal/link/report_retry_test.go +++ b/internal/link/report_retry_test.go @@ -20,85 +20,85 @@ func (a *saidTo) Nack(_ uint64, _ bool, requeue bool) error { return nil } func (a *saidTo) Reject(uint64, bool) error { return nil } +func (a *saidTo) unsettled() bool { return !a.acked && !a.nacked } type heardWith struct{ err error } func (h heardWith) Heard(context.Context, Report) error { return h.err } -func aReport(t *testing.T, to *saidTo) amqp.Delivery { +// switchable answers with whatever it is set to — the store away, then back. +type switchable struct{ err error } + +func (h *switchable) Heard(context.Context, Report) error { return h.err } + +var tag uint64 + +func aReport(t *testing.T, to *saidTo, node, declared string) amqp.Delivery { t.Helper() - body, err := json.Marshal(Report{Node: "anchor", Applied: []string{"store"}}) + body, err := json.Marshal(Report{Node: node, Declared: declared, Applied: []string{"store"}}) if err != nil { t.Fatal(err) } - return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body} + tag++ + return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body, DeliveryTag: tag} } -// A report the store could not take right now goes back to the broker to be asked again; one the -// store answered no to is acknowledged, or it would come back for ever (novox/hq issue 082). -func TestAReportTheStoreCouldNotTakeIsKeptAndOneItRefusedIsNot(t *testing.T) { - quiet := log.New(io.Discard, "", 0) +func quiet() *log.Logger { return log.New(io.Discard, "", 0) } - notNow := &saidTo{} - s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet, again: 1} - s.handleReport(context.Background(), aReport(t, notNow)) - if notNow.acked || !notNow.nacked || !notNow.requeued { - t.Fatalf("a report the store could not take yet was not handed back to be asked again: %+v", notNow) +// A report the store could not take right now is held, unsettled, and recorded when the store is +// back; one the store answered no to is acknowledged; one recorded is acknowledged (issue 082, 083). +func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) { + store := &switchable{err: errors.Join(ErrTryAgain, errors.New("starting up"))} + s := &Server{listener: store, log: quiet()} + held := &saidTo{} + s.handleReport(context.Background(), aReport(t, held, "anchor", "d1")) + if !held.unsettled() || len(s.parked) != 1 { + t.Fatalf("a report the store could not take was not held: %+v, %d held", held, len(s.parked)) + } + store.err = nil + s.retryHeld(context.Background()) + if !held.acked || len(s.parked) != 0 { + t.Fatalf("a held report was not recorded once the store was back: %+v, %d held", held, len(s.parked)) } refused := &saidTo{} - s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet, again: 1} - s.handleReport(context.Background(), aReport(t, refused)) + s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet()} + s.handleReport(context.Background(), aReport(t, refused, "anchor", "d1")) if !refused.acked || refused.nacked { - t.Fatalf("a report the store answered no to was not acknowledged, so it would spin: %+v", refused) + t.Fatalf("a report the store answered no to was not acknowledged: %+v", refused) } +} - recorded := &saidTo{} - s = &Server{listener: heardWith{}, log: quiet} - s.handleReport(context.Background(), aReport(t, recorded)) - if !recorded.acked || recorded.nacked { - t.Fatalf("a recorded report was not acknowledged: %+v", recorded) +// A newer report from the same node supersedes one of its reports still held: recorded after the +// newer, the older would overwrite what the node is doing now. +func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) { + s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet()} + older, newer, other := &saidTo{}, &saidTo{}, &saidTo{} + s.handleReport(context.Background(), aReport(t, older, "anchor", "d1")) + s.handleReport(context.Background(), aReport(t, other, "laptop", "d7")) + s.handleReport(context.Background(), aReport(t, newer, "anchor", "d2")) + if !older.acked { + t.Fatalf("the older report was not set aside by the newer: %+v", older) + } + if !newer.unsettled() || !other.unsettled() || len(s.parked) != 2 { + t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held", + newer, other, len(s.parked)) } } // A store that has not come back within the bound is not restarting: the report is let go, loudly, -// rather than holding every enrolment and report behind it for ever (novox/hq issue 082, review). +// rather than held for ever. func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) { - quiet := log.New(io.Discard, "", 0) s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))}, - log: quiet, again: 1, giveUp: time.Millisecond} - - first := &saidTo{} - s.handleReport(context.Background(), aReport(t, first)) - if !first.requeued { - t.Fatalf("the first failure was not handed back: %+v", first) + log: quiet(), giveUp: time.Millisecond} + held := &saidTo{} + s.handleReport(context.Background(), aReport(t, held, "anchor", "d1")) + if !held.unsettled() { + t.Fatalf("the first failure was not held: %+v", held) } time.Sleep(5 * time.Millisecond) - later := &saidTo{} - s.handleReport(context.Background(), aReport(t, later)) - if !later.acked || later.nacked { - t.Fatalf("a report past the bound was not let go, so it would hold the queue for ever: %+v", later) - } - // Let go, and forgotten: the same report arriving fresh starts a new clock. - again := &saidTo{} - s.handleReport(context.Background(), aReport(t, again)) - if !again.requeued { - t.Fatalf("a report let go was not forgotten, so its next arrival is given no chance: %+v", again) - } -} - -// Shutting down does not wait out the pause. -func TestAReportBeingTriedAgainDoesNotHoldUpShutdown(t *testing.T) { - quiet := log.New(io.Discard, "", 0) - s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, - log: quiet, again: time.Hour} - ctx, cancel := context.WithCancel(context.Background()) - cancel() - done := make(chan struct{}) - go func() { s.handleReport(ctx, aReport(t, &saidTo{})); close(done) }() - select { - case <-done: - case <-time.After(5 * time.Second): - t.Fatal("a cancelled context still waited out the pause") + s.retryHeld(context.Background()) + if !held.acked || len(s.parked) != 0 { + t.Fatalf("a report past the bound was not let go: %+v, %d held", held, len(s.parked)) } } diff --git a/internal/link/serve.go b/internal/link/serve.go index 3c55b8b..9b8f3de 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -55,29 +55,40 @@ type Server struct { log *log.Logger upgrader Upgrader replayer Replayer - // again is how long a report the store could not take waits before it is handed back to - // the broker; zero means TryAgainAfter. giveUp is how long one report is kept trying before - // it is let go; zero means GiveUpAfter. waiting is when each report still trying first failed. - again time.Duration - giveUp time.Duration - waiting map[string]time.Time + // Messages the store could not take right now, held unacknowledged and tried again on a + // ticker, by subject (novox/hq issues 082, 083). again is the ticker's interval, zero meaning + // TryAgainAfter; giveUp is how long one is kept, zero meaning GiveUpAfter. + again time.Duration + giveUp time.Duration + parked map[string]*held +} + +// held is one message the store could not take, kept to be tried again. +type held struct { + delivery amqp.Delivery + retry func(context.Context, amqp.Delivery) + what string + first time.Time } // ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and // asking again, as when the store is restarting (novox/hq issue 082). var ErrTryAgain = errors.New("not now, try again") -// TryAgainAfter is the pause before a report the store could not take goes back to the broker. -// The consumer takes one message at a time, so without it a restarting store would be asked in -// a tight loop; a store comes back in seconds, and a report a few seconds late is still current. +// TryAgainAfter is how often messages the store could not take are tried again. A store comes +// back in seconds, and a report a few seconds late is still current. const TryAgainAfter = 2 * time.Second -// GiveUpAfter bounds how long one report holds the queue. The consumer takes one message at a -// time, so a report being tried again holds every enrolment, build result and other report -// behind it; a store that has not come back in this long is not restarting, and holding the -// mesh's control queue for it would turn one lost report into a mesh that answers nothing. +// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this +// long is not restarting, and the message is let go with a line saying it was lost. const GiveUpAfter = 2 * time.Minute +// Prefetch is how many messages the broker hands the control plane before it has settled them. +// More than one because a message the store could not take is held, unsettled, while the loop goes +// on answering others — an enrolment above all, which a host is waiting on (novox/hq issue 083). +// Bounded, because what is held is also what the broker has not kept on its own disk as pending. +const Prefetch = 64 + // Records tells the server where to keep build results. // // Set after Connect rather than passed to it, because a control plane that only publishes — the @@ -206,10 +217,11 @@ func (s *Server) Close() { // receive half of what it expects — a fault this project has already had, between a module's // daemon and its capability server. func (s *Server) Serve(ctx context.Context) error { - // Prefetch of one. The control plane writes to a database per message, and a burst of - // enrolments delivered all at once would be held in memory rather than left on the broker, - // which is the one place they survive a restart. - if err := s.channel.Qos(1, 0, false); err != nil { + // A bounded prefetch rather than one. The loop still takes messages one at a time; what the + // prefetch buys is that a message the store could not take can be held while the loop goes on + // to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything held + // goes back to the broker if the control plane stops, because nothing held is acknowledged. + if err := s.channel.Qos(Prefetch, 0, false); err != nil { return err } @@ -250,10 +262,19 @@ func (s *Server) Serve(ctx context.Context) error { s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded) } + again := s.again + if again == 0 { + again = TryAgainAfter + } + ticker := time.NewTicker(again) + defer ticker.Stop() + for { select { case <-ctx.Done(): return nil + case <-ticker.C: + s.retryHeld(ctx) case delivery, ok := <-catchups: if !ok { if catchups != nil { @@ -334,13 +355,17 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { _ = delivery.Reject(false) return } + // A node's newer report supersedes one of its older reports still held: the older is its + // past, and recorded after the newer it would overwrite what the node is doing now. + subject := "report " + report.Node + s.supersede(subject, delivery) if s.listener != nil { if err := s.listener.Heard(context.Background(), report); err != nil { - // Kept, not acknowledged, while the store cannot take it: the node reports an apply + // Held, not acknowledged, while the store cannot take it: the node reports an apply // once, and a report lost here is a node the mesh never hears from again — the store // restarting under the adoption that node just applied lost exactly that (issue 082). what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) - if s.tryLater(ctx, delivery, what, err) { + if s.tryLater(delivery, subject, what, err, s.handle) { return } // Said rather than swallowed. A report the mesh heard and failed to write down is a @@ -357,59 +382,82 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { default: s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied)) } - s.settled(delivery) + s.settled(subject, delivery) _ = delivery.Ack(false) } -// tryLater hands a message the store could not take right now back to the broker, to be asked -// again after a pause, and says whether it did (novox/hq issues 082, 083). +// tryLater holds a message the store could not take right now, to be tried again on the ticker, +// and says whether it did (novox/hq issues 082, 083). // -// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error -// the inventory reads as an outage. Anything else is an answer, and is left to the caller to -// settle. One message holds its queue at most giveUpAfter: the consumer takes one message at a -// time, so a message tried again holds everything behind it, and a store that has not come back -// in that long is not restarting. Past it, the message is let go with a line saying it was lost. -func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, what string, err error) bool { +// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error the +// inventory reads as an outage. Anything else is an answer, and is left to the caller to settle. +// Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a +// host is waiting on is answered while a report waits for the store. One message is held at most +// giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it. +func (s *Server) tryLater(delivery amqp.Delivery, subject, what string, err error, + retry func(context.Context, amqp.Delivery)) bool { if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) { return false } - key := waitingKey(delivery) - if s.waiting == nil { - s.waiting = map[string]time.Time{} + if s.parked == nil { + s.parked = map[string]*held{} } - first, seen := s.waiting[key] - if !seen { - first = time.Now() - s.waiting[key] = first + h, ok := s.parked[subject] + if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag { + h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()} + s.parked[subject] = h + s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err) + return true } - if time.Since(first) >= s.giveUpAfter() { - delete(s.waiting, key) - s.log.Printf("LOST %s: the store has not come back in %s, and the queue cannot wait longer: %v", - what, s.giveUpAfter(), err) + if time.Since(h.first) >= s.giveUpAfter() { + delete(s.parked, subject) + s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err) return false } - again := s.again - if again == 0 { - again = TryAgainAfter - } - s.log.Printf("could not keep %s yet, and will again in %s: %v", what, again, err) - select { - case <-ctx.Done(): - case <-time.After(again): - } - _ = delivery.Nack(false, true) return true } -// settled forgets a message's time spent waiting, once it has been handled either way. -func (s *Server) settled(delivery amqp.Delivery) { delete(s.waiting, waitingKey(delivery)) } +// supersede drops a message held for a subject when a newer one for it arrives: the older is +// acknowledged, because acting on it after the newer would undo the newer. +func (s *Server) supersede(subject string, newer amqp.Delivery) { + h, ok := s.parked[subject] + if !ok || h.delivery.DeliveryTag == newer.DeliveryTag { + return + } + delete(s.parked, subject) + s.log.Printf("set aside %s: a newer one arrived", h.what) + _ = h.delivery.Ack(false) +} -// waitingKey is a message by its content: the same message handed back is the same key. -func waitingKey(delivery amqp.Delivery) string { - sum := sha256.Sum256(append([]byte(delivery.RoutingKey+"\x00"), delivery.Body...)) +// settled forgets a message once it has been handled either way. +func (s *Server) settled(subject string, delivery amqp.Delivery) { + if h, ok := s.parked[subject]; ok && h.delivery.DeliveryTag == delivery.DeliveryTag { + delete(s.parked, subject) + } +} + +// digest names a message by its content. +func digest(body []byte) string { + sum := sha256.Sum256(body) return hex.EncodeToString(sum[:]) } +// retryHeld tries every held message again. Each handler holds it again, settles it, or lets it +// go past the bound. +func (s *Server) retryHeld(ctx context.Context) { + for _, h := range s.snapshot() { + h.retry(ctx, h.delivery) + } +} + +func (s *Server) snapshot() []*held { + out := make([]*held, 0, len(s.parked)) + for _, h := range s.parked { + out = append(out, h) + } + return out +} + func (s *Server) giveUpAfter() time.Duration { if s.giveUp == 0 { return GiveUpAfter @@ -445,9 +493,8 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { s.reply(ctx, delivery, reply) // Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the - // request on the broker rather than having consumed it silently. Enrolment is idempotent - // only in the sense that the token is spent — a redelivery gets the refusal, which is - // correct and visible, where a lost request is neither. + // request on the broker rather than having consumed it silently. Asked again by the same + // presenter, an enrolment finishes: the token is held for its key and spent last (issue 083). _ = delivery.Ack(false) } @@ -493,18 +540,20 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { _ = delivery.Reject(false) return } + // Each build result its own subject: none supersedes another, and recording one twice is + // harmless — the build is kept by its id. + subject := "build " + digest(delivery.Body) if err := s.recorder.Built(ctx, result); err != nil { - // Kept while the store cannot take it: a build result lost here is never announced, and - // recording one twice is harmless — the build is kept by its id (issue 083). - if s.tryLater(ctx, delivery, fmt.Sprintf("a build result from %s", result.On), err) { + // Held while the store cannot take it: a build result lost here is never announced (083). + if s.tryLater(delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) { return } s.log.Printf("cannot keep a build result from %s: %v", result.On, err) - s.settled(delivery) + s.settled(subject, delivery) _ = delivery.Reject(false) return } - s.settled(delivery) + s.settled(subject, delivery) switch { case result.Failed != "": s.log.Printf("%s could not build %s", result.On, result.Repository) @@ -518,13 +567,16 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { // // Acknowledged after the work. A replay that fails for a reason other than the store is not one // that succeeds by being handed the same request again, so that is acknowledged and said; but a -// store that could not be read right now is asked again after a pause, bounded, rather than the -// request lost until the catalogue next restarts (issue 083). +// store that could not be read right now is held and asked again, bounded, rather than the +// request lost until the catalogue next restarts (issue 083). One request stands for all: a newer +// one supersedes one still held. func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { - requeued := false + const subject = "catch-up" + s.supersede(subject, delivery) + holding := false defer func() { - if !requeued { - s.settled(delivery) + if !holding { + s.settled(subject, delivery) _ = delivery.Ack(false) } }() @@ -534,8 +586,8 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { } announcements, err := s.replayer.Announceable(ctx) if err != nil { - if s.tryLater(ctx, delivery, "a catalogue's request to catch up", err) { - requeued = true + if s.tryLater(delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) { + holding = true return } s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err) @@ -560,19 +612,24 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { // // **Acknowledged whatever happens, but one thing.** A failure here is usually the control plane // being unable to act on an upgrade — a machine that cannot be resolved, a broker that will not -// take a declaration — and none of those get better by being handed the same message again. -// Requeuing those would put a poison message at the head of a durable queue and stop every -// upgrade behind it. The one exception is the store unreachable for the moment, which does get -// better: that is asked again after a pause, for a bounded time (novox/hq issue 083). +// take a declaration — and none of those get better by being handed the same message again. The +// one exception is the upgrader saying the store could not be read for the moment (ErrTryAgain): +// that is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts +// here, not an error that merely looks like an outage — a push that timed out on the second +// machine is not asked again, or the first would be pushed every few seconds for two minutes. +// A newer move of the same module supersedes one still held. func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { - requeued := false + var u Upgraded + _ = json.Unmarshal(delivery.Body, &u) + subject := "upgrade " + u.Module + s.supersede(subject, delivery) + holding := false defer func() { - if !requeued { - s.settled(delivery) + if !holding { + s.settled(subject, delivery) _ = delivery.Ack(false) } }() - var u Upgraded if err := json.Unmarshal(delivery.Body, &u); err != nil { s.log.Printf("an upgrade announcement could not be read: %v", err) return @@ -582,11 +639,9 @@ func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { return } if err := s.upgrader.Upgraded(ctx, u); err != nil { - // A store that could not be read right now is asked again, bounded: acting on an upgrade - // twice pushes the same declarations twice, which converges (issue 083). Any other - // failure is acknowledged, as before — requeued, it would stop every upgrade behind it. - if s.tryLater(ctx, delivery, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err) { - requeued = true + if errors.Is(err, ErrTryAgain) && + s.tryLater(delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { + holding = true return } s.log.Printf("%s moved to %s and the mesh could not act on it: %v", diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go index d894335..7fe9ffd 100644 --- a/internal/link/store_window_test.go +++ b/internal/link/store_window_test.go @@ -42,10 +42,13 @@ func a(t *testing.T, to *settledAs, key string, v any) amqp.Delivery { if err != nil { t.Fatal(err) } - return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body} + tag++ + return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body, DeliveryTag: tag} } -func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0), again: 1} } +func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0)} } + +func (a *settledAs) held() bool { return !a.acked && !a.nacked && !a.rejected } // A build result the store could not take right now is handed back; one it refused is rejected, // as before; one it kept is acknowledged. @@ -56,7 +59,7 @@ func TestABuildResultWaitsOutARestartingStore(t *testing.T) { err error want func(*settledAs) bool }{ - {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked && !s.rejected }}, + {"restarting", restarting, func(s *settledAs) bool { return s.held() }}, {"refused", errors.New("no such module"), func(s *settledAs) bool { return s.rejected && !s.nacked }}, {"kept", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, } { @@ -79,7 +82,8 @@ func TestAnUpgradeWaitsOutARestartingStoreAndNothingElse(t *testing.T) { err error want func(*settledAs) bool }{ - {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked }}, + {"the store away, said by the upgrader", errors.Join(ErrTryAgain, restarting), func(s *settledAs) bool { return s.held() }}, + {"a push that timed out", context.DeadlineExceeded, func(s *settledAs) bool { return s.acked && !s.nacked }}, {"cannot act", errors.New("anchor cannot be resolved"), func(s *settledAs) bool { return s.acked && !s.nacked }}, {"acted", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, } { @@ -101,7 +105,7 @@ func TestACatchUpWaitsOutARestartingStore(t *testing.T) { err error want func(*settledAs) bool }{ - {"restarting", restarting, func(s *settledAs) bool { return s.requeued && !s.acked }}, + {"restarting", restarting, func(s *settledAs) bool { return s.held() }}, {"unreadable", errors.New("a build row is malformed"), func(s *settledAs) bool { return s.acked && !s.nacked }}, {"nothing to replay", nil, func(s *settledAs) bool { return s.acked && !s.nacked }}, } { From 4567fa666caae25215b7c9e59160f213f828b9e9 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 14:33:13 +0200 Subject: [PATCH 3/5] =?UTF-8?q?Review=20of=20083:=20finishing=20an=20enrol?= =?UTF-8?q?ment=20whose=20token=20was=20spent=20takes=20proof=20of=20the?= =?UTF-8?q?=20key's=20private=20half,=20a=20live=20lease=20and=20a=20first?= =?UTF-8?q?=20delivery=20=E2=80=94=20a=20public=20key=20alone=20cannot=20r?= =?UTF-8?q?eplay=20a=20spent=20token;=20shutdown=20leaves=20held=20message?= =?UTF-8?q?s=20for=20the=20broker;=20identical=20builds=20supersede;=20wha?= =?UTF-8?q?t=20is=20held=20leaves=20room=20in=20the=20prefetch?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/inventory/inventory_test.go | 35 +++++++++----- internal/inventory/nodes.go | 17 ++++--- internal/link/enrol_claim_test.go | 71 +++++++++++++++++++++++++++- internal/link/enrolment.go | 14 +++++- internal/link/protocol.go | 21 ++++++++ internal/link/serve.go | 34 +++++++++++-- internal/link/store_window_test.go | 47 ++++++++++++++++++ 7 files changed, 214 insertions(+), 25 deletions(-) diff --git a/internal/inventory/inventory_test.go b/internal/inventory/inventory_test.go index 7ae1532..daabd80 100644 --- a/internal/inventory/inventory_test.go +++ b/internal/inventory/inventory_test.go @@ -397,14 +397,14 @@ func TestATokenIsClaimedByOnePresenterAndSpentOnlyByIt(t *testing.T) { if err != nil { t.Fatal(err) } - node, err := inv.Claim(ctx, issued.Secret, "key-a") + node, err := inv.Claim(ctx, issued.Secret, "key-a", false) if err != nil || node.Name != "laptop" { t.Fatalf("a fresh token was not claimed for its node: %v %q", err, node.Name) } - if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { + if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); err != nil { t.Fatalf("the presenter holding the claim could not claim again after an interruption: %v", err) } - if _, err := inv.Claim(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenInUse) { + if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); !errors.Is(err, ErrTokenInUse) { t.Fatalf("a second presenter was not held off while the claim is live: %v", err) } if err := inv.Spend(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenRefused) { @@ -414,18 +414,31 @@ func TestATokenIsClaimedByOnePresenterAndSpentOnlyByIt(t *testing.T) { t.Fatalf("the presenter holding the claim could not spend it: %v", err) } // Spent: nobody else may claim it, ever. - if _, err := inv.Claim(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenRefused) { + if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); !errors.Is(err, ErrTokenRefused) { t.Fatalf("a spent token was claimed by another presenter: %v", err) } - // But the presenter that spent it may, and spend it again: its spend reached the store and the - // answer did not reach the node, which asked again — refusing it would lock out a machine the - // mesh holds as enrolled. - if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { - t.Fatalf("the presenter whose answer was lost after its spend was refused: %v", err) + // Nor the presenter that spent it, without proof it holds the key: a public key is no secret. + if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("a spent token was claimed again with no proof of the key: %v", err) + } + // With it, and inside the lease, it may, and spend it again: its spend reached the store and + // the answer did not reach the node, which asked again — refusing it would lock out a machine + // the mesh holds as enrolled. + if _, err := inv.Claim(ctx, issued.Secret, "key-a", true); err != nil { + t.Fatalf("the proven presenter whose answer was lost after its spend was refused: %v", err) } if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil { t.Fatalf("spending again by the same presenter failed: %v", err) } + // And not once the lease is over: then a spent token is spent to everyone, proof or not. + if _, err := inv.store.Pool().Exec(ctx, + `update enrolment_token set claimed_until = now() - interval '1 second' where secret = $1`, + hashSecret(issued.Secret)); err != nil { + t.Fatal(err) + } + if _, err := inv.Claim(ctx, issued.Secret, "key-a", true); !errors.Is(err, ErrTokenRefused) { + t.Fatalf("a spent token was claimed again after its lease: %v", err) + } } // A claim lapses: a host that gave up and was started over, with keys of its own, is not held off @@ -440,7 +453,7 @@ func TestAClaimThatLapsedCanBeTakenByAnotherPresenter(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := inv.Claim(ctx, issued.Secret, "key-a"); err != nil { + if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); err != nil { t.Fatal(err) } if _, err := inv.store.Pool().Exec(ctx, @@ -448,7 +461,7 @@ func TestAClaimThatLapsedCanBeTakenByAnotherPresenter(t *testing.T) { hashSecret(issued.Secret)); err != nil { t.Fatal(err) } - if _, err := inv.Claim(ctx, issued.Secret, "key-b"); err != nil { + if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); err != nil { t.Fatalf("a lapsed claim held off a new presenter: %v", err) } if err := inv.Spend(ctx, issued.Secret, "key-a"); !errors.Is(err, ErrTokenRefused) { diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index b3792c6..2b822ce 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -216,17 +216,20 @@ var ErrTokenInUse = errors.New("the token is being used by another enrolment") // of a lease, and says which node it enrols. The same presenter may claim it again, as may anyone // once the lease has lapsed; nothing is spent until Spend (novox/hq 04-ISSUES/083). // -// A token this same presenter already spent is claimed again too: its spend reached the store and -// the answer did not reach the node, which asked again. Refusing it then would lock out a machine -// the mesh holds as enrolled — with the key it is still presenting. -func (i *Inventory) Claim(ctx context.Context, secret, by string) (Node, error) { +// A token this same presenter already spent may be claimed again when `again` says so — the +// caller has proof the presenter holds the key's private half — and only while its claim's lease +// is live: its spend reached the store and the answer did not reach the node, which asked again. +// Refusing it then would lock out a machine the mesh holds as enrolled. Without the proof a spent +// token stays spent to everyone, as ADR 0004 says. +func (i *Inventory) Claim(ctx context.Context, secret, by string, again bool) (Node, error) { var id string err := i.store.Pool().QueryRow(ctx, - `update enrolment_token set claimed_by = $2, claimed_until = now() + $3::interval + `update enrolment_token set claimed_by = $2, + claimed_until = case when redeemed is null then now() + $3::interval else claimed_until end where secret = $1 and expires > now() and ((redeemed is null and (claimed_by is null or claimed_by = $2 or claimed_until < now())) - or (redeemed is not null and claimed_by = $2)) - returning node`, hashSecret(secret), by, ClaimLease.String()).Scan(&id) + or (redeemed is not null and $4 and claimed_by = $2 and claimed_until > now())) + returning node`, hashSecret(secret), by, ClaimLease.String(), again).Scan(&id) if errors.Is(err, pgx.ErrNoRows) { // Unusable, or held by someone else — told apart, because the second passes. var held bool diff --git a/internal/link/enrol_claim_test.go b/internal/link/enrol_claim_test.go index bf373f3..0d4eafd 100644 --- a/internal/link/enrol_claim_test.go +++ b/internal/link/enrol_claim_test.go @@ -3,6 +3,8 @@ package link_test import ( "crypto/ed25519" "crypto/rand" + "crypto/sha256" + "encoding/hex" "errors" "testing" "time" @@ -23,7 +25,7 @@ func TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err := inv.Claim(ctx, issued.Secret, "another enrolment's key"); err != nil { + if _, err := inv.Claim(ctx, issued.Secret, "another enrolment's key", false); err != nil { t.Fatal(err) } public, _, err := ed25519.GenerateKey(rand.Reader) @@ -42,3 +44,70 @@ func TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain(t *testing.T) { t.Fatalf("a token that cannot be used was not refused outright: %v", err) } } + +// **A spent token cannot be replayed with a node's public key** (novox/hq issue 083, on review). A +// public key is no secret: finishing an enrolment whose token that key spent takes proof the +// presenter holds its private half, and a proof made with another key — or for another request — +// is refused outright. +func TestASpentTokenCannotBeReplayedWithAPublicKeyAlone(t *testing.T) { + inv := inventory.ForTest(t) + ctx := t.Context() + if _, err := inv.AddNode(ctx, "laptop"); err != nil { + t.Fatal(err) + } + issued, err := inv.IssueToken(ctx, "laptop", time.Hour) + if err != nil { + t.Fatal(err) + } + victim, victimPrivate, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + // The victim's enrolment spent the token. + by := claimantOf(victim) + if _, err := inv.Claim(ctx, issued.Secret, by, false); err != nil { + t.Fatal(err) + } + if err := inv.Spend(ctx, issued.Secret, by); err != nil { + t.Fatal(err) + } + + // Someone with the leaked token and the victim's public key, and no proof. + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's"}) + if err == nil || errors.Is(err, link.ErrTryAgain) { + t.Fatalf("a spent token was taken again with a public key alone: %v", err) + } + // With a proof made by another key. + _, forger, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + forged := ed25519.Sign(forger, link.EnrolProof(issued.Secret, victim, "", "the attacker's", "")) + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's", Proof: forged}) + if err == nil { + t.Fatal("a spent token was taken again with a proof made by another key") + } + // With the victim's own proof, but for another request — keys swapped in. + theirs := ed25519.Sign(victimPrivate, link.EnrolProof(issued.Secret, victim, "", "the victim's", "")) + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's", Proof: theirs}) + if err == nil { + t.Fatal("a spent token was taken again with a proof made for another request") + } + // And the victim's own, proven request is not honoured when the broker redelivered it: the + // first delivery may already have been answered. + _, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{ + Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the victim's", + Proof: theirs, Redelivered: true}) + if err == nil { + t.Fatal("a redelivered request finished an enrolment already spent") + } +} + +// claimantOf is the claimant the enrolment derives from a key, recomputed here. +func claimantOf(public ed25519.PublicKey) string { + sum := sha256.Sum256(public) + return hex.EncodeToString(sum[:]) +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index b17ddf2..52a185d 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -60,8 +60,20 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol // token, so the two cannot be one transaction; a failure between them used to leave a spent // token and a node with no key, which the host — making new keys on every attempt — could not // recover from. Every write below overwrites, so an attempt made again is safe. + // A proof that does not verify is refused outright: it was made with another key, or for + // another request. One that verifies lets this presenter finish an enrolment whose token it + // already spent — never a request the broker handed over a second time, which may already + // have been answered. + proven := false + if len(request.Proof) > 0 { + if !ed25519.Verify(public, EnrolProof(secret, public, request.OverlayKey, request.SealingKey, + request.ServingKey), request.Proof) { + return EnrolReply{}, errors.New("the enrolment's proof does not match the key it presents") + } + proven = true + } by := claimant(public) - node, err := e.Inventory.Claim(ctx, secret, by) + node, err := e.Inventory.Claim(ctx, secret, by, proven && !request.Redelivered) if err != nil { return EnrolReply{}, err } diff --git a/internal/link/protocol.go b/internal/link/protocol.go index 25e24fa..ba6dd7c 100644 --- a/internal/link/protocol.go +++ b/internal/link/protocol.go @@ -6,6 +6,8 @@ // traffic between them, each receiving half of what it expects. That has happened here before. package link +import "encoding/base64" + // Exchange is where nodes publish everything they have to say. const Exchange = "mesh" @@ -58,6 +60,18 @@ type EnrolRequest struct { // Profile is what this machine can be asked to do. The control plane cannot decide what a // node should run without it, so it arrives with enrolment rather than being asked for after. Profile map[string]any `json:"profile,omitempty"` + + // Proof is the node's identity key signing EnrolProof over this request: that the presenter + // holds the private half of PublicKey, not only knows the public one. Required to finish an + // enrolment whose token this key already spent — the case of an answer lost after the spend — + // because a public key is no secret, and without it anyone holding a leaked token and a + // node's public key could replay the spent token (novox/hq issue 083, on review). + Proof []byte `json:"proof,omitempty"` + + // Redelivered is set by the control plane, never sent: the broker handed this request over a + // second time. Such a request does not finish an enrolment already spent — the first time may + // have answered, and the node holds what it was told. + Redelivered bool `json:"-"` } // Signed is a declaration and the signature over it. @@ -137,3 +151,10 @@ type EnrolReply struct { // token can fail — unknown, spent, expired — so that guessing learns nothing. Refusal string `json:"refusal,omitempty"` } + +// EnrolProof is what a node signs with its identity key when it enrols: the token and every key it +// presents, so a proof cannot be moved to another request. +func EnrolProof(secret string, public []byte, overlay, sealing, serving string) []byte { + return []byte("novox-mesh-enrol\x00" + secret + "\x00" + base64.StdEncoding.EncodeToString(public) + + "\x00" + overlay + "\x00" + sealing + "\x00" + serving) +} diff --git a/internal/link/serve.go b/internal/link/serve.go index 9b8f3de..f3aa49e 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -89,6 +89,10 @@ const GiveUpAfter = 2 * time.Minute // Bounded, because what is held is also what the broker has not kept on its own disk as pending. const Prefetch = 64 +// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to +// answer — an enrolment above all — while others wait for the store. +const PrefetchHeadroom = 8 + // Records tells the server where to keep build results. // // Set after Connect rather than passed to it, because a control plane that only publishes — the @@ -274,6 +278,9 @@ func (s *Server) Serve(ctx context.Context) error { case <-ctx.Done(): return nil case <-ticker.C: + if ctx.Err() != nil { + return nil + } s.retryHeld(ctx) case delivery, ok := <-catchups: if !ok { @@ -365,7 +372,7 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { // once, and a report lost here is a node the mesh never hears from again — the store // restarting under the adoption that node just applied lost exactly that (issue 082). what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) - if s.tryLater(delivery, subject, what, err, s.handle) { + if s.tryLater(ctx, delivery, subject, what, err, s.handle) { return } // Said rather than swallowed. A report the mesh heard and failed to write down is a @@ -394,8 +401,13 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) { // Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a // host is waiting on is answered while a report waits for the store. One message is held at most // giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it. -func (s *Server) tryLater(delivery amqp.Delivery, subject, what string, err error, +func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, subject, what string, err error, retry func(context.Context, amqp.Delivery)) bool { + // Shutting down: nothing is settled. Unsettled, the broker hands the message to whatever + // consumes next — a cancelled context is not an answer about the message (issue 083, review). + if ctx.Err() != nil { + return true + } if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) { return false } @@ -404,6 +416,13 @@ func (s *Server) tryLater(delivery amqp.Delivery, subject, what string, err erro } h, ok := s.parked[subject] if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag { + // Held no further than the prefetch leaves room: past it, the broker would hand the loop + // nothing new — enrolments included — until something held was let go. + if !ok && len(s.parked) >= Prefetch-PrefetchHeadroom { + s.log.Printf("LOST %s: %d messages are already held for the store, and holding more "+ + "would stop the queue: %v", what, len(s.parked), err) + return false + } h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()} s.parked[subject] = h s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err) @@ -446,6 +465,9 @@ func digest(body []byte) string { // go past the bound. func (s *Server) retryHeld(ctx context.Context) { for _, h := range s.snapshot() { + if ctx.Err() != nil { + return + } h.retry(ctx, h.delivery) } } @@ -472,6 +494,7 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { if err := json.Unmarshal(delivery.Body, &request); err != nil { s.log.Printf("an enrolment request could not be read: %v", err) } else { + request.Redelivered = delivery.Redelivered accepted, err := s.enroller.Enrol(ctx, request) switch { case errors.Is(err, ErrTryAgain): @@ -543,9 +566,10 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) { // Each build result its own subject: none supersedes another, and recording one twice is // harmless — the build is kept by its id. subject := "build " + digest(delivery.Body) + s.supersede(subject, delivery) if err := s.recorder.Built(ctx, result); err != nil { // Held while the store cannot take it: a build result lost here is never announced (083). - if s.tryLater(delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) { + if s.tryLater(ctx, delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) { return } s.log.Printf("cannot keep a build result from %s: %v", result.On, err) @@ -586,7 +610,7 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) { } announcements, err := s.replayer.Announceable(ctx) if err != nil { - if s.tryLater(delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) { + if s.tryLater(ctx, delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) { holding = true return } @@ -640,7 +664,7 @@ func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { } if err := s.upgrader.Upgraded(ctx, u); err != nil { if errors.Is(err, ErrTryAgain) && - s.tryLater(delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { + s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { holding = true return } diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go index 7fe9ffd..931eeaf 100644 --- a/internal/link/store_window_test.go +++ b/internal/link/store_window_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "io" "log" "testing" @@ -118,3 +119,49 @@ func TestACatchUpWaitsOutARestartingStore(t *testing.T) { } } } + +// Shutting down is not an answer about a message: one handled with a cancelled context is left +// unsettled, for the broker to hand to whatever consumes next (issue 083, review). +func TestAMessageHandledDuringShutdownIsLeftForTheBroker(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + s := quietServer() + s.recorder = recordsWith{err: context.Canceled} + to := &settledAs{} + s.handleBuilt(ctx, a(t, to, KeyBuilt, BuildResult{On: "anchor", Repository: "/r", Commit: "abc"})) + if !to.held() { + t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to) + } +} + +// Two identical build results: the newer sets the older aside rather than leaving it unsettled +// for ever, holding a place in the prefetch. +func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) { + s := quietServer() + s.recorder = recordsWith{err: restarting} + built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"} + first, second := &settledAs{}, &settledAs{} + s.handleBuilt(context.Background(), a(t, first, KeyBuilt, built)) + s.handleBuilt(context.Background(), a(t, second, KeyBuilt, built)) + if !first.acked || !second.held() || len(s.parked) != 1 { + t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held", + *first, *second, len(s.parked)) + } +} + +// What is held stops short of the prefetch, so the loop always has room to answer an enrolment. +func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) { + s := quietServer() + s.recorder = recordsWith{err: restarting} + var last *settledAs + for i := 0; i < Prefetch; i++ { + last = &settledAs{} + s.handleBuilt(context.Background(), a(t, last, KeyBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)})) + } + if len(s.parked) != Prefetch-PrefetchHeadroom { + t.Fatalf("%d messages were held; the ceiling is %d", len(s.parked), Prefetch-PrefetchHeadroom) + } + if last.held() { + t.Fatalf("a message past the ceiling was held: %+v", *last) + } +} From 32b8af6a9b26cc98e77b47d121cb8e286378e9a1 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 14:33:33 +0200 Subject: [PATCH 4/5] The enrolment proof's message has a known answer, repeated in the host's test --- internal/link/enrol_proof_test.go | 13 +++++++++++++ 1 file changed, 13 insertions(+) create mode 100644 internal/link/enrol_proof_test.go diff --git a/internal/link/enrol_proof_test.go b/internal/link/enrol_proof_test.go new file mode 100644 index 0000000..727197f --- /dev/null +++ b/internal/link/enrol_proof_test.go @@ -0,0 +1,13 @@ +package link + +import "testing" + +// The bytes a node signs when it enrols, built here to check its signature. The host builds the +// same bytes in another repository; this known answer is repeated in its test, so the two cannot +// drift apart without one of them failing (novox/hq issue 083). +func TestWhatAnEnrollingNodeSignsIsFixed(t *testing.T) { + got := string(EnrolProof("s", []byte{1, 2, 3}, "o", "e", "v")) + if want := "novox-mesh-enrol\x00s\x00AQID\x00o\x00e\x00v"; got != want { + t.Fatalf("the enrolment proof's message changed: %q, want %q", got, want) + } +} From 4f3b4e6014bdba5a2d9c809c94f5273943ef8f8a Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 22 Sep 2026 14:39:15 +0200 Subject: [PATCH 5/5] Review of 083: an upgrade handled during shutdown is left for the broker; an enrolment cut short by shutdown is told to try again; a refused redelivery says what it probably is --- internal/link/enrolment.go | 3 ++- internal/link/serve.go | 12 ++++++++++++ internal/link/store_window_test.go | 14 ++++++++++++++ 3 files changed, 28 insertions(+), 1 deletion(-) diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index 52a185d..f39bc52 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -45,7 +45,8 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol // A store that could not be asked right now, or a token another presenter holds for the // moment, is "not now": the node asks again with the same request (novox/hq issue 083). defer func() { - if inventory.Unreachable(err) || errors.Is(err, inventory.ErrTokenInUse) { + if inventory.Unreachable(err) || errors.Is(err, inventory.ErrTokenInUse) || + errors.Is(err, context.Canceled) { err = fmt.Errorf("%w: %w", ErrTryAgain, err) } }() diff --git a/internal/link/serve.go b/internal/link/serve.go index f3aa49e..3bc2d76 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -503,6 +503,13 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { // this answer — decides when to ask, and the queue behind it moves (issue 083). reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"} s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err) + case err != nil && request.Redelivered: + // Said as what it most likely is: the broker handed this request over again after + // the control plane stopped mid-answer, and an enrolment already spent is not + // finished a second time. The node may need a new token. + s.log.Printf("refusing a redelivered enrolment for %q — it may have finished before "+ + "the control plane stopped, and if the node did not get its answer it needs a new "+ + "token: %v", request.Node, err) case err != nil: // Logged in full here, where an operator can see it; sent back as one refusal, so // that somebody guessing learns nothing from which reason came back. @@ -663,6 +670,11 @@ func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) { return } if err := s.upgrader.Upgraded(ctx, u); err != nil { + // Shutting down is not an answer about the announcement: left for the broker. + if ctx.Err() != nil { + holding = true + return + } if errors.Is(err, ErrTryAgain) && s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) { holding = true diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go index 931eeaf..f0fcd2d 100644 --- a/internal/link/store_window_test.go +++ b/internal/link/store_window_test.go @@ -165,3 +165,17 @@ func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) { t.Fatalf("a message past the ceiling was held: %+v", *last) } } + +// An upgrade handled during shutdown is left for the broker too — the upgrader's error is the +// cancelled context, which is no answer about the announcement. +func TestAnUpgradeHandledDuringShutdownIsLeftForTheBroker(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + s := quietServer() + s.upgrader = upgradesWith{err: context.Canceled} + to := &settledAs{} + s.upgraded(ctx, a(t, to, "upgraded", Upgraded{Module: "gitea", Commit: "abcdef0123"})) + if !to.held() { + t.Fatalf("an upgrade handled during shutdown was settled, and so lost: %+v", *to) + } +}