Merge pull request 'A consumer's max-deliveries condition has one watcher and counts what was given up (issue 440)' (#214) from fix/440-one-key-one-watcher into main
This commit was merged in pull request #214.
This commit is contained in:
@@ -93,3 +93,88 @@ func TestAConditionOfThePresentCountsItsLooks(t *testing.T) {
|
||||
t.Fatalf("observed %d time(s), %d evidence lines, last %s", c.Observations, len(c.Evidence), c.LastObserved)
|
||||
}
|
||||
}
|
||||
|
||||
// **One key, one watcher** (novox/hq issue 440). A consumer's max-deliveries condition is said from
|
||||
// DEAD_LETTERS while it holds a message the consumer gave up on (issue 330). A max-deliveries advisory
|
||||
// without a token names the same key; read beside it, each look added an observation and a line of
|
||||
// evidence through the dead-letter half, and the two summaries overwrote each other on every look. The
|
||||
// dead-letter row alone says that key, and counts the messages given up on, never the looks.
|
||||
func TestAHeldDeadLetterIsCountedByWhatWasGivenUpNotByTheLooks(t *testing.T) {
|
||||
gaveUp := time.Date(2026, 10, 10, 21, 4, 0, 0, time.UTC)
|
||||
now := gaveUp.Add(20 * time.Second)
|
||||
store := conditions.NewInMemory()
|
||||
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store,
|
||||
Teller: &conditions.Told{}, Now: func() time.Time { return now }})
|
||||
t.Cleanup(func() { k.Close(t.Context()) })
|
||||
|
||||
const key = "bus.EVENTS.media_sonarr.max-deliveries"
|
||||
held := map[string]int{"EVENTS.media_sonarr": 1}
|
||||
newest := map[string]time.Time{"EVENTS.media_sonarr": gaveUp}
|
||||
advisory := 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: gaveUp, Last: gaveUp, Count: 1}
|
||||
look := func() conditions.Condition {
|
||||
t.Helper()
|
||||
f := &signalFacts{now: now, host: "novox", advisories: []link.Advisory{advisory},
|
||||
deadLetters: held, deadLettersNewest: newest}
|
||||
if err := k.Reconcile(t.Context(), "S9", watchAdvisories(f)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
c, found, err := conditions.ReadOne(t.Context(), store, key)
|
||||
if err != nil || !found {
|
||||
t.Fatalf("%s: found %v, %v", key, found, err)
|
||||
}
|
||||
if !strings.Contains(c.Summary, "dead-letters verb") {
|
||||
t.Fatalf("the summary is not the dead-letter row's: %q", c.Summary)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
var c conditions.Condition
|
||||
for range 5 {
|
||||
c = look()
|
||||
now = now.Add(30 * time.Second)
|
||||
}
|
||||
if c.Observations != 1 || len(c.Evidence) != 1 || !c.LastObserved.Equal(gaveUp) {
|
||||
t.Fatalf("one message given up on, looked at five times, reads as observed %d time(s), last %s, "+
|
||||
"with %d evidence lines: %+v", c.Observations, c.LastObserved, len(c.Evidence), c.Evidence)
|
||||
}
|
||||
|
||||
// The consumer gives up on a second message between two looks: two given up on, one more line.
|
||||
again := now.Add(-10 * time.Second)
|
||||
held["EVENTS.media_sonarr"], newest["EVENTS.media_sonarr"] = 2, again
|
||||
c = look()
|
||||
if c.Observations != 2 || len(c.Evidence) != 2 || !c.LastObserved.Equal(again) {
|
||||
t.Fatalf("two given up on read as observed %d time(s), last %s, evidence %+v", c.Observations,
|
||||
c.LastObserved, c.Evidence)
|
||||
}
|
||||
now = now.Add(30 * time.Second)
|
||||
if c = look(); c.Observations != 2 || len(c.Evidence) != 2 {
|
||||
t.Fatalf("a look with nothing new counted: observed %d time(s), evidence %+v", c.Observations, c.Evidence)
|
||||
}
|
||||
}
|
||||
|
||||
// The advisory still says the max-deliveries it alone knows: a message given up on that could not be
|
||||
// kept, under a key of its own (issue 330), which issue 440's change leaves as it was.
|
||||
func TestAMessageThatCouldNotBeKeptIsStillSaidFromTheAdvisory(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", Token: link.AdvisoryNotKept,
|
||||
Said: "sonarr on media gave up on message 7, and it could not be kept", First: at, Last: at, Count: 1}}}
|
||||
said := watchAdvisories(f)
|
||||
if len(said) != 1 || said[0].Key() != "bus.EVENTS.media_sonarr.not-kept" {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -523,6 +523,14 @@ func watchAdvisories(f *signalFacts) []conditions.Observation {
|
||||
if a.Kind == link.AdvisoryConsumerLost && !f.lostConsumers[a.Stream+"."+a.Consumer] {
|
||||
continue // it exists again, or the mesh no longer expects it: a removal, not a loss
|
||||
}
|
||||
if a.Kind == link.AdvisoryMaxDeliveries && a.Token == "" {
|
||||
// A consumer's max-deliveries key is the dead-letter row's (novox/hq issue 330): said while
|
||||
// DEAD_LETTERS holds what it gave up on, cleared when that is delivered again or dropped. The
|
||||
// listener records no such advisory today; one read here as well added a look to the count and
|
||||
// overwrote the row's words on every look (novox/hq issue 440). Only one that could not be kept
|
||||
// (AdvisoryNotKept), which DEAD_LETTERS cannot say, is said from here.
|
||||
continue
|
||||
}
|
||||
severity := conditions.Warning
|
||||
times := ""
|
||||
if a.Count > 1 {
|
||||
@@ -568,7 +576,17 @@ 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 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
|
||||
}
|
||||
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: key, Kind: link.AdvisoryMaxDeliveries,
|
||||
Happened: happened, Times: times,
|
||||
Machine: consumerMachine(stream, consumer), Severity: conditions.Warning,
|
||||
Summary: fmt.Sprintf("%s gave up on %s; %s kept in %s until delivered again or dropped, with why, "+
|
||||
"through the controller's dead-letters verb", link.ConsumerInWords(stream, consumer), messages,
|
||||
|
||||
@@ -100,14 +100,15 @@ var suppressions = map[string]suppression{
|
||||
inside: func(f *signalFacts) { f.standings = []conditions.Condition{standingSaid(f.now.Add(-29 * time.Minute))} },
|
||||
past: func(f *signalFacts) { f.standings = []conditions.Condition{standingSaid(f.now.Add(-31 * time.Minute))} },
|
||||
},
|
||||
// A message given up on and not kept: the one max-deliveries advisory S9 says itself (issue 440).
|
||||
"S9": {
|
||||
inside: func(f *signalFacts) {
|
||||
f.advisories = []link.Advisory{{Kind: link.AdvisoryMaxDeliveries, ID: "EVENTS.anchor_shop", Said: "gave up",
|
||||
First: f.now.Add(-2 * time.Hour), Last: f.now.Add(-61 * time.Minute), Count: 1}}
|
||||
f.advisories = []link.Advisory{{Kind: link.AdvisoryMaxDeliveries, ID: "EVENTS.anchor_shop",
|
||||
Token: link.AdvisoryNotKept, Said: "gave up, not kept", First: f.now.Add(-2 * time.Hour), Last: f.now.Add(-61 * time.Minute), Count: 1}}
|
||||
},
|
||||
past: func(f *signalFacts) {
|
||||
f.advisories = []link.Advisory{{Kind: link.AdvisoryMaxDeliveries, ID: "EVENTS.anchor_shop", Said: "gave up",
|
||||
First: f.now.Add(-2 * time.Hour), Last: f.now.Add(-59 * time.Minute), Count: 1}}
|
||||
f.advisories = []link.Advisory{{Kind: link.AdvisoryMaxDeliveries, ID: "EVENTS.anchor_shop",
|
||||
Token: link.AdvisoryNotKept, Said: "gave up, not kept", First: f.now.Add(-2 * time.Hour), Last: f.now.Add(-59 * time.Minute), Count: 1}}
|
||||
},
|
||||
},
|
||||
"S10": {
|
||||
|
||||
@@ -79,8 +79,11 @@ type signalFacts struct {
|
||||
lostConsumers map[string]bool
|
||||
// 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
|
||||
advisoriesErr error
|
||||
deadLetters map[string]int
|
||||
// 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
|
||||
|
||||
selfCheck selfCheckFacts
|
||||
|
||||
@@ -350,6 +353,16 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
|
||||
if f.advisoriesErr == nil && w.js != nil {
|
||||
f.deadLetters, f.advisoriesErr = link.HeldDeadLetters(w.js.Context())
|
||||
}
|
||||
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()
|
||||
return f
|
||||
|
||||
@@ -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"`
|
||||
@@ -276,6 +278,10 @@ func (o Observation) check() error {
|
||||
return fmt.Errorf("the condition %s says nothing", o.Key())
|
||||
case strings.TrimSpace(o.Source) == "":
|
||||
return fmt.Errorf("the condition %s does not say what raised it", o.Key())
|
||||
case o.Times != 0 && o.Happened.IsZero():
|
||||
// Times is what a record counts beside when it last happened: alone, the raise would ignore it
|
||||
// and an update count it, so the count would mean two things (novox/hq issue 440).
|
||||
return fmt.Errorf("the condition %s says how often it happened (%d) but not when", o.Key(), o.Times)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -494,3 +494,19 @@ func TestASecondReopeningKeepsBothGaps(t *testing.T) {
|
||||
t.Fatalf("open at the wrong moments: %+v", g)
|
||||
}
|
||||
}
|
||||
|
||||
// **Times is how often a record says it happened, never alone** (novox/hq issue 440). The raise read
|
||||
// Times only beside Happened while an update read it always, so a source setting Times alone would have
|
||||
// jumped the count on its second look and not its first. The keeper refuses it, saying why.
|
||||
func TestTimesWithoutWhenItHappenedIsRefused(t *testing.T) {
|
||||
k, store, _, _ := keeper(t)
|
||||
o := Observation{Scope: ScopeMachine, ID: "ace", Kind: "silent", Machine: "ace", Severity: Warning,
|
||||
Summary: "ace has not been heard from", Source: "S1", Times: 5}
|
||||
_, err := k.Observe(t.Context(), o)
|
||||
if err == nil || !strings.Contains(err.Error(), "when") {
|
||||
t.Fatalf("Times without Happened was taken: %v", err)
|
||||
}
|
||||
if _, found, _ := ReadOne(t.Context(), store, o.Key()); found {
|
||||
t.Fatal("a refused observation raised its condition")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -167,6 +167,29 @@ func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) {
|
||||
return held, nil
|
||||
}
|
||||
|
||||
// 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 {
|
||||
stream, consumer, _ := strings.Cut(key, ".")
|
||||
raw, err := js.GetLastMsg(broker.DeadLettersStream, broker.DeadLetterSubject(stream, consumer))
|
||||
if errors.Is(err, nats.ErrMsgNotFound) {
|
||||
continue // delivered again or dropped since it was counted
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the newest dead letter of %s cannot be read: %w", key, err)
|
||||
}
|
||||
newest[key] = raw.Time.UTC()
|
||||
}
|
||||
return newest, nil
|
||||
}
|
||||
|
||||
// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them (all when most is not
|
||||
// positive); for one consumer when consumer names one (its name, or `<stream>.<consumer>`). The total is
|
||||
// the stream's own count per consumer, so it is right however few are read.
|
||||
|
||||
@@ -334,3 +334,50 @@ func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) {
|
||||
}
|
||||
eventually(t, "a report heard while the notices cannot be taken", func() bool { return held.count() == 1 })
|
||||
}
|
||||
|
||||
// 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)
|
||||
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)
|
||||
}
|
||||
}
|
||||
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
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
// 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 !after["EVENTS.media_radarr"].Equal(before["EVENTS.media_radarr"]) {
|
||||
t.Fatalf("radarr's newest moved: before %v, after %v", before, after)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user