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