From 4d6d6262afae56f8133101c26a33a3ade7358963 Mon Sep 17 00:00:00 2001 From: Joshua Chi Date: Mon, 13 Apr 2026 15:54:33 +0800 Subject: [PATCH] fix #308: follower and learner respond success only after snapshot apply completes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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). --- d-engine-core/src/raft_role/follower_state.rs | 40 ++---- .../src/raft_role/follower_state_test.rs | 132 ++++++++++++++++++ d-engine-core/src/raft_role/learner_state.rs | 84 ++++++----- .../src/raft_role/learner_state_test.rs | 68 +++++++++ d-engine-server/src/node/mod.rs | 6 +- 5 files changed, 255 insertions(+), 75 deletions(-) diff --git a/d-engine-core/src/raft_role/follower_state.rs b/d-engine-core/src/raft_role/follower_state.rs index b42e5110..6ea20f94 100644 --- a/d-engine-core/src/raft_role/follower_state.rs +++ b/d-engine-core/src/raft_role/follower_state.rs @@ -10,7 +10,6 @@ use d_engine_proto::server::election::VoteResponse; use d_engine_proto::server::storage::SnapshotAck; use d_engine_proto::server::storage::SnapshotResponse; -use d_engine_proto::server::storage::snapshot_ack::ChunkStatus; use tokio::sync::mpsc; use tokio::time::Instant; use tonic::Status; @@ -334,31 +333,9 @@ impl RaftRoleState for FollowerState { } RaftEvent::InstallSnapshotChunk(stream, sender) => { - // Create ACK channel (follower sends ACKs to leader) - let (ack_tx, mut ack_rx) = mpsc::channel::(32); - - // Spawn ACK handler to send final response. - // process_snapshot_stream sends one ACK per chunk; drain all of them and - // use the last one as the final status once the sender side is dropped. - tokio::spawn(async move { - let mut last_ack: Option = None; - while let Some(ack) = ack_rx.recv().await { - last_ack = Some(ack); - } - match last_ack { - Some(final_ack) => { - let response = SnapshotResponse { - term: my_term, - success: final_ack.status == (ChunkStatus::Accepted as i32), - next_chunk: final_ack.next_requested, - }; - let _ = sender.send(Ok(response)); - } - None => { - let _ = sender.send(Err(Status::internal("No ACK received"))); - } - } - }); + // ack_tx is used internally by process_snapshot_stream for per-chunk + // validation only; _ack_rx is intentionally discarded in push mode. + let (ack_tx, _ack_rx) = mpsc::channel::(32); let snap_result = ctx .handlers @@ -370,6 +347,17 @@ impl RaftRoleState for FollowerState { &ctx.node_config.raft.snapshot, ) .await; + + // Raft §7: respond success only after snapshot is fully applied. + // Using last-chunk ACK status (as before) would report success even when + // apply_snapshot_from_file fails, causing the leader to advance match_index + // prematurely and stop retrying — leaving the follower permanently behind (#308). + let _ = sender.send(Ok(SnapshotResponse { + term: my_term, + success: snap_result.is_ok(), + next_chunk: 0, + })); + match snap_result { Err(e) => { // Transient failure: leader may have crashed mid-transfer. diff --git a/d-engine-core/src/raft_role/follower_state_test.rs b/d-engine-core/src/raft_role/follower_state_test.rs index 0db2efa1..196f1558 100644 --- a/d-engine-core/src/raft_role/follower_state_test.rs +++ b/d-engine-core/src/raft_role/follower_state_test.rs @@ -1,3 +1,4 @@ +use crate::test_utils::create_test_snapshot_stream; use d_engine_proto::client::ClientReadRequest; use d_engine_proto::client::ClientWriteRequest; use d_engine_proto::client::ReadConsistencyPolicy; @@ -17,6 +18,9 @@ use d_engine_proto::server::election::VoteResponse; use d_engine_proto::server::election::VotedFor; use d_engine_proto::server::replication::AppendEntriesRequest; use d_engine_proto::server::replication::AppendEntriesResponse; +use d_engine_proto::server::storage::SnapshotAck; +use d_engine_proto::server::storage::SnapshotChunk; +use d_engine_proto::server::storage::snapshot_ack::ChunkStatus; use std::sync::Arc; use tonic::Code; use tonic::Status; @@ -2479,3 +2483,131 @@ async fn test_follower_commit_index_and_ack_both_sent_immediately() { let response = resp_rx.try_recv().expect("ACK must be sent immediately"); assert!(response.unwrap().is_success()); } + +// ============================================================================ +// InstallSnapshotChunk Tests +// ============================================================================ + +/// Follower reports success when snapshot is fully transferred and applied. +/// +/// # Given +/// - apply_snapshot_stream_from_leader returns Ok(()) +/// +/// # When +/// - Leader pushes a snapshot (InstallSnapshotChunk event) +/// +/// # Then +/// - Response success: true +#[tokio::test] +async fn test_follower_install_snapshot_reports_success_when_apply_succeeds() { + let (_graceful_tx, graceful_rx) = watch::channel(()); + let (mut context, _temp_dir) = mock_raft_context_with_temp(graceful_rx, None); + + let mut sm_handler = MockStateMachineHandler::new(); + sm_handler.expect_apply_snapshot_stream_from_leader().once().returning( + |_term, _stream, ack_tx, _config| { + let _ = ack_tx.try_send(SnapshotAck { + seq: 0, + status: ChunkStatus::Accepted as i32, + next_requested: 1, + }); + Ok(()) + }, + ); + sm_handler.expect_get_latest_snapshot_metadata().returning(|| None); + context.handlers.state_machine_handler = Arc::new(sm_handler); + + let mut state = + FollowerState::::new(1, context.node_config.clone(), None, None); + + let stream = create_test_snapshot_stream(vec![SnapshotChunk::default()]); + let (resp_tx, mut resp_rx) = MaybeCloneOneshot::new(); + let (role_tx, _role_rx) = mpsc::unbounded_channel(); + + state + .handle_raft_event( + RaftEvent::InstallSnapshotChunk(Box::new(stream), resp_tx), + &context, + role_tx, + ) + .await + .unwrap(); + + let response = tokio::time::timeout(std::time::Duration::from_secs(2), resp_rx.recv()) + .await + .expect("response must arrive within 2s") + .expect("recv must not fail") + .expect("response must be Ok"); + + assert!( + response.success, + "Follower must report success when apply succeeds" + ); +} + +/// Follower must NOT report success when apply fails after transfer completes (#308). +/// +/// # Raft §7 + #308 +/// The previous implementation derived success from the last per-chunk ACK status. +/// When all chunks are received (last ACK = Accepted) but apply_snapshot_from_file +/// then fails, the spawned ACK-handler still sends success:true — causing the leader +/// to advance match_index and stop retrying, leaving the follower permanently behind. +/// +/// # Given +/// - apply_snapshot_stream_from_leader: sends Accepted ACK (transfer succeeded), +/// then returns Err (apply_snapshot_from_file failed) +/// +/// # When +/// - Leader pushes a snapshot (InstallSnapshotChunk event) +/// +/// # Then +/// - Response MUST be success: false +#[tokio::test] +async fn test_follower_install_snapshot_reports_failure_when_apply_fails_after_transfer() { + let (_graceful_tx, graceful_rx) = watch::channel(()); + let (mut context, _temp_dir) = mock_raft_context_with_temp(graceful_rx, None); + + let mut sm_handler = MockStateMachineHandler::new(); + sm_handler.expect_apply_snapshot_stream_from_leader().once().returning( + |_term, _stream, ack_tx, _config| { + // Transfer phase succeeds: all chunks accepted + let _ = ack_tx.try_send(SnapshotAck { + seq: 0, + status: ChunkStatus::Accepted as i32, + next_requested: 1, + }); + // Apply phase fails (apply_snapshot_from_file returned Err) + Err(crate::Error::Fatal( + "apply_snapshot_from_file failed".into(), + )) + }, + ); + context.handlers.state_machine_handler = Arc::new(sm_handler); + + let mut state = + FollowerState::::new(1, context.node_config.clone(), None, None); + + let stream = create_test_snapshot_stream(vec![SnapshotChunk::default()]); + let (resp_tx, mut resp_rx) = MaybeCloneOneshot::new(); + let (role_tx, _role_rx) = mpsc::unbounded_channel(); + + // Follower absorbs the error and continues (does not propagate) + let _ = state + .handle_raft_event( + RaftEvent::InstallSnapshotChunk(Box::new(stream), resp_tx), + &context, + role_tx, + ) + .await; + + let response = tokio::time::timeout(std::time::Duration::from_secs(2), resp_rx.recv()) + .await + .expect("response must arrive within 2s") + .expect("recv must not fail") + .expect("response must be Ok(SnapshotResponse)"); + + assert!( + !response.success, + "Follower must NOT report success when apply failed after transfer (got success:true — #308 bug)" + ); +} diff --git a/d-engine-core/src/raft_role/learner_state.rs b/d-engine-core/src/raft_role/learner_state.rs index e5673d4e..6a206ceb 100644 --- a/d-engine-core/src/raft_role/learner_state.rs +++ b/d-engine-core/src/raft_role/learner_state.rs @@ -15,7 +15,6 @@ use d_engine_proto::server::election::VoteResponse; use d_engine_proto::server::election::VotedFor; use d_engine_proto::server::storage::SnapshotAck; use d_engine_proto::server::storage::SnapshotResponse; -use d_engine_proto::server::storage::snapshot_ack::ChunkStatus; use tokio::sync::mpsc::{self}; use tokio::time::Instant; use tonic::Status; @@ -287,34 +286,11 @@ impl RaftRoleState for LearnerState { } RaftEvent::InstallSnapshotChunk(stream, sender) => { - // Create ACK channel (follower sends ACKs to leader) - let (ack_tx, mut ack_rx) = mpsc::channel::(32); + // ack_tx is used internally by process_snapshot_stream for per-chunk + // validation only; _ack_rx is intentionally discarded in push mode. + let (ack_tx, _ack_rx) = mpsc::channel::(32); - // Spawn ACK handler to send final response. - // Drain ALL ACKs and use the last one — the channel closes only after - // apply_snapshot_stream_from_leader returns (success or error), so the - // last ACK reflects the true final outcome of the full snapshot install. - tokio::spawn(async move { - let mut last_ack: Option = None; - while let Some(ack) = ack_rx.recv().await { - last_ack = Some(ack); - } - match last_ack { - Some(final_ack) => { - let response = SnapshotResponse { - term: my_term, - success: final_ack.status == (ChunkStatus::Accepted as i32), - next_chunk: final_ack.next_requested, - }; - let _ = sender.send(Ok(response)); - } - None => { - let _ = sender.send(Err(Status::internal("ACK channel closed"))); - } - } - }); - - if let Err(e) = ctx + let snap_result = ctx .handlers .state_machine_handler .apply_snapshot_stream_from_leader( @@ -323,24 +299,42 @@ impl RaftRoleState for LearnerState { ack_tx, &ctx.node_config.raft.snapshot, ) - .await - { - error!(?e, "Learner handle RaftEvent::InstallSnapshotChunk"); - return Err(e); - } + .await; - // Advance raft log purge boundary to snapshot's last_included so that - // last_log_id() returns the correct position after snapshot install. - if let Some(metadata) = ctx.state_machine_handler().get_latest_snapshot_metadata() - && let Some(last_included) = metadata.last_included - { - if let Err(e) = ctx.raft_log().purge_logs_up_to(last_included).await { - error!(?e, "Failed to set raft log boundary after snapshot install"); - } else { - info!( - ?last_included, - "Learner raft log boundary set after InstallSnapshotChunk" - ); + // Raft §7: respond success only after snapshot is fully applied. + // Using last-chunk ACK status (as before) would report success even when + // apply_snapshot_from_file fails, causing the leader to advance match_index + // prematurely and stop retrying — leaving the learner permanently behind (#308). + let _ = sender.send(Ok(SnapshotResponse { + term: my_term, + success: snap_result.is_ok(), + next_chunk: 0, + })); + + match snap_result { + Err(e) => { + error!(?e, "Learner handle RaftEvent::InstallSnapshotChunk"); + return Err(e); + } + Ok(()) => { + // Advance raft log purge boundary to snapshot's last_included so that + // last_log_id() returns the correct position after snapshot install. + if let Some(metadata) = + ctx.state_machine_handler().get_latest_snapshot_metadata() + && let Some(last_included) = metadata.last_included + { + if let Err(e) = ctx.raft_log().purge_logs_up_to(last_included).await { + error!( + ?e, + "Failed to set raft log boundary after snapshot install" + ); + } else { + info!( + ?last_included, + "Learner raft log boundary set after InstallSnapshotChunk" + ); + } + } } } } diff --git a/d-engine-core/src/raft_role/learner_state_test.rs b/d-engine-core/src/raft_role/learner_state_test.rs index 20d76bc7..bf4e7962 100644 --- a/d-engine-core/src/raft_role/learner_state_test.rs +++ b/d-engine-core/src/raft_role/learner_state_test.rs @@ -1886,3 +1886,71 @@ async fn test_learner_install_snapshot_does_not_report_success_on_mid_chunk_fail // If no response arrived or it was an error — acceptable here. // The primary assertion: success:true must never be sent on failure. } + +/// Learner must NOT report success when apply fails after transfer succeeds (#308). +/// +/// # Raft §7 + #308 +/// The existing test above covers chunk-level failures (Failed ACK sent before Err). +/// This test covers the distinct #308 scenario: ALL chunks are accepted (transfer ok), +/// but apply_snapshot_from_file then fails. The last per-chunk ACK is still Accepted, +/// so the buggy ACK-handler sends success:true — causing the leader to advance +/// match_index and stop retrying, leaving the learner permanently behind. +/// +/// # Given +/// - apply_snapshot_stream_from_leader: sends Accepted ACKs (transfer succeeded), +/// then returns Err (apply_snapshot_from_file failed) +/// +/// # When +/// - Leader pushes a snapshot (InstallSnapshotChunk event) +/// +/// # Then +/// - Response MUST be success: false +#[tokio::test] +async fn test_learner_install_snapshot_reports_failure_when_apply_fails_after_transfer() { + let (_graceful_tx, graceful_rx) = watch::channel(()); + let (mut context, _temp_dir) = mock_raft_context_with_temp(graceful_rx, None); + + let mut sm_handler = MockStateMachineHandler::new(); + sm_handler.expect_apply_snapshot_stream_from_leader().once().returning( + |_term, _stream, ack_tx, _config| { + // Transfer phase succeeds: all chunks accepted + let _ = ack_tx.try_send(SnapshotAck { + seq: 0, + status: ChunkStatus::Accepted as i32, + next_requested: 1, + }); + // Apply phase fails (apply_snapshot_from_file returned Err) + Err(crate::Error::Fatal( + "apply_snapshot_from_file failed".into(), + )) + }, + ); + context.handlers.state_machine_handler = Arc::new(sm_handler); + + let mut state = LearnerState::::new(1, context.node_config.clone()); + state.update_current_term(2); + + let stream = create_test_snapshot_stream(vec![SnapshotChunk::default()]); + let (resp_tx, mut resp_rx) = MaybeCloneOneshot::new(); + let (role_tx, _role_rx) = mpsc::unbounded_channel(); + + // Learner propagates apply errors — ignore the return value here + let _ = state + .handle_raft_event( + RaftEvent::InstallSnapshotChunk(Box::new(stream), resp_tx), + &context, + role_tx, + ) + .await; + + let response = tokio::time::timeout(std::time::Duration::from_secs(2), resp_rx.recv()) + .await + .expect("response must arrive within 2s") + .expect("recv must not fail") + .expect("response must be Ok(SnapshotResponse)"); + + assert!( + !response.success, + "Learner must NOT report success when apply failed after transfer (got success:true — #308 bug)" + ); +} diff --git a/d-engine-server/src/node/mod.rs b/d-engine-server/src/node/mod.rs index eab92400..0a524b36 100644 --- a/d-engine-server/src/node/mod.rs +++ b/d-engine-server/src/node/mod.rs @@ -169,10 +169,8 @@ where // Note: IO thread is closed inside Raft::run() on shutdown before returning. self.start_raft_loop().await?; - // Shutdown order is reverse of startup order (TiKV convention): - // Raft loop has exited → no more applies will be enqueued. - // Join sm-worker thread so its Arc clone is dropped before we return. - // This ensures RocksDB LOCK is released before the caller can reopen the DB. + // Shutdown in reverse startup order: join sm-worker thread first so its + // Arc clone is dropped before we return, releasing the RocksDB LOCK. let handle = self.sm_worker_handle.lock().unwrap().take(); if let Some(handle) = handle { tokio::task::spawn_blocking(move || {