A build ask its seat's worker gave up on could only be dropped, though the controller already holds the publish on that seat's accepts and the stream API to remove the original. Asks to other seats stay refused: delivering them needs a grant ADR 0264 withholds, in the asker's name ADR 0259 protects.
415 lines
18 KiB
Go
415 lines
18 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"`
|
|
// 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). The seat's
|
|
// stream keeps it until a holder pulls, as it keeps any ask while the seat has none. 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):
|
|
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 {
|
|
// An ask, back in its seat's work queue: the original, given up on, is never acknowledged and
|
|
// would stay beside its copy until the stream's age drops it (novox/hq issue 334). Removed first,
|
|
// so a copy that then cannot be published leaves the kept one to try again, and never two. One
|
|
// the queue no longer holds has aged out, which is no reason to refuse.
|
|
if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
|
|
return d, to, fmt.Errorf("the ask dead letter %d was kept from, message %d of %s, could not be "+
|
|
"removed, so it was not delivered again: %w; it is still kept", id, d.Sequence, d.Stream, err)
|
|
}
|
|
}
|
|
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 {
|
|
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
|
|
|
|
// 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())
|
|
}
|