Files
mesh-controller/internal/link/deadletters_ask_test.go
T
jschoubben ef551fdfb6
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
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.
2026-10-11 03:24:52 +02:00

218 lines
8.1 KiB
Go

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