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