Keep what a consumer gives up on until a person delivers it again or drops it
Design 25 promised a dead-letter stream that did not exist: a message a consumer gave up on stayed only in its source, which drops it after a week, and its condition cleared when the advisories stopped (hq issue 330, ADR 0264).
This commit is contained in:
@@ -34,6 +34,9 @@ const (
|
||||
AdvisoryMaxDeliveries = "max-deliveries"
|
||||
AdvisoryRefused = "refused"
|
||||
AdvisoryConsumerLost = "consumer-lost"
|
||||
// AdvisoryNotKept is the token of a max-deliveries advisory whose message could not be kept in
|
||||
// DEAD_LETTERS (novox/hq issue 330): said from the log, since the stream does not hold it.
|
||||
AdvisoryNotKept = "not-kept"
|
||||
)
|
||||
|
||||
// Advisory is one thing the bus said, kept as its newest word and how often it was said.
|
||||
@@ -43,6 +46,8 @@ type Advisory struct {
|
||||
ID string
|
||||
// Stream and Consumer are the consumer it is about, when it is about one.
|
||||
Stream, Consumer string
|
||||
// Token is the condition key's last part, when it is not Kind.
|
||||
Token string
|
||||
// Said is the newest saying, in the mesh's words.
|
||||
Said string
|
||||
First, Last time.Time
|
||||
@@ -62,7 +67,7 @@ var Advisories = &AdvisoryLog{seen: map[string]*Advisory{}}
|
||||
func (l *AdvisoryLog) Heard(a Advisory, at time.Time) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
key := a.Kind + "/" + a.ID
|
||||
key := a.Kind + "/" + a.ID + "/" + a.Token
|
||||
if had, ok := l.seen[key]; ok {
|
||||
had.Last, had.Said, had.Count = at, a.Said, had.Count+1
|
||||
return
|
||||
@@ -107,8 +112,9 @@ func ReadAdvisory(subject string, body []byte) (Advisory, bool) {
|
||||
switch {
|
||||
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."):
|
||||
return Advisory{Kind: AdvisoryMaxDeliveries, ID: id, Stream: a.Stream, Consumer: a.Consumer,
|
||||
Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: it will not be delivered "+
|
||||
"again, and what it asked for was not done", who, a.StreamSeq, a.Deliveries)}, true
|
||||
Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: what it asked for was not "+
|
||||
"done, and the controller keeps it in %s until it is delivered again or dropped", who, a.StreamSeq,
|
||||
a.Deliveries, broker.DeadLettersStream)}, true
|
||||
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.DELETED."):
|
||||
if !MeshNamed(a.Stream, a.Consumer) {
|
||||
// A reader's own consumer, gone when it finished — every watch of a bucket and every
|
||||
@@ -175,7 +181,12 @@ func (s *Server) HearAdvisories(logf func(string, ...any)) (func(), error) {
|
||||
for _, subject := range broker.BusAdvisories {
|
||||
sub, err := conn.Subscribe(subject, func(m *nats.Msg) {
|
||||
if a, ok := ReadAdvisory(m.Subject, m.Data); ok {
|
||||
Advisories.Heard(a, time.Now())
|
||||
// A message given up on is said from DEAD_LETTERS, where it is kept, for as long as it is
|
||||
// kept (novox/hq issue 330) — not for an hour after the server said it; one that could not
|
||||
// be kept is recorded by whoever tried.
|
||||
if a.Kind != AdvisoryMaxDeliveries {
|
||||
Advisories.Heard(a, time.Now())
|
||||
}
|
||||
logf("the bus says: %s", a.Said)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -37,7 +37,8 @@ func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) {
|
||||
[]byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`))
|
||||
if !ok || a.Kind != AdvisoryMaxDeliveries || a.ID != "EVENTS.anchor_shop" ||
|
||||
a.Said != "how shop on anchor hears what it consumes handed message 7 over 5 times and gave up on it: "+
|
||||
"it will not be delivered again, and what it asked for was not done" {
|
||||
"what it asked for was not done, and the controller keeps it in DEAD_LETTERS until it is delivered "+
|
||||
"again or dropped" {
|
||||
t.Fatalf("%+v", a)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,302 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// What a consumer gave up on, kept until a person acts on it (novox/hq issue 330, design 25 §3).
|
||||
//
|
||||
// A durable consumer hands a message over as often as it may, gives up on it and says so on
|
||||
// `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES`. The server's notice is captured in its own stream
|
||||
// (broker.DeadLetterNoticesStream), so one said while no controller listens waits for the next. The
|
||||
// serving controller takes each notice, fetches the message it names by its sequence, while the source
|
||||
// stream still holds it, and keeps a copy in DEAD_LETTERS with what the notice said: the consumer, the
|
||||
// subject, how often it was handed over and when it was given up. Until then a message given up on was
|
||||
// kept only by its source, which drops an event after a week, and nothing said which.
|
||||
//
|
||||
// The copy's headers are the message's own, but for the server's (`Nats-…`), which describe the first
|
||||
// publish and would refuse or de-duplicate the copy; they are kept under DeadHeader names.
|
||||
|
||||
// The headers a kept message carries beside its own.
|
||||
const (
|
||||
DeadHeaderStream = "Mesh-Dead-Stream"
|
||||
DeadHeaderConsumer = "Mesh-Dead-Consumer"
|
||||
DeadHeaderSequence = "Mesh-Dead-Sequence"
|
||||
DeadHeaderSubject = "Mesh-Dead-Subject"
|
||||
DeadHeaderDeliveries = "Mesh-Dead-Deliveries"
|
||||
DeadHeaderGaveUp = "Mesh-Dead-Gave-Up"
|
||||
DeadHeaderPublished = "Mesh-Dead-Published"
|
||||
DeadHeaderLost = "Mesh-Dead-Lost"
|
||||
// DeadHeaderMsgID is the message's own de-duplication id, kept, since the copy carries its own.
|
||||
DeadHeaderMsgID = "Mesh-Dead-Msg-Id"
|
||||
// AgainHeader marks a message delivered again, with the kept copy's id it came from.
|
||||
AgainHeader = "Mesh-Delivered-Again"
|
||||
)
|
||||
|
||||
// DeadLetter is one message a consumer gave up on, as DEAD_LETTERS keeps it.
|
||||
type DeadLetter struct {
|
||||
// ID is its sequence in DEAD_LETTERS: what the verb names it by.
|
||||
ID uint64 `json:"id"`
|
||||
Stream string `json:"stream"`
|
||||
Consumer string `json:"consumer"`
|
||||
// Who is the consumer in the mesh's words: whose, and for what.
|
||||
Who string `json:"who"`
|
||||
Subject string `json:"subject"`
|
||||
Sequence uint64 `json:"sequence"`
|
||||
Deliveries uint64 `json:"deliveries"`
|
||||
GaveUp time.Time `json:"gave_up"`
|
||||
Published time.Time `json:"published,omitzero"`
|
||||
// Lost says why the message itself is not kept: its source no longer held it when the notice was
|
||||
// taken. The record of it is kept all the same, so it is said and dropped, never silently missing.
|
||||
Lost string `json:"lost,omitempty"`
|
||||
Size int `json:"size"`
|
||||
// Body and Headers are the message itself, given only for one dead letter asked by its id.
|
||||
Body string `json:"body,omitempty"`
|
||||
Headers map[string]string `json:"headers,omitempty"`
|
||||
}
|
||||
|
||||
// maxDeliveries is the part of the server's notice a kept message is made from.
|
||||
type maxDeliveries struct {
|
||||
Stream string `json:"stream"`
|
||||
Consumer string `json:"consumer"`
|
||||
StreamSeq uint64 `json:"stream_seq"`
|
||||
Deliveries uint64 `json:"deliveries"`
|
||||
Timestamp time.Time `json:"timestamp"`
|
||||
}
|
||||
|
||||
// KeepDeadLetter keeps the message one maximum-deliveries notice is about. An error leaves nothing
|
||||
// kept, and the notice is to be taken again: the source still holds the message, or a full
|
||||
// DEAD_LETTERS has room again. A source that no longer holds it is no error: the record is kept with
|
||||
// why, so it is said.
|
||||
func KeepDeadLetter(js nats.JetStreamContext, notice []byte) (DeadLetter, error) {
|
||||
var a maxDeliveries
|
||||
if err := json.Unmarshal(notice, &a); err != nil || a.Stream == "" || a.Consumer == "" || a.StreamSeq == 0 {
|
||||
return DeadLetter{}, fmt.Errorf("not a maximum-deliveries notice the mesh can read: %.200s", notice)
|
||||
}
|
||||
gaveUp := a.Timestamp
|
||||
if gaveUp.IsZero() {
|
||||
gaveUp = time.Now()
|
||||
}
|
||||
kept := &nats.Msg{Subject: broker.DeadLetterSubject(a.Stream, a.Consumer), Header: nats.Header{}}
|
||||
raw, err := js.GetMsg(a.Stream, a.StreamSeq)
|
||||
switch {
|
||||
case err == nil:
|
||||
for k, v := range raw.Header {
|
||||
if strings.HasPrefix(k, "Nats-") {
|
||||
continue
|
||||
}
|
||||
kept.Header[k] = v
|
||||
}
|
||||
if id := raw.Header.Get(nats.MsgIdHdr); id != "" {
|
||||
kept.Header.Set(DeadHeaderMsgID, id)
|
||||
}
|
||||
kept.Header.Set(DeadHeaderSubject, raw.Subject)
|
||||
kept.Header.Set(DeadHeaderPublished, raw.Time.UTC().Format(time.RFC3339Nano))
|
||||
kept.Data = raw.Data
|
||||
case errors.Is(err, nats.ErrMsgNotFound) || errors.Is(err, nats.ErrStreamNotFound):
|
||||
kept.Header.Set(DeadHeaderLost, fmt.Sprintf("%s no longer held message %d when the notice was taken: %v",
|
||||
a.Stream, a.StreamSeq, err))
|
||||
default:
|
||||
return DeadLetter{}, fmt.Errorf("message %d of %s could not be read: %w", a.StreamSeq, a.Stream, err)
|
||||
}
|
||||
kept.Header.Set(DeadHeaderStream, a.Stream)
|
||||
kept.Header.Set(DeadHeaderConsumer, a.Consumer)
|
||||
kept.Header.Set(DeadHeaderSequence, strconv.FormatUint(a.StreamSeq, 10))
|
||||
kept.Header.Set(DeadHeaderDeliveries, strconv.FormatUint(a.Deliveries, 10))
|
||||
kept.Header.Set(DeadHeaderGaveUp, gaveUp.UTC().Format(time.RFC3339Nano))
|
||||
// One copy per message given up, however often its notice is taken: a controller that copied and
|
||||
// stopped before acknowledging, or two controllers each taking it.
|
||||
kept.Header.Set(nats.MsgIdHdr, fmt.Sprintf("%s.%s.%d", a.Stream, a.Consumer, a.StreamSeq))
|
||||
ack, err := js.PublishMsg(kept)
|
||||
if err != nil {
|
||||
return DeadLetter{}, fmt.Errorf("message %d of %s could not be kept in %s: %w", a.StreamSeq, a.Stream,
|
||||
broker.DeadLettersStream, err)
|
||||
}
|
||||
return deadLetterOf(ack.Sequence, kept.Subject, kept.Header, kept.Data, false), nil
|
||||
}
|
||||
|
||||
// deadLetterOf reads a kept message.
|
||||
func deadLetterOf(id uint64, subject string, h nats.Header, data []byte, whole bool) DeadLetter {
|
||||
d := DeadLetter{ID: id, Stream: h.Get(DeadHeaderStream), Consumer: h.Get(DeadHeaderConsumer),
|
||||
Subject: h.Get(DeadHeaderSubject), Lost: h.Get(DeadHeaderLost), Size: len(data)}
|
||||
if d.Stream == "" || d.Consumer == "" {
|
||||
d.Stream, d.Consumer, _ = broker.DeadLetterOf(subject)
|
||||
}
|
||||
d.Who = ConsumerInWords(d.Stream, d.Consumer)
|
||||
d.Sequence, _ = strconv.ParseUint(h.Get(DeadHeaderSequence), 10, 64)
|
||||
d.Deliveries, _ = strconv.ParseUint(h.Get(DeadHeaderDeliveries), 10, 64)
|
||||
d.GaveUp, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderGaveUp))
|
||||
d.Published, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderPublished))
|
||||
if whole {
|
||||
d.Body = string(data)
|
||||
d.Headers = map[string]string{}
|
||||
for k := range h {
|
||||
if !strings.HasPrefix(k, "Mesh-Dead-") && !strings.HasPrefix(k, "Nats-") {
|
||||
d.Headers[k] = h.Get(k)
|
||||
}
|
||||
}
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
// HeldDeadLetters is how many messages DEAD_LETTERS holds for each consumer, by `<stream>.<consumer>`.
|
||||
func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) {
|
||||
info, err := js.StreamInfo(broker.DeadLettersStream, &nats.StreamInfoRequest{SubjectsFilter: ">"})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err)
|
||||
}
|
||||
held := map[string]int{}
|
||||
for subject, n := range info.State.Subjects {
|
||||
if stream, consumer, ok := broker.DeadLetterOf(subject); ok && n > 0 {
|
||||
held[stream+"."+consumer] += int(n)
|
||||
}
|
||||
}
|
||||
return held, nil
|
||||
}
|
||||
|
||||
// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them; for one consumer when
|
||||
// consumer names one (its name, or `<stream>.<consumer>`).
|
||||
func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLetter, int, error) {
|
||||
info, err := js.StreamInfo(broker.DeadLettersStream)
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err)
|
||||
}
|
||||
var out []DeadLetter
|
||||
total := 0
|
||||
for seq := info.State.LastSeq; seq >= info.State.FirstSeq && seq > 0; seq-- {
|
||||
raw, err := js.GetMsg(broker.DeadLettersStream, seq)
|
||||
if errors.Is(err, nats.ErrMsgNotFound) {
|
||||
continue // delivered again or dropped
|
||||
}
|
||||
if err != nil {
|
||||
return out, total, fmt.Errorf("dead letter %d cannot be read: %w", seq, err)
|
||||
}
|
||||
d := deadLetterOf(seq, raw.Subject, raw.Header, raw.Data, false)
|
||||
if consumer != "" && consumer != d.Consumer && consumer != d.Stream+"."+d.Consumer {
|
||||
continue
|
||||
}
|
||||
total++
|
||||
if most <= 0 || len(out) < most {
|
||||
out = append(out, d)
|
||||
}
|
||||
}
|
||||
return out, total, nil
|
||||
}
|
||||
|
||||
// ErrNoDeadLetter is a dead letter's id DEAD_LETTERS does not hold.
|
||||
var ErrNoDeadLetter = errors.New("no such dead letter")
|
||||
|
||||
// DeadLetterNamed is one kept message whole: what it said, and the headers it said it with.
|
||||
func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
||||
raw, err := js.GetMsg(broker.DeadLettersStream, id)
|
||||
if errors.Is(err, nats.ErrMsgNotFound) {
|
||||
return DeadLetter{}, fmt.Errorf("%w: %s holds no message %d — delivered again or dropped already, "+
|
||||
"or never kept; dead-letters lists what it holds", ErrNoDeadLetter, broker.DeadLettersStream, id)
|
||||
}
|
||||
if err != nil {
|
||||
return DeadLetter{}, fmt.Errorf("dead letter %d cannot be read: %w", id, err)
|
||||
}
|
||||
return deadLetterOf(id, raw.Subject, raw.Header, raw.Data, true), nil
|
||||
}
|
||||
|
||||
// AgainTo is where a kept message is delivered again so that only the consumer that gave it up gets
|
||||
// it: an event under that consumer's own again subject on EVENTS; an ask back onto its seat's queue,
|
||||
// whose one worker is that consumer. Any other stream's message is refused, with why.
|
||||
func AgainTo(d DeadLetter) (string, error) {
|
||||
switch {
|
||||
case d.Lost != "":
|
||||
return "", fmt.Errorf("dead letter %d holds no message to deliver: %s. Drop it", d.ID, d.Lost)
|
||||
case d.Subject == "":
|
||||
return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+
|
||||
"delivered again. Drop it", d.ID)
|
||||
case d.Stream == broker.EventsStream:
|
||||
return broker.AgainSubject(d.Consumer, d.Subject), nil
|
||||
case strings.HasPrefix(d.Stream, "SEAT_") && strings.HasSuffix(d.Consumer, "_worker"):
|
||||
return d.Subject, nil
|
||||
}
|
||||
return "", fmt.Errorf("dead letter %d is from %s, which nothing delivers again to one consumer alone: "+
|
||||
"publishing it again would reach every consumer of its subject. Drop it, and have its sender say it again",
|
||||
d.ID, d.Stream)
|
||||
}
|
||||
|
||||
// DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it
|
||||
// from DEAD_LETTERS. The message carries its own headers and AgainHeader; its de-duplication id is the
|
||||
// kept copy's, so asking twice delivers it once.
|
||||
func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, error) {
|
||||
d, err := DeadLetterNamed(js, id)
|
||||
if err != nil {
|
||||
return d, "", err
|
||||
}
|
||||
to, err := AgainTo(d)
|
||||
if err != nil {
|
||||
return d, "", err
|
||||
}
|
||||
again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)}
|
||||
for k, v := range d.Headers {
|
||||
again.Header.Set(k, v)
|
||||
}
|
||||
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
|
||||
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
|
||||
if _, err := js.PublishMsg(again); err != nil {
|
||||
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err)
|
||||
}
|
||||
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
|
||||
return d, to, fmt.Errorf("dead letter %d was delivered again on %s and could not be removed from %s: %w",
|
||||
id, to, broker.DeadLettersStream, err)
|
||||
}
|
||||
return d, to, nil
|
||||
}
|
||||
|
||||
// DropDeadLetter removes a kept message for good.
|
||||
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
||||
d, err := DeadLetterNamed(js, id)
|
||||
if err != nil {
|
||||
return d, err
|
||||
}
|
||||
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil {
|
||||
return d, fmt.Errorf("dead letter %d could not be dropped: %w", id, err)
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
|
||||
// noticeRetry is how long a notice whose message could not be kept waits before it is taken again.
|
||||
var noticeRetry = time.Minute
|
||||
|
||||
// keepGivenUp takes the server's maximum-deliveries notices off their stream and keeps the message
|
||||
// each is about, for as long as the subscription stands. A notice that could not be kept is offered
|
||||
// again after noticeRetry, and said: in the log, and as the consumer's max-deliveries condition.
|
||||
func keepGivenUp(js nats.JetStreamContext, logger *log.Logger) (*nats.Subscription, error) {
|
||||
return js.Subscribe("", func(m *nats.Msg) {
|
||||
d, err := KeepDeadLetter(js, m.Data)
|
||||
if err != nil {
|
||||
logger.Printf("a message a consumer gave up on could NOT be kept: %v", err)
|
||||
var a maxDeliveries
|
||||
if json.Unmarshal(m.Data, &a) == nil && a.Stream != "" && a.Consumer != "" {
|
||||
Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: a.Stream + "." + a.Consumer,
|
||||
Stream: a.Stream, Consumer: a.Consumer, Token: AdvisoryNotKept,
|
||||
Said: fmt.Sprintf("%s gave up on message %d, and it could not be kept: %v",
|
||||
ConsumerInWords(a.Stream, a.Consumer), a.StreamSeq, err)}, time.Now())
|
||||
}
|
||||
_ = m.NakWithDelay(noticeRetry)
|
||||
return
|
||||
}
|
||||
if d.Lost != "" {
|
||||
logger.Printf("%s gave up on message %d of %s, which its stream no longer held; kept as dead letter %d "+
|
||||
"without it", d.Who, d.Sequence, d.Stream, d.ID)
|
||||
} else {
|
||||
logger.Printf("%s gave up on message %d of %s (%s) after %d deliveries; kept as dead letter %d",
|
||||
d.Who, d.Sequence, d.Stream, d.Subject, d.Deliveries, d.ID)
|
||||
}
|
||||
_ = m.Ack()
|
||||
}, nats.Bind(broker.DeadLetterNoticesStream, broker.ControllerName), nats.ManualAck())
|
||||
}
|
||||
@@ -0,0 +1,264 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"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}` || whole.Headers["x-node"] != "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")
|
||||
}
|
||||
}
|
||||
|
||||
// A message from a stream no consumer alone can be given again on is refused, with why.
|
||||
func TestOnlyAnEventOrAnAskIsDeliveredAgain(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
d DeadLetter
|
||||
want string
|
||||
}{
|
||||
{DeadLetter{ID: 1, Stream: "EVENTS", Consumer: "a_b", Subject: "mesh.mod.x.event.y"}, "mesh.again.a_b.mod.x.event.y"},
|
||||
{DeadLetter{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker",
|
||||
Subject: "mesh.seat.telegram-sender.accept.send"}, "mesh.seat.telegram-sender.accept.send"},
|
||||
} {
|
||||
if to, err := AgainTo(c.d); err != nil || to != c.want {
|
||||
t.Errorf("%+v: %s %v", c.d, to, err)
|
||||
}
|
||||
}
|
||||
if _, err := AgainTo(DeadLetter{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"}); err == nil {
|
||||
t.Error("a bucket's message was given a subject to be delivered again on")
|
||||
}
|
||||
}
|
||||
|
||||
// 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 })
|
||||
}
|
||||
@@ -94,6 +94,14 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
|
||||
holding.Store(true)
|
||||
defer holding.Store(false)
|
||||
|
||||
// What a consumer gave up on, kept by the controller acting now (novox/hq issue 330): the server's
|
||||
// notices wait in their stream for it, so one said while no controller listened is not lost.
|
||||
kept, err := keepGivenUp(js, log.Default())
|
||||
if err != nil {
|
||||
return fmt.Errorf("taking the notices of messages consumers gave up on: %w", err)
|
||||
}
|
||||
defer func() { _ = kept.Unsubscribe() }()
|
||||
|
||||
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
|
||||
// they are their own guarantee: a lost one is the next one.
|
||||
beats := make(chan *nats.Msg, Prefetch)
|
||||
@@ -167,6 +175,11 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
|
||||
func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control),
|
||||
msg *nats.Msg, streamed bool) {
|
||||
|
||||
// An event the controller gave up on and a person delivered again (novox/hq issue 330) arrives under
|
||||
// the controller's own again subject, and is the event it was: acted on as first published.
|
||||
if original, ok := broker.OriginalOfAgain(msg.Subject); ok {
|
||||
msg.Subject = original
|
||||
}
|
||||
kind, known := kindOfSubject(msg.Subject)
|
||||
if !known {
|
||||
if streamed {
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user