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 `.`. 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 (all when most is not // positive); for one consumer when consumer names one (its name, or `.`). 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. Any other // stream's message is refused, with why: an ask given up on by a seat's worker is not delivered again yet // (novox/hq issue 330's follow-up), since publishing it again leaves the original stuck in the queue. 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 d.Stream != broker.EventsStream: return "", fmt.Errorf("dead letter %d is from %s, and only an event is delivered again: publishing it "+ "again would reach every consumer of its subject, or leave the original in its queue. 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. 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 } 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()) }