Repository navigation
fix #308: snapshot install success driven by apply result not transfer ACK - #360
Conversation
…ply completes Leader advanced match_index based on the last per-chunk ACK (Accepted), not on whether apply_snapshot_from_file actually succeeded. When apply failed after a successful transfer, the spawned ACK-drain task still sent success:true — causing the leader to stop retrying and leaving the node permanently behind. Remove the spawned ACK-drain task in FollowerState and LearnerState; send SnapshotResponse directly from snap_result.is_ok() so the leader only advances match_index when the snapshot is fully applied (Raft §7).
📝 WalkthroughWalkthroughThis PR removes the spawned ACK-draining task from snapshot installation handling in followers and learners. Instead of deriving Changes
Sequence DiagramsequenceDiagram
participant Leader
participant Follower
participant ApplyTask as apply_snapshot_stream<br/>from_leader
participant RocksDB
rect rgba(0, 100, 150, 0.5)
Note over Leader,Follower: OLD FLOW: ACK-based success
Leader->>Follower: InstallSnapshotChunk
Follower->>ApplyTask: Call with ack_tx
ApplyTask->>ApplyTask: Send per-chunk SnapshotAck<br/>(ChunkStatus::Accepted)
Follower->>Follower: Spawn task to drain ACKs<br/>from ack_rx
ApplyTask->>RocksDB: Apply chunks
ApplyTask-->>Follower: Return result
Follower->>Follower: Derive success from<br/>last ACK status
Follower->>Leader: SnapshotResponse<br/>(success from ACK)
end
rect rgba(150, 100, 0, 0.5)
Note over Leader,Follower: NEW FLOW: Result-based success
Leader->>Follower: InstallSnapshotChunk
Follower->>ApplyTask: Call with ack_tx
ApplyTask->>ApplyTask: Send per-chunk SnapshotAck<br/>(for validation only)
ApplyTask->>RocksDB: Apply chunks
ApplyTask-->>Follower: Return Ok()/Err
Follower->>Follower: Derive success from<br/>result.is_ok()
Follower->>Leader: SnapshotResponse<br/>(success from result)
end
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 minutes Possibly related issues
Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
d-engine-core/src/raft_role/learner_state.rs (1)
289-312:⚠️ Potential issue | 🔴 CriticalThe learner path has the same 32-chunk deadlock.
Line 291 creates a bounded ACK channel but never drains
ack_rx. Becaused-engine-core/src/state_machine_handler/snapshot_stream_processor.rs:1-20awaitstx.send(ack)for each chunk, a pushed snapshot with more than 32 chunks will block mid-transfer andapply_snapshot_stream_from_leaderwill never return. This needs the same fix as the follower path: keep draining the channel, or pass a sink that cannot backpressure.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@d-engine-core/src/raft_role/learner_state.rs` around lines 289 - 312, The bounded ACK channel created with let (ack_tx, _ack_rx) = mpsc::channel::<SnapshotAck>(32) can block apply_snapshot_stream_from_leader because _ack_rx is never drained; fix by consuming _ack_rx (or switch to a non-blocking sink) similarly to the follower path: spawn a background task to loop and drain/ignore received SnapshotAck values from _ack_rx so ack_tx never backpressures apply_snapshot_stream_from_leader, leaving the rest of the logic (use of ack_tx, call to apply_snapshot_stream_from_leader, and sending the SnapshotResponse via sender.send) unchanged.d-engine-core/src/raft_role/follower_state.rs (1)
336-359:⚠️ Potential issue | 🔴 CriticalThe ACK channel will block on snapshots larger than 32 chunks.
Line 338 discards
_ack_rx, butprocess_snapshot_stream(line 963 in default_state_machine_handler.rs) sends an ACK for every successfully processed chunk viaack_tx.send(ack).await.map_err(...). With a 32-slot bounded channel and no consumer, the buffer fills after 32 chunks. On chunk 33+, the bounded channel fails to send because the receiver is dropped—causing the entire snapshot install to fail. Tests only use snapshots with ≤4 chunks, so this goes undetected.Either spawn a drain task for the channel, or refactor
apply_snapshot_stream_from_leader/process_snapshot_streamto skip or accept a no-op ACK sink in push mode instead of a bounded channel.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@d-engine-core/src/raft_role/follower_state.rs` around lines 336 - 359, The current implementation creates a bounded mpsc channel (mpsc::channel::<SnapshotAck>(32)) and drops the receiver (_ack_rx), which will block apply_snapshot_stream_from_leader/process_snapshot_stream once more than 32 chunk ACKs are sent; fix by ensuring ACKs are drained: either spawn a background task that consumes _ack_rx and silently ignores incoming SnapshotAck values (so ack_tx.send(...).await never blocks), or change the sink passed to apply_snapshot_stream_from_leader to a no-op/discard sink instead of a bounded channel (refactor process_snapshot_stream to accept a Box<dyn Sink> or similar). Ensure the fix targets the channel creation and usage around ack_tx/_ack_rx and update apply_snapshot_stream_from_leader/process_snapshot_stream wiring accordingly.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@d-engine-core/src/raft_role/follower_state.rs`:
- Around line 336-359: The current implementation creates a bounded mpsc channel
(mpsc::channel::<SnapshotAck>(32)) and drops the receiver (_ack_rx), which will
block apply_snapshot_stream_from_leader/process_snapshot_stream once more than
32 chunk ACKs are sent; fix by ensuring ACKs are drained: either spawn a
background task that consumes _ack_rx and silently ignores incoming SnapshotAck
values (so ack_tx.send(...).await never blocks), or change the sink passed to
apply_snapshot_stream_from_leader to a no-op/discard sink instead of a bounded
channel (refactor process_snapshot_stream to accept a Box<dyn Sink> or similar).
Ensure the fix targets the channel creation and usage around ack_tx/_ack_rx and
update apply_snapshot_stream_from_leader/process_snapshot_stream wiring
accordingly.
In `@d-engine-core/src/raft_role/learner_state.rs`:
- Around line 289-312: The bounded ACK channel created with let (ack_tx,
_ack_rx) = mpsc::channel::<SnapshotAck>(32) can block
apply_snapshot_stream_from_leader because _ack_rx is never drained; fix by
consuming _ack_rx (or switch to a non-blocking sink) similarly to the follower
path: spawn a background task to loop and drain/ignore received SnapshotAck
values from _ack_rx so ack_tx never backpressures
apply_snapshot_stream_from_leader, leaving the rest of the logic (use of ack_tx,
call to apply_snapshot_stream_from_leader, and sending the SnapshotResponse via
sender.send) unchanged.
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 446b2d0c-dfd2-4917-9838-b919422481d6
📒 Files selected for processing (5)
d-engine-core/src/raft_role/follower_state.rsd-engine-core/src/raft_role/follower_state_test.rsd-engine-core/src/raft_role/learner_state.rsd-engine-core/src/raft_role/learner_state_test.rsd-engine-server/src/node/mod.rs
Codecov Report✅ All modified and coverable lines are covered by tests. 📢 Thoughts on this report? Let us know! |
What Does This PR Do?
Fixes a bug where follower/learner nodes reported snapshot install success
based on whether all chunks were received, not whether the snapshot was
actually applied — causing the leader to advance
match_indexprematurelyand stop retrying, leaving the node permanently behind.
Type:
Why Is This Needed?
Bug: When
apply_snapshot_from_filefailed after a successful chunktransfer, the spawned ACK-drain task saw the last per-chunk ACK
(
ChunkStatus::Accepted) and sentsuccess:trueto the leader. Theleader then called
init_peers_next_index_and_match_index, advancingmatch_indexto its ownlast_entry_id. Withmatch_indexcaught up,the heartbeat loop never retried the snapshot — leaving the follower
permanently behind with no recovery path.
Fix: Remove the spawned ACK-drain task in
FollowerStateandLearnerState. SendSnapshotResponsedirectly afterapply_snapshot_stream_from_leaderreturns, withsuccess = snap_result.is_ok(). The leader now only advancesmatch_indexwhen the snapshot is fully applied (Raft §7). No newretry timer needed — the existing heartbeat loop handles retries
naturally once
match_indexstays behind the purge boundary.Checklist
Required:
make testpassesTesting
How tested:
test_follower_install_snapshot_reports_failure_when_apply_fails_after_transfer— fails without fix, passes after (leader may advance match_index on snapshot transfer complete without waiting for follower apply confirmation #308 exact repro)test_follower_install_snapshot_reports_success_when_apply_succeeds— happy pathtest_learner_install_snapshot_reports_failure_when_apply_fails_after_transfer— fails without fix, passes after (leader may advance match_index on snapshot transfer complete without waiting for follower apply confirmation #308 exact repro)For bug fixes:
Does This Follow d-engine's Principles?
Reviewer Notes
Core change is 2 files × ~15 lines each. The spawned task and its
ACK-drain loop are removed entirely;
sender.send(...)is now a directcall in the same execution path as
snap_result.The
_ack_rxdiscard is intentional:ack_txis still passed intoprocess_snapshot_streamfor per-chunk internal validation; the receiverside is simply not needed at the call site in push mode.
Estimated review complexity:
Summary by CodeRabbit
Release Notes
Bug Fixes
Tests