package link import ( "testing" "time" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "github.com/novox/mesh-controller/internal/broker" ) // novox/hq issue 330, replayed with only what the link had before its fix, so it can be laid over the // older commit. On 2026-10-06 a media server's event consumer gave up on several messages after five // deliveries each. The controller raised a condition for each run and cleared it within minutes; the // messages were kept only by EVENTS, which drops an event after a week, and nothing said which they // were. Design 25 promised a dead-letter stream that does not exist. A consumer that never acknowledges // an event: once it gives up, the event is kept with its consumer and how often it was handed over. func TestReplay330(t *testing.T) { js := aBus(t) consumer, ok := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr", Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"}) if !ok { t.Fatal("a module that consumes got no consumer") } if err := js.EnsureConsumer(consumer); err != nil { t.Fatal(err) } _, stop := servingOn(t, js, &counted{}) defer stop() const subject = "mesh.mod.gitea.event.pull.merged" if _, err := js.Context().Publish(subject, []byte(`{"repository":"novox/media"}`)); err != nil { t.Fatal(err) } // The module's handler fails every time, as the media server's did. api, err := jetstream.New(js.Conn()) if err != nil { t.Fatal(err) } reading, err := api.Consumer(t.Context(), broker.EventsStream, consumer.Name) if err != nil { t.Fatal(err) } for handed := 0; handed < consumer.MaxDeliver; handed++ { batch, err := reading.Fetch(1, jetstream.FetchMaxWait(5*time.Second)) if err != nil { t.Fatal(err) } n := 0 for msg := range batch.Messages() { n++ _ = msg.Nak() } if n != 1 { t.Fatalf("handed over %d times, then nothing: %v", handed, batch.Error()) } } // And it goes on reading, as a module's runtime does: the server gives the event up when it would // hand it over a sixth time. if batch, err := reading.Fetch(1, jetstream.FetchMaxWait(time.Second)); err == nil { for msg := range batch.Messages() { t.Fatalf("handed over a sixth time: %s", msg.Subject()) } } kept := "mesh.events.dead." + broker.EventsStream + "." + consumer.Name var held *nats.RawStreamMsg deadline := time.Now().Add(10 * time.Second) for time.Now().Before(deadline) && held == nil { if stream, err := js.Context().StreamNameBySubject(kept); err == nil { held, _ = js.Context().GetLastMsg(stream, kept) } if held == nil { time.Sleep(50 * time.Millisecond) } } if held == nil { t.Fatalf("%s gave up on the event and nothing on the bus keeps it under %s", consumer.Name, kept) } if string(held.Data) != `{"repository":"novox/media"}` { t.Errorf("kept %q, not the event", held.Data) } for header, want := range map[string]string{"Mesh-Dead-Consumer": consumer.Name, "Mesh-Dead-Stream": "EVENTS", "Mesh-Dead-Subject": subject, "Mesh-Dead-Deliveries": "5"} { if got := held.Header.Get(header); got != want { t.Errorf("the kept event's %s is %q, not %q", header, got, want) } } if held.Header.Get("Mesh-Dead-Gave-Up") == "" { t.Error("the kept event does not say when it was given up") } }