Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
66de6774e0 | ||
|
|
c401bbff2d |
@@ -839,7 +839,10 @@ func ConsumerOf(node, module string) string {
|
||||
}
|
||||
|
||||
// NakDelay is how long an event a handler failed waits before it is offered again: a transient cause
|
||||
// gets another attempt, a permanent one exhausts the consumer's max-deliver rather than spinning.
|
||||
// gets another attempt, a permanent one exhausts the consumer's max-deliver rather than spinning. The
|
||||
// consumer then gives the event up, and the controller keeps it in DEAD_LETTERS until a person delivers
|
||||
// it again — under `mesh.again.<consumer>.…`, from which keyOf reads the same key — or drops it
|
||||
// (novox/hq issue 330).
|
||||
var NakDelay = 5 * time.Second
|
||||
|
||||
// ConsumeAs reads a module's durable consumer and hands each event to deliver (novox/hq ADR 0198):
|
||||
|
||||
@@ -46,3 +46,17 @@ func TestWhatWaitingCannotFixIsFinal(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// An event given up on and delivered again by the controller arrives under `mesh.again.<consumer>.` and
|
||||
// the event's own subject without its `mesh.` (novox/hq issue 330): the handler sees the key it saw the
|
||||
// first time, for a module's event and a seat's.
|
||||
func TestAnEventDeliveredAgainHasItsOwnKey(t *testing.T) {
|
||||
for original, again := range map[string]string{
|
||||
"mesh.mod.gitea.event.pull.merged": "mesh.again.ace_sonarr.mod.gitea.event.pull.merged",
|
||||
"mesh.seat.node-build-agent.event.built": "mesh.again.ace_sonarr.seat.node-build-agent.event.built",
|
||||
} {
|
||||
if keyOf(again) != keyOf(original) {
|
||||
t.Errorf("%s reads as %q, the event as %q", again, keyOf(again), keyOf(original))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -498,9 +498,9 @@ async function deliver(
|
||||
try {
|
||||
env = toEnvelope<unknown>(msg);
|
||||
} catch {
|
||||
// Unparseable: acknowledge it. Redelivering a message no version of this code can read is
|
||||
// an infinite loop, and the stream's dead-letter is for handlers that fail, not for bytes
|
||||
// that were never an envelope.
|
||||
// Unparseable: terminate it. Redelivering a message no version of this code can read is
|
||||
// an infinite loop, and what the controller keeps of a message given up on (DEAD_LETTERS,
|
||||
// novox/hq issue 330) is for handlers that fail, not for bytes that were never an envelope.
|
||||
msg.term();
|
||||
return;
|
||||
}
|
||||
@@ -517,8 +517,11 @@ async function deliver(
|
||||
msg.ack();
|
||||
} catch {
|
||||
// Negative-acknowledge with a delay, so a handler failing on a transient cause gets another
|
||||
// attempt, and one failing permanently exhausts max-deliver and dead-letters rather than
|
||||
// spinning. The consumer's limits are the controller's; this only says "not done".
|
||||
// attempt, and one failing permanently exhausts max-deliver rather than spinning: the consumer
|
||||
// gives the event up, and the controller keeps it in DEAD_LETTERS until a person delivers it
|
||||
// again or drops it (novox/hq issue 330). Delivered again, it arrives under
|
||||
// `mesh.again.<consumer>.…`, and keyFromSubject reads the same key from it. The consumer's
|
||||
// limits are the controller's; this only says "not done".
|
||||
msg.nak(5_000);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user