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

This commit is contained in:
2026-09-22 14:23:03 +02:00
parent 1a41b88ed3
commit a3b7e830c8
7 changed files with 280 additions and 184 deletions
+137 -82
View File
@@ -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",