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) + } + } +}