From f510b46319250e59fa7e09f4702ec5e97e809904 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 03:18:34 +0200 Subject: [PATCH 1/2] Deliver again an ask the controller made, removing the original from its queue (issue 334) A build ask its seat's worker gave up on could only be dropped, though the controller already holds the publish on that seat's accepts and the stream API to remove the original. Asks to other seats stay refused: delivering them needs a grant ADR 0264 withholds, in the asker's name ADR 0259 protects. --- internal/broker/nats.go | 15 ++++ internal/link/deadletters.go | 33 +++++++-- internal/link/deadletters_ask_test.go | 102 ++++++++++++++++++++++++++ 3 files changed, 143 insertions(+), 7 deletions(-) create mode 100644 internal/link/deadletters_ask_test.go diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 09966256..d46caf7f 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -163,6 +163,21 @@ type Principal struct { // goes with the retired seat row. var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"} +// TheControllersAsk says whether a message of stream, published on subject, is an ask the controller +// itself makes: one on the accept subject of a seat in seatsTheControllerAsks, kept in that seat's own +// work queue. Such an ask, given up on by the seat's worker, can be delivered again with the authority +// the controller already holds — the publish on that seat's accepts and the stream API to remove the +// original — and no other can (novox/hq issue 334, ADR 0264's consequences: no grant over the seats' +// queues). +func TheControllersAsk(stream, subject string) bool { + for _, seat := range seatsTheControllerAsks { + if stream == seatStreamName(seat) && strings.HasPrefix(subject, "mesh.seat."+seat+".accept.") { + return true + } + } + return false +} + // SeatVerb is one verb of one seat, on every machine holding it. type SeatVerb struct{ Seat, Verb string } diff --git a/internal/link/deadletters.go b/internal/link/deadletters.go index d6471b66..2b6c3ec5 100644 --- a/internal/link/deadletters.go +++ b/internal/link/deadletters.go @@ -250,9 +250,15 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) { // 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. +// filters that subject, so a message is never let go as delivered while nobody receives it. +// +// An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's +// work queue has one worker per subject, so only the worker that gave it up takes it, and DeliverAgain +// removes the original from the queue first, so there are never two (novox/hq issue 334). The seat's +// stream keeps it until a holder pulls, as it keeps any ask while the seat has none. Any other stream's +// message is refused, with why: an ask to a seat the controller does not ask would need a publish it is +// not granted, in its asker's name (ADR 0264's consequences, ADR 0259 §3), and a stream that is neither +// would reach every consumer of its subject. func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) { switch { case d.Lost != "": @@ -260,10 +266,13 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) { case d.Subject == "": return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+ "delivered again. Drop it", d.ID) + case broker.TheControllersAsk(d.Stream, d.Subject): + return d.Subject, nil case d.Stream != broker.EventsStream: - return "", fmt.Errorf("dead letter %d is from %s, and only an event 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) + return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+ + "delivered again: publishing it again would reach every consumer of its subject, or need a grant the "+ + "controller does not hold to speak for its asker. Drop it, and have its sender say it again", + d.ID, d.Stream) } info, err := js.ConsumerInfo(d.Stream, d.Consumer) if errors.Is(err, nats.ErrConsumerNotFound) { @@ -284,7 +293,7 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) { } // 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 +// from DEAD_LETTERS; an ask's original is removed from its work queue before. The message carries its own headers and AgainHeader; its de-duplication id is the // kept copy's, so asking twice delivers it once. func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, error) { d, err := DeadLetterNamed(js, id) @@ -295,6 +304,16 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro if err != nil { return d, "", err } + if d.Stream != broker.EventsStream { + // An ask, back in its seat's work queue: the original, given up on, is never acknowledged and + // would stay beside its copy until the stream's age drops it (novox/hq issue 334). Removed first, + // so a copy that then cannot be published leaves the kept one to try again, and never two. One + // the queue no longer holds has aged out, which is no reason to refuse. + if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) { + return d, to, fmt.Errorf("the ask dead letter %d was kept from, message %d of %s, could not be "+ + "removed, so it was not delivered again: %w; it is still kept", id, d.Sequence, d.Stream, err) + } + } again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)} for k, v := range d.Headers { again.Header[k] = append([]string(nil), v...) diff --git a/internal/link/deadletters_ask_test.go b/internal/link/deadletters_ask_test.go new file mode 100644 index 00000000..5fd180e8 --- /dev/null +++ b/internal/link/deadletters_ask_test.go @@ -0,0 +1,102 @@ +package link + +import ( + "errors" + "strings" + "testing" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// What a seat's worker gave up on, when the ask is one the controller itself makes (novox/hq issue 334). +// +// The controller already publishes on the accept subjects of the seats it asks (broker's +// seatsTheControllerAsks) and reaches the stream API, so delivering its own ask again needs no grant it +// does not hold: the original is removed from the work queue by its sequence, and the kept copy is +// published on the ask's own subject, where the seat's one worker takes it. An ask to any other seat is +// still refused, and stays kept: delivering it would need a publish the controller is not granted +// (ADR 0264's consequences), in the asker's name (ADR 0259 §3). + +// theBuildWorker is the build seat's worker as a holder pulls from it. +func theBuildWorker(t *testing.T, js *broker.JetStream) jetstream.Consumer { + t.Helper() + api, err := jetstream.New(js.Conn()) + if err != nil { + t.Fatal(err) + } + stream := "SEAT_NODE_BUILD_AGENT" + worker, err := api.Consumer(t.Context(), stream, stream+"_worker") + if err != nil { + t.Fatal(err) + } + return worker +} + +func TestAnAskTheControllerMadeIsDeliveredAgainToTheSeatsWorker(t *testing.T) { + js := aBusWithTheBuildRole(t) + keeping(t, js) + worker := theBuildWorker(t, js) + + subject := BuildWorkOf(TheBuildMachine) + ask := &nats.Msg{Subject: subject, Data: []byte(`{"module":"x"}`), Header: nats.Header{}} + ask.Header.Set(nats.MsgIdHdr, "build-1") + if _, err := js.Context().PublishMsg(ask); err != nil { + t.Fatal(err) + } + d := givenUp(t, js, worker) + if d.Stream != "SEAT_NODE_BUILD_AGENT" || d.Subject != subject || d.Lost != "" { + t.Fatalf("kept as %+v", d) + } + if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil { + t.Fatalf("the original is not in the work queue before it is delivered again: %v", err) + } + + _, to, err := DeliverAgain(js.Context(), d.ID) + if err != nil { + t.Fatal(err) + } + if to != subject { + t.Fatalf("delivered again on %s, not the ask's own subject %s", to, subject) + } + if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) { + t.Fatalf("the original is still in the work queue beside its copy: %v", err) + } + again := next(t, worker) + if again == nil { + t.Fatal("the seat's worker was not handed the ask again") + } + if string(again.Data()) != `{"module":"x"}` || again.Headers().Get(AgainHeader) == "" { + t.Fatalf("handed again as %s %v", again.Data(), again.Headers()) + } + _ = again.Ack() + if m := next(t, worker); m != nil { + t.Fatalf("the worker 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) + } + if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) { + t.Fatalf("a second delivery answered %v", err) + } +} + +// An ask to a seat the controller does not ask is refused, says why, and nothing is done. +func TestAnAskTheControllerDidNotMakeIsStillOnlyDropped(t *testing.T) { + js := aBus(t) + for _, d := range []DeadLetter{ + {ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", + Subject: "mesh.seat.telegram-sender.accept.send"}, + // The subject of a seat the controller asks, on a stream that is not that seat's queue. + {ID: 3, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", + Subject: "mesh.seat.node-build-agent.accept.build"}, + } { + if to, err := AgainTo(js.Context(), d); err == nil { + t.Errorf("%s on %s was given %s to be delivered again on", d.Subject, d.Stream, to) + } else if !strings.Contains(err.Error(), "Drop it") { + t.Errorf("refused without saying what to do: %v", err) + } + } +} From ef551fdfb6b29d4a9b68ff330bf6b930d3158bd1 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 03:24:52 +0200 Subject: [PATCH 2/2] Remove an ask's original only when its sequence still holds it, and only with a worker (issue 334) A seat's queue made again numbers from one, so an old dead letter's sequence can name a live ask; deleting by number alone would drop it silently. An ask with no worker would wait unseen while counted as delivered. --- cmd/mesh-controller/deadletters.go | 6 +- internal/link/deadletters.go | 57 ++++++++++--- internal/link/deadletters_ask_test.go | 115 ++++++++++++++++++++++++++ 3 files changed, 167 insertions(+), 11 deletions(-) diff --git a/cmd/mesh-controller/deadletters.go b/cmd/mesh-controller/deadletters.go index c01dd2c4..74f22cb0 100644 --- a/cmd/mesh-controller/deadletters.go +++ b/cmd/mesh-controller/deadletters.go @@ -144,11 +144,15 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64, } switch act { case "deliver": - _, to, err := link.DeliverAgain(on.js, id) + delivered, to, err := link.DeliverAgain(on.js, id) if err != nil { return nil, err } answer["delivered_on"] = to + if delivered.Original != "" { + // What became of the ask in its seat's queue (novox/hq issue 334): removed, or left, and why. + answer["original"] = delivered.Original + } 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": diff --git a/internal/link/deadletters.go b/internal/link/deadletters.go index 2b6c3ec5..0b4f41fa 100644 --- a/internal/link/deadletters.go +++ b/internal/link/deadletters.go @@ -62,6 +62,9 @@ type DeadLetter struct { // taken. The record of it is kept all the same, so it is said and dropped, never silently missing. Lost string `json:"lost,omitempty"` Size int `json:"size"` + // Original says what became of the ask it was kept from, in its seat's work queue, when it was + // delivered again (novox/hq issue 334). + Original string `json:"original,omitempty"` // Body and Headers are the message itself, given only for one dead letter asked by its id; a header // with several values keeps them all. Body string `json:"body,omitempty"` @@ -254,8 +257,8 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) { // // An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's // work queue has one worker per subject, so only the worker that gave it up takes it, and DeliverAgain -// removes the original from the queue first, so there are never two (novox/hq issue 334). The seat's -// stream keeps it until a holder pulls, as it keeps any ask while the seat has none. Any other stream's +// removes the original from the queue first, so there are never two (novox/hq issue 334); only while the +// worker that gave it up is on the bus. Any other stream's // message is refused, with why: an ask to a seat the controller does not ask would need a publish it is // not granted, in its asker's name (ADR 0264's consequences, ADR 0259 §3), and a stream that is neither // would reach every consumer of its subject. @@ -267,6 +270,15 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) { return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+ "delivered again. Drop it", d.ID) case broker.TheControllersAsk(d.Stream, d.Subject): + // The seat's worker takes it, and only while it is on the bus: without one the queue would keep + // the ask for a holder that may never come, and it would be let go as delivered meanwhile. + if _, err := js.ConsumerInfo(d.Stream, d.Consumer); errors.Is(err, nats.ErrConsumerNotFound) { + return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so the ask "+ + "would only wait in %s for a holder. Nothing was done, and it is still kept: deliver it again once "+ + "the seat has a holder, or drop it", d.ID, d.Who, d.Stream) + } else if err != nil { + return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err) + } return d.Subject, nil case d.Stream != broker.EventsStream: return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+ @@ -305,14 +317,7 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro return d, "", err } if d.Stream != broker.EventsStream { - // An ask, back in its seat's work queue: the original, given up on, is never acknowledged and - // would stay beside its copy until the stream's age drops it (novox/hq issue 334). Removed first, - // so a copy that then cannot be published leaves the kept one to try again, and never two. One - // the queue no longer holds has aged out, which is no reason to refuse. - if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) { - return d, to, fmt.Errorf("the ask dead letter %d was kept from, message %d of %s, could not be "+ - "removed, so it was not delivered again: %w; it is still kept", id, d.Sequence, d.Stream, err) - } + d.Original = removeOriginal(js, d) } again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)} for k, v := range d.Headers { @@ -321,6 +326,10 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro again.Header.Set(AgainHeader, strconv.FormatUint(id, 10)) again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10)) if _, err := js.PublishMsg(again); err != nil { + if d.Original != "" { + return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept, so "+ + "delivering it again tries once more. Its original: %s", id, to, err, d.Original) + } return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err) } if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) { @@ -330,6 +339,34 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro return d, to, nil } +// removeOriginal takes the ask a dead letter was kept from out of its seat's work queue, and says what +// became of it. Given up on, the original is never acknowledged and would stay beside its copy until the +// stream's age drops it (novox/hq issue 334). It is removed before the copy is published, so a copy that +// cannot be published leaves the kept one to try again, and never two. +// +// **Only when the queue still holds that same message**: its subject and the time it was stored are the +// dead letter's. A seat's stream deleted and made again — the build handover deletes one (builds.go) — +// numbers from one again, and an old dead letter's sequence may then name another, live ask; deleting by +// the number alone would drop that one silently. Anything else is said, never an error: the original is +// gone or is not this one, and the copy is the only one there will be. +func removeOriginal(js nats.JetStreamContext, d DeadLetter) string { + held, err := js.GetMsg(d.Stream, d.Sequence) + switch { + case errors.Is(err, nats.ErrMsgNotFound): + return fmt.Sprintf("%s no longer held message %d, so there was nothing to remove", d.Stream, d.Sequence) + case err != nil: + return fmt.Sprintf("message %d of %s could not be read, so it was left as it is: %v", d.Sequence, d.Stream, err) + case d.Published.IsZero() || held.Subject != d.Subject || !held.Time.Equal(d.Published): + return fmt.Sprintf("message %d of %s is another message now (%s, stored %s), so it was left as it is; "+ + "the one given up on is gone", d.Sequence, d.Stream, held.Subject, held.Time.UTC().Format(time.RFC3339)) + } + if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) { + return fmt.Sprintf("message %d of %s could not be removed, so it stays beside its copy until the "+ + "stream's age drops it; nothing delivers it again: %v", d.Sequence, d.Stream, err) + } + return fmt.Sprintf("message %d of %s, the one given up on, was removed from the queue", d.Sequence, d.Stream) +} + // DropDeadLetter removes a kept message for good. func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) { d, err := DeadLetterNamed(js, id) diff --git a/internal/link/deadletters_ask_test.go b/internal/link/deadletters_ask_test.go index 5fd180e8..2def6f3e 100644 --- a/internal/link/deadletters_ask_test.go +++ b/internal/link/deadletters_ask_test.go @@ -100,3 +100,118 @@ func TestAnAskTheControllerDidNotMakeIsStillOnlyDropped(t *testing.T) { } } } + +// aBuildAskGivenUp publishes one build ask and lets the worker give it up. +func aBuildAskGivenUp(t *testing.T, js *broker.JetStream, worker jetstream.Consumer, body string) DeadLetter { + t.Helper() + if _, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(body)); err != nil { + t.Fatal(err) + } + return givenUp(t, js, worker) +} + +// A copy that cannot be published: the original is already out of the queue, the dead letter is kept and +// the answer says both; asked again once it can be published, it is delivered once. +func TestAnAskWhoseCopyIsRefusedStaysKeptAndIsDeliveredOnTheNextTry(t *testing.T) { + js := aBusWithTheBuildRole(t) + keeping(t, js) + worker := theBuildWorker(t, js) + d := aBuildAskGivenUp(t, js, worker, `{"module":"x"}`) + + // The queue stops taking the ask's subject, so the copy's publish is refused by the server. + info, err := js.Context().StreamInfo(d.Stream) + if err != nil { + t.Fatal(err) + } + cfg := info.Config + taking := cfg.Subjects + cfg.Subjects = []string{"mesh.seat." + TheBuildMachine + ".accept.nothing"} + if _, err := js.Context().UpdateStream(&cfg); err != nil { + t.Fatal(err) + } + _, _, err = DeliverAgain(js.Context(), d.ID) + if err == nil || !strings.Contains(err.Error(), "still kept") || !strings.Contains(err.Error(), "was removed") { + t.Fatalf("a refused copy answered %v", err) + } + if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) { + t.Fatalf("the original is still in the queue: %v", err) + } + if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil { + t.Fatalf("a refused copy let the dead letter go: %v", err) + } + + cfg.Subjects = taking + if _, err := js.Context().UpdateStream(&cfg); err != nil { + t.Fatal(err) + } + retried, _, err := DeliverAgain(js.Context(), d.ID) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(retried.Original, "no longer held") { + t.Fatalf("the retry said of the original: %q", retried.Original) + } + if again := next(t, worker); again == nil || string(again.Data()) != `{"module":"x"}` { + t.Fatal("the retry did not hand the ask to the worker") + } else { + _ = again.Ack() + } + if m := next(t, worker); m != nil { + t.Fatalf("handed twice: %s", m.Data()) + } +} + +// A queue made again numbers from one: the old dead letter's sequence then names a live ask, which is left +// alone. +func TestAnAskWhoseSequenceNamesAnotherMessageLeavesThatOneAlone(t *testing.T) { + js := aBusWithTheBuildRole(t) + keeping(t, js) + d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"old"}`) + + // The build handover's way: the seat's queue deleted and made again, with its worker. + if err := js.Context().DeleteStream(d.Stream); err != nil { + t.Fatal(err) + } + if err := broker.RaiseSeats(js, []broker.DeclaredSeat{{Name: TheBuildMachine, Accepts: []string{"build"}, + Emits: []string{"built"}}}, map[string]broker.Holder{TheBuildMachine: {Node: "anchor", Module: "builder"}}); err != nil { + t.Fatal(err) + } + live, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(`{"module":"live"}`)) + if err != nil { + t.Fatal(err) + } + if live.Sequence != d.Sequence { + t.Fatalf("the live ask is message %d, the dead letter names %d: the test does not set up the collision", + live.Sequence, d.Sequence) + } + + delivered, _, err := DeliverAgain(js.Context(), d.ID) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(delivered.Original, "another message") { + t.Fatalf("the answer said of the original: %q", delivered.Original) + } + if held, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil || string(held.Data) != `{"module":"live"}` { + t.Fatalf("the live ask at that sequence was touched: %v", err) + } +} + +// No worker on the seat's queue: refused, kept, and nothing published. +func TestAnAskIsNotDeliveredAgainWhileTheSeatHasNoWorker(t *testing.T) { + js := aBusWithTheBuildRole(t) + keeping(t, js) + d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"x"}`) + if err := js.Context().DeleteConsumer(d.Stream, d.Consumer); err != nil { + t.Fatal(err) + } + if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "wait in") { + t.Fatalf("delivered with no worker: %v", err) + } + if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil { + t.Fatalf("a refusal removed the original: %v", err) + } + if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil { + t.Fatalf("a refusal let the dead letter go: %v", err) + } +}