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.
218 lines
8.1 KiB
Go
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)
|
|
}
|
|
}
|