Files
mesh-controller/internal/link/deadletters.go
T
jschoubben ef551fdfb6
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
Remove an ask's original only when its sequence still holds it, and only with a worker (issue 334)
A seat's queue made again numbers from one, so an old dead letter's sequence
can name a live ask; deleting by number alone would drop it silently. An ask
with no worker would wait unseen while counted as delivered.
2026-10-11 03:24:52 +02:00

452 lines
21 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"slices"
"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"`
// Original says what became of the ask it was kept from, in its seat's work queue, when it was
// delivered again (novox/hq issue 334).
Original string `json:"original,omitempty"`
// Body and Headers are the message itself, given only for one dead letter asked by its id; a header
// with several values keeps them all.
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, v := range h {
if !strings.HasPrefix(k, "Mesh-Dead-") && !strings.HasPrefix(k, "Nats-") {
d.Headers[k] = append([]string(nil), v...)
}
}
}
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
}
// 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.
func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLetter, int, error) {
held, err := HeldDeadLetters(js)
if err != nil {
return nil, 0, err
}
total := 0
for key, n := range held {
if consumer == "" || key == consumer || strings.HasSuffix(key, "."+consumer) {
total += n
}
}
if total == 0 {
return nil, 0, nil
}
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
for seq := info.State.LastSeq; seq >= info.State.FirstSeq && seq > 0; seq-- {
if most > 0 && len(out) >= most || len(out) >= total {
break
}
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
}
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 — and only when the consumer exists and
// filters that subject, so a message is never let go as delivered while nobody receives it.
//
// An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's
// work queue has one worker per subject, so only the worker that gave it up takes it, and DeliverAgain
// removes the original from the queue first, so there are never two (novox/hq issue 334); only while the
// worker that gave it up is on the bus. Any other stream's
// message is refused, with why: an ask to a seat the controller does not ask would need a publish it is
// not granted, in its asker's name (ADR 0264's consequences, ADR 0259 §3), and a stream that is neither
// would reach every consumer of its subject.
func AgainTo(js nats.JetStreamContext, 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 broker.TheControllersAsk(d.Stream, d.Subject):
// The seat's worker takes it, and only while it is on the bus: without one the queue would keep
// the ask for a holder that may never come, and it would be let go as delivered meanwhile.
if _, err := js.ConsumerInfo(d.Stream, d.Consumer); errors.Is(err, nats.ErrConsumerNotFound) {
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so the ask "+
"would only wait in %s for a holder. Nothing was done, and it is still kept: deliver it again once "+
"the seat has a holder, or drop it", d.ID, d.Who, d.Stream)
} else if err != nil {
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
}
return d.Subject, nil
case d.Stream != broker.EventsStream:
return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+
"delivered again: publishing it again would reach every consumer of its subject, or need a grant the "+
"controller does not hold to speak for its asker. Drop it, and have its sender say it again",
d.ID, d.Stream)
}
info, err := js.ConsumerInfo(d.Stream, d.Consumer)
if errors.Is(err, nats.ErrConsumerNotFound) {
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so nothing would "+
"receive it. Drop it, or deliver it again once the module is assigned there again", d.ID, d.Who)
}
if err != nil {
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
}
want := broker.AgainFilter(d.Consumer)
filters := append([]string{info.Config.FilterSubject}, info.Config.FilterSubjects...)
if !slices.Contains(filters, want) {
return "", fmt.Errorf("%s does not yet take events delivered again (it does not filter %s): the controller "+
"sets that at the next send to its machine. Nothing was done, and dead letter %d is still kept",
d.Who, want, d.ID)
}
return broker.AgainSubject(d.Consumer, d.Subject), nil
}
// DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it
// from DEAD_LETTERS; an ask's original is removed from its work queue before. 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(js, d)
if err != nil {
return d, "", err
}
if d.Stream != broker.EventsStream {
d.Original = removeOriginal(js, d)
}
again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)}
for k, v := range d.Headers {
again.Header[k] = append([]string(nil), 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 {
if d.Original != "" {
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept, so "+
"delivering it again tries once more. Its original: %s", id, to, err, d.Original)
}
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
}
// removeOriginal takes the ask a dead letter was kept from out of its seat's work queue, and says what
// became of it. Given up on, the original is never acknowledged and would stay beside its copy until the
// stream's age drops it (novox/hq issue 334). It is removed before the copy is published, so a copy that
// cannot be published leaves the kept one to try again, and never two.
//
// **Only when the queue still holds that same message**: its subject and the time it was stored are the
// dead letter's. A seat's stream deleted and made again — the build handover deletes one (builds.go) —
// numbers from one again, and an old dead letter's sequence may then name another, live ask; deleting by
// the number alone would drop that one silently. Anything else is said, never an error: the original is
// gone or is not this one, and the copy is the only one there will be.
func removeOriginal(js nats.JetStreamContext, d DeadLetter) string {
held, err := js.GetMsg(d.Stream, d.Sequence)
switch {
case errors.Is(err, nats.ErrMsgNotFound):
return fmt.Sprintf("%s no longer held message %d, so there was nothing to remove", d.Stream, d.Sequence)
case err != nil:
return fmt.Sprintf("message %d of %s could not be read, so it was left as it is: %v", d.Sequence, d.Stream, err)
case d.Published.IsZero() || held.Subject != d.Subject || !held.Time.Equal(d.Published):
return fmt.Sprintf("message %d of %s is another message now (%s, stored %s), so it was left as it is; "+
"the one given up on is gone", d.Sequence, d.Stream, held.Subject, held.Time.UTC().Format(time.RFC3339))
}
if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
return fmt.Sprintf("message %d of %s could not be removed, so it stays beside its copy until the "+
"stream's age drops it; nothing delivers it again: %v", d.Sequence, d.Stream, err)
}
return fmt.Sprintf("message %d of %s, the one given up on, was removed from the queue", d.Sequence, d.Stream)
}
// 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
// keepingGivenUp keeps taking the notices until ctx ends; the returned function stops it. A subscription
// that cannot be made is said — in the log, and as the max-deliveries condition of the notices themselves
// — and tried again every noticeRetry, so it never stops the controller serving.
func keepingGivenUp(ctx context.Context, bus *broker.JetStream, logger *log.Logger) func() {
ctx, cancel := context.WithCancel(ctx)
done := make(chan struct{})
go func() {
defer close(done)
for {
sub, err := func() (*nats.Subscription, error) {
if err := bus.EnsureConsumer(broker.NoticesConsumer()); err != nil {
return nil, err
}
return keepGivenUp(bus.Context(), logger)
}()
if err == nil {
<-ctx.Done()
_ = sub.Unsubscribe()
return
}
logger.Printf("the notices of messages consumers gave up on cannot be taken, so none is kept until "+
"they can: %v", err)
Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: broker.DeadLetterNoticesStream + "." +
broker.ControllerName, Stream: broker.DeadLetterNoticesStream, Consumer: broker.ControllerName,
Token: AdvisoryNotKept, Said: "the controller cannot take the notices of messages consumers gave " +
"up on, so none is kept: " + err.Error()}, time.Now())
select {
case <-ctx.Done():
return
case <-time.After(noticeRetry):
}
}
}()
return func() {
cancel()
<-done
}
}
// 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())
}