Count a dead letter by when it was kept, so one kept late still counts (issue 440)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered

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.
This commit is contained in:
2026-10-11 02:36:58 +02:00
parent f83fcdc15f
commit cfbbad8209
6 changed files with 74 additions and 35 deletions
@@ -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)
}
}
+5 -3
View File
@@ -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
+9 -2
View File
@@ -80,8 +80,8 @@ type signalFacts struct {
// deadLetters are how many messages DEAD_LETTERS holds per consumer, by `<stream>.<consumer>`
// (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()
+3 -1
View File
@@ -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"`
+8 -9
View File
@@ -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
// `<stream>.<consumer>`: 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
// `<stream>.<consumer>`: 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
}
+37 -20
View File
@@ -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)
}
}