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) + } +}