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.
This commit is contained in:
@@ -144,11 +144,15 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64,
|
|||||||
}
|
}
|
||||||
switch act {
|
switch act {
|
||||||
case "deliver":
|
case "deliver":
|
||||||
_, to, err := link.DeliverAgain(on.js, id)
|
delivered, to, err := link.DeliverAgain(on.js, id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
answer["delivered_on"] = to
|
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",
|
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))
|
id, consumerWho(d.Stream, d.Consumer))
|
||||||
case "drop":
|
case "drop":
|
||||||
|
|||||||
@@ -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.
|
// taken. The record of it is kept all the same, so it is said and dropped, never silently missing.
|
||||||
Lost string `json:"lost,omitempty"`
|
Lost string `json:"lost,omitempty"`
|
||||||
Size int `json:"size"`
|
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
|
// 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.
|
// with several values keeps them all.
|
||||||
Body string `json:"body,omitempty"`
|
Body string `json:"body,omitempty"`
|
||||||
@@ -254,8 +257,8 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
|||||||
//
|
//
|
||||||
// An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's
|
// 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
|
// 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
|
// removes the original from the queue first, so there are never two (novox/hq issue 334); only while the
|
||||||
// stream keeps it until a holder pulls, as it keeps any ask while the seat has none. Any other stream's
|
// 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
|
// 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
|
// 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.
|
// would reach every consumer of its subject.
|
||||||
@@ -267,6 +270,15 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
|
|||||||
return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+
|
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)
|
"delivered again. Drop it", d.ID)
|
||||||
case broker.TheControllersAsk(d.Stream, d.Subject):
|
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
|
return d.Subject, nil
|
||||||
case d.Stream != broker.EventsStream:
|
case d.Stream != broker.EventsStream:
|
||||||
return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+
|
return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+
|
||||||
@@ -305,14 +317,7 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
|||||||
return d, "", err
|
return d, "", err
|
||||||
}
|
}
|
||||||
if d.Stream != broker.EventsStream {
|
if d.Stream != broker.EventsStream {
|
||||||
// An ask, back in its seat's work queue: the original, given up on, is never acknowledged and
|
d.Original = removeOriginal(js, d)
|
||||||
// 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)}
|
again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)}
|
||||||
for k, v := range d.Headers {
|
for k, v := range d.Headers {
|
||||||
@@ -321,6 +326,10 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
|||||||
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
|
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
|
||||||
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
|
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
|
||||||
if _, err := js.PublishMsg(again); err != nil {
|
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)
|
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) {
|
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
|
||||||
@@ -330,6 +339,34 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
|||||||
return d, to, nil
|
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.
|
// DropDeadLetter removes a kept message for good.
|
||||||
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
||||||
d, err := DeadLetterNamed(js, id)
|
d, err := DeadLetterNamed(js, id)
|
||||||
|
|||||||
@@ -100,3 +100,118 @@ func TestAnAskTheControllerDidNotMakeIsStillOnlyDropped(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user