Compare commits

...
Author SHA1 Message Date
mesh-admin 0a2e58c070 Merge pull request 'Deliver again an ask the controller made, removing the original from its queue (issue 334, part)' (#219) from fix/334-a-given-up-ask-can-be-delivered-again into main 2026-10-11 02:36:26 +00:00
mesh-admin 45b4ae92bc Merge pull request 'A provider whose wait fails its check lists who waits on it (issue 450)' (#221) from fix/450-a-provider-whose-wait-fails-lists-its-waiters into main 2026-10-11 02:19:02 +00:00
jschoubben ed45cc6415 Check a provider's waits before saying who waits on it (issue 450)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
A provider whose wait fails the check raises unhealthy, but the holding read
its stored statement with the wait unchecked: sayWaiters built the
needs-operator key, found it not open and listed no held consumer. The
holding now reads each statement as judged (ADR 0283 decision 3), in
sayWaiters and in the provider's state its consumers are held by.
2026-10-11 03:29:10 +02:00
jschoubben ef551fdfb6 Remove an ask's original only when its sequence still holds it, and only with a worker (issue 334)
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
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
jschoubben f510b46319 Deliver again an ask the controller made, removing the original from its queue (issue 334)
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 superseded: a newer head of the same pull request
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.
2026-10-11 03:18:34 +02:00
8 changed files with 409 additions and 14 deletions
+5 -1
View File
@@ -144,11 +144,15 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64,
}
switch act {
case "deliver":
_, to, err := link.DeliverAgain(on.js, id)
delivered, to, err := link.DeliverAgain(on.js, id)
if err != nil {
return nil, err
}
answer["delivered_on"] = to
if delivered.Original != "" {
// What became of the ask in its seat's queue (novox/hq issue 334): removed, or left, and why.
answer["original"] = delivered.Original
}
answer["done"] = fmt.Sprintf("dead letter %d was delivered again to %s, and nobody else; it is no longer kept",
id, consumerWho(d.Stream, d.Consumer))
case "drop":
@@ -0,0 +1,66 @@
package main
import (
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A provider whose wait fails the controller's check lists who waits on it (novox/hq issue 450).
//
// A wait that does not check out is judged unhealthy (ADR 0283 decision 3), so the provider raises its unhealthy
// condition and its consumers are held under it (ADR 0240 rule 5). What is stored is what the machine said, the
// wait unchecked: read as said, the provider seems to wait for the operator, and the held consumers were never
// listed on the condition that is open.
// refusedWait is a wait for a secret the database's manifest does not declare: it fails the check.
var refusedWait = inventory.Wait{Part: "postgres-database", Secret: "certificate", What: "the database's certificate"}
// judgedRefused is the provider's resource as the controller judges it: unhealthy, saying why.
func judgedRefused() inventory.ResourceHealth {
r := waitingDatabaseFor(refusedWait)
r.State, r.Waits = link.StateUnhealthy, nil
r.Reason = "says it waits for the secret certificate, which db does not declare"
return r
}
// The consumer's statement arrives after the provider's, so it is the consumer's judging (sayWaiters) that must list
// it at the provider's unhealthy condition, made urgent by who waits on it.
func TestAProviderWhoseWaitFailedItsCheckListsWhoWaitsOnIt(t *testing.T) {
k, _ := withConditionsInMemory(t)
stored := waitingDatabaseFor(refusedWait)
open := judgeBoth(t, k, []inventory.ResourceHealth{judgedRefused()},
map[string]inventory.NodeHealth{"anchor": {Node: "anchor", Resources: []inventory.ResourceHealth{stored}}},
shopFailingBeside(stored))
provider := conditionOf(open, moduleUnhealthyKey("db", "anchor"))
if len(open) != 1 || provider == nil {
t.Fatalf("a provider whose wait failed its check and a consumer of it raised %v; want the provider's "+
"unhealthy alone", openKeysOf(open))
}
if said := provider.Evidence[0].Said; !strings.Contains(said, "shop on laptop") {
t.Fatalf("the provider's unhealthy condition does not list shop on laptop as waiting on it: %s", said)
}
if provider.Severity != conditions.Urgent {
t.Fatalf("the provider's unhealthy condition is %s with a consumer waiting on it; want urgent", provider.Severity)
}
}
// A provider that waits for a secret already given is unhealthy to its consumers too, before its own condition opens:
// they are held under it as under any unhealthy provider, never as waiting for the operator, and its condition, when
// open, is said as unhealthy with who waits on it.
func TestAProviderWaitingForASecretAlreadyGivenIsUnhealthyToItsConsumers(t *testing.T) {
given := holdingOf(nil, shopFailingBeside(waitingDatabase()), shopOnTheDatabase)
given.waits.given["db@anchor"] = map[string]time.Time{"licence": time.Now().Add(-time.Hour)}
by, held := given.heldWith("laptop", "shop", []inventory.ResourceHealth{failingConsumer("shop")})
if !held || by.provider != theDatabase || len(by.waits) > 0 {
t.Fatalf("shop under a database waiting for a licence already given: held %v under %v with waits %v; want "+
"held under db on anchor as unhealthy, no waits", held, by.provider, by.waits)
}
if _, uncovered := given.waitingUncovered("laptop", "shop", []inventory.ResourceHealth{failingConsumer("shop")}); uncovered {
t.Fatalf("a provider whose wait failed its check is said as waiting for another part")
}
}
+3 -1
View File
@@ -326,8 +326,10 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
// sayWaiters observes a provider's open condition again, with who waits on it, from its machine's newest
// statement. Nothing when its condition is not open: it is raised by its own statements, on its own looks.
func sayWaiters(ctx context.Context, k *conditions.Keeper, hold *holding, p catalogue.Chosen, now time.Time) error {
// Its statement as it was judged, its waits checked (novox/hq issue 450): a wait that failed the check is
// unhealthy, so the condition built here is the one its own statement raised, and who waits on it is listed.
var rs []inventory.ResourceHealth
for _, r := range hold.healths[p.Node].Resources {
for _, r := range hold.checked(p.Node) {
if r.Module == p.Module && (r.State == link.StateUnhealthy || r.State == link.StateWaiting) {
rs = append(rs, r)
}
+27 -1
View File
@@ -37,6 +37,32 @@ type holding struct {
shelf map[string]catalogue.Manifest
// providers memoises providerFor by machine, consumer and provision.
providers map[string]providerLookup
// waits is what the waits of a statement are checked against, read as they are asked for (novox/hq issue 450).
waits operatorWaitFacts
}
// checked is a machine's newest statement as the controller judges it: each wait checked (ADR 0283 decision 3), so
// a wait that does not check out reads as unhealthy, as it did when the statement was judged (novox/hq issue 450).
// What is stored is what the machine said, the waits unchecked; read as stored, a provider whose wait failed its
// check seems to wait for the operator while its unhealthy condition is open.
func (h *holding) checked(machine string) []inventory.ResourceHealth {
rs := h.healths[machine].Resources
mods := waitingModules(rs)
if len(mods) == 0 {
return rs
}
if h.waits.manifests == nil {
h.waits.manifests = map[string]catalogue.Manifest{}
}
for _, module := range mods {
if _, has := h.waits.manifests[module]; !has {
if m, ok := h.manifestOf(module); ok {
h.waits.manifests[module] = m
}
}
}
readWaitFacts(h.ctx, h.inv, machine, mods, nil, &h.waits)
return checkWaiting(machine, rs, h.waits)
}
type providerLookup struct {
@@ -125,7 +151,7 @@ func (h *holding) providerState(p catalogue.Chosen) (unhealthy, waiting bool, wa
return true, false, nil
}
}
for _, r := range h.healths[p.Node].Resources {
for _, r := range h.checked(p.Node) {
if r.Module != p.Module {
continue
}
+13 -4
View File
@@ -41,12 +41,21 @@ func failingConsumer(module string) inventory.ResourceHealth {
State: link.StateUnhealthy, Reason: "http /health on web: answered 500", Check: "http", Needs: "postgres-database"}
}
// dbSecrets is what the database's waits name: own secrets issued outside the mesh, so a wait for one checks out
// while nobody gave it (ADR 0283 decision 3, novox/hq issue 450).
var dbSecrets = catalogue.OwnSecrets{
"licence": {Path: "/s/licence", IssuedBy: catalogue.IssuedOutside},
"backup-key": {Path: "/s/backup-key", IssuedBy: catalogue.IssuedOutside},
}
// holdingOf is one reading of the record without a store: the newest statements, the open conditions, the
// catalogue, and each consumer's provider already looked up, as providerFor memoises it.
// catalogue, what was given on the provider's machine (nothing), and each consumer's provider already looked up,
// as providerFor memoises it.
func holdingOf(open []conditions.Condition, healths map[string]inventory.NodeHealth, bound map[[3]string]catalogue.Chosen) *holding {
h := &holding{ctx: context.Background(), healths: healths, open: open, providers: map[string]providerLookup{},
shelf: map[string]catalogue.Manifest{"db": {Module: "db", Version: "1",
Provides: []catalogue.Offer{{Name: "postgres-database", Scope: catalogue.ScopeMesh}}}}}
shelf: map[string]catalogue.Manifest{"db": {Module: "db", Version: "1", OwnSecrets: dbSecrets,
Provides: []catalogue.Offer{{Name: "postgres-database", Scope: catalogue.ScopeMesh}}}},
waits: operatorWaitFacts{given: map[string]map[string]time.Time{"db@anchor": {}}}}
for k, p := range bound {
h.providers[k[0]+"\x00"+k[1]+"\x00"+k[2]] = providerLookup{p, true}
}
@@ -75,7 +84,7 @@ func TestAProviderWaitingForTheOperatorHoldsOnlyTheConsumersOfTheWaitingPart(t *
healthyDB.State, healthyDB.Waits = link.StateHealthy, nil
byCredential := holdingOf(nil, shopFailingBeside(waitingDatabaseFor(inventory.Wait{Part: "the server",
Secret: "licence", What: "the licence key"})), shopOnTheDatabase)
byCredential.shelf["db"] = catalogue.Manifest{Module: "db", Version: "1", Provides: []catalogue.Offer{{
byCredential.shelf["db"] = catalogue.Manifest{Module: "db", Version: "1", OwnSecrets: dbSecrets, Provides: []catalogue.Offer{{
Name: "postgres-database", Credential: &catalogue.OfferCredential{Own: "licence"}}}}
for _, c := range []struct {
name string
+15
View File
@@ -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 }
+63 -7
View File
@@ -62,6 +62,9 @@ type DeadLetter struct {
// taken. The record of it is kept all the same, so it is said and dropped, never silently missing.
Lost string `json:"lost,omitempty"`
Size int `json:"size"`
// Original says what became of the ask it was kept from, in its seat's work queue, when it was
// delivered again (novox/hq issue 334).
Original string `json:"original,omitempty"`
// Body and Headers are the message itself, given only for one dead letter asked by its id; a header
// with several values keeps them all.
Body string `json:"body,omitempty"`
@@ -250,9 +253,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); only while the
// worker that gave it up is on the bus. 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 +269,22 @@ 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):
// The seat's worker takes it, and only while it is on the bus: without one the queue would keep
// the ask for a holder that may never come, and it would be let go as delivered meanwhile.
if _, err := js.ConsumerInfo(d.Stream, d.Consumer); errors.Is(err, nats.ErrConsumerNotFound) {
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so the ask "+
"would only wait in %s for a holder. Nothing was done, and it is still kept: deliver it again once "+
"the seat has a holder, or drop it", d.ID, d.Who, d.Stream)
} else if err != nil {
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
}
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 +305,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 +316,9 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
if err != nil {
return d, "", err
}
if d.Stream != broker.EventsStream {
d.Original = removeOriginal(js, d)
}
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...)
@@ -302,6 +326,10 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
if _, err := js.PublishMsg(again); err != nil {
if d.Original != "" {
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept, so "+
"delivering it again tries once more. Its original: %s", id, to, err, d.Original)
}
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err)
}
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
@@ -311,6 +339,34 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
return d, to, nil
}
// removeOriginal takes the ask a dead letter was kept from out of its seat's work queue, and says what
// became of it. Given up on, the original is never acknowledged and would stay beside its copy until the
// stream's age drops it (novox/hq issue 334). It is removed before the copy is published, so a copy that
// cannot be published leaves the kept one to try again, and never two.
//
// **Only when the queue still holds that same message**: its subject and the time it was stored are the
// dead letter's. A seat's stream deleted and made again — the build handover deletes one (builds.go) —
// numbers from one again, and an old dead letter's sequence may then name another, live ask; deleting by
// the number alone would drop that one silently. Anything else is said, never an error: the original is
// gone or is not this one, and the copy is the only one there will be.
func removeOriginal(js nats.JetStreamContext, d DeadLetter) string {
held, err := js.GetMsg(d.Stream, d.Sequence)
switch {
case errors.Is(err, nats.ErrMsgNotFound):
return fmt.Sprintf("%s no longer held message %d, so there was nothing to remove", d.Stream, d.Sequence)
case err != nil:
return fmt.Sprintf("message %d of %s could not be read, so it was left as it is: %v", d.Sequence, d.Stream, err)
case d.Published.IsZero() || held.Subject != d.Subject || !held.Time.Equal(d.Published):
return fmt.Sprintf("message %d of %s is another message now (%s, stored %s), so it was left as it is; "+
"the one given up on is gone", d.Sequence, d.Stream, held.Subject, held.Time.UTC().Format(time.RFC3339))
}
if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
return fmt.Sprintf("message %d of %s could not be removed, so it stays beside its copy until the "+
"stream's age drops it; nothing delivers it again: %v", d.Sequence, d.Stream, err)
}
return fmt.Sprintf("message %d of %s, the one given up on, was removed from the queue", d.Sequence, d.Stream)
}
// DropDeadLetter removes a kept message for good.
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
d, err := DeadLetterNamed(js, id)
+217
View File
@@ -0,0 +1,217 @@
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)
}
}