Conversation
| RedisModule_Log(mr_staticCtx, | ||
| "warning", | ||
| "Rejecting duplicate execution id %s", | ||
| e->idStr); |
There was a problem hiding this comment.
I don't think that rejecting a dupped id is a wise choice. The acceptance of dups is probably by design, following the principle of idempotency in distributed systems.
There was a problem hiding this comment.
Agreed. Commit 3c38e0b preserves idempotency: when the ID already exists, the receiver ACKs the retransmission, frees the newly deserialized duplicate, and keeps the existing execution.
| /* The shard id is stable across a process restart. Starting the sequence at | ||
| * zero would therefore reuse ids while remote shards can still hold an | ||
| * execution created by the previous process. Keep the wire format intact | ||
| * and randomize the 64-bit sequence origin for each module load. */ |
There was a problem hiding this comment.
The comment is confusing (e.g., it mentions the shard id, which gets on board only in SetId()) and redundant.
There was a problem hiding this comment.
Removed the confusing initialization comment. The code now explicitly initializes a nonzero 32-bit execution epoch in the high bits and keeps the incrementing counter in the low bits.
| * execution created by the previous process. Keep the wire format intact | ||
| * and randomize the 64-bit sequence origin for each module load. */ | ||
| RedisModule_GetRandomBytes((unsigned char *)&mrCtx.lastExecutionId, | ||
| sizeof(mrCtx.lastExecutionId)); |
There was a problem hiding this comment.
I agreed that resetting the execution ID on restart risks collisions, but:
- Should reproduce first: I want to understand why other safeguards (e.g., dmc closing the sconn to the dead shard) didn't prevent this.
- Keep the counter: it lets us trace execution ids in post-mortems, and fully random 64-bit ids lose that. Suggest 32 random high bits (per shard start) + 32-bit counter in the low bits.
There was a problem hiding this comment.
Updated as suggested: a nonzero random 32-bit restart epoch in the high bits and a 32-bit counter in the low bits, so IDs remain traceable. The reproduction leaves an unevenwork execution in the surviving shard, kills only its initiator, restarts it, and then runs readallkeys. Before the fix, the first post-restart ID collides and the surviving shard uses the stale pipeline; the revised focused two-shard test passes. DMC closing its connection to the dead shard does not clear the LibMR execution dictionary in the surviving shard, which retains that execution until max-idle cleanup.
|
|
||
|
|
||
| @MRTestDecorator(skipTest=Defaults.num_shards == 1) | ||
| def testExecutionIdsDoNotCollideAfterShardRestart(env, conn): |
There was a problem hiding this comment.
I don't understand how this test is relevant to anything...
There was a problem hiding this comment.
I renamed the helper and added a test docstring plus comments describing the exact failure. unevenwork keeps the old execution alive on the surviving shard; after the initiator restarts, readallkeys detects whether its first execution ID collides with that retained state. The old implementation times out or runs the stale pipeline; commit 3c38e0b passes this focused two-shard regression.
|
Validated the fix with an A/B test against the exact restart-collision trigger. Test sequence:
Results using the same test binary and harness:
Command used: PYTHON="$PWD/.venv/bin/python" \
tests/mr_test_module/pytests/run_tests.sh \
-t tests/mr_test_module/pytests/test_execution_id_restart.py \
--env oss-cluster --shards-count 2This isolates the root cause: the surviving shard retains an old execution while the restarted initiator resets its local counter. With the old code, the new command can reuse that ID and the survivor executes/returns data for the stale pipeline, which is the mechanism that produced mixed MGET/MRANGE records and incomplete RESP replies in the Enterprise failure. With the restart epoch in the high 32 bits, the new execution has a different ID; normal retransmission of the same execution remains idempotent. The normal CI coverage on Redis 7.2 and 7.4 also passes. The unstable jobs still fail during the known RLTest OSS-cluster startup issue (MOD-18751), before reaching this regression test. |
A shard restart resets LibMR's process-local execution counter while the shard ID remains stable. Remote shards can still hold an in-flight execution from the previous process, so a new command can reuse the same
(shard ID, counter)ID. The receiver ignoredmr_dictAddduplicate failure and ACKed the new execution, causing later invoke/result messages to operate on the stale command and mix record types across distributed commands.Seed the existing counter with random bytes at each module load, preserving the wire format while making IDs unique across restarts. Also reject duplicate IDs instead of ACKing them as a defensive invariant.
The regression test leaves an execution active on a peer, restarts its initiator, and issues a different distributed command. Before this change the second command times out; with this change it completes and returns the expected key.
Validation:
End-to-end TimeSeries proof used two shards and 200 labeled series. On the affected LibMR, an MGET execution remained on the remote shard while its initiator restarted. The first post-restart MRANGE reused the exact same ID; instrumentation recorded
existing=ShardMgetMapper incoming=ShardSeriesMapper, and the MRANGE client timed out on the resulting incomplete response. With this change, MGET and MRANGE received different IDs, no duplicate was recorded, and MRANGE returned all 200 series after the normal abandoned-execution timeout window.Note
Medium Risk
Changes core distributed execution identity and receive-path handling; incorrect behavior could still cause cross-command mix-ups or stuck executions in cluster mode.
Overview
Fixes execution ID collisions when a LibMR initiator shard restarts while peers still hold an in-flight execution for the same
(shard ID, counter)pair. A restarted process used to reset its local counter to zero, so the next distributed command could reuse an ID still registered on another shard and route later messages to the stale pipeline.MR_Initnow seedslastExecutionIdwith a random non-zero 32-bit epoch in the high half of the counter (wire format unchanged; only the starting point moves).MR_ReceivedExecutiontreats a failedmr_dictAddas an idempotent retransmit: ACK without replacing the running execution and free the duplicate payload.Adds a two-shard regression test that leaves
unevenworkactive on a peer, restarts the initiator, and asserts a subsequentreadallkeyscompletes with the expected key.Reviewed by Cursor Bugbot for commit e989512. Bugbot is set up for automated code reviews on this repo. Configure here.