From f510b46319250e59fa7e09f4702ec5e97e809904 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 11 Oct 2026 03:18:34 +0200 Subject: [PATCH] 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) + } + } +}