Merge pull request 'Keep what a consumer gives up on until a person delivers it again or drops it (hq issue 330, ADR 0264)' (#159) from fix/330-a-message-given-up-on-is-kept into main

This commit was merged in pull request #159.
This commit is contained in:
2026-10-08 19:07:33 +00:00
26 changed files with 1454 additions and 28 deletions
+166
View File
@@ -0,0 +1,166 @@
package main
import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"sync/atomic"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
// What a consumer gave up on, answered by the serving controller (novox/hq issue 330).
//
// **In this process, on its own connection**: DEAD_LETTERS is read and changed on the bus, and the
// serving controller is already on it. A verb run as a fresh process would open a connection of its own
// for each call (novox/hq issue 327).
// deadLettersOn is the serving controller's bus, set when it starts watching the mesh; nil in a
// process that serves nothing.
var deadLettersOn atomic.Pointer[busHandles]
// busHandles are the serving controller's connection and JetStream handle.
type busHandles struct {
conn *nats.Conn
js nats.JetStreamContext
}
// defaultDeadLetters is how many the list says when not asked for more.
const defaultDeadLetters = 50
// causeDeadLetter is the cause a delivery or drop gives when the caller gives none.
const causeDeadLetter = "dead-letter"
// deadLettersAnswer is what `dead-letters` answers: the list, one whole, or what came of delivering one
// again or dropping it.
func deadLettersAnswer(ctx context.Context, a *verbArguments) (any, error) {
on := deadLettersOn.Load()
if on == nil {
return nil, errors.New("this controller is not serving, so it does not read DEAD_LETTERS: ask again, " +
"and the serving controller answers")
}
// The shape first, and every argument it reads; one given beside it is refused before anything is
// done, as every verb refuses what it would pass over (novox/hq issue 244).
var deliver, drop, why, cause, idText, consumer, limit string
switch {
case a.given["deliver"] != "" || a.given["drop"] != "":
deliver, drop, why, cause = a.str("deliver"), a.str("drop"), a.str("why"), a.str("cause")
case a.given["id"] != "":
idText = a.str("id")
default:
consumer, limit = a.str("consumer"), a.str("limit")
}
if unused := a.unused(); len(unused) > 0 {
return nil, fmt.Errorf("dead-letters did not use %s together with %s, and an argument a verb would pass "+
"over is refused: nothing was done", quoteAll(unused), quoteAll(a.usedGiven()))
}
switch {
case deliver != "" && drop != "":
return nil, errors.New("dead-letters delivers one again or drops one, not both. Nothing was done")
case deliver != "" || drop != "":
act, text := "deliver", deliver
if drop != "" {
act, text = "drop", drop
}
if strings.TrimSpace(why) == "" {
return nil, fmt.Errorf("dead-letters %s is a hand act, and says why: why is required and recorded in "+
"the hand-act log (novox/hq to-be 45 §7). Nothing was done", act)
}
id, err := deadLetterID(text)
if err != nil {
return nil, err
}
return actOnDeadLetter(ctx, on, act, id, why, cause)
case idText != "":
id, err := deadLetterID(idText)
if err != nil {
return nil, err
}
return link.DeadLetterNamed(on.js, id)
}
most := defaultDeadLetters
if limit != "" {
n, err := strconv.Atoi(limit)
if err != nil || n <= 0 {
return nil, fmt.Errorf("limit is a number of dead letters, not %q", limit)
}
most = n
}
held, total, err := link.DeadLetters(on.js, consumer, most)
if err != nil {
return nil, err
}
answer := map[string]any{"dead_letters": held, "held": total,
"note": "newest first; with id, one whole; deliver or drop one with why"}
if total == 0 {
answer["note"] = "no consumer gave up on a message that is still kept"
}
return answer, nil
}
// deadLetterID is a dead letter's id as a caller wrote it.
func deadLetterID(text string) (uint64, error) {
id, err := strconv.ParseUint(strings.TrimSpace(text), 10, 64)
if err != nil || id == 0 {
return 0, fmt.Errorf("a dead letter's id is its number in %s, as dead-letters lists it, not %q",
broker.DeadLettersStream, text)
}
return id, nil
}
// actOnDeadLetter delivers one again or drops it, recorded in the hand-act log before it is done. A log
// that cannot be written is said, and the act still happens: the log is never the reason a person's act
// is refused (handacts.go).
func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64, why, cause string) (any, error) {
d, err := link.DeadLetterNamed(on.js, id)
if err != nil {
return nil, err
}
if act == "deliver" {
// Refused before it is recorded: an act that cannot be done is not an act.
if _, err := link.AgainTo(on.js, d); err != nil {
return nil, err
}
}
if cause == "" {
cause = causeDeadLetter
}
by := link.CallerIn(ctx)
if by == "" {
by = "a seat call whose caller the bus did not name"
}
answer := map[string]any{"dead_letter": d.ID, "consumer": d.Who, "subject": d.Subject}
recorded, logErr := link.RecordHandAct(ctx, on.conn, link.HandAct{Verb: "dead-letters " + act,
Args: []string{strconv.FormatUint(id, 10), d.Stream + "." + d.Consumer}, Why: why, Cause: cause,
Condition: conditions.Key(conditions.ScopeBus, d.Stream+"."+d.Consumer, link.AdvisoryMaxDeliveries),
By: by + ", through the " + catalogue.ControllerSeatName + " seat"})
if logErr != nil {
answer["unrecorded"] = "the hand-act log could not be written, and the act was done all the same: " + logErr.Error()
} else {
answer["recorded"] = recorded.ID
}
switch act {
case "deliver":
_, to, err := link.DeliverAgain(on.js, id)
if err != nil {
return nil, err
}
answer["delivered_on"] = to
answer["done"] = fmt.Sprintf("dead letter %d was delivered again to %s, and nobody else; it is no longer kept",
id, consumerWho(d.Stream, d.Consumer))
case "drop":
if _, err := link.DropDeadLetter(on.js, id); err != nil {
return nil, err
}
answer["done"] = fmt.Sprintf("dead letter %d, which %s gave up on, was dropped for good", id,
consumerWho(d.Stream, d.Consumer))
}
return answer, nil
}
+132
View File
@@ -0,0 +1,132 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
"github.com/novox/mesh-controller/internal/testbus"
)
// The serving controller's bus, with the mesh's streams, for the verb to read and act on.
func servingDeadLetters(t *testing.T) *broker.JetStream {
t.Helper()
js, err := broker.Dial(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
if err := broker.AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
if err := js.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
before := deadLettersOn.Load()
deadLettersOn.Store(&busHandles{conn: js.Conn(), js: js.Context()})
t.Cleanup(func() { deadLettersOn.Store(before) })
return js
}
func askDeadLetters(t *testing.T, args map[string]any) (any, error) {
t.Helper()
a, err := readArguments("dead-letters", args)
if err != nil {
return nil, err
}
return deadLettersAnswer(context.Background(), a)
}
func TestDeadLettersListsDropsAndRecordsWhy(t *testing.T) {
js := servingDeadLetters(t)
if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{"n":1}`)); err != nil {
t.Fatal(err)
}
d, err := link.KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":1,"deliveries":5}`))
if err != nil {
t.Fatal(err)
}
answer, err := askDeadLetters(t, map[string]any{})
if err != nil {
t.Fatal(err)
}
listed := answer.(map[string]any)
if listed["held"] != 1 || len(listed["dead_letters"].([]link.DeadLetter)) != 1 {
t.Fatalf("listed %v", listed)
}
// Refused before anything is done: an act without why, an argument the shape passes over, both acts.
for _, args := range []map[string]any{
{"drop": "1"},
{"drop": "1", "why": "x", "limit": "3"},
{"drop": "1", "deliver": "1", "why": "x"},
{"why": "x"},
{"id": "nought"},
} {
if _, err := askDeadLetters(t, args); err == nil {
t.Errorf("%v was done", args)
}
}
if _, total, _ := link.DeadLetters(js.Context(), "", 0); total != 1 {
t.Fatalf("a refused call changed what is kept: %d left", total)
}
done, err := askDeadLetters(t, map[string]any{"drop": "1", "why": "the media server took the download in by hand"})
if err != nil {
t.Fatal(err)
}
if said := done.(map[string]any); said["recorded"] == nil || !strings.Contains(said["done"].(string), "dropped") {
t.Fatalf("answered %v", said)
}
acts, err := link.HandActs(context.Background(), js.Conn(), time.Now().Add(-time.Minute))
if err != nil || len(acts) != 1 || acts[0].Verb != "dead-letters drop" || acts[0].Cause != causeDeadLetter ||
!strings.Contains(acts[0].Condition, "media_sonarr") {
t.Fatalf("recorded %+v (%v)", acts, err)
}
if !personsDecision(acts[0]) {
t.Error("dropping a dead letter is counted as a repair, so S15 would want a healer for it")
}
if _, err := askDeadLetters(t, map[string]any{"id": "1"}); err == nil {
t.Errorf("dead letter %d is still answered after it was dropped", d.ID)
}
}
// Open while DEAD_LETTERS holds a message for the consumer, in words the operator reads in one pass:
// what is held, and where to act.
func TestAConsumersDeadLettersAreSaidUntilActedOn(t *testing.T) {
f := &signalFacts{now: time.Now(), deadLetters: map[string]int{"EVENTS.media_sonarr": 4, "EVENTS.controller": 1}}
found := watchDeadLetters(f)
if len(found) != 2 {
t.Fatalf("said %d conditions", len(found))
}
for _, o := range found {
if o.Kind != "max-deliveries" || o.Severity != conditions.Warning || o.Needs == "" {
t.Errorf("%+v", o)
}
if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Needs: o.Needs,
Explanation: o.Explanation, Resolved: o.Resolved}, "media"); !ok {
t.Errorf("%q is not plain: %s", o.Headline, why)
}
if !strings.Contains(o.Needs, "mesh MCP server") {
t.Errorf("does not say where to act: %q", o.Needs)
}
}
sonarr := found[1]
if sonarr.ID != "EVENTS.media_sonarr" || sonarr.Machine != "media" ||
sonarr.Headline != "Sonarr on media could not handle 4 messages" ||
!strings.Contains(sonarr.Summary, "DEAD_LETTERS") {
t.Errorf("%+v", sonarr)
}
if found[0].Headline != "The controller could not handle a message" {
t.Errorf("%q", found[0].Headline)
}
// None held, none said: it clears when they are delivered again or dropped.
if left := watchDeadLetters(&signalFacts{now: time.Now()}); len(left) != 0 {
t.Fatalf("%v", left)
}
}
+5
View File
@@ -82,6 +82,11 @@ var handActVerbs = []handActVerb{
{Verb: "retire approve", Decision: "nothing is retired past its bound without a person (ADR 0230)"}, {Verb: "retire approve", Decision: "nothing is retired past its bound without a person (ADR 0230)"},
{Verb: "retire reject", Decision: "keeping a consumer active is a person's word (ADR 0230)"}, {Verb: "retire reject", Decision: "keeping a consumer active is a person's word (ADR 0230)"},
{Verb: "cleanup delete", Decision: "nothing retired is deleted without a person (ADR 0230)"}, {Verb: "cleanup delete", Decision: "nothing retired is deleted without a person (ADR 0230)"},
// What becomes of a message a consumer gave up on (novox/hq issue 330): kept until a person says.
{Verb: "dead-letters deliver", Decision: "a message a consumer gave up on is delivered again only on a " +
"person's word (issue 330)"},
{Verb: "dead-letters drop", Decision: "a message a consumer gave up on is let go only on a person's word " +
"(issue 330)"},
// The sweep run on a person's word rather than after a build: the same decision the records make, at // The sweep run on a person's word rather than after a build: the same decision the records make, at
// a moment the person chose (ADR 0251) — never a repair. // a moment the person chose (ADR 0251) — never a repair.
{Verb: "collect", Decision: "letting the store go of what the records keep for no reason, now rather " + {Verb: "collect", Decision: "letting the store go of what the records keep for no reason, now rather " +
+5 -3
View File
@@ -544,9 +544,11 @@ var plainWordings = map[string]func(conditions.Observation) words{
Resolved: "The listener keeps up again"} Resolved: "The listener keeps up again"}
}), }),
"max-deliveries": worded(func(o conditions.Observation) words { "max-deliveries": worded(func(o conditions.Observation) words {
return words{Headline: "A message could not be handled", return words{Headline: "A listener gave up on messages",
Explanation: "The bus gave up on a message after trying to hand it over too many times.", Needs: "deliver them again or drop them, from the mesh MCP server.",
Resolved: "Resolved: messages are handled again"} Explanation: "A listener on the bus could not handle messages after several tries, so what they asked " +
"for was not done. The mesh keeps them until you deliver them again or drop them.",
Resolved: "Resolved: the messages given up on were delivered again or dropped"}
}), }),
"refused": worded(func(o conditions.Observation) words { "refused": worded(func(o conditions.Observation) words {
return words{Headline: "The bus refuses some messages", return words{Headline: "The bus refuses some messages",
+10
View File
@@ -875,6 +875,16 @@ func streamDiffers(want broker.Stream, have nats.StreamConfig) string {
if perSubject != 0 && have.MaxMsgsPerSubject != perSubject { if perSubject != 0 && have.MaxMsgsPerSubject != perSubject {
differs = append(differs, fmt.Sprintf("keeps %d per subject, defined %d", have.MaxMsgsPerSubject, perSubject)) differs = append(differs, fmt.Sprintf("keeps %d per subject, defined %d", have.MaxMsgsPerSubject, perSubject))
} }
if want.MaxBytes > 0 && have.MaxBytes != want.MaxBytes {
differs = append(differs, fmt.Sprintf("holds up to %d bytes, defined %d", have.MaxBytes, want.MaxBytes))
}
if want.DuplicatesSeconds > 0 && have.Duplicates != time.Duration(want.DuplicatesSeconds)*time.Second {
differs = append(differs, fmt.Sprintf("keeps one of a message id for %s, defined %s", have.Duplicates,
time.Duration(want.DuplicatesSeconds)*time.Second))
}
if want.DiscardNew && have.Discard != nats.DiscardNew {
differs = append(differs, "drops what it holds when full, defined to refuse what comes next")
}
return strings.Join(differs, "; ") return strings.Join(differs, "; ")
} }
+9 -1
View File
@@ -250,6 +250,8 @@ func (a *verbArguments) commandLine() ([]string, error) {
return argv, nil return argv, nil
case "tools": case "tools":
return nil, errors.New("tools is answered from the records, not by a command") return nil, errors.New("tools is answered from the records, not by a command")
case "dead-letters":
return nil, errors.New("dead-letters is answered by the serving controller, on its own connection, not by a command")
case "status": case "status":
return []string{"status", "--json"}, nil return []string{"status", "--json"}, nil
case "nodes": case "nodes":
@@ -982,6 +984,9 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) {
if verb == "calls" { if verb == "calls" {
return callsAnswer(link.Calls, a.given["call"]) return callsAnswer(link.Calls, a.given["call"])
} }
if verb == "dead-letters" {
return deadLettersAnswer(ctx, a)
}
if verb == "doctor" { if verb == "doctor" {
// From the serving controller, which runs the self-check and hears the signals // From the serving controller, which runs the self-check and hears the signals
// (novox/hq to-be 45 §4): the last verdict at once, or a run now. // (novox/hq to-be 45 §4): the last verdict at once, or a run now.
@@ -1064,7 +1069,10 @@ func actsOnAPlan(args map[string]any) bool {
// inProcess are the verbs answered by this process rather than by a command it runs: `tools` from // inProcess are the verbs answered by this process rather than by a command it runs: `tools` from
// the records, `calls` from what this process served. // the records, `calls` from what this process served.
var inProcess = map[string]bool{"tools": true, "calls": true, "doctor": true} var inProcess = map[string]bool{"tools": true, "calls": true, "doctor": true,
// What a consumer gave up on, read and changed on the serving controller's own connection (novox/hq
// issue 330).
"dead-letters": true}
// answersFirst is a command line whose caller is answered before it runs: a push, by its verb or // answersFirst is a command line whose caller is answered before it runs: a push, by its verb or
// through `command`. A push sends the machine holding the bus first when its user list changed, the // through `command`. A push sends the machine holding the bus first when its user list changed, the
+87 -4
View File
@@ -154,8 +154,9 @@ var signalsTable = []signalRow{
}}, }},
{Row: "S9", Signal: "bus advisories: maximum deliveries, consumer deleted; the controller's own slow " + {Row: "S9", Signal: "bus advisories: maximum deliveries, consumer deleted; the controller's own slow " +
"consumer and refused subjects", Emitter: "bus server's advisory subjects; the controller's connection", "consumer and refused subjects", Emitter: "bus server's advisory subjects; the controller's connection",
Trigger: "any", Bound: "any occurrence; clears after an hour without another, and a deleted consumer " + Trigger: "any", Bound: "any occurrence; clears after an hour without another, a deleted consumer " +
"once it exists again or the mesh no longer expects it", "once it exists again or the mesh no longer expects it, and a message given up on once DEAD_LETTERS " +
"no longer holds it (novox/hq issue 330)",
Kind: "slow-consumer, max-deliveries, refused, consumer-lost", Severity: conditions.Warning, Phase: 1, Kind: "slow-consumer, max-deliveries, refused, consumer-lost", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.advisoriesErr }, watch: watchAdvisories, needs: func(f *signalFacts) error { return f.advisoriesErr }, watch: watchAdvisories,
newest: func(f *signalFacts) time.Time { newest: func(f *signalFacts) time.Time {
@@ -485,12 +486,94 @@ func watchAdvisories(f *signalFacts) []conditions.Observation {
if a.ID == "controller" { if a.ID == "controller" {
machine = f.host machine = f.host
} }
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind, o := conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind, Token: a.Token,
Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said}) Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said}
if a.Kind == link.AdvisoryMaxDeliveries && a.Token == link.AdvisoryNotKept {
o.Machine = consumerMachine(a.Stream, a.Consumer)
o.Headline = clip(conditions.Capital(fmt.Sprintf("%s gave up on a message, not kept",
consumerWho(a.Stream, a.Consumer))), 60)
o.Explanation = "A listener on the bus could not handle a message, and the mesh could not keep it " +
"for you yet. It tries again every minute."
o.Resolved = "Resolved: the message is kept"
}
out = append(out, o)
}
return append(out, watchDeadLetters(f)...)
}
// watchDeadLetters says each consumer that DEAD_LETTERS holds a message for (novox/hq issue 330): open
// while it holds any, so it clears when they are delivered again or dropped, never because the server
// stopped saying it.
func watchDeadLetters(f *signalFacts) []conditions.Observation {
keys := make([]string, 0, len(f.deadLetters))
for k := range f.deadLetters {
keys = append(keys, k)
}
sort.Strings(keys)
var out []conditions.Observation
for _, key := range keys {
n := f.deadLetters[key]
stream, consumer, _ := strings.Cut(key, ".")
messages, them := "a message", "it"
if n > 1 {
messages, them = fmt.Sprintf("%d messages", n), "them"
}
who := consumerWho(stream, consumer)
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: key, Kind: link.AdvisoryMaxDeliveries,
Machine: consumerMachine(stream, consumer), Severity: conditions.Warning,
Summary: fmt.Sprintf("%s gave up on %s; %s kept in %s until delivered again or dropped, with why, "+
"through the controller's dead-letters verb", link.ConsumerInWords(stream, consumer), messages,
map[bool]string{true: "they are", false: "it is"}[n > 1], broker.DeadLettersStream),
Said: fmt.Sprintf("%d held for %s", n, key),
Headline: clip(conditions.Capital(fmt.Sprintf("%s could not handle %s", who, messages)), 60),
Needs: fmt.Sprintf("deliver %s again or drop %s, from the mesh MCP server.", them, them),
Explanation: conditions.Capital(fmt.Sprintf("%s was handed %s several times and gave up, so what %s "+
"asked for was not done. The mesh keeps %s until you deliver %s again or drop %s.", who, messages,
them, them, them, them)),
Resolved: "Resolved: the messages it gave up on were delivered again or dropped"})
} }
return out return out
} }
// consumerWho is a durable consumer's holder as the operator says it: a module on its machine, the
// controller, or a seat's holders.
func consumerWho(stream, consumer string) string {
switch {
case consumer == broker.ControllerName:
return "the controller"
case strings.HasPrefix(stream, "SEAT_") && strings.HasSuffix(consumer, "_worker"):
seat := strings.ToLower(strings.ReplaceAll(strings.TrimSuffix(strings.TrimPrefix(consumer, "SEAT_"), "_worker"), "_", "-"))
return "the holder of " + seat
case stream == broker.EventsStream:
if node, module, ok := strings.Cut(consumer, "_"); ok {
return module + " on " + node
}
}
return "a listener on the bus"
}
// consumerMachine is the machine a module's consumer is on; empty for the others.
func consumerMachine(stream, consumer string) string {
if stream == broker.EventsStream {
if node, _, ok := strings.Cut(consumer, "_"); ok {
return node
}
}
return ""
}
// clip is words at most n characters long, cut at a word.
func clip(s string, n int) string {
if len(s) <= n {
return s
}
cut := s[:n]
if i := strings.LastIndex(cut, " "); i > 0 {
cut = cut[:i]
}
return cut
}
func watchSelfCheck(f *signalFacts) []conditions.Observation { func watchSelfCheck(f *signalFacts) []conditions.Observation {
every := f.selfCheck.every every := f.selfCheck.every
if every <= 0 { if every <= 0 {
+10
View File
@@ -75,6 +75,9 @@ type signalFacts struct {
advisories []link.Advisory advisories []link.Advisory
lostConsumers map[string]bool lostConsumers map[string]bool
// deadLetters are how many messages DEAD_LETTERS holds per consumer, by `<stream>.<consumer>`
// (novox/hq issue 330): each consumer's max-deliveries condition is open while it holds any.
deadLetters map[string]int
advisoriesErr error advisoriesErr error
selfCheck selfCheckFacts selfCheck selfCheckFacts
@@ -332,6 +335,9 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
} }
f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet)) f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet))
f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories) f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories)
if f.advisoriesErr == nil && w.js != nil {
f.deadLetters, f.advisoriesErr = link.HeldDeadLetters(w.js.Context())
}
f.handActs, f.handActsErr = w.gatherHandActs(ctx, now) f.handActs, f.handActsErr = w.gatherHandActs(ctx, now)
f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last() f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last()
return f return f
@@ -654,6 +660,8 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li
return nil return nil
} }
conditionsFrom = keeper conditionsFrom = keeper
// What a consumer gave up on is read and acted on by the verb in this process (novox/hq issue 330).
deadLettersOn.Store(&busHandles{conn: server.JetStream().Conn(), js: server.JetStream().Context()})
logf := func(format string, args ...any) { fmt.Printf(format+"\n", args...) } logf := func(format string, args ...any) { fmt.Printf(format+"\n", args...) }
stopHearing, err := server.HearAdvisories(logf) stopHearing, err := server.HearAdvisories(logf)
if err != nil { if err != nil {
@@ -680,6 +688,8 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li
return func() { return func() {
stop() stop()
stopHearing() stopHearing()
// No longer serving: the verb answers that it does not read DEAD_LETTERS here (novox/hq issue 330).
deadLettersOn.Store(nil)
flushing, cancel := context.WithTimeout(context.Background(), 10*time.Second) flushing, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel() defer cancel()
keeper.Close(flushing) keeper.Close(flushing)
+7 -3
View File
@@ -38,7 +38,8 @@ type Consumer struct {
Push bool Push bool
// AckWaitSeconds before an unacknowledged delivery is redelivered. // AckWaitSeconds before an unacknowledged delivery is redelivered.
AckWaitSeconds int AckWaitSeconds int
// MaxDeliver before the message is dead-lettered; zero for the mesh's default. // MaxDeliver is how often a message is handed over before the consumer gives it up; zero for no
// bound. What a consumer gives up is kept in DEAD_LETTERS by the controller (novox/hq issue 330).
MaxDeliver int MaxDeliver int
// MaxAckPending is how many deliveries the server lets stand unacknowledged at once; zero for // MaxAckPending is how many deliveries the server lets stand unacknowledged at once; zero for
// the server's default, which is many. **One, for a consumer handled one at a time** // the server's default, which is many. **One, for a consumer handled one at a time**
@@ -145,13 +146,16 @@ func ConsumerFor(p Principal) (Consumer, bool) {
return Consumer{}, false return Consumer{}, false
} }
sort.Strings(filters) sort.Strings(filters)
// And the events given up on and delivered again to this consumer alone (novox/hq issue 330).
filters = append(filters, AgainFilter(consumerDurable(p)))
return Consumer{ return Consumer{
Name: consumerDurable(p), Name: consumerDurable(p),
Stream: consumerStream(p), Stream: consumerStream(p),
Filters: filters, Filters: filters,
AckWaitSeconds: 30, AckWaitSeconds: 30,
MaxDeliver: 5, MaxDeliver: 5,
Why: "what " + p.Module + " declared it consumes; after max-deliver it dead-letters", Why: "what " + p.Module + " declared it consumes; after max-deliver it gives an event up, and the " +
"controller keeps it in DEAD_LETTERS until a person delivers it again or drops it",
}, true }, true
} }
@@ -234,7 +238,7 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
// //
// **No max-deliver, and a long ack wait.** A declaration is settled only after the node has applied // **No max-deliver, and a long ack wait.** A declaration is settled only after the node has applied
// it and reported, which is minutes on a machine pulling images; and a declaration the mesh cannot // it and reported, which is minutes on a machine pulling images; and a declaration the mesh cannot
// get a node to accept is not one to dead-letter, because the stream keeps only the newest per node // get a node to accept is not one to give up on, because the stream keeps only the newest per node
// anyway — so there is exactly one message per node to redeliver, for as long as that node is away. // anyway — so there is exactly one message per node to redeliver, for as long as that node is away.
func NodeConsumer(node string) Consumer { func NodeConsumer(node string) Consumer {
return Consumer{ return Consumer{
+3 -2
View File
@@ -61,8 +61,9 @@ func TestAModuleGetsOneConsumerCarryingEveryFilter(t *testing.T) {
if !ok { if !ok {
t.Fatal("a module that consumes got no consumer") t.Fatal("a module that consumes got no consumer")
} }
if len(c.Filters) != 2 { // Both, and its own share of what is delivered again (novox/hq issue 330).
t.Fatalf("expected both subjects as filters, got %v", c.Filters) if len(c.Filters) != 3 || c.Filters[2] != "mesh.again.one_audit.>" {
t.Fatalf("expected both subjects and its own again filter, got %v", c.Filters)
} }
perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit", perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"}) Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
+9
View File
@@ -154,6 +154,15 @@ func (j *JetStream) EnsureStream(s Stream) error {
Description: s.Why, Description: s.Why,
} }
want.AllowDirect = s.Direct want.AllowDirect = s.Direct
if s.MaxBytes > 0 {
want.MaxBytes = s.MaxBytes
}
if s.DiscardNew {
want.Discard = nats.DiscardNew
}
if s.DuplicatesSeconds > 0 {
want.Duplicates = time.Duration(s.DuplicatesSeconds) * time.Second
}
if s.Retention == RetentionLastPerSubject { if s.Retention == RetentionLastPerSubject {
// Last-per-subject is a limits stream with one message kept per subject, not a // Last-per-subject is a limits stream with one message kept per subject, not a
// retention policy of its own — the state shape, spelled the way the server spells it. // retention policy of its own — the state shape, spelled the way the server spells it.
+5
View File
@@ -412,6 +412,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// both in the mesh's own account; the controller says each as a condition in the mesh's words. // both in the mesh's own account; the controller says each as a condition in the mesh's words.
// Named, not `$JS.EVENT.>`: the other advisories are every API call the mesh makes. // Named, not `$JS.EVENT.>`: the other advisories are every API call the mesh makes.
sub = append(sub, BusAdvisories...) sub = append(sub, BusAdvisories...)
// **And what a consumer gave up on, kept and delivered again** (novox/hq issue 330): the notice
// acknowledged once the message is copied, the copy kept, and an event delivered again to the one
// consumer that gave it up. An ask to a seat is not delivered again, so no seat's queue is granted.
pub = append(pub, "$JS.ACK."+DeadLetterNoticesStream+"."+ControllerName+".>", deadLetterPrefix+">",
againPrefix+">")
case KindPerson: case KindPerson:
// Tools, and nothing else. Every subject a person may publish is a tool call; a person // Tools, and nothing else. Every subject a person may publish is a tool call; a person
+134 -6
View File
@@ -3,11 +3,13 @@ package broker
import ( import (
"fmt" "fmt"
"sort" "sort"
"strings"
) )
// The mesh's own streams. // The mesh's own streams.
// //
// **These four and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118). An earlier // **These and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118; the two that keep what a
// consumer gave up on added for issue 330). An earlier
// reading had the controller create *every* stream at genesis, from a fixed set. That is only the // reading had the controller create *every* stream at genesis, from a fixed set. That is only the
// mesh's own half: a seat's streams are created when the module declaring it is registered, and a // mesh's own half: a seat's streams are created when the module declaring it is registered, and a
// module's durable consumers when it is assigned — neither of which has happened at genesis. What // module's durable consumers when it is assigned — neither of which has happened at genesis. What
@@ -50,11 +52,80 @@ type Stream struct {
// Direct lets a client read a subject's last message without a consumer, which is how a // Direct lets a client read a subject's last message without a consumer, which is how a
// runtime reads its own membership with no JetStream API beyond one request (ADR 0160). // runtime reads its own membership with no JetStream API beyond one request (ADR 0160).
Direct bool Direct bool
// MaxBytes bounds the stream's size, zero for unbounded. With DiscardNew a full stream refuses
// what comes next rather than dropping what it holds: the publisher is told, and says so.
MaxBytes int64
DiscardNew bool
// DuplicatesSeconds is the window in which a message id published twice is kept once; zero for the
// server's default (two minutes).
DuplicatesSeconds int
} }
// AssignmentsStream holds every assignment's membership, the newest per subject. // AssignmentsStream holds every assignment's membership, the newest per subject.
const AssignmentsStream = "ASSIGNMENTS" const AssignmentsStream = "ASSIGNMENTS"
// What a durable consumer gave up on is kept (novox/hq issue 330, design 25 §3).
//
// **The server says it and keeps it; the controller copies it.** A consumer that handed a message over
// as often as it may stops offering it and publishes `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES` with
// the stream and the message's sequence. DeadLetterNoticesStream captures those advisories as the
// server publishes them, so one said while no controller listens is still there when one starts. The
// controller's consumer on it fetches the given-up message by its sequence, while the source stream
// still holds it, and keeps a copy in DeadLettersStream under DeadLetterSubject, with the consumer,
// the subject, how often it was handed over and when it was given up. It stays there until a person
// delivers it again or drops it, with why; a condition is open for as long as it does.
const (
DeadLetterNoticesStream = "DEAD_LETTER_NOTICES"
DeadLettersStream = "DEAD_LETTERS"
// MaxDeliveriesAdvisories is the subject the server says a given-up message on.
MaxDeliveriesAdvisories = "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>"
deadLetterPrefix = "mesh.events.dead."
againPrefix = "mesh.again."
// DeadLettersBytes bounds the kept copies. Full, the stream refuses the next copy, which is said
// as the consumer's condition; it never drops one it holds.
DeadLettersBytes = 256 << 20
)
// DeadLetterSubject is where a message one consumer gave up on is kept: `mesh.events.dead.<stream>.<consumer>`.
func DeadLetterSubject(stream, consumer string) string {
return deadLetterPrefix + stream + "." + consumer
}
// DeadLetterOf is the stream and consumer a kept message's subject names; false for any other subject.
func DeadLetterOf(subject string) (stream, consumer string, ok bool) {
rest, found := strings.CutPrefix(subject, deadLetterPrefix)
if !found {
return "", "", false
}
stream, consumer, ok = strings.Cut(rest, ".")
return stream, consumer, ok && stream != "" && consumer != "" && !strings.Contains(consumer, ".")
}
// AgainSubject is where an event given up on is delivered again to the one consumer that gave it up,
// and to nobody else: `mesh.again.<consumer>.` and the original subject without its `mesh.`. Every
// consumer on EVENTS filters its own (AgainFilter), and a module's runtime reads the event's key from
// the tokens around `.event.`, so the handler sees the same key it saw the first time.
func AgainSubject(consumer, original string) string {
return againPrefix + consumer + "." + strings.TrimPrefix(original, "mesh.")
}
// AgainFilter is the one consumer's share of the subjects events are delivered again on.
func AgainFilter(consumer string) string { return againPrefix + consumer + ".>" }
// OriginalOfAgain is the subject an event delivered again was first published on; false for a subject
// that is not one delivered again.
func OriginalOfAgain(subject string) (string, bool) {
rest, found := strings.CutPrefix(subject, againPrefix)
if !found {
return "", false
}
_, original, ok := strings.Cut(rest, ".")
if !ok || original == "" {
return "", false
}
return "mesh." + original, true
}
// MeshStreams is the foundation set, in the order a person reads it. // MeshStreams is the foundation set, in the order a person reads it.
// //
// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live // **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live
@@ -97,7 +168,9 @@ func MeshStreams() []Stream {
// A seat's own events ride here too: they are 1:many like any event, and the // A seat's own events ride here too: they are 1:many like any event, and the
// `event` token keeps them clear of both the seat's work queue (`accept`) and its // `event` token keeps them clear of both the seat's work queue (`accept`) and its
// tools (`tool`), which must not be persisted. // tools (`tool`), which must not be persisted.
Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>"}, // And an event given up on, delivered again to the one consumer that gave it up
// (novox/hq issue 330): under `mesh.again.<consumer>.`, which only that consumer filters.
Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>", againPrefix + ">"},
Retention: RetentionLimits, Retention: RetentionLimits,
MaxAge: 7 * 24 * 60 * 60, MaxAge: 7 * 24 * 60 * 60,
MaxMsgsPerSubject: 10000, MaxMsgsPerSubject: 10000,
@@ -105,6 +178,26 @@ func MeshStreams() []Stream {
"excluded by the event token; per-subject caps keep a noisy emitter from " + "excluded by the event token; per-subject caps keep a noisy emitter from " +
"evicting a quiet one without splitting the stream", "evicting a quiet one without splitting the stream",
}, },
{
Name: DeadLetterNoticesStream,
Subjects: []string{MaxDeliveriesAdvisories},
Retention: RetentionWorkQueue,
MaxAge: 7 * 24 * 60 * 60,
Why: "the server's word that a consumer gave up on a message, kept until the controller has " +
"copied the message into DEAD_LETTERS (novox/hq issue 330); a week, the longest the source " +
"streams keep what they are about",
},
{
Name: DeadLettersStream,
Subjects: []string{deadLetterPrefix + ">"},
Retention: RetentionLimits,
MaxBytes: DeadLettersBytes,
DiscardNew: true,
DuplicatesSeconds: 24 * 60 * 60,
Why: "every message a consumer gave up on, with its consumer, subject, deliveries and when, kept " +
"until a person delivers it again or drops it with why (novox/hq issue 330); no age, and full " +
"it refuses the next copy rather than drop one it holds",
},
} }
} }
@@ -298,7 +391,7 @@ const EventsStream = "EVENTS"
// //
// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's, // **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's,
// not the server's (window.go): a message is held with a nak-and-delay until the controller either // not the server's (window.go): a message is held with a nak-and-delay until the controller either
// takes it or gives up and says so. A max-deliver here would dead-letter a push that was being // takes it or gives up and says so. A max-deliver here would give up on a push that was being
// held through a store restart — the exact message the stream exists to protect — some minutes // held through a store restart — the exact message the stream exists to protect — some minutes
// before the controller had finished deciding about it. // before the controller had finished deciding about it.
func MeshConsumers() []Consumer { func MeshConsumers() []Consumer {
@@ -314,7 +407,7 @@ func MeshConsumers() []Consumer {
{ {
Name: ControllerName, Name: ControllerName,
Stream: "EVENTS", Stream: "EVENTS",
Filters: ControllerFollows, Filters: append(append([]string(nil), ControllerFollows...), AgainFilter(ControllerName)),
Push: true, Push: true,
AckWaitSeconds: 30, AckWaitSeconds: 30,
MaxDeliver: 5, MaxDeliver: 5,
@@ -331,12 +424,47 @@ func MeshConsumers() []Consumer {
Resettable: "what it drops is caught up: merges by the catch-up pass (issue 266), build outcomes " + Resettable: "what it drops is caught up: merges by the catch-up pass (issue 266), build outcomes " +
"from the build records (issue 214), a provider's failing word said again (ADR 0224)", "from the build records (issue 214), a provider's failing word said again (ADR 0224)",
Why: "the events the mesh's own controller reacts to, one at a time; after " + Why: "the events the mesh's own controller reacts to, one at a time; after " +
"max-deliver it dead-letters, because an announcement it cannot act on will not " + "max-deliver it gives the event up, and the controller keeps it in DEAD_LETTERS until " +
"become actionable", "a person delivers it again or drops it",
},
// What the server said a consumer gave up on (novox/hq issue 330): copied into DEAD_LETTERS and
// acknowledged. No max-deliver: a notice the controller could not copy is offered again, and said.
{
Name: ControllerName,
Stream: DeadLetterNoticesStream,
Push: true,
AckWaitSeconds: 30,
Why: "the controller copies each message a consumer gave up on into DEAD_LETTERS; no max-deliver, " +
"because a notice it gave up on would lose the message it is about",
}, },
} }
} }
// NoticesConsumer is the controller's consumer on DEAD_LETTER_NOTICES (novox/hq issue 330).
func NoticesConsumer() Consumer {
for _, c := range MeshConsumers() {
if c.Stream == DeadLetterNoticesStream {
return c
}
}
panic("the mesh's consumers carry none on " + DeadLetterNoticesStream)
}
// AssertServingConsumers are the controller's own consumers its serving cannot go without: all but the
// one on DEAD_LETTER_NOTICES, which the keeper of dead letters asserts and retries by itself, so a fault
// there never stops the controller serving (novox/hq issue 330).
func AssertServingConsumers(e Ensurer) error {
for _, c := range MeshConsumers() {
if c.Stream == DeadLetterNoticesStream {
continue
}
if err := e.EnsureConsumer(c); err != nil {
return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err)
}
}
return nil
}
// Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason // Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason
// Asserter is. // Asserter is.
type Ensurer interface { type Ensurer interface {
+4
View File
@@ -127,6 +127,10 @@ func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) {
"NODES": RetentionLastPerSubject, "NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits, "EVENTS": RetentionLimits,
"ASSIGNMENTS": RetentionLastPerSubject, "ASSIGNMENTS": RetentionLastPerSubject,
// What a consumer gave up on (novox/hq issue 330): the server's notice taken once, the message
// kept until somebody acts.
"DEAD_LETTER_NOTICES": RetentionWorkQueue,
"DEAD_LETTERS": RetentionLimits,
} }
got := map[string]Retention{} got := map[string]Retention{}
for _, s := range MeshStreams() { for _, s := range MeshStreams() {
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] }
subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
+5
View File
@@ -133,6 +133,11 @@ var WritersTable = []WriterRow{
Others: "—"}, Others: "—"},
{State: "the facts snapshot", Writer: "controller", KeptIn: "the artifact store, facts/latest", {State: "the facts snapshot", Writer: "controller", KeptIn: "the artifact store, facts/latest",
Others: "the build seat reads"}, Others: "the build seat reads"},
// What a consumer gave up on (novox/hq issue 330): copied by the controller from the server's notice,
// and delivered again or dropped only through its verb, with why.
{State: "a message a consumer gave up on", Writer: "controller", KeptIn: "the bus, the stream " + DeadLettersStream,
Others: "read, delivered again or dropped through the controller's dead-letters verb",
Subjects: []string{deadLetterPrefix + ">", againPrefix + ">"}, Writes: isController},
} }
// CheckWriters refuses a grant that lets a principal publish on a subject the writers table gives // CheckWriters refuses a grant that lets a principal publish on a subject the writers table gives
+1
View File
@@ -30,6 +30,7 @@ var designRows = []string{
"a provider's standing", "a provider's standing",
"the operator-channel's open messages", "the operator-channel's open messages",
"the facts snapshot", "the facts snapshot",
"a message a consumer gave up on",
} }
func TestTheWritersTableIsTheDesigns(t *testing.T) { func TestTheWritersTableIsTheDesigns(t *testing.T) {
+15
View File
@@ -424,6 +424,21 @@ var ControllerVerbs = []Verb{
"why": "with consumer or older-than: why — required, and recorded in the hand-act log", "why": "with consumer or older-than: why — required, and recorded in the hand-act log",
"cause": "with consumer or older-than: the cause in a word (cleanup-waiting when absent)", "cause": "with consumer or older-than: the cause in a word (cleanup-waiting when absent)",
}, nil, "confirm")}, }, nil, "confirm")},
// What a consumer gave up on (novox/hq issue 330): kept in DEAD_LETTERS until a person acts on it.
{Name: "dead-letters", Description: "Every message a consumer on the bus gave up on after handing it over " +
"as often as it may, kept in DEAD_LETTERS: whose consumer, the subject, how often it was handed over and " +
"when it was given up, newest first. With id: that one whole, with what it said. With deliver: hand it " +
"again to the consumer that gave it up, and nobody else; with drop: let it go for good. Delivering and " +
"dropping are hand acts, which say why (novox/hq issue 330).",
Input: schema(map[string]string{
"id": "a dead letter's id, as the list gives it: that one whole, with its message",
"consumer": "a consumer's name, or <stream>.<consumer>: only what that one gave up on",
"limit": "how many to list, newest first (default 50); only when listing",
"deliver": "a dead letter's id: deliver it again to the consumer that gave it up (needs why)",
"drop": "a dead letter's id: drop it for good (needs why)",
"why": "with deliver or drop: why — required, and recorded in the hand-act log",
"cause": "with deliver or drop: the cause in a word (dead-letter when absent)",
}, nil)},
// The data every machine declares (novox/hq ADR 0233). // The data every machine declares (novox/hq ADR 0233).
{Name: "data", Description: "Every item of data every machine declares, as the self-check last measured it: " + {Name: "data", Description: "Every item of data every machine declares, as the self-check last measured it: " +
"its class (irreplaceable, rebuildable, cache), where it is, its size, its newest write, its newest good " + "its class (irreplaceable, rebuildable, cache), where it is, its size, its newest write, its newest good " +
+15 -4
View File
@@ -34,6 +34,9 @@ const (
AdvisoryMaxDeliveries = "max-deliveries" AdvisoryMaxDeliveries = "max-deliveries"
AdvisoryRefused = "refused" AdvisoryRefused = "refused"
AdvisoryConsumerLost = "consumer-lost" 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. // 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 ID string
// Stream and Consumer are the consumer it is about, when it is about one. // Stream and Consumer are the consumer it is about, when it is about one.
Stream, Consumer string 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 is the newest saying, in the mesh's words.
Said string Said string
First, Last time.Time First, Last time.Time
@@ -62,7 +67,7 @@ var Advisories = &AdvisoryLog{seen: map[string]*Advisory{}}
func (l *AdvisoryLog) Heard(a Advisory, at time.Time) { func (l *AdvisoryLog) Heard(a Advisory, at time.Time) {
l.mu.Lock() l.mu.Lock()
defer l.mu.Unlock() defer l.mu.Unlock()
key := a.Kind + "/" + a.ID key := a.Kind + "/" + a.ID + "/" + a.Token
if had, ok := l.seen[key]; ok { if had, ok := l.seen[key]; ok {
had.Last, had.Said, had.Count = at, a.Said, had.Count+1 had.Last, had.Said, had.Count = at, a.Said, had.Count+1
return return
@@ -107,8 +112,9 @@ func ReadAdvisory(subject string, body []byte) (Advisory, bool) {
switch { switch {
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."): case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."):
return Advisory{Kind: AdvisoryMaxDeliveries, ID: id, Stream: a.Stream, Consumer: a.Consumer, 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 "+ Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: what it asked for was not "+
"again, and what it asked for was not done", who, a.StreamSeq, a.Deliveries)}, true "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."): case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.DELETED."):
if !MeshNamed(a.Stream, a.Consumer) { if !MeshNamed(a.Stream, a.Consumer) {
// A reader's own consumer, gone when it finished — every watch of a bucket and every // 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 { for _, subject := range broker.BusAdvisories {
sub, err := conn.Subscribe(subject, func(m *nats.Msg) { sub, err := conn.Subscribe(subject, func(m *nats.Msg) {
if a, ok := ReadAdvisory(m.Subject, m.Data); ok { 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) logf("the bus says: %s", a.Said)
} }
}) })
+2 -1
View File
@@ -37,7 +37,8 @@ func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) {
[]byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`)) []byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`))
if !ok || a.Kind != AdvisoryMaxDeliveries || a.ID != "EVENTS.anchor_shop" || 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: "+ 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) t.Fatalf("%+v", a)
} }
} }
+5 -2
View File
@@ -8,13 +8,16 @@ import (
// JetStream connection the caller has already raised the streams on. Nothing is declared here — // JetStream connection the caller has already raised the streams on. Nothing is declared here —
// the streams and the controller's consumers are asserted by Raise, before anything is served. // the streams and the controller's consumers are asserted by Raise, before anything is served.
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server { func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
logger := newLog()
inbound := Nats(js).(*natsInbound)
inbound.log = logger
return &Server{ return &Server{
inbound: Nats(js), inbound: inbound,
bus: OverNATS{JS: js.Context(), Conn: js.Conn()}, bus: OverNATS{JS: js.Context(), Conn: js.Conn()},
js: js, js: js,
enroller: enroller, enroller: enroller,
listener: listener, listener: listener,
log: newLog(), log: logger,
} }
} }
+372
View File
@@ -0,0 +1,372 @@
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
}
// 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. 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())
}
+336
View File
@@ -0,0 +1,336 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"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}` || len(whole.Headers["x-node"]) != 1 || whole.Headers["x-node"][0] != "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")
}
}
// Only an event is delivered again, and only to a consumer that is on the bus and filters its again
// subject: a dead letter is never let go as delivered while nobody receives it.
func TestADeadLetterIsDeliveredAgainOnlyWhereItIsReceived(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)
// A consumer as it was before this fix: its filters do not take what is delivered again.
c, _ := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr",
Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"})
c.Filters = c.Filters[:len(c.Filters)-1]
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "does not yet take") {
t.Fatalf("delivered to a consumer that does not filter its again subject: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// The consumer gone: nothing would receive it.
if err := js.Context().DeleteConsumer(broker.EventsStream, c.Name); err != nil {
t.Fatal(err)
}
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "no longer on the bus") {
t.Fatalf("delivered to a consumer that is gone: %v", err)
}
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
t.Fatalf("a refused delivery let the dead letter go: %v", err)
}
// Not an event: refused, whatever stream it is from.
for _, other := range []DeadLetter{
{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", Subject: "mesh.seat.telegram-sender.accept.send"},
{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"},
} {
if _, err := AgainTo(js.Context(), other); err == nil {
t.Errorf("%s was given a subject to be delivered again on", other.Stream)
}
}
}
// A total says how many are held, however few the list carries.
func TestTheListSaysTheTotalAndStopsAtItsLimit(t *testing.T) {
js := aBus(t)
for seq := 1; seq <= 5; seq++ {
if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_sonarr",`+
`"stream_seq":%d,"deliveries":5}`, seq))); err != nil {
t.Fatal(err)
}
}
list, total, err := DeadLetters(js.Context(), "media_sonarr", 2)
if err != nil || total != 5 || len(list) != 2 || list[0].Sequence != 5 {
t.Fatalf("%d of %d (%v): %+v", len(list), total, err, list)
}
if _, total, _ := DeadLetters(js.Context(), "another_one", 2); total != 0 {
t.Fatalf("another consumer's total is %d", total)
}
}
// 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 })
}
// The notices cannot be taken (their stream is gone): said as a condition, and the controller serves on.
func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) {
js := aBus(t)
if err := js.Context().DeleteStream(broker.DeadLetterNoticesStream); err != nil {
t.Fatal(err)
}
stop := keepingGivenUp(context.Background(), js, quiet())
defer stop()
eventually(t, "the failure said", func() bool {
for _, a := range Advisories.Since(time.Now().Add(-time.Minute)) {
if a.Token == AdvisoryNotKept && a.Stream == broker.DeadLetterNoticesStream {
return true
}
}
return false
})
held := &counted{}
_, stopServing := servingOn(t, js, held)
defer stopServing()
if _, err := js.Context().Publish("mesh.control.anchor.report", []byte(`{"node":"anchor"}`)); err != nil {
t.Fatal(err)
}
eventually(t, "a report heard while the notices cannot be taken", func() bool { return held.count() == 1 })
}
+22 -1
View File
@@ -41,6 +41,15 @@ type natsInbound struct {
// that restarts loses these and starts the window again, which is correct — it is holding // that restarts loses these and starts the window again, which is correct — it is holding
// nothing, and the messages are all still on the server. // nothing, and the messages are all still on the server.
since map[uint64]time.Time since map[uint64]time.Time
// log is where it says what it could not do; the standard logger when nobody gave one.
log *log.Logger
}
func (n *natsInbound) logger() *log.Logger {
if n.log != nil {
return n.log
}
return log.Default()
} }
// Nats is the consume side of the bus being built. // Nats is the consume side of the bus being built.
@@ -70,7 +79,7 @@ func (n *natsInbound) Close() {}
// which message is held, and since when — is read and written without a lock because the AMQP loop // which message is held, and since when — is read and written without a lock because the AMQP loop
// never had two. A second goroutine would make that wrong in a way no test would catch. // never had two. A second goroutine would make that wrong in a way no test would catch.
func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error { func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
if err := broker.AssertMeshConsumers(n.js); err != nil { if err := broker.AssertServingConsumers(n.js); err != nil {
return err return err
} }
js, conn := n.js.Context(), n.js.Conn() js, conn := n.js.Context(), n.js.Conn()
@@ -94,6 +103,13 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
holding.Store(true) holding.Store(true)
defer holding.Store(false) 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.
// **Never the reason the controller stops serving**: a notice it cannot take yet waits in its stream,
// and the failure is said and tried again every minute.
stopKeeping := keepingGivenUp(ctx, n.js, n.logger())
defer stopKeeping()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // 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. // they are their own guarantee: a lost one is the next one.
beats := make(chan *nats.Msg, Prefetch) beats := make(chan *nats.Msg, Prefetch)
@@ -167,6 +183,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), func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control),
msg *nats.Msg, streamed bool) { 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) kind, known := kindOfSubject(msg.Subject)
if !known { if !known {
if streamed { if streamed {
+93
View File
@@ -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")
}
}
+1
View File
@@ -71,6 +71,7 @@
"bus", "bus",
"retire", "retire",
"cleanup", "cleanup",
"dead-letters",
"data", "data",
"build", "build",
"artifacts", "artifacts",