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)

This commit is contained in:
2026-09-22 14:06:28 +02:00
parent b9cdb96a90
commit 1a41b88ed3
8 changed files with 433 additions and 59 deletions
+104 -43
View File
@@ -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)
}