Review of 083: finishing an enrolment whose token was spent takes proof of the key's private half, a live lease and a first delivery — a public key alone cannot replay a spent token; shutdown leaves held messages for the broker; identical builds supersede; what is held leaves room in the prefetch

This commit is contained in:
2026-09-22 14:33:13 +02:00
parent a3b7e830c8
commit 4567fa666c
7 changed files with 214 additions and 25 deletions
+29 -5
View File
@@ -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
}