diff --git a/cmd/mesh-controller/deadletters.go b/cmd/mesh-controller/deadletters.go new file mode 100644 index 00000000..d929bb9a --- /dev/null +++ b/cmd/mesh-controller/deadletters.go @@ -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(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 +} diff --git a/cmd/mesh-controller/deadletters_test.go b/cmd/mesh-controller/deadletters_test.go new file mode 100644 index 00000000..5b3aef7a --- /dev/null +++ b/cmd/mesh-controller/deadletters_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 19a662b8..5d5ad5d7 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -82,6 +82,11 @@ var handActVerbs = []handActVerb{ {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: "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 // 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 " + diff --git a/cmd/mesh-controller/plain_words.go b/cmd/mesh-controller/plain_words.go index 78d75c1e..0b333b26 100644 --- a/cmd/mesh-controller/plain_words.go +++ b/cmd/mesh-controller/plain_words.go @@ -544,9 +544,11 @@ var plainWordings = map[string]func(conditions.Observation) words{ Resolved: "The listener keeps up again"} }), "max-deliveries": worded(func(o conditions.Observation) words { - return words{Headline: "A message could not be handled", - Explanation: "The bus gave up on a message after trying to hand it over too many times.", - Resolved: "Resolved: messages are handled again"} + return words{Headline: "A listener gave up on messages", + Needs: "deliver them again or drop them, from the mesh MCP server.", + 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 { return words{Headline: "The bus refuses some messages", diff --git a/cmd/mesh-controller/probes.go b/cmd/mesh-controller/probes.go index 598e97b9..1f19b088 100644 --- a/cmd/mesh-controller/probes.go +++ b/cmd/mesh-controller/probes.go @@ -875,6 +875,12 @@ func streamDiffers(want broker.Stream, have nats.StreamConfig) string { if perSubject != 0 && 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.DiscardNew && have.Discard != nats.DiscardNew { + differs = append(differs, "drops what it holds when full, defined to refuse what comes next") + } return strings.Join(differs, "; ") } diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index e134cf3e..58aed4a9 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -250,6 +250,8 @@ func (a *verbArguments) commandLine() ([]string, error) { return argv, nil case "tools": 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": return []string{"status", "--json"}, nil case "nodes": @@ -982,6 +984,9 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { if verb == "calls" { return callsAnswer(link.Calls, a.given["call"]) } + if verb == "dead-letters" { + return deadLettersAnswer(ctx, a) + } if verb == "doctor" { // 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. @@ -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 // 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 // through `command`. A push sends the machine holding the bus first when its user list changed, the diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index 875c4940..a337b0f1 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -154,8 +154,9 @@ var signalsTable = []signalRow{ }}, {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", - Trigger: "any", Bound: "any occurrence; clears after an hour without another, and a deleted consumer " + - "once it exists again or the mesh no longer expects it", + 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, 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, needs: func(f *signalFacts) error { return f.advisoriesErr }, watch: watchAdvisories, newest: func(f *signalFacts) time.Time { @@ -485,12 +486,94 @@ func watchAdvisories(f *signalFacts) []conditions.Observation { if a.ID == "controller" { machine = f.host } - out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind, - Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said}) + 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} + 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 } +// 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 { every := f.selfCheck.every if every <= 0 { diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 6f73c860..1bede623 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -75,6 +75,9 @@ type signalFacts struct { advisories []link.Advisory lostConsumers map[string]bool + // deadLetters are how many messages DEAD_LETTERS holds per consumer, by `.` + // (novox/hq issue 330): each consumer's max-deliveries condition is open while it holds any. + deadLetters map[string]int advisoriesErr error selfCheck selfCheckFacts @@ -332,6 +335,9 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts { } f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet)) 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.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last() return f @@ -654,6 +660,8 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li return nil } 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...) } stopHearing, err := server.HearAdvisories(logf) if err != nil { diff --git a/internal/broker/derived.go b/internal/broker/derived.go index a0bad9c7..28831526 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -38,7 +38,8 @@ type Consumer struct { Push bool // AckWaitSeconds before an unacknowledged delivery is redelivered. 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 // 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** @@ -145,13 +146,16 @@ func ConsumerFor(p Principal) (Consumer, bool) { return Consumer{}, false } 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{ Name: consumerDurable(p), Stream: consumerStream(p), Filters: filters, AckWaitSeconds: 30, 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 } @@ -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 // 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. func NodeConsumer(node string) Consumer { return Consumer{ diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go index c7257d5c..fc541705 100644 --- a/internal/broker/derived_test.go +++ b/internal/broker/derived_test.go @@ -61,8 +61,9 @@ func TestAModuleGetsOneConsumerCarryingEveryFilter(t *testing.T) { if !ok { t.Fatal("a module that consumes got no consumer") } - if len(c.Filters) != 2 { - t.Fatalf("expected both subjects as filters, got %v", c.Filters) + // Both, and its own share of what is delivered again (novox/hq issue 330). + 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", Consumes: []string{"shop.order.placed"}, PasswordHash: "x"}) diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 198ee68c..7480f578 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -154,6 +154,15 @@ func (j *JetStream) EnsureStream(s Stream) error { Description: s.Why, } 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 { // 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. diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 4e022964..d313237a 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -412,6 +412,12 @@ 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. // Named, not `$JS.EVENT.>`: the other advisories are every API call the mesh makes. 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 a message delivered again — an + // event to the one consumer that gave it up, an ask back onto its seat's queue, whose one worker + // is the consumer that gave it up. + pub = append(pub, "$JS.ACK."+DeadLetterNoticesStream+"."+ControllerName+".>", deadLetterPrefix+">", + againPrefix+">", "mesh.seat.*.accept.>") case KindPerson: // Tools, and nothing else. Every subject a person may publish is a tool call; a person diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 4d0a91b5..2e514e47 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -3,11 +3,13 @@ package broker import ( "fmt" "sort" + "strings" ) // 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 // 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 @@ -50,11 +52,80 @@ type Stream struct { // 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). 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. 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..`. +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..` 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. // // **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 // `event` token keeps them clear of both the seat's work queue (`accept`) and its // 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..`, which only that consumer filters. + Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>", againPrefix + ">"}, Retention: RetentionLimits, MaxAge: 7 * 24 * 60 * 60, MaxMsgsPerSubject: 10000, @@ -105,6 +178,26 @@ func MeshStreams() []Stream { "excluded by the event token; per-subject caps keep a noisy emitter from " + "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, // 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 // before the controller had finished deciding about it. func MeshConsumers() []Consumer { @@ -314,7 +407,7 @@ func MeshConsumers() []Consumer { { Name: ControllerName, Stream: "EVENTS", - Filters: ControllerFollows, + Filters: append(append([]string(nil), ControllerFollows...), AgainFilter(ControllerName)), Push: true, AckWaitSeconds: 30, MaxDeliver: 5, @@ -331,8 +424,18 @@ func MeshConsumers() []Consumer { 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)", 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 " + - "become actionable", + "max-deliver it gives the event up, and the controller keeps it in DEAD_LETTERS until " + + "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", }, } } diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index 9a2592ec..7b736172 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -127,6 +127,10 @@ func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) { "NODES": RetentionLastPerSubject, "EVENTS": RetentionLimits, "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{} for _, s := range MeshStreams() { diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 2891cc56..5f31c221 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { 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.*.accept.>", "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"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/broker/writers.go b/internal/broker/writers.go index 9c18ed15..70cb099b 100644 --- a/internal/broker/writers.go +++ b/internal/broker/writers.go @@ -133,6 +133,11 @@ var WritersTable = []WriterRow{ Others: "—"}, {State: "the facts snapshot", Writer: "controller", KeptIn: "the artifact store, facts/latest", 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 diff --git a/internal/broker/writers_test.go b/internal/broker/writers_test.go index 32d53dc2..72aba5e4 100644 --- a/internal/broker/writers_test.go +++ b/internal/broker/writers_test.go @@ -30,6 +30,7 @@ var designRows = []string{ "a provider's standing", "the operator-channel's open messages", "the facts snapshot", + "a message a consumer gave up on", } func TestTheWritersTableIsTheDesigns(t *testing.T) { diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index de1faf8a..332cb4ea 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -424,6 +424,21 @@ var ControllerVerbs = []Verb{ "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)", }, 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 .: 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). {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 " + diff --git a/internal/link/advisories.go b/internal/link/advisories.go index 36fc6bb0..013f2b5f 100644 --- a/internal/link/advisories.go +++ b/internal/link/advisories.go @@ -34,6 +34,9 @@ const ( AdvisoryMaxDeliveries = "max-deliveries" AdvisoryRefused = "refused" AdvisoryConsumerLost = "consumer-lost" + // AdvisoryNotKept is the token of a max-deliveries advisory whose message could not be kept in + // DEAD_LETTERS (novox/hq issue 330): said from the log, since the stream does not hold it. + AdvisoryNotKept = "not-kept" ) // Advisory is one thing the bus said, kept as its newest word and how often it was said. @@ -43,6 +46,8 @@ type Advisory struct { ID string // Stream and Consumer are the consumer it is about, when it is about one. Stream, Consumer string + // Token is the condition key's last part, when it is not Kind. + Token string // Said is the newest saying, in the mesh's words. Said string First, Last time.Time @@ -62,7 +67,7 @@ var Advisories = &AdvisoryLog{seen: map[string]*Advisory{}} func (l *AdvisoryLog) Heard(a Advisory, at time.Time) { l.mu.Lock() defer l.mu.Unlock() - key := a.Kind + "/" + a.ID + key := a.Kind + "/" + a.ID + "/" + a.Token if had, ok := l.seen[key]; ok { had.Last, had.Said, had.Count = at, a.Said, had.Count+1 return @@ -107,8 +112,9 @@ func ReadAdvisory(subject string, body []byte) (Advisory, bool) { switch { case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."): return Advisory{Kind: AdvisoryMaxDeliveries, ID: id, Stream: a.Stream, Consumer: a.Consumer, - Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: it will not be delivered "+ - "again, and what it asked for was not done", who, a.StreamSeq, a.Deliveries)}, true + Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: what it asked for was not "+ + "done, and the controller keeps it in %s until it is delivered again or dropped", who, a.StreamSeq, + a.Deliveries, broker.DeadLettersStream)}, true case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.DELETED."): if !MeshNamed(a.Stream, a.Consumer) { // A reader's own consumer, gone when it finished — every watch of a bucket and every @@ -175,7 +181,12 @@ func (s *Server) HearAdvisories(logf func(string, ...any)) (func(), error) { for _, subject := range broker.BusAdvisories { sub, err := conn.Subscribe(subject, func(m *nats.Msg) { if a, ok := ReadAdvisory(m.Subject, m.Data); ok { - Advisories.Heard(a, time.Now()) + // A message given up on is said from DEAD_LETTERS, where it is kept, for as long as it is + // kept (novox/hq issue 330) — not for an hour after the server said it; one that could not + // be kept is recorded by whoever tried. + if a.Kind != AdvisoryMaxDeliveries { + Advisories.Heard(a, time.Now()) + } logf("the bus says: %s", a.Said) } }) diff --git a/internal/link/advisories_test.go b/internal/link/advisories_test.go index f121b8db..56d174f2 100644 --- a/internal/link/advisories_test.go +++ b/internal/link/advisories_test.go @@ -37,7 +37,8 @@ func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) { []byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`)) if !ok || a.Kind != AdvisoryMaxDeliveries || a.ID != "EVENTS.anchor_shop" || a.Said != "how shop on anchor hears what it consumes handed message 7 over 5 times and gave up on it: "+ - "it will not be delivered again, and what it asked for was not done" { + "what it asked for was not done, and the controller keeps it in DEAD_LETTERS until it is delivered "+ + "again or dropped" { t.Fatalf("%+v", a) } } diff --git a/internal/link/deadletters.go b/internal/link/deadletters.go new file mode 100644 index 00000000..6cc6fd10 --- /dev/null +++ b/internal/link/deadletters.go @@ -0,0 +1,302 @@ +package link + +import ( + "encoding/json" + "errors" + "fmt" + "log" + "strconv" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// What a consumer gave up on, kept until a person acts on it (novox/hq issue 330, design 25 §3). +// +// A durable consumer hands a message over as often as it may, gives up on it and says so on +// `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES`. The server's notice is captured in its own stream +// (broker.DeadLetterNoticesStream), so one said while no controller listens waits for the next. The +// serving controller takes each notice, fetches the message it names by its sequence, while the source +// stream still holds it, and keeps a copy in DEAD_LETTERS with what the notice said: the consumer, the +// subject, how often it was handed over and when it was given up. Until then a message given up on was +// kept only by its source, which drops an event after a week, and nothing said which. +// +// The copy's headers are the message's own, but for the server's (`Nats-…`), which describe the first +// publish and would refuse or de-duplicate the copy; they are kept under DeadHeader names. + +// The headers a kept message carries beside its own. +const ( + DeadHeaderStream = "Mesh-Dead-Stream" + DeadHeaderConsumer = "Mesh-Dead-Consumer" + DeadHeaderSequence = "Mesh-Dead-Sequence" + DeadHeaderSubject = "Mesh-Dead-Subject" + DeadHeaderDeliveries = "Mesh-Dead-Deliveries" + DeadHeaderGaveUp = "Mesh-Dead-Gave-Up" + DeadHeaderPublished = "Mesh-Dead-Published" + DeadHeaderLost = "Mesh-Dead-Lost" + // DeadHeaderMsgID is the message's own de-duplication id, kept, since the copy carries its own. + DeadHeaderMsgID = "Mesh-Dead-Msg-Id" + // AgainHeader marks a message delivered again, with the kept copy's id it came from. + AgainHeader = "Mesh-Delivered-Again" +) + +// DeadLetter is one message a consumer gave up on, as DEAD_LETTERS keeps it. +type DeadLetter struct { + // ID is its sequence in DEAD_LETTERS: what the verb names it by. + ID uint64 `json:"id"` + Stream string `json:"stream"` + Consumer string `json:"consumer"` + // Who is the consumer in the mesh's words: whose, and for what. + Who string `json:"who"` + Subject string `json:"subject"` + Sequence uint64 `json:"sequence"` + Deliveries uint64 `json:"deliveries"` + GaveUp time.Time `json:"gave_up"` + Published time.Time `json:"published,omitzero"` + // Lost says why the message itself is not kept: its source no longer held it when the notice was + // taken. The record of it is kept all the same, so it is said and dropped, never silently missing. + Lost string `json:"lost,omitempty"` + Size int `json:"size"` + // Body and Headers are the message itself, given only for one dead letter asked by its id. + Body string `json:"body,omitempty"` + Headers map[string]string `json:"headers,omitempty"` +} + +// maxDeliveries is the part of the server's notice a kept message is made from. +type maxDeliveries struct { + Stream string `json:"stream"` + Consumer string `json:"consumer"` + StreamSeq uint64 `json:"stream_seq"` + Deliveries uint64 `json:"deliveries"` + Timestamp time.Time `json:"timestamp"` +} + +// KeepDeadLetter keeps the message one maximum-deliveries notice is about. An error leaves nothing +// kept, and the notice is to be taken again: the source still holds the message, or a full +// DEAD_LETTERS has room again. A source that no longer holds it is no error: the record is kept with +// why, so it is said. +func KeepDeadLetter(js nats.JetStreamContext, notice []byte) (DeadLetter, error) { + var a maxDeliveries + if err := json.Unmarshal(notice, &a); err != nil || a.Stream == "" || a.Consumer == "" || a.StreamSeq == 0 { + return DeadLetter{}, fmt.Errorf("not a maximum-deliveries notice the mesh can read: %.200s", notice) + } + gaveUp := a.Timestamp + if gaveUp.IsZero() { + gaveUp = time.Now() + } + kept := &nats.Msg{Subject: broker.DeadLetterSubject(a.Stream, a.Consumer), Header: nats.Header{}} + raw, err := js.GetMsg(a.Stream, a.StreamSeq) + switch { + case err == nil: + for k, v := range raw.Header { + if strings.HasPrefix(k, "Nats-") { + continue + } + kept.Header[k] = v + } + if id := raw.Header.Get(nats.MsgIdHdr); id != "" { + kept.Header.Set(DeadHeaderMsgID, id) + } + kept.Header.Set(DeadHeaderSubject, raw.Subject) + kept.Header.Set(DeadHeaderPublished, raw.Time.UTC().Format(time.RFC3339Nano)) + kept.Data = raw.Data + case errors.Is(err, nats.ErrMsgNotFound) || errors.Is(err, nats.ErrStreamNotFound): + kept.Header.Set(DeadHeaderLost, fmt.Sprintf("%s no longer held message %d when the notice was taken: %v", + a.Stream, a.StreamSeq, err)) + default: + return DeadLetter{}, fmt.Errorf("message %d of %s could not be read: %w", a.StreamSeq, a.Stream, err) + } + kept.Header.Set(DeadHeaderStream, a.Stream) + kept.Header.Set(DeadHeaderConsumer, a.Consumer) + kept.Header.Set(DeadHeaderSequence, strconv.FormatUint(a.StreamSeq, 10)) + kept.Header.Set(DeadHeaderDeliveries, strconv.FormatUint(a.Deliveries, 10)) + kept.Header.Set(DeadHeaderGaveUp, gaveUp.UTC().Format(time.RFC3339Nano)) + // One copy per message given up, however often its notice is taken: a controller that copied and + // stopped before acknowledging, or two controllers each taking it. + kept.Header.Set(nats.MsgIdHdr, fmt.Sprintf("%s.%s.%d", a.Stream, a.Consumer, a.StreamSeq)) + ack, err := js.PublishMsg(kept) + if err != nil { + return DeadLetter{}, fmt.Errorf("message %d of %s could not be kept in %s: %w", a.StreamSeq, a.Stream, + broker.DeadLettersStream, err) + } + return deadLetterOf(ack.Sequence, kept.Subject, kept.Header, kept.Data, false), nil +} + +// deadLetterOf reads a kept message. +func deadLetterOf(id uint64, subject string, h nats.Header, data []byte, whole bool) DeadLetter { + d := DeadLetter{ID: id, Stream: h.Get(DeadHeaderStream), Consumer: h.Get(DeadHeaderConsumer), + Subject: h.Get(DeadHeaderSubject), Lost: h.Get(DeadHeaderLost), Size: len(data)} + if d.Stream == "" || d.Consumer == "" { + d.Stream, d.Consumer, _ = broker.DeadLetterOf(subject) + } + d.Who = ConsumerInWords(d.Stream, d.Consumer) + d.Sequence, _ = strconv.ParseUint(h.Get(DeadHeaderSequence), 10, 64) + d.Deliveries, _ = strconv.ParseUint(h.Get(DeadHeaderDeliveries), 10, 64) + d.GaveUp, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderGaveUp)) + d.Published, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderPublished)) + if whole { + d.Body = string(data) + d.Headers = map[string]string{} + for k := range h { + if !strings.HasPrefix(k, "Mesh-Dead-") && !strings.HasPrefix(k, "Nats-") { + d.Headers[k] = h.Get(k) + } + } + } + return d +} + +// HeldDeadLetters is how many messages DEAD_LETTERS holds for each consumer, by `.`. +func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) { + info, err := js.StreamInfo(broker.DeadLettersStream, &nats.StreamInfoRequest{SubjectsFilter: ">"}) + if err != nil { + return nil, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err) + } + held := map[string]int{} + for subject, n := range info.State.Subjects { + if stream, consumer, ok := broker.DeadLetterOf(subject); ok && n > 0 { + held[stream+"."+consumer] += int(n) + } + } + return held, nil +} + +// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them; for one consumer when +// consumer names one (its name, or `.`). +func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLetter, int, error) { + info, err := js.StreamInfo(broker.DeadLettersStream) + if err != nil { + return nil, 0, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err) + } + var out []DeadLetter + total := 0 + for seq := info.State.LastSeq; seq >= info.State.FirstSeq && seq > 0; seq-- { + raw, err := js.GetMsg(broker.DeadLettersStream, seq) + if errors.Is(err, nats.ErrMsgNotFound) { + continue // delivered again or dropped + } + if err != nil { + return out, total, fmt.Errorf("dead letter %d cannot be read: %w", seq, err) + } + d := deadLetterOf(seq, raw.Subject, raw.Header, raw.Data, false) + if consumer != "" && consumer != d.Consumer && consumer != d.Stream+"."+d.Consumer { + continue + } + total++ + if most <= 0 || len(out) < most { + out = append(out, d) + } + } + return out, total, nil +} + +// ErrNoDeadLetter is a dead letter's id DEAD_LETTERS does not hold. +var ErrNoDeadLetter = errors.New("no such dead letter") + +// DeadLetterNamed is one kept message whole: what it said, and the headers it said it with. +func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) { + raw, err := js.GetMsg(broker.DeadLettersStream, id) + if errors.Is(err, nats.ErrMsgNotFound) { + return DeadLetter{}, fmt.Errorf("%w: %s holds no message %d — delivered again or dropped already, "+ + "or never kept; dead-letters lists what it holds", ErrNoDeadLetter, broker.DeadLettersStream, id) + } + if err != nil { + return DeadLetter{}, fmt.Errorf("dead letter %d cannot be read: %w", id, err) + } + return deadLetterOf(id, raw.Subject, raw.Header, raw.Data, true), nil +} + +// AgainTo is where a kept message is delivered again so that only the consumer that gave it up gets +// it: an event under that consumer's own again subject on EVENTS; an ask back onto its seat's queue, +// whose one worker is that consumer. Any other stream's message is refused, with why. +func AgainTo(d DeadLetter) (string, error) { + switch { + case d.Lost != "": + return "", fmt.Errorf("dead letter %d holds no message to deliver: %s. Drop it", d.ID, d.Lost) + case d.Subject == "": + return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+ + "delivered again. Drop it", d.ID) + case d.Stream == broker.EventsStream: + return broker.AgainSubject(d.Consumer, d.Subject), nil + case strings.HasPrefix(d.Stream, "SEAT_") && strings.HasSuffix(d.Consumer, "_worker"): + return d.Subject, nil + } + return "", fmt.Errorf("dead letter %d is from %s, which nothing delivers again to one consumer alone: "+ + "publishing it again would reach every consumer of its subject. Drop it, and have its sender say it again", + d.ID, d.Stream) +} + +// DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it +// from DEAD_LETTERS. The message carries its own headers and AgainHeader; its de-duplication id is the +// kept copy's, so asking twice delivers it once. +func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, error) { + d, err := DeadLetterNamed(js, id) + if err != nil { + return d, "", err + } + to, err := AgainTo(d) + if err != nil { + return d, "", err + } + again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)} + for k, v := range d.Headers { + again.Header.Set(k, v) + } + again.Header.Set(AgainHeader, strconv.FormatUint(id, 10)) + again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10)) + if _, err := js.PublishMsg(again); err != nil { + return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err) + } + if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) { + return d, to, fmt.Errorf("dead letter %d was delivered again on %s and could not be removed from %s: %w", + id, to, broker.DeadLettersStream, err) + } + return d, to, nil +} + +// DropDeadLetter removes a kept message for good. +func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) { + d, err := DeadLetterNamed(js, id) + if err != nil { + return d, err + } + if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil { + return d, fmt.Errorf("dead letter %d could not be dropped: %w", id, err) + } + return d, nil +} + +// noticeRetry is how long a notice whose message could not be kept waits before it is taken again. +var noticeRetry = time.Minute + +// keepGivenUp takes the server's maximum-deliveries notices off their stream and keeps the message +// each is about, for as long as the subscription stands. A notice that could not be kept is offered +// again after noticeRetry, and said: in the log, and as the consumer's max-deliveries condition. +func keepGivenUp(js nats.JetStreamContext, logger *log.Logger) (*nats.Subscription, error) { + return js.Subscribe("", func(m *nats.Msg) { + d, err := KeepDeadLetter(js, m.Data) + if err != nil { + logger.Printf("a message a consumer gave up on could NOT be kept: %v", err) + var a maxDeliveries + if json.Unmarshal(m.Data, &a) == nil && a.Stream != "" && a.Consumer != "" { + Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: a.Stream + "." + a.Consumer, + Stream: a.Stream, Consumer: a.Consumer, Token: AdvisoryNotKept, + Said: fmt.Sprintf("%s gave up on message %d, and it could not be kept: %v", + ConsumerInWords(a.Stream, a.Consumer), a.StreamSeq, err)}, time.Now()) + } + _ = m.NakWithDelay(noticeRetry) + return + } + if d.Lost != "" { + logger.Printf("%s gave up on message %d of %s, which its stream no longer held; kept as dead letter %d "+ + "without it", d.Who, d.Sequence, d.Stream, d.ID) + } else { + logger.Printf("%s gave up on message %d of %s (%s) after %d deliveries; kept as dead letter %d", + d.Who, d.Sequence, d.Stream, d.Subject, d.Deliveries, d.ID) + } + _ = m.Ack() + }, nats.Bind(broker.DeadLetterNoticesStream, broker.ControllerName), nats.ManualAck()) +} diff --git a/internal/link/deadletters_test.go b/internal/link/deadletters_test.go new file mode 100644 index 00000000..ea013d09 --- /dev/null +++ b/internal/link/deadletters_test.go @@ -0,0 +1,264 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "io" + "log" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// A module's consumer as the controller makes it, on a bus whose streams the controller asserted. +func aModuleConsumer(t *testing.T, js *broker.JetStream, node, module string, consumes ...string) jetstream.Consumer { + t.Helper() + c, ok := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: node, Module: module, + Consumes: consumes, PasswordHash: "x"}) + if !ok { + t.Fatal("a module that consumes got no consumer") + } + if err := js.EnsureConsumer(c); err != nil { + t.Fatal(err) + } + api, err := jetstream.New(js.Conn()) + if err != nil { + t.Fatal(err) + } + reading, err := api.Consumer(t.Context(), broker.EventsStream, c.Name) + if err != nil { + t.Fatal(err) + } + return reading +} + +// next is the one message a consumer hands over within a second, or nil. +func next(t *testing.T, c jetstream.Consumer) jetstream.Msg { + t.Helper() + batch, err := c.Fetch(1, jetstream.FetchMaxWait(time.Second)) + if err != nil { + t.Fatal(err) + } + for msg := range batch.Messages() { + return msg + } + return nil +} + +// keeping takes the notices as the serving controller does, until the test ends. +func keeping(t *testing.T, js *broker.JetStream) { + t.Helper() + if err := broker.AssertMeshConsumers(js); err != nil { + t.Fatal(err) + } + sub, err := keepGivenUp(js.Context(), log.New(io.Discard, "", 0)) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sub.Unsubscribe() }) +} + +// givenUp hands one event to a consumer as often as it may, failing each time, and waits for it kept. +func givenUp(t *testing.T, js *broker.JetStream, c jetstream.Consumer) DeadLetter { + t.Helper() + for { + msg := next(t, c) + if msg == nil { + break + } + _ = msg.Nak() + } + var kept []DeadLetter + eventually(t, "the event kept", func() bool { + kept, _, _ = DeadLetters(js.Context(), c.CachedInfo().Name, 0) + return len(kept) == 1 + }) + return kept[0] +} + +// **Delivered again to the consumer that gave it up, and nobody else**: another module consuming the +// same event handled it the first time and must not handle it twice. +func TestADeadLetterIsDeliveredAgainToItsConsumerAlone(t *testing.T) { + js := aBus(t) + keeping(t, js) + failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged") + handling := aModuleConsumer(t, js, "media", "radarr", "gitea.pull.merged") + + const subject = "mesh.mod.gitea.event.pull.merged" + msg := &nats.Msg{Subject: subject, Data: []byte(`{"n":1}`), Header: nats.Header{}} + msg.Header.Set("x-node", "forge") + msg.Header.Set(nats.MsgIdHdr, "event-1") + if _, err := js.Context().PublishMsg(msg); err != nil { + t.Fatal(err) + } + if m := next(t, handling); m == nil { + t.Fatal("the other module was not handed the event") + } else { + _ = m.Ack() + } + d := givenUp(t, js, failing) + if d.Consumer != "media_sonarr" || d.Stream != "EVENTS" || d.Subject != subject || d.Deliveries != 5 || + d.GaveUp.IsZero() || d.Published.IsZero() || d.Lost != "" { + t.Fatalf("kept as %+v", d) + } + held, err := HeldDeadLetters(js.Context()) + if err != nil || held["EVENTS.media_sonarr"] != 1 { + t.Fatalf("held %v (%v)", held, err) + } + + whole, err := DeadLetterNamed(js.Context(), d.ID) + if err != nil || whole.Body != `{"n":1}` || whole.Headers["x-node"] != "forge" { + t.Fatalf("one asked whole is %+v (%v)", whole, err) + } + + _, to, err := DeliverAgain(js.Context(), d.ID) + if err != nil { + t.Fatal(err) + } + if to != "mesh.again.media_sonarr.mod.gitea.event.pull.merged" { + t.Fatalf("delivered again on %s", to) + } + again := next(t, failing) + if again == nil { + t.Fatal("the consumer that gave it up was not handed it again") + } + if string(again.Data()) != `{"n":1}` || again.Headers().Get("x-node") != "forge" || + again.Headers().Get(AgainHeader) == "" { + t.Fatalf("handed again as %s %v", again.Data(), again.Headers()) + } + if original, ok := broker.OriginalOfAgain(again.Subject()); !ok || original != subject { + t.Fatalf("delivered again on %s, which does not say the event's own subject", again.Subject()) + } + _ = again.Ack() + if m := next(t, handling); m != nil { + t.Fatalf("the module that handled it was handed it twice: %s", m.Subject()) + } + if _, err := DeadLetterNamed(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) { + t.Fatalf("still kept after it was delivered again: %v", err) + } + held, _ = HeldDeadLetters(js.Context()) + if held["EVENTS.media_sonarr"] != 0 { + t.Fatalf("still held: %v", held) + } + // Asked twice, delivered once: the second finds nothing kept. + if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) { + t.Fatalf("a second delivery answered %v", err) + } +} + +func TestADeadLetterDroppedIsGone(t *testing.T) { + js := aBus(t) + keeping(t, js) + failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged") + if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil { + t.Fatal(err) + } + d := givenUp(t, js, failing) + if _, err := DropDeadLetter(js.Context(), d.ID); err != nil { + t.Fatal(err) + } + if left, total, err := DeadLetters(js.Context(), "", 0); err != nil || total != 0 || len(left) != 0 { + t.Fatalf("after the drop: %v %d %v", left, total, err) + } + if m := next(t, failing); m != nil { + t.Fatalf("a dropped event was handed over: %s", m.Subject()) + } +} + +// A notice whose message its stream no longer holds is kept as a record that says so — said, and +// dropped by a person, never silently missing — and cannot be delivered again. +func TestANoticeWhoseMessageIsGoneIsKeptAndSaysSo(t *testing.T) { + js := aBus(t) + d, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+ + `"deliveries":5,"timestamp":"2026-10-06T15:51:00Z"}`)) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(d.Lost, "no longer held message 42") || d.Consumer != "media_sonarr" { + t.Fatalf("kept as %+v", d) + } + if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "Drop it") { + t.Fatalf("a record without its message was delivered again: %v", err) + } + // The same notice taken twice is kept once. + if _, err := KeepDeadLetter(js.Context(), []byte(`{"stream":"EVENTS","consumer":"media_sonarr","stream_seq":42,`+ + `"deliveries":5}`)); err != nil { + t.Fatal(err) + } + if _, total, _ := DeadLetters(js.Context(), "media_sonarr", 0); total != 1 { + t.Fatalf("one notice taken twice is kept %d times", total) + } +} + +// What only the giving-up consumer filters is its own: a subject delivered again reads back as the +// event's, and a module's runtime reads the same key from it (the tokens around `.event.`). +func TestASubjectDeliveredAgainSaysTheEventsOwn(t *testing.T) { + for _, original := range []string{"mesh.mod.gitea.event.pull.merged", "mesh.seat.node-build-agent.event.built"} { + again := broker.AgainSubject("ace_sonarr", original) + if back, ok := broker.OriginalOfAgain(again); !ok || back != original { + t.Errorf("%s reads back as %s", again, back) + } + key := func(subject string) string { + before, event, _ := strings.Cut(subject, ".event.") + parts := strings.Split(before, ".") + return parts[len(parts)-1] + "." + event + } + if key(again) != key(original) { + t.Errorf("%s reads as key %s, the event as %s", again, key(again), key(original)) + } + } + if _, ok := broker.OriginalOfAgain("mesh.mod.gitea.event.pull.merged"); ok { + t.Error("an event's own subject reads as delivered again") + } +} + +// A message from a stream no consumer alone can be given again on is refused, with why. +func TestOnlyAnEventOrAnAskIsDeliveredAgain(t *testing.T) { + for _, c := range []struct { + d DeadLetter + want string + }{ + {DeadLetter{ID: 1, Stream: "EVENTS", Consumer: "a_b", Subject: "mesh.mod.x.event.y"}, "mesh.again.a_b.mod.x.event.y"}, + {DeadLetter{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", + Subject: "mesh.seat.telegram-sender.accept.send"}, "mesh.seat.telegram-sender.accept.send"}, + } { + if to, err := AgainTo(c.d); err != nil || to != c.want { + t.Errorf("%+v: %s %v", c.d, to, err) + } + } + if _, err := AgainTo(DeadLetter{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"}); err == nil { + t.Error("a bucket's message was given a subject to be delivered again on") + } +} + +// An event the controller itself gave up on, delivered again by a person, is acted on as the event it +// was: it arrives under the controller's again subject, which its consumer filters. +func TestTheControllerActsOnAnEventDeliveredAgain(t *testing.T) { + js := aBus(t) + told := &toldAbout{} + s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()} + if err := s.Follows(told); err != nil { + t.Fatal(err) + } + if err := s.Answers(replaysWith{}); err != nil { + t.Fatal(err) + } + ctx, stop := context.WithCancel(context.Background()) + defer stop() + go func() { _ = s.Serve(ctx) }() + eventually(t, "the controller's event consumer being made", func() bool { + _, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName) + return err == nil + }) + moved, _ := json.Marshal(Upgraded{Module: "gitea", Commit: "abcdef0123"}) + if _, err := js.Context().Publish(broker.AgainSubject(broker.ControllerName, broker.ControllerFollows[0]), moved); err != nil { + t.Fatal(err) + } + eventually(t, "the upgrade delivered again reaching the controller", func() bool { return told.count() == 1 }) +} diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index f3d9d5e0..14f6eb16 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -94,6 +94,14 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con holding.Store(true) defer holding.Store(false) + // What a consumer gave up on, kept by the controller acting now (novox/hq issue 330): the server's + // notices wait in their stream for it, so one said while no controller listened is not lost. + kept, err := keepGivenUp(js, log.Default()) + if err != nil { + return fmt.Errorf("taking the notices of messages consumers gave up on: %w", err) + } + defer func() { _ = kept.Unsubscribe() }() + // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // they are their own guarantee: a lost one is the next one. beats := make(chan *nats.Msg, Prefetch) @@ -167,6 +175,11 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control), msg *nats.Msg, streamed bool) { + // An event the controller gave up on and a person delivered again (novox/hq issue 330) arrives under + // the controller's own again subject, and is the event it was: acted on as first published. + if original, ok := broker.OriginalOfAgain(msg.Subject); ok { + msg.Subject = original + } kind, known := kindOfSubject(msg.Subject) if !known { if streamed { diff --git a/internal/link/replay330_test.go b/internal/link/replay330_test.go new file mode 100644 index 00000000..ac25b362 --- /dev/null +++ b/internal/link/replay330_test.go @@ -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") + } +} diff --git a/module.json b/module.json index d5390683..ca6f2374 100644 --- a/module.json +++ b/module.json @@ -71,6 +71,7 @@ "bus", "retire", "cleanup", + "dead-letters", "data", "build", "artifacts",