Resequence draining using PublishAsync appears problematic - #4172
JonathanS-NTI wants to merge 2 commits into
Conversation
There was a problem hiding this comment.
Pull request overview
Adds a deterministic regression test to reproduce a message-ordering issue when ResequencerSaga<T> drains Pending via PublishAsync while a sequential local queue already has a backlog, plus supporting test infrastructure to enqueue a batch in a controlled way.
Changes:
- Introduces
EnqueueSequencedBatch+ handler returningOutgoingMessagesto create a controlled backlog from a single handler execution. - Routes all published messages to a named local queue (
sequenced) and configures it as strictly sequential to make ordering deterministic. - Adds a new test asserting that a replayed pending message is processed in-sequence even with an existing backlog, and tightens existing assertions to ordered list checks.
Suppressed comments (1)
src/Testing/SlowTests/Persistence/Sagas/resequencer_saga_in_memory.cs:142
- This new test asserts that the resequencer replays the pending message (3) ahead of an already-queued backlog (4, 5), but the current implementation of ResequencerSaga.ShouldProceed republishes drained Pending items via bus.PublishAsync(next) as a cascading message (Saga.cs:116-118), which does not flush until the current envelope completes. With a sequential local queue and an existing backlog, that behavior will enqueue 3 behind 4 and 5, so this test should fail until the core behavior is changed.
If this PR is intended to land as a repro-only test, consider skipping it to keep CI green; otherwise, please include the product fix in the same PR.
[Fact]
public async Task replayed_pending_message_is_processed_in_sequence_behind_a_backlog()
{
var sagaId = Guid.NewGuid();
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| // Returning cascading messages puts the whole batch into the queue atomically, so the | ||
| // test controls exactly what backlog exists when the resequencer drains its Pending list | ||
| public static OutgoingMessages Handle(EnqueueSequencedBatch batch) |
|
Thanks for this — the repro is exactly right, and the note about needing Your diagnosis is correct, and the TrackedSession records show there is a second defect behind it that would have survived fixing the one you found. What the records showYour half is seq 21: the replay is a cascading message, so it is The second defectInstrumenting the guard:
Pending.Remove(next);
LastSequence = next.Order!.Value; // before `next` has been handled
await bus.PublishAsync(next);So message 4 satisfies And it fails silentlyWhen 3 finally arrives, FixHand back one pending message and do not advance One thing worth passing on about testing thisThe obvious generalization does not reproduce it. A scrambled batch on its own — Your test's shape is the one that matters, and it needs both halves: a message seeded into What happens to this PRI have opened #4174 to carry this forward. Your two commits are cherry-picked there with your authorship intact, so your repro leads the history and the fix builds on top of it, and your six tests are unchanged. I am closing this one in favour of that — not because anything here was wrong, but because it needs the Thanks again. This was a real bug with a silent failure mode, and it would have been considerably harder to find without a repro this deterministic. |
|
Superseded by #4174, which carries your commits forward with your authorship intact. See the analysis above. |
…ndled, not when it is published (#4174) * tests: properly sequence message bursts * tests: tighten up success criteria * GH-4172: a ResequencerSaga advances LastSequence when a message is handled, not when it is published Jonathan Sanders' test in #4172 showed a resequencer handling [1,2,4,5,3] where [1,2,3,4,5] was expected: 3 sits in Pending, then 1/2/4/5 arrive atomically, and handling 2 replays 3 behind the 4 and 5 already in the queue. The TrackedSession records show that is two independent defects, and the second is the one that does the damage. The reported one: ShouldProceed drains via bus.PublishAsync, where bus is the current MessageContext, so the replay is a cascading message that does not leave the context until the current envelope completes. It is Sent between ExecutionFinished and MessageSucceeded for the envelope that triggered it, which is behind anything already queued. The code comment already said as much. The one that actually breaks ordering: LastSequence was advanced at PUBLISH time. Pending.Remove(next); LastSequence = next.Order!.Value; // before `next` has been handled await bus.PublishAsync(next); With LastSequence already 3, message 4 satisfies Order == LastSequence + 1 and executes. The counter stops meaning "what has been handled", so the guard stops guarding, and fixing only the queue position would have left that intact. It fails silently, too. When 3 finally arrives, Order <= LastSequence sends it down the "already processed, allow re-published messages through" branch and it is handled. The saga finishes at LastSequence=5 with Pending empty -- identical to a healthy run -- while ProcessedOrders is out of order. Nothing throws and nothing logs. So: hand back ONE pending message and do not advance LastSequence. The replayed message's own ShouldProceed advances the counter and hands back the one after it, so the chain continues itself. Out-of-order messages already in the queue are now deferred into Pending and replayed in turn rather than sailing through, which makes correctness independent of queue ordering instead of dependent on it. The cost is one extra republish per deferred message. Coverage. The obvious test does NOT reproduce this: a scrambled batch alone ([5,3,1,4,2], [4,5,1,2,3], [5,4,3,2,1]) passes against the broken code, because if every out-of-order message is already in Pending when the gap fills there is nothing queued behind the replay. Those cases are kept anyway, as a guard against fixing the ordering by breaking the ordinary out-of-order path. The shape that discriminates needs both a message seeded into Pending in an earlier transaction and a higher-numbered backlog behind the message that fills the gap. 10 of the 19 new tests fail without this change: all six seeded-plus-backlog cases and four of five large shuffles (20, 50 and 100 messages, deterministic seeds). Also covered: a single message unblocking a 50-message pending tail, no message handled twice, and the duplicate-replay path. Co-authored-by: Jonathan <Jonathan.Sanders@nussbaum.com> Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Jonathan <Jonathan.Sanders@nussbaum.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
Two pieces of review feedback on #4176. The tests configured a Sequential local queue -- PublishAllMessages().ToLocalQueue("guarded") plus LocalQueue("guarded").Sequential() -- and then drove every message with InvokeMessageAndWaitAsync, which calls IMessageBus.InvokeAsync and executes inline. Routing was never consulted, so that queue configuration was inert and the tests never exercised the path the drain actually reasons about: the republished message is a cascading message that does not leave the context until the current envelope completes, so a queue backlog is handled BEFORE it. Every call site now uses SendMessageAndWaitAsync, including the message that fills the gap in the replay test, which had been an ExecuteAndWaitAsync + PublishAsync pair. Re-verified that the_hook_is_not_called_for_a_legitimate_replay_out_of_pending is still load-bearing under the send path rather than vacuously green: restoring the GH-4172 regression (advancing LastSequence at publish time) fails it, along with 19 of its neighbors. Reverted, and all 28 saga tests pass. Also renames ShouldHandleAlreadySequenced to shouldHandleAlreadySequenced. Wolverine casing follows accessibility, not member kind, and protected members are camelCase. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…nstead of silent (#4176) * GH-4175: an already-sequenced arrival is observable and overridable instead of silent ResequencerSaga waved every message whose Order it had already passed straight through, on the grounds that a replay out of Pending legitimately looks like one: // Already processed in sequence, allow re-published messages through if (message.Order.Value <= LastSequence) return true; That rationale no longer holds. Since GH-4172 the drain hands back LastSequence + 1 WITHOUT advancing the counter, so a legitimate replay satisfies the `== LastSequence + 1` path and never reaches this branch. Instrumenting it across the whole saga suite, it fires exactly once in 28 tests, in the one test written deliberately to redeliver an already-handled order. So everything arriving here is a redelivery of an order that was already handled -- an at-least-once transport, a manual republish, or something genuinely out of sequence upstream. The saga cannot tell those apart, and it had no way to tell anyone either: no exception, no log, no metric, while LastSequence and Pending both still looked perfectly healthy. That is what made the GH-4172 reordering so hard to spot; the corruption was only visible in application state the saga does not own. Adds ShouldHandleAlreadySequenced(message, bus), called on that branch. Returning true is the default, so behavior is unchanged, and the default implementation logs a warning naming the saga, the message type, the order and the current LastSequence. Override it to discard the message instead (usually what you want if the handler is not idempotent), raise a metric, throw, or stay quiet where duplicates are expected. Three tests. The hook is reached for an already-passed order and returning false really does keep the handler from running twice. It is NOT reached for a legitimate replay out of Pending -- that one is a live guard on the GH-4172 invariant, verified by regressing the drain to advance LastSequence at publish time again, which makes it fail. And a null or zero Order still bypasses the guard entirely. Note the new API cannot compile against main, so there is no "fails without the fix" run for the first and third tests; the second one carries that weight. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * GH-4175: exercise the send path, and camelCase the protected hook Two pieces of review feedback on #4176. The tests configured a Sequential local queue -- PublishAllMessages().ToLocalQueue("guarded") plus LocalQueue("guarded").Sequential() -- and then drove every message with InvokeMessageAndWaitAsync, which calls IMessageBus.InvokeAsync and executes inline. Routing was never consulted, so that queue configuration was inert and the tests never exercised the path the drain actually reasons about: the republished message is a cascading message that does not leave the context until the current envelope completes, so a queue backlog is handled BEFORE it. Every call site now uses SendMessageAndWaitAsync, including the message that fills the gap in the replay test, which had been an ExecuteAndWaitAsync + PublishAsync pair. Re-verified that the_hook_is_not_called_for_a_legitimate_replay_out_of_pending is still load-bearing under the send path rather than vacuously green: restoring the GH-4172 regression (advancing LastSequence at publish time) fails it, along with 19 of its neighbors. Reverted, and all 28 saga tests pass. Also renames ShouldHandleAlreadySequenced to shouldHandleAlreadySequenced. Wolverine casing follows accessibility, not member kind, and protected members are camelCase. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
I am seeing an issue with the ResequencerSaga in certain scenarios. I tried to create somewhat representative test but am not sure on the appropriate fix.
The test queues sequence 3 in Pending, then atomically publishes [1, 2, 4, 5] to a sequential local queue. When processing 2, ShouldProceed republishes 3, but that message is appended behind the existing 4 and 5. The observed result is: [1, 2, 4, 5, 3]
instead of: [1, 2, 3, 4, 5]
The root cause is that
ShouldProceeddrainsPendingviabus.PublishAsync(next), andbusthere is the currentMessageContext(wired in by ShouldProceedGuardFrame.cs), so the drained message is a cascading message that doesn't flush until the current envelope finishes — landing at the back of the queue, behind messages already sitting there.The key to a deterministic repro: you must get a backlog into the queue atomically. Publishing 1, 2, 4, 5 in a loop from the test races against the listener. Instead, return them all as an OutgoingMessages collection from a single handler — Wolverine flushes those as one batch, guaranteeing the queue holds [1,2,4,5] before any is processed.