From cfbbad8209aac36be2728dc024309a2667df1152 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 02:36:58 +0200 Subject: [PATCH] Count a dead letter by when it was kept, so one kept late still counts (issue 440) The give-up time is not newer for every letter kept: a notice that could not be kept is offered again a minute later, so a later give-up can be kept first and the one after it read as not fresh. The stored time rises with the stream's sequence. A consumer whose letters vanish between the two reads is left out of the look instead of counted as a look, and the skip of a token-less max-deliveries advisory now has a test of its own. --- .../advisories_counted_test.go | 12 ++++ cmd/mesh-controller/signals.go | 8 ++- cmd/mesh-controller/watchdogs.go | 11 +++- internal/conditions/condition.go | 4 +- internal/link/deadletters.go | 17 +++--- internal/link/deadletters_test.go | 57 ++++++++++++------- 6 files changed, 74 insertions(+), 35 deletions(-) diff --git a/cmd/mesh-controller/advisories_counted_test.go b/cmd/mesh-controller/advisories_counted_test.go index 5fcd3164..0a999eb8 100644 --- a/cmd/mesh-controller/advisories_counted_test.go +++ b/cmd/mesh-controller/advisories_counted_test.go @@ -166,3 +166,15 @@ func TestAMessageThatCouldNotBeKeptIsStillSaidFromTheAdvisory(t *testing.T) { t.Fatalf("%+v", said) } } + +// A max-deliveries advisory without a token names a consumer's dead-letter key, which only DEAD_LETTERS +// says (novox/hq issues 330 and 440): with nothing held, the advisory alone raises nothing. +func TestAMaxDeliveriesAdvisoryAloneRaisesNothing(t *testing.T) { + at := time.Date(2026, 10, 10, 21, 4, 0, 0, time.UTC) + f := &signalFacts{now: at, host: "novox", advisories: []link.Advisory{{Kind: link.AdvisoryMaxDeliveries, + ID: "EVENTS.media_sonarr", Stream: "EVENTS", Consumer: "media_sonarr", + Said: "sonarr on media handed message 7 over 5 times and gave up on it", First: at, Last: at, Count: 1}}} + if said := watchAdvisories(f); len(said) != 0 { + t.Fatalf("said %+v", said) + } +} diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index 3f9cbb07..ab3e8a93 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -576,9 +576,11 @@ func watchDeadLetters(f *signalFacts) []conditions.Observation { messages, them = fmt.Sprintf("%d messages", n), "them" } who := consumerWho(stream, consumer) - // DEAD_LETTERS is a record of what was given up on: the condition counts the messages and says when - // the newest was given up, never how often the controller looked (novox/hq issues 402 and 440). - // Without that time it counts its looks, as a source of the present does. + // DEAD_LETTERS is a record of what was given up on: the condition says when the newest held message + // was kept, and its count is a running total of the messages given up on while it is open, never + // how often the controller looked (novox/hq issues 402 and 440). It is at least the number held now, + // and does not go down when one is delivered again or dropped. Without that time it counts its + // looks, as a source of the present does. happened, times := f.deadLettersNewest[key], 0 if !happened.IsZero() { times = n diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 9d8dd999..35c2c951 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -80,8 +80,8 @@ type signalFacts struct { // deadLetters are how many messages DEAD_LETTERS holds per consumer, by `.` // (novox/hq issue 330): each consumer's max-deliveries condition is open while it holds any. deadLetters map[string]int - // deadLettersNewest is when each of those consumers last gave up on a message DEAD_LETTERS still - // holds: the condition counts the messages given up on, not the looks (novox/hq issue 440). + // deadLettersNewest is when DEAD_LETTERS kept the newest message it holds for each of those + // consumers: the condition counts the messages given up on, not the looks (novox/hq issue 440). deadLettersNewest map[string]time.Time advisoriesErr error @@ -355,6 +355,13 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts { } if f.advisoriesErr == nil && len(f.deadLetters) > 0 { f.deadLettersNewest, f.advisoriesErr = link.NewestDeadLetters(w.js.Context(), f.deadLetters) + for key := range f.deadLetters { + if _, ok := f.deadLettersNewest[key]; !ok && f.advisoriesErr == nil { + // Delivered again or dropped between the two reads: left out of this look, rather than + // observed without a time and counted as a look (novox/hq issue 440). + delete(f.deadLetters, key) + } + } } f.handActs, f.handActsErr = w.gatherHandActs(ctx, now) f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last() diff --git a/internal/conditions/condition.go b/internal/conditions/condition.go index d7779745..8bf9a195 100644 --- a/internal/conditions/condition.go +++ b/internal/conditions/condition.go @@ -154,7 +154,9 @@ type Condition struct { Gaps []Gap `json:"gaps,omitempty"` // Observations is how many times it was observed since raised: a look that sees it, for a source // that looks at the present; for one that reads a record of what happened (Observation.Happened), - // how many times it happened, and LastObserved when it last did (novox/hq issue 402). + // how many times it happened, and LastObserved when it last did (novox/hq issue 402). That count is a + // running total since raised, never lowered: a consumer's dead letters count each message given up on + // while the condition is open, not the number held now (novox/hq issue 440). Observations int `json:"observations"` // Count is how many times it has been raised, a reopening within ReopenWithin counted. Count int `json:"count"` diff --git a/internal/link/deadletters.go b/internal/link/deadletters.go index cc7fcfc4..d6471b66 100644 --- a/internal/link/deadletters.go +++ b/internal/link/deadletters.go @@ -167,10 +167,13 @@ func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) { return held, nil } -// NewestDeadLetters is when each consumer of held last gave up on a message DEAD_LETTERS still holds, by -// `.`: the newest kept, by when the server said it gave up, or when it was kept if that -// was not said. The consumer's condition counts the messages given up on by it, never how often the -// controller looked at them (novox/hq issue 440). +// NewestDeadLetters is when DEAD_LETTERS kept the newest message it holds for each consumer of held, by +// `.`: the stored time of its last message, which rises with the stream's sequence. A +// consumer whose letters were all delivered again or dropped since held was read is left out. The +// consumer's condition counts the messages given up on by it, never how often the controller looked at +// them (novox/hq issue 440), so it needs a time that is newer for every message newly kept. The time the +// server said it gave up is not one: a notice that could not be kept is offered again a minute later +// (noticeRetry), so a later give-up can be kept first, and the one kept after it carries an older time. func NewestDeadLetters(js nats.JetStreamContext, held map[string]int) (map[string]time.Time, error) { newest := map[string]time.Time{} for key := range held { @@ -182,11 +185,7 @@ func NewestDeadLetters(js nats.JetStreamContext, held map[string]int) (map[strin if err != nil { return nil, fmt.Errorf("the newest dead letter of %s cannot be read: %w", key, err) } - at := deadLetterOf(raw.Sequence, raw.Subject, raw.Header, nil, false).GaveUp - if at.IsZero() { - at = raw.Time - } - newest[key] = at.UTC() + newest[key] = raw.Time.UTC() } return newest, nil } diff --git a/internal/link/deadletters_test.go b/internal/link/deadletters_test.go index a7008a88..f575c4ce 100644 --- a/internal/link/deadletters_test.go +++ b/internal/link/deadletters_test.go @@ -335,32 +335,49 @@ func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) { eventually(t, "a report heard while the notices cannot be taken", func() bool { return held.count() == 1 }) } -// When each consumer last gave up on a message DEAD_LETTERS still holds, for the condition that counts -// what was given up on rather than how often it was looked at (novox/hq issue 440): the newest kept, by -// when the server said it gave up. -func TestTheNewestHeldDeadLetterSaysWhenItWasGivenUp(t *testing.T) { +// When DEAD_LETTERS kept each consumer's newest message, for the condition that counts what was given up +// on rather than how often it was looked at (novox/hq issue 440). It rises with every message kept, even +// one the consumer gave up on before the last: a notice that could not be kept is offered again a minute +// later, so a later give-up can be kept first, and the one kept after it must still count. +func TestTheNewestHeldDeadLetterIsWhenItWasKept(t *testing.T) { js := aBus(t) - first := time.Date(2026, 10, 10, 21, 4, 0, 0, time.UTC) - for seq, at := range map[int]time.Time{1: first, 2: first.Add(time.Minute)} { - if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_sonarr",`+ - `"stream_seq":%d,"deliveries":5,"timestamp":%q}`, seq, at.Format(time.RFC3339Nano)))); err != nil { + gaveUp := time.Date(2026, 10, 10, 21, 4, 0, 0, time.UTC) + keep := func(consumer string, seq int, at time.Time) { + t.Helper() + if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":%q,`+ + `"stream_seq":%d,"deliveries":5,"timestamp":%q}`, consumer, seq, at.Format(time.RFC3339Nano)))); err != nil { t.Fatal(err) } } - if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_radarr",`+ - `"stream_seq":3,"deliveries":5,"timestamp":%q}`, first.Format(time.RFC3339Nano)))); err != nil { - t.Fatal(err) + look := func() map[string]time.Time { + t.Helper() + held, err := HeldDeadLetters(js.Context()) + if err != nil { + t.Fatal(err) + } + newest, err := NewestDeadLetters(js.Context(), held) + if err != nil { + t.Fatal(err) + } + return newest } - held, err := HeldDeadLetters(js.Context()) - if err != nil { - t.Fatal(err) + + // Sonarr gave up at 21:05 and that is kept first; radarr's letter is kept after it. + keep("media_sonarr", 2, gaveUp.Add(time.Minute)) + keep("media_radarr", 3, gaveUp) + before := look() + if len(before) != 2 || before["EVENTS.media_sonarr"].IsZero() || + !before["EVENTS.media_radarr"].After(before["EVENTS.media_sonarr"]) { + t.Fatalf("%v", before) } - newest, err := NewestDeadLetters(js.Context(), held) - if err != nil { - t.Fatal(err) + // Sonarr's give-up of 21:04, kept late, is newer than anything the condition has seen. + keep("media_sonarr", 1, gaveUp) + after := look() + if !after["EVENTS.media_sonarr"].After(before["EVENTS.media_sonarr"]) || + !after["EVENTS.media_sonarr"].After(before["EVENTS.media_radarr"]) { + t.Fatalf("a letter kept late does not read as newer: before %v, after %v", before, after) } - if len(newest) != 2 || !newest["EVENTS.media_sonarr"].Equal(first.Add(time.Minute)) || - !newest["EVENTS.media_radarr"].Equal(first) { - t.Fatalf("%v", newest) + if !after["EVENTS.media_radarr"].Equal(before["EVENTS.media_radarr"]) { + t.Fatalf("radarr's newest moved: before %v, after %v", before, after) } }