Files
mesh-controller/internal/link/deadletters_test.go
T
jschoubben f83fcdc15f
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 superseded: a newer head of the same pull request
Give a consumer's max-deliveries key one watcher, counting what was given up (issue 440)
The dead-letter row and a max-deliveries advisory without a token both said
bus.<stream>.<consumer>.max-deliveries: each look added an observation and a
line of evidence through the row, and the two overwrote each other's words.
The row owns the key since issue 330, so the advisory watcher leaves it, and
the row now says when the newest held message was given up, so a letter held
for a day no longer reads as observed every 30 seconds. Times without Happened
is refused, since the raise ignored it and an update counted it.
2026-10-11 02:30:09 +02:00

367 lines
13 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// A module's consumer as the controller makes it, on a bus whose streams the controller asserted.
func aModuleConsumer(t *testing.T, js *broker.JetStream, node, module string, consumes ...string) jetstream.Consumer {
t.Helper()
c, ok := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: node, Module: module,
Consumes: consumes, PasswordHash: "x"})
if !ok {
t.Fatal("a module that consumes got no consumer")
}
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
reading, err := api.Consumer(t.Context(), broker.EventsStream, c.Name)
if err != nil {
t.Fatal(err)
}
return reading
}
// next is the one message a consumer hands over within a second, or nil.
func next(t *testing.T, c jetstream.Consumer) jetstream.Msg {
t.Helper()
batch, err := c.Fetch(1, jetstream.FetchMaxWait(time.Second))
if err != nil {
t.Fatal(err)
}
for msg := range batch.Messages() {
return msg
}
return nil
}
// keeping takes the notices as the serving controller does, until the test ends.
func keeping(t *testing.T, js *broker.JetStream) {
t.Helper()
if err := broker.AssertMeshConsumers(js); err != nil {
t.Fatal(err)
}
sub, err := keepGivenUp(js.Context(), log.New(io.Discard, "", 0))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
}
// givenUp hands one event to a consumer as often as it may, failing each time, and waits for it kept.
func givenUp(t *testing.T, js *broker.JetStream, c jetstream.Consumer) DeadLetter {
t.Helper()
for {
msg := next(t, c)
if msg == nil {
break
}
_ = msg.Nak()
}
var kept []DeadLetter
eventually(t, "the event kept", func() bool {
kept, _, _ = DeadLetters(js.Context(), c.CachedInfo().Name, 0)
return len(kept) == 1
})
return kept[0]
}
// **Delivered again to the consumer that gave it up, and nobody else**: another module consuming the
// same event handled it the first time and must not handle it twice.
func TestADeadLetterIsDeliveredAgainToItsConsumerAlone(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
handling := aModuleConsumer(t, js, "media", "radarr", "gitea.pull.merged")
const subject = "mesh.mod.gitea.event.pull.merged"
msg := &nats.Msg{Subject: subject, Data: []byte(`{"n":1}`), Header: nats.Header{}}
msg.Header.Set("x-node", "forge")
msg.Header.Set(nats.MsgIdHdr, "event-1")
if _, err := js.Context().PublishMsg(msg); err != nil {
t.Fatal(err)
}
if m := next(t, handling); m == nil {
t.Fatal("the other module was not handed the event")
} else {
_ = m.Ack()
}
d := givenUp(t, js, failing)
if d.Consumer != "media_sonarr" || d.Stream != "EVENTS" || d.Subject != subject || d.Deliveries != 5 ||
d.GaveUp.IsZero() || d.Published.IsZero() || d.Lost != "" {
t.Fatalf("kept as %+v", d)
}
held, err := HeldDeadLetters(js.Context())
if err != nil || held["EVENTS.media_sonarr"] != 1 {
t.Fatalf("held %v (%v)", held, err)
}
whole, err := DeadLetterNamed(js.Context(), d.ID)
if err != nil || whole.Body != `{"n":1}` || len(whole.Headers["x-node"]) != 1 || whole.Headers["x-node"][0] != "forge" {
t.Fatalf("one asked whole is %+v (%v)", whole, err)
}
_, to, err := DeliverAgain(js.Context(), d.ID)
if err != nil {
t.Fatal(err)
}
if to != "mesh.again.media_sonarr.mod.gitea.event.pull.merged" {
t.Fatalf("delivered again on %s", to)
}
again := next(t, failing)
if again == nil {
t.Fatal("the consumer that gave it up was not handed it again")
}
if string(again.Data()) != `{"n":1}` || again.Headers().Get("x-node") != "forge" ||
again.Headers().Get(AgainHeader) == "" {
t.Fatalf("handed again as %s %v", again.Data(), again.Headers())
}
if original, ok := broker.OriginalOfAgain(again.Subject()); !ok || original != subject {
t.Fatalf("delivered again on %s, which does not say the event's own subject", again.Subject())
}
_ = again.Ack()
if m := next(t, handling); m != nil {
t.Fatalf("the module that handled it was handed it twice: %s", m.Subject())
}
if _, err := DeadLetterNamed(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
t.Fatalf("still kept after it was delivered again: %v", err)
}
held, _ = HeldDeadLetters(js.Context())
if held["EVENTS.media_sonarr"] != 0 {
t.Fatalf("still held: %v", held)
}
// Asked twice, delivered once: the second finds nothing kept.
if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
t.Fatalf("a second delivery answered %v", err)
}
}
func TestADeadLetterDroppedIsGone(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil {
t.Fatal(err)
}
d := givenUp(t, js, failing)
if _, err := DropDeadLetter(js.Context(), d.ID); err != nil {
t.Fatal(err)
}
if left, total, err := DeadLetters(js.Context(), "", 0); err != nil || total != 0 || len(left) != 0 {
t.Fatalf("after the drop: %v %d %v", left, total, err)
}
if m := next(t, failing); m != nil {
t.Fatalf("a dropped event was handed over: %s", m.Subject())
}
}
// A notice whose message its stream no longer holds is kept as a record that says so — said, and
// dropped by a person, never silently missing — and cannot be delivered again.
func TestANoticeWhoseMessageIsGoneIsKeptAndSaysSo(t *testing.T) {
js := aBus(t)
d, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+
`"deliveries":5,"timestamp":"2026-10-06T15:51:00Z"}`))
if err != nil {
t.Fatal(err)
}
if !strings.Contains(d.Lost, "no longer held message 42") || d.Consumer != "media_sonarr" {
t.Fatalf("kept as %+v", d)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "Drop it") {
t.Fatalf("a record without its message was delivered again: %v", err)
}
// The same notice taken twice is kept once.
if _, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+
`"deliveries":5}`)); err != nil {
t.Fatal(err)
}
if _, total, _ := DeadLetters(js.Context(), "media_sonarr", 0); total != 1 {
t.Fatalf("one notice taken twice is kept %d times", total)
}
}
// What only the giving-up consumer filters is its own: a subject delivered again reads back as the
// event's, and a module's runtime reads the same key from it (the tokens around `.event.`).
func TestASubjectDeliveredAgainSaysTheEventsOwn(t *testing.T) {
for _, original := range []string{"mesh.mod.gitea.event.pull.merged", "mesh.seat.node-build-agent.event.built"} {
again := broker.AgainSubject("ace_sonarr", original)
if back, ok := broker.OriginalOfAgain(again); !ok || back != original {
t.Errorf("%s reads back as %s", again, back)
}
key := func(subject string) string {
before, event, _ := strings.Cut(subject, ".event.")
parts := strings.Split(before, ".")
return parts[len(parts)-1] + "." + event
}
if key(again) != key(original) {
t.Errorf("%s reads as key %s, the event as %s", again, key(again), key(original))
}
}
if _, ok := broker.OriginalOfAgain("mesh.mod.gitea.event.pull.merged"); ok {
t.Error("an event's own subject reads as delivered again")
}
}
// Only an event is delivered again, and only to a consumer that is on the bus and filters its again
// subject: a dead letter is never let go as delivered while nobody receives it.
func TestADeadLetterIsDeliveredAgainOnlyWhereItIsReceived(t *testing.T) {
js := aBus(t)
keeping(t, js)
failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged")
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil {
t.Fatal(err)
}
d := givenUp(t, js, failing)
// A consumer as it was before this fix: its filters do not take what is delivered again.
c, _ := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr",
Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"})
c.Filters = c.Filters[:len(c.Filters)-1]
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "does not yet take") {
t.Fatalf("delivered to a consumer that does not filter its again subject: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// The consumer gone: nothing would receive it.
if err := js.Context().DeleteConsumer(broker.EventsStream, c.Name); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "no longer on the bus") {
t.Fatalf("delivered to a consumer that is gone: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// Not an event: refused, whatever stream it is from.
for _, other := range []DeadLetter{
{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", Subject: "mesh.seat.telegram-sender.accept.send"},
{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"},
} {
if _, err := AgainTo(js.Context(), other); err == nil {
t.Errorf("%s was given a subject to be delivered again on", other.Stream)
}
}
}
// A total says how many are held, however few the list carries.
func TestTheListSaysTheTotalAndStopsAtItsLimit(t *testing.T) {
js := aBus(t)
for seq := 1; seq <= 5; seq++ {
if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_sonarr",`+
`"stream_seq":%d,"deliveries":5}`, seq))); err != nil {
t.Fatal(err)
}
}
list, total, err := DeadLetters(js.Context(), "media_sonarr", 2)
if err != nil || total != 5 || len(list) != 2 || list[0].Sequence != 5 {
t.Fatalf("%d of %d (%v): %+v", len(list), total, err, list)
}
if _, total, _ := DeadLetters(js.Context(), "another_one", 2); total != 0 {
t.Fatalf("another consumer's total is %d", total)
}
}
// An event the controller itself gave up on, delivered again by a person, is acted on as the event it
// was: it arrives under the controller's again subject, which its consumer filters.
func TestTheControllerActsOnAnEventDeliveredAgain(t *testing.T) {
js := aBus(t)
told := &toldAbout{}
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()}
if err := s.Follows(told); err != nil {
t.Fatal(err)
}
if err := s.Answers(replaysWith{}); err != nil {
t.Fatal(err)
}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
eventually(t, "the controller's event consumer being made", func() bool {
_, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName)
return err == nil
})
moved, _ := json.Marshal(Upgraded{Module: "gitea", Commit: "abcdef0123"})
if _, err := js.Context().Publish(broker.AgainSubject(broker.ControllerName, broker.ControllerFollows[0]), moved); err != nil {
t.Fatal(err)
}
eventually(t, "the upgrade delivered again reaching the controller", func() bool { return told.count() == 1 })
}
// The notices cannot be taken (their stream is gone): said as a condition, and the controller serves on.
func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) {
js := aBus(t)
if err := js.Context().DeleteStream(broker.DeadLetterNoticesStream); err != nil {
t.Fatal(err)
}
stop := keepingGivenUp(context.Background(), js, quiet())
defer stop()
eventually(t, "the failure said", func() bool {
for _, a := range Advisories.Since(time.Now().Add(-time.Minute)) {
if a.Token == AdvisoryNotKept && a.Stream == broker.DeadLetterNoticesStream {
return true
}
}
return false
})
held := &counted{}
_, stopServing := servingOn(t, js, held)
defer stopServing()
if _, err := js.Context().Publish("mesh.control.anchor.report", []byte(`{"node":"anchor"}`)); err != nil {
t.Fatal(err)
}
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) {
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 {
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)
}
held, err := HeldDeadLetters(js.Context())
if err != nil {
t.Fatal(err)
}
newest, err := NewestDeadLetters(js.Context(), held)
if err != nil {
t.Fatal(err)
}
if len(newest) != 2 || !newest["EVENTS.media_sonarr"].Equal(first.Add(time.Minute)) ||
!newest["EVENTS.media_radarr"].Equal(first) {
t.Fatalf("%v", newest)
}
}