From e5fc7eb130b52f8564abbfe6245bc00b9f463039 Mon Sep 17 00:00:00 2001 From: Zhang Yanpo Date: Tue, 29 Sep 2026 19:30:55 +0800 Subject: [PATCH 1/2] feat: replication: add Trigger::reset_backoff() # Summary Applications can end the active AppendEntries backoff for selected replication targets, or every target when no node ids are given. # Details A follower may be reachable again while its replication stream still waits after an RPC failure. Calling `reset_backoff()` before `Trigger::transfer_leader()` lets replication resume without that delay. A reset ends the current or next wait of the active backoff and drops its `Backoff` iterator. `BackoffState::rank` remains until an RPC succeeds, so another failure can start backoff again. A reset sent before backoff starts does not affect the later backoff. Snapshot transfer has a separate retry backoff. The reset receiver goes to `ReplicationCore::spawn()`. It stays out of `ReplicationContext`, which is also used by `SnapshotTransmitter`. This avoids giving snapshot tasks a receiver whose sender is dropped. `ExternalCommandName` gains `ResetBackoff` for metrics. - Fix: #2143 --- openraft/src/core/raft_core.rs | 28 ++- .../src/core/raft_msg/external_command.rs | 14 ++ openraft/src/core/raft_msg/raft_msg_name.rs | 2 + openraft/src/raft/trigger.rs | 15 ++ openraft/src/replication/backoff_consumer.rs | 103 ++++++++++-- openraft/src/replication/backoff_state.rs | 21 +-- openraft/src/replication/mod.rs | 21 ++- .../src/replication/replication_handle.rs | 18 ++ openraft/src/replication/stream_state.rs | 26 +-- .../stream_state/stream_state_test.rs | 3 +- tests/tests/replication/main.rs | 1 + tests/tests/replication/t52_reset_backoff.rs | 159 ++++++++++++++++++ 12 files changed, 360 insertions(+), 51 deletions(-) create mode 100644 tests/tests/replication/t52_reset_backoff.rs diff --git a/openraft/src/core/raft_core.rs b/openraft/src/core/raft_core.rs index 02aa6766b..f520774eb 100644 --- a/openraft/src/core/raft_core.rs +++ b/openraft/src/core/raft_core.rs @@ -929,6 +929,22 @@ where self.engine.snapshot_handler().trigger_snapshot() } + /// Reset the backoff of the replication streams to `to`, or of every stream if `to` is empty. + pub(crate) fn trigger_reset_backoff(&mut self, to: BatchOf) { + if to.is_empty() { + for (_node_id, repl_handle) in self.replications.iter() { + repl_handle.reset_backoff(); + } + return; + } + + for node_id in to { + if let Some(repl_handle) = self.replications.get(&node_id) { + repl_handle.reset_backoff(); + } + } + } + /// Trigger routine actions that need to be checked after processing messages. /// /// This is called in the main event loop after processing messages and running engine commands. @@ -1120,10 +1136,12 @@ where }; let (replicate_tx, replicate_rx) = C::watch_channel(Replicate::default()); + let (backoff_reset_tx, backoff_reset_rx) = C::watch_channel(()); let event_watcher = self.new_event_watcher(replicate_rx); - let (mut replication_handle, replication_context) = self.new_replication(leader_vote, prog, replicate_tx); + let (mut replication_handle, replication_context) = + self.new_replication(leader_vote, prog, replicate_tx, backoff_reset_tx); let progress = replication_progress::ReplicationProgress { local_committed: self.engine.state.local_committed().cloned(), @@ -1136,6 +1154,7 @@ where network, self.log_store.get_log_reader().await, event_watcher, + backoff_reset_rx, tracing::span!(parent: &self.span, Level::DEBUG, "replication", id=display(&self.id), target=display(&prog.target)), ); @@ -1149,12 +1168,13 @@ where leader_vote: CommittedVoteOf, prog: &TargetProgress, replicate_tx: WatchSenderOf>, + backoff_reset_tx: WatchSenderOf, ) -> (ReplicationHandle, ReplicationContext) { let (cancel_tx, cancel_rx) = C::watch_channel(()); let context = self.new_replication_context(leader_vote, prog, cancel_rx); - let handle = ReplicationHandle::new(prog.progress.data.stream_id, replicate_tx, cancel_tx); + let handle = ReplicationHandle::new(prog.progress.data.stream_id, replicate_tx, cancel_tx, backoff_reset_tx); (handle, context) } @@ -1856,6 +1876,9 @@ where ExternalCommand::TriggerTransferLeader { to } => { self.engine.trigger_transfer_leader(to); } + ExternalCommand::ResetBackoff { to } => { + self.trigger_reset_backoff(to); + } ExternalCommand::AllowNextRevert { to, allow, tx } => { // let res = match self.engine.try_leader_handler() { @@ -2305,6 +2328,7 @@ where // Drop sender to notify the task to shutdown drop(s.replicate_tx); drop(s.cancel_tx); + drop(s.backoff_reset_tx); let target = target.clone(); #[allow(clippy::let_underscore_future)] diff --git a/openraft/src/core/raft_msg/external_command.rs b/openraft/src/core/raft_msg/external_command.rs index c743469f2..bbe322012 100644 --- a/openraft/src/core/raft_msg/external_command.rs +++ b/openraft/src/core/raft_msg/external_command.rs @@ -4,12 +4,15 @@ use std::fmt; use std::sync::Arc; use display_more::DisplayOptionExt; +use display_more::DisplaySliceExt; use crate::RaftTypeConfig; +use crate::batch::Batch; use crate::core::raft_msg::ExternalCommandName; use crate::core::raft_msg::ResultSender; use crate::errors::AllowNextRevertError; use crate::metrics::MetricsRecorder; +use crate::type_config::alias::BatchOf; use crate::type_config::alias::LogIdOf; use crate::type_config::alias::VoteOf; @@ -52,6 +55,9 @@ where C: RaftTypeConfig /// Submit a command to inform RaftCore to transfer leadership to the specified node. TriggerTransferLeader { to: C::NodeId }, + /// Reset the backoff state of the replication to the specified nodes. + ResetBackoff { to: BatchOf }, + /// Allow or not the next revert of the replication to the specified node. AllowNextRevert { to: C::NodeId, @@ -98,6 +104,7 @@ where C: RaftTypeConfig ExternalCommand::Snapshot => ExternalCommandName::Snapshot, ExternalCommand::PurgeLog { .. } => ExternalCommandName::PurgeLog, ExternalCommand::TriggerTransferLeader { .. } => ExternalCommandName::TriggerTransferLeader, + ExternalCommand::ResetBackoff { .. } => ExternalCommandName::ResetBackoff, ExternalCommand::AllowNextRevert { .. } => ExternalCommandName::AllowNextRevert, ExternalCommand::SetMetricsRecorder { .. } => ExternalCommandName::SetMetricsRecorder, ExternalCommand::RefreshServerState { .. } => ExternalCommandName::RefreshServerState, @@ -133,6 +140,13 @@ where C: RaftTypeConfig ExternalCommand::TriggerTransferLeader { to } => { write!(f, "TriggerTransferLeader: to {}", to) } + ExternalCommand::ResetBackoff { to } => { + if to.is_empty() { + write!(f, "ResetBackoff: to all") + } else { + write!(f, "ResetBackoff: to {}", to.as_ref().display()) + } + } ExternalCommand::AllowNextRevert { to, allow, .. } => { write!( f, diff --git a/openraft/src/core/raft_msg/raft_msg_name.rs b/openraft/src/core/raft_msg/raft_msg_name.rs index 9b78f8dcf..e82db8328 100644 --- a/openraft/src/core/raft_msg/raft_msg_name.rs +++ b/openraft/src/core/raft_msg/raft_msg_name.rs @@ -2,6 +2,7 @@ use openraft_macros::VariantName; use openraft_macros::since; /// Enum representing the name of each `ExternalCommand` variant. +#[since(version = "0.10.0", change = "added ResetBackoff variant")] #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] #[derive(VariantName)] #[variant_name(prefix = "Ext::")] @@ -11,6 +12,7 @@ pub enum ExternalCommandName { Snapshot, PurgeLog, TriggerTransferLeader, + ResetBackoff, AllowNextRevert, SetMetricsRecorder, RefreshServerState, diff --git a/openraft/src/raft/trigger.rs b/openraft/src/raft/trigger.rs index 0d39cfb1d..bef460359 100644 --- a/openraft/src/raft/trigger.rs +++ b/openraft/src/raft/trigger.rs @@ -3,11 +3,13 @@ use openraft_macros::since; use crate::RaftTypeConfig; +use crate::batch::Batch; use crate::core::raft_msg::external_command::ExternalCommand; use crate::errors::AllowNextRevertError; use crate::errors::Fatal; use crate::raft::RaftInner; use crate::type_config::TypeConfigExt; +use crate::type_config::alias::BatchOf; use crate::type_config::alias::LogIdOf; use crate::type_config::alias::VoteOf; @@ -107,6 +109,19 @@ where C: RaftTypeConfig self.raft_inner.send_external_command(ExternalCommand::TriggerTransferLeader { to }).await } + /// Request to reset the current replication backoff for the specified nodes. + /// + /// A reset ends the current or next wait of the backoff that is active when it arrives. + /// A backoff that starts later is not affected. Snapshot transfer retries are not affected + /// either. + /// + /// If no node ids are specified, this applies to all replication targets. + #[since(version = "0.10.0")] + pub async fn reset_backoff(&self, to: impl IntoIterator) -> Result<(), Fatal> { + let to = BatchOf::::of(to); + self.raft_inner.send_external_command(ExternalCommand::ResetBackoff { to }).await + } + /// Request the RaftCore to allow to reset replication for a specific node when log revert is /// detected. /// diff --git a/openraft/src/replication/backoff_consumer.rs b/openraft/src/replication/backoff_consumer.rs index 77b65ac07..a19fb6294 100644 --- a/openraft/src/replication/backoff_consumer.rs +++ b/openraft/src/replication/backoff_consumer.rs @@ -1,31 +1,44 @@ -//! Read-only consumer handle for the shared [`Backoff`] owned by +//! Consumer handle for the shared [`Backoff`] owned by //! [`BackoffState`](crate::replication::backoff_state::BackoffState). //! -//! The request-stream generator in -//! [`StreamState`](crate::replication::stream_state::StreamState) samples the next -//! delay before emitting each AppendEntries request. It does not — and must not — -//! enable, disable, or replace the backoff itself; only `BackoffState` does that, -//! maintaining the invariants described in -//! [issue #1723](https://github.com/databendlabs/openraft/issues/1723). +//! The request stream samples a delay before each AppendEntries request and clears +//! the active iterator on reset. `BackoffState` tracks the error rank and enables +//! new iterators after RPC errors. +use std::fmt; use std::sync::Arc; use std::sync::Mutex; use std::time::Duration; +use futures_util::FutureExt; +use rt::OptionalSend; +use rt::WatchReceiver; + +use crate::RaftTypeConfig; use crate::network::Backoff; use crate::replication::EXHAUSTED_BACKOFF_DELAY; +use crate::type_config::TypeConfigExt; +use crate::type_config::alias::WatchReceiverOf; -/// Read-only handle to the backoff shared with the request-stream generator. +/// Handle for sampling delays and clearing the active backoff on reset. /// -/// Exposes exactly one operation — [`next_delay`](Self::next_delay) — so the consumer -/// cannot enable, disable, or replace the backoff. Constructed only by +/// Shares the iterator with +/// [`BackoffState`](crate::replication::backoff_state::BackoffState), which controls when a new +/// iterator is enabled. Constructed by /// [`BackoffState::consumer`](crate::replication::backoff_state::BackoffState::consumer). #[derive(Clone)] -pub(crate) struct BackoffConsumer { +pub(crate) struct BackoffConsumer +where C: RaftTypeConfig +{ pub(crate) inner: Arc>>, + + /// Backoff is reset and the sleep should be aborted. + pub(crate) reset_rx: WatchReceiverOf, } -impl BackoffConsumer { +impl BackoffConsumer +where C: RaftTypeConfig +{ /// Returns the next delay to wait before emitting the next request, or `None` /// if backoff is not currently enabled. Advances the iterator when enabled. pub(crate) fn next_delay(&self) -> Option { @@ -33,14 +46,52 @@ impl BackoffConsumer { let backoff = guard.as_mut()?; Some(backoff.next().unwrap_or(EXHAUSTED_BACKOFF_DELAY)) } + + /// Waits for the next backoff delay if backoff is enabled, or returns immediately. + /// + /// The wait ends early when `cancel` resolves or a reset arrives. A reset also drops the + /// iterator, so later requests are not delayed until `BackoffState` enables a new one. + pub(crate) async fn backoff_if_enabled(&mut self, cancel: Fu) + where + Fu: Future + OptionalSend, + T: fmt::Debug, + { + let Some(sleep_duration) = self.next_delay() else { + return; + }; + + let sleep = C::sleep(sleep_duration); + let reset = self.reset_rx.changed(); + + tracing::debug!("backoff timeout: {:?}", sleep_duration); + + futures_util::select! { + _ = sleep.fuse() => { + tracing::debug!("backoff timeout"); + } + cancel_res = cancel.fuse() => { + tracing::info!("Replication Stream is canceled, res: {:?}, when:(backoff_if_enabled:wait-for-changed)", cancel_res); + } + reset_res = reset.fuse() => { + tracing::info!("Backoff is reset, res: {:?}, when:(backoff_if_enabled:wait-for-changed)", reset_res); + // once reset is received, clear the backoff state. + self.inner.lock().unwrap().take(); + } + } + } } #[cfg(test)] mod tests { use std::time::Duration; + use futures_util::FutureExt; + use rt::WatchSender; + + use crate::engine::testing::UTConfig; use crate::network::Backoff; use crate::replication::backoff_state::BackoffState; + use crate::type_config::TypeConfigExt; fn constant_200ms_backoff() -> Backoff { Backoff::new(std::iter::repeat(Duration::from_millis(200))) @@ -50,8 +101,9 @@ mod tests { /// `BackoffState`, and returns `None` once the state clears it. #[test] fn samples_delay_only_when_enabled() { + let (_tx, rx) = UTConfig::<()>::watch_channel(()); let mut state = BackoffState::new(); - let consumer = state.consumer(); + let consumer = state.consumer::(rx); assert_eq!(consumer.next_delay(), None, "disabled: no delay"); @@ -68,4 +120,29 @@ mod tests { assert_eq!(consumer.next_delay(), None, "after success: no delay"); } + + /// A reset during an active wait stops the delay and clears the iterator. + #[test] + fn reset_interrupts_active_backoff() { + UTConfig::<()>::run(async { + let (tx, rx) = UTConfig::<()>::watch_channel(()); + let mut state = BackoffState::new(); + let mut consumer = state.consumer::(rx); + + state.on_error(100); + state.reconcile(constant_200ms_backoff); + + let wait = consumer.backoff_if_enabled(std::future::pending::<()>()); + let mut wait = Box::pin(wait); + let pending = wait.as_mut().now_or_never(); + assert!(pending.is_none(), "backoff should wait before reset"); + + tx.send(()).unwrap(); + let result = UTConfig::<()>::timeout(Duration::from_millis(100), wait).await; + assert!(result.is_ok(), "reset should interrupt the 200ms wait"); + + let next_delay = consumer.next_delay(); + assert_eq!(next_delay, None); + }); + } } diff --git a/openraft/src/replication/backoff_state.rs b/openraft/src/replication/backoff_state.rs index baae712ba..cd87213df 100644 --- a/openraft/src/replication/backoff_state.rs +++ b/openraft/src/replication/backoff_state.rs @@ -4,9 +4,8 @@ //! should throttle outgoing requests after transient RPC errors: //! //! - `rank`: accumulated error weight, reset on success. -//! - `inner`: the active [`Backoff`] iterator. Handed to the request-stream generator as a -//! [`BackoffConsumer`] so the consumer can sample the next delay without being able to enable or -//! disable the backoff itself. +//! - `inner`: the active [`Backoff`] iterator. Shared with the request-stream generator through a +//! [`BackoffConsumer`], which samples delays and clears the iterator on an explicit reset. //! //! Both pieces must be cleared on success; otherwise a stale [`Backoff`] left over //! from a prior error throttles every subsequent request in the same stream session @@ -20,6 +19,7 @@ use crate::RaftTypeConfig; use crate::errors::RPCError; use crate::network::Backoff; use crate::replication::backoff_consumer::BackoffConsumer; +use crate::type_config::alias::WatchReceiverOf; /// Error-rank threshold above which backoff is enabled. const BACKOFF_RANK_THRESHOLD: u64 = 20; @@ -27,8 +27,8 @@ const BACKOFF_RANK_THRESHOLD: u64 = 20; /// Coordinates rank accumulation and the shared [`Backoff`] iterator for a /// replication session. /// -/// Owned by `ReplicationCore`. The request-stream generator receives a read-only -/// [`BackoffConsumer`] that can only sample delays. +/// Owned by `ReplicationCore`. The request-stream generator receives a +/// [`BackoffConsumer`] that samples delays and handles resets. pub(crate) struct BackoffState { rank: u64, inner: Arc>>, @@ -42,10 +42,12 @@ impl BackoffState { } } - /// Returns a consumer handle that can only sample the current backoff delay. - pub(crate) fn consumer(&self) -> BackoffConsumer { + /// Returns a consumer handle for sampling delays and handling resets. + pub(crate) fn consumer(&self, reset_rx: WatchReceiverOf) -> BackoffConsumer + where C: RaftTypeConfig { BackoffConsumer { inner: self.inner.clone(), + reset_rx, } } @@ -55,9 +57,8 @@ impl BackoffState { /// stream session is active, so this is the only opportunity to drop the stale /// [`Backoff`] before the next RPC in the same session reads it. /// - /// Fast path: `rank == 0` implies `inner` is already `None` (invariant maintained - /// by every writer in this module), so we skip the lock on the common hot path of - /// successive successes. + /// Fast path: `rank == 0` implies `inner` is already `None`, so we skip the + /// lock on the common hot path of successive successes. pub(crate) fn on_success(&mut self) { if self.rank == 0 { return; diff --git a/openraft/src/replication/mod.rs b/openraft/src/replication/mod.rs index 8f4541662..e49cec41c 100644 --- a/openraft/src/replication/mod.rs +++ b/openraft/src/replication/mod.rs @@ -29,7 +29,7 @@ use replication_progress::ReplicationProgress; pub(crate) use replication_session_id::ReplicationSessionId; pub(crate) use response::Progress; -/// Fallback delay used when a [`Backoff`](crate::network::Backoff) iterator is exhausted. +/// Fallback delay used when a [`Backoff`] iterator is exhausted. pub(crate) const EXHAUSTED_BACKOFF_DELAY: Duration = Duration::from_millis(500); use response::ReplicationResult; use stream_state::StreamState; @@ -45,6 +45,7 @@ use crate::display_ext::display_instant::DisplayInstantExt; use crate::errors::RPCError; use crate::errors::ReplicationClosed; use crate::log_id_range::LogIdRange; +use crate::network::Backoff; use crate::network::NetBackoff; use crate::network::NetStreamAppend; use crate::network::RPCOption; @@ -64,6 +65,7 @@ use crate::type_config::alias::InstantOf; use crate::type_config::alias::JoinHandleOf; use crate::type_config::alias::LogIdOf; use crate::type_config::alias::MutexOf; +use crate::type_config::alias::WatchReceiverOf; use crate::type_config::async_runtime::mpsc::MpscSender; /// A task responsible for sending replication events to a target follower in the Raft cluster. @@ -125,6 +127,7 @@ where network: N::Network, log_reader: LS::LogReader, event_watcher: EventWatcher, + backoff_reset_rx: WatchReceiverOf, span: tracing::Span, ) -> JoinHandleOf> { tracing::debug!( @@ -136,6 +139,7 @@ where ); let backoff_state = BackoffState::new(); + let backoff_consumer = backoff_state.consumer(backoff_reset_rx); let this = Self { replication_context: replication_context.clone(), @@ -145,7 +149,7 @@ where log_reader, payload: None, inflight_id: None, - backoff_consumer: backoff_state.consumer(), + backoff_consumer, })), inflight_id: None, event_watcher, @@ -212,7 +216,7 @@ where // to `network`, which is needed to construct the `Backoff` iterator. // If the network returns None, fall back to the policy configured in `Config::backoff`. let config = self.replication_context.config.clone(); - self.backoff_state.reconcile(|| network.backoff().unwrap_or_else(|| config.build_backoff())); + self.reconcile_backoff(|| network.backoff().unwrap_or_else(|| config.build_backoff())).await; let payload = self.select_next_payload().await?; @@ -220,6 +224,17 @@ where } } + async fn reconcile_backoff(&mut self, backoff_factory: impl FnOnce() -> Backoff) { + let mut stream_state = self.stream_state.lock().await; + + self.backoff_state.reconcile(|| { + let backoff = backoff_factory(); + // Discard resets from before this backoff; later resets remain pending. + stream_state.backoff_consumer.reset_rx.borrow_and_update(); + backoff + }); + } + /// Check the two reasons this task must stop before opening another stream: the leader /// dropped its handle, or a different leader has taken over locally. fn ensure_still_leading(&mut self) -> Result<(), ReplicationClosed> { diff --git a/openraft/src/replication/replication_handle.rs b/openraft/src/replication/replication_handle.rs index 2942cfae4..0dc34df51 100644 --- a/openraft/src/replication/replication_handle.rs +++ b/openraft/src/replication/replication_handle.rs @@ -1,3 +1,5 @@ +use rt::WatchSender; + use crate::RaftTypeConfig; use crate::errors::ReplicationClosed; use crate::progress::stream_id::StreamId; @@ -19,6 +21,9 @@ where C: RaftTypeConfig /// Sender for the cancellation signal; dropping this stops replication. pub(crate) cancel_tx: WatchSenderOf, + /// Sender for the backoff reset signal; sending on it ends the active backoff of replication. + pub(crate) backoff_reset_tx: WatchSenderOf, + /// The spawn handle of the `ReplicationCore` task. pub(crate) join_handle: Option>>, @@ -33,6 +38,7 @@ where C: RaftTypeConfig stream_id: StreamId, replicate_tx: WatchSenderOf>, cancel_tx: WatchSenderOf, + backoff_reset_tx: WatchSenderOf, ) -> Self { Self { stream_id, @@ -40,6 +46,18 @@ where C: RaftTypeConfig replicate_tx, snapshot_transmit_handle: None, cancel_tx, + backoff_reset_tx, + } + } + + /// Signal the replication task to end its active backoff. + pub(crate) fn reset_backoff(&self) { + let send_res = self.backoff_reset_tx.send(()); + if send_res.is_err() { + tracing::warn!( + "Failed to reset backoff for replication stream {}: channel closed", + self.stream_id + ); } } } diff --git a/openraft/src/replication/stream_state.rs b/openraft/src/replication/stream_state.rs index a4e44720f..07e4b02f8 100644 --- a/openraft/src/replication/stream_state.rs +++ b/openraft/src/replication/stream_state.rs @@ -49,11 +49,9 @@ where pub(crate) inflight_id: Option, - /// Read-only handle to the shared backoff state, sampled before each request. - /// - /// The consumer can only query the next delay; only `ReplicationCore` (via its - /// owned `BackoffState`) enables or clears the backoff. - pub(crate) backoff_consumer: BackoffConsumer, + /// Handle that samples backoff delays and clears the current backoff on reset. + /// `ReplicationCore` enables new backoff after RPC errors. + pub(crate) backoff_consumer: BackoffConsumer, } impl StreamState @@ -181,23 +179,7 @@ where /// Waits for the backoff duration if backoff is enabled, or returns immediately. async fn backoff_if_enabled(&mut self) { - let Some(sleep_duration) = self.backoff_consumer.next_delay() else { - return; - }; - - let sleep = C::sleep(sleep_duration); - let cancel = self.replication_context.cancel_rx.changed(); - - tracing::debug!("backoff timeout: {:?}", sleep_duration); - - futures_util::select! { - _ = sleep.fuse() => { - tracing::debug!("backoff timeout"); - } - cancel_res = cancel.fuse() => { - tracing::info!("Replication Stream is canceled, res: {:?}, when:(backoff_if_enabled:wait-for-changed)", cancel_res); - } - } + self.backoff_consumer.backoff_if_enabled(self.replication_context.cancel_rx.changed()).await; } /// Advances the payload after generating a request through `sent`. diff --git a/openraft/src/replication/stream_state/stream_state_test.rs b/openraft/src/replication/stream_state/stream_state_test.rs index 5ed3c6b87..01d898f3e 100644 --- a/openraft/src/replication/stream_state/stream_state_test.rs +++ b/openraft/src/replication/stream_state/stream_state_test.rs @@ -103,6 +103,7 @@ fn probe_stops_after_storage_limited_prefix() { let (_committed_tx, committed_rx) = UTConfig::<()>::watch_channel(None); let (_io_tx, io_rx) = UTConfig::<()>::watch_channel(IOId::new_log_io(vote.clone(), range.last)); let (_cancel_tx, cancel_rx) = UTConfig::<()>::watch_channel(()); + let (_backoff_reset_tx, backoff_reset_rx) = UTConfig::<()>::watch_channel(()); let (tx_notify, _rx_notify) = UTConfig::<()>::mpsc(1); StreamState:: { @@ -131,7 +132,7 @@ fn probe_stops_after_storage_limited_prefix() { }, payload: Some(Payload::Probe { log_id_range: range }), inflight_id: Some(inflight_id), - backoff_consumer: BackoffState::new().consumer(), + backoff_consumer: BackoffState::new().consumer(backoff_reset_rx), } }; diff --git a/tests/tests/replication/main.rs b/tests/tests/replication/main.rs index 84a132fdd..07de023a0 100644 --- a/tests/tests/replication/main.rs +++ b/tests/tests/replication/main.rs @@ -10,6 +10,7 @@ mod t20_empty_log_entries; mod t50_append_entries_backoff; mod t50_append_entries_backoff_rejoin; mod t51_backoff_cleared_after_success; +mod t52_reset_backoff; mod t60_feature_loosen_follower_log_revert; mod t61_allow_follower_log_revert; mod t62_follower_clear_restart_recover; diff --git a/tests/tests/replication/t52_reset_backoff.rs b/tests/tests/replication/t52_reset_backoff.rs new file mode 100644 index 000000000..28d09ff59 --- /dev/null +++ b/tests/tests/replication/t52_reset_backoff.rs @@ -0,0 +1,159 @@ +use std::sync::Arc; +use std::sync::atomic::AtomicU64; +use std::sync::atomic::Ordering; +use std::time::Duration; + +use anyhow::Result; +use maplit::btreeset; +use openraft::Config; +use openraft::RPCTypes; +use openraft::async_runtime::WatchReceiver; +use openraft::errors::RPCError; +use openraft::errors::Unreachable; +use openraft::type_config::TypeConfigExt; +use openraft_memstore::TypeConfig; + +use crate::fixtures::RaftRouter; +use crate::fixtures::ut_harness; + +/// A constant 5-second backoff delay, much longer than `timeout()`. +const BACKOFF: &str = "5s ...5s"; + +/// `Trigger::reset_backoff()` ends the active backoff wait, so a reachable follower catches up +/// before the backoff delay ends. +#[tracing::instrument] +#[test_harness::test(harness = ut_harness)] +async fn reset_backoff_ends_active_backoff() -> Result<()> { + let config = Arc::new( + Config { + enable_tick: false, + backoff: BACKOFF.to_string(), + ..Default::default() + } + .validate()?, + ); + + let mut router = RaftRouter::new(config.clone()); + + tracing::info!("--- bring up a 3-node cluster"); + let mut log_index = router.new_cluster(btreeset! {0,1,2}, btreeset! {}).await?; + + let n0 = router.get_raft_handle(&0)?; + let n2 = router.get_raft_handle(&2)?; + + let failures_remaining = Arc::new(AtomicU64::new(0)); + { + let failures_remaining = failures_remaining.clone(); + router + .set_rpc_pre_hook(RPCTypes::AppendEntries, move |_router, _req, _from, to| { + let should_fail = to == 2 && failures_remaining.load(Ordering::SeqCst) > 0; + let res = if should_fail { + failures_remaining.fetch_sub(1, Ordering::SeqCst); + Err(RPCError::Unreachable(Unreachable::::from_string( + "injected", + ))) + } else { + Ok(()) + }; + Box::pin(futures::future::ready(res)) + }) + .await; + } + + tracing::info!(log_index, "--- fail one AppendEntries to node-2 to start its backoff"); + { + failures_remaining.store(1, Ordering::SeqCst); + + router.client_request(0, "foo", 1).await?; + log_index += 1; + } + + tracing::info!(log_index, "--- node-2 is reachable but the backoff keeps it behind"); + { + TypeConfig::sleep(Duration::from_millis(500)).await; + + let failures_remaining = failures_remaining.load(Ordering::SeqCst); + assert_eq!(failures_remaining, 0, "the injected failure is consumed"); + + let n2_metrics = n2.metrics(); + let last_log_index = n2_metrics.borrow_watched().last_log_index; + assert_eq!( + last_log_index, + Some(log_index - 1), + "node-2 has not received the new log" + ); + } + + tracing::info!(log_index, "--- reset the backoff, node-2 catches up within timeout()"); + { + n0.trigger().reset_backoff([2]).await?; + + router.wait(&2, timeout()).applied_index(Some(log_index), "node-2 catches up after reset").await?; + } + + Ok(()) +} + +/// A reset sent while no backoff is active is dropped when the next backoff starts, so it does +/// not skip the first wait of that backoff. +#[tracing::instrument] +#[test_harness::test(harness = ut_harness)] +async fn reset_backoff_before_backoff_is_dropped() -> Result<()> { + let config = Arc::new( + Config { + enable_tick: false, + backoff: BACKOFF.to_string(), + ..Default::default() + } + .validate()?, + ); + + let mut router = RaftRouter::new(config.clone()); + + tracing::info!("--- bring up a 3-node cluster"); + let mut log_index = router.new_cluster(btreeset! {0,1,2}, btreeset! {}).await?; + + let n0 = router.get_raft_handle(&0)?; + + tracing::info!(log_index, "--- reset the backoff to node-2 while its stream is healthy"); + { + n0.trigger().reset_backoff([2]).await?; + } + + let attempts = Arc::new(AtomicU64::new(0)); + + tracing::info!(log_index, "--- cut off node-2 and write a log to start its backoff"); + { + let attempts = attempts.clone(); + router + .set_rpc_pre_hook(RPCTypes::AppendEntries, move |_router, _req, _from, to| { + let res = if to == 2 { + attempts.fetch_add(1, Ordering::SeqCst); + Err(RPCError::Unreachable(Unreachable::::from_string( + "injected", + ))) + } else { + Ok(()) + }; + Box::pin(futures::future::ready(res)) + }) + .await; + + router.client_request(0, "foo", 1).await?; + log_index += 1; + } + + tracing::info!(log_index, "--- the earlier reset must not skip the first backoff wait"); + { + TypeConfig::sleep(Duration::from_millis(1_000)).await; + + let attempts = attempts.load(Ordering::SeqCst); + assert_eq!(attempts, 1, "node-2 is retried only after the backoff delay"); + } + + Ok(()) +} + +fn timeout() -> Option { + Some(Duration::from_millis(1_000)) +} From b2fd62c30a6db9eb8baedd82675fdd65f2322524 Mon Sep 17 00:00:00 2001 From: Zhang Yanpo Date: Tue, 29 Sep 2026 20:50:29 +0800 Subject: [PATCH 2/2] feat: replication: reset target backoff on transfer # Summary Leadership transfer now ends the target's active AppendEntries backoff so a restarted node can catch up before the transfer request times out. # Details A reachable transfer target may still be waiting through a retry delay after an earlier RPC failure. `RaftCore` signals that target's reset channel before it starts the leadership transfer. `Config::reset_backoff_on_transfer_leader` is a new `Option`. `None` and `Some(true)` enable the reset; `Some(false)` keeps the existing backoff delay. Configs written before this field was added also enable the reset. --- openraft/src/config/config.rs | 19 +++++ openraft/src/config/config_clap_test.rs | 23 ++++++ openraft/src/config/config_test.rs | 14 ++++ openraft/src/core/raft_core.rs | 4 + openraft/src/raft/trigger.rs | 5 ++ tests/tests/replication/t52_reset_backoff.rs | 83 ++++++++++++++++++++ 6 files changed, 148 insertions(+) diff --git a/openraft/src/config/config.rs b/openraft/src/config/config.rs index e5a5c0197..455d30a5c 100644 --- a/openraft/src/config/config.rs +++ b/openraft/src/config/config.rs @@ -178,6 +178,7 @@ impl SnapshotPolicy { /// /// [`Raft::new`]: crate::Raft::new #[since] +#[since(version = "0.10.0", change = "added reset_backoff_on_transfer_leader option")] #[since(version = "0.10.0", change = "added opt-in quorum-loss inactivity setting")] #[derive(Clone, Debug)] #[cfg_attr(feature = "clap", derive(Parser))] @@ -522,6 +523,18 @@ pub struct Config { #[cfg_attr(feature = "clap", clap(long, default_value = DEFAULTS.backoff))] pub backoff: String, + /// Whether transferring leadership resets the target's active replication backoff. + /// + /// `None` (the default) enables the reset. Set this to `Some(false)` to keep the target's + /// current backoff delay during a leadership transfer. + #[since(version = "0.10.0")] + #[cfg_attr(feature = "clap", clap(long, + action = clap::ArgAction::Set, + num_args = 0..=1, + default_missing_value = "true" + ))] + pub reset_backoff_on_transfer_leader: Option, + /// Whether to allow to reset the replication progress to `None`, when the /// follower's log is found reverted to an early state. **Do not enable this in production** /// unless you know what you are doing. @@ -611,6 +624,7 @@ impl Default for Config { quorum_loss_probe_interval: DEFAULTS.quorum_loss_probe_interval, enable_pre_vote: DEFAULTS.enable_pre_vote, backoff: DEFAULTS.backoff.to_string(), + reset_backoff_on_transfer_leader: None, allow_log_reversion: None, enable_leader_restore: None, } @@ -660,6 +674,11 @@ impl Config { self.enable_leader_restore.unwrap_or(true) } + /// Whether a leadership transfer resets the target's active replication backoff. + pub(crate) fn get_reset_backoff_on_transfer_leader(&self) -> bool { + self.reset_backoff_on_transfer_leader.unwrap_or(true) + } + /// Whether a follower runs a Pre-Vote round before starting a real election. /// /// Evaluates the [`enable_pre_vote`](Self::enable_pre_vote) option: `None` is treated as diff --git a/openraft/src/config/config_clap_test.rs b/openraft/src/config/config_clap_test.rs index 1da81ed7f..a483d77df 100644 --- a/openraft/src/config/config_clap_test.rs +++ b/openraft/src/config/config_clap_test.rs @@ -215,6 +215,29 @@ fn test_config_enable_pre_vote() -> anyhow::Result<()> { Ok(()) } +#[test] +fn test_config_reset_backoff_on_transfer_leader() -> anyhow::Result<()> { + let config = Config::build(&["foo"])?; + assert_eq!(None, config.reset_backoff_on_transfer_leader); + let enabled = config.get_reset_backoff_on_transfer_leader(); + assert!(enabled); + + let config = Config::build(&["foo", "--reset-backoff-on-transfer-leader=false"])?; + assert_eq!(Some(false), config.reset_backoff_on_transfer_leader); + let enabled = config.get_reset_backoff_on_transfer_leader(); + assert!(!enabled); + + let config = Config::build(&["foo", "--reset-backoff-on-transfer-leader=true"])?; + assert_eq!(Some(true), config.reset_backoff_on_transfer_leader); + let enabled = config.get_reset_backoff_on_transfer_leader(); + assert!(enabled); + + let config = Config::build(&["foo", "--reset-backoff-on-transfer-leader"])?; + assert_eq!(Some(true), config.reset_backoff_on_transfer_leader); + + Ok(()) +} + #[test] fn test_config_allow_log_reversion() -> anyhow::Result<()> { let config = Config::build(&["foo", "--allow-log-reversion=false"])?; diff --git a/openraft/src/config/config_test.rs b/openraft/src/config/config_test.rs index 0ffa916be..a7caaeaa0 100644 --- a/openraft/src/config/config_test.rs +++ b/openraft/src/config/config_test.rs @@ -75,6 +75,20 @@ fn test_removed_leader_step_down_serde_default() -> anyhow::Result<()> { Ok(()) } +#[cfg(feature = "serde")] +#[test] +fn test_reset_backoff_on_transfer_leader_serde_default() -> anyhow::Result<()> { + let mut value = serde_json::to_value(Config::default())?; + value.as_object_mut().unwrap().remove("reset_backoff_on_transfer_leader"); + + let config: Config = serde_json::from_value(value)?; + assert_eq!(None, config.reset_backoff_on_transfer_leader); + let enabled = config.get_reset_backoff_on_transfer_leader(); + assert!(enabled); + + Ok(()) +} + #[test] fn test_invalid_election_timeout_config_produces_expected_error() { let config = Config { diff --git a/openraft/src/core/raft_core.rs b/openraft/src/core/raft_core.rs index f520774eb..cd17c0873 100644 --- a/openraft/src/core/raft_core.rs +++ b/openraft/src/core/raft_core.rs @@ -1874,6 +1874,10 @@ where self.engine.trigger_purge_log(upto); } ExternalCommand::TriggerTransferLeader { to } => { + if self.config.get_reset_backoff_on_transfer_leader() { + let targets = BatchOf::::of([to.clone()]); + self.trigger_reset_backoff(targets); + } self.engine.trigger_transfer_leader(to); } ExternalCommand::ResetBackoff { to } => { diff --git a/openraft/src/raft/trigger.rs b/openraft/src/raft/trigger.rs index bef460359..f4ca06912 100644 --- a/openraft/src/raft/trigger.rs +++ b/openraft/src/raft/trigger.rs @@ -104,7 +104,12 @@ where C: RaftTypeConfig /// Submit a command to inform RaftCore to transfer leadership to the specified node. /// + /// By default, this also resets the target's active replication backoff. Set + /// [`Config::reset_backoff_on_transfer_leader`] to `Some(false)` to disable the reset. + /// /// If this node is not a Leader, it is just ignored. + /// + /// [`Config::reset_backoff_on_transfer_leader`]: crate::Config::reset_backoff_on_transfer_leader pub async fn transfer_leader(&self, to: C::NodeId) -> Result<(), Fatal> { self.raft_inner.send_external_command(ExternalCommand::TriggerTransferLeader { to }).await } diff --git a/tests/tests/replication/t52_reset_backoff.rs b/tests/tests/replication/t52_reset_backoff.rs index 28d09ff59..b32983491 100644 --- a/tests/tests/replication/t52_reset_backoff.rs +++ b/tests/tests/replication/t52_reset_backoff.rs @@ -154,6 +154,89 @@ async fn reset_backoff_before_backoff_is_dropped() -> Result<()> { Ok(()) } +/// A leadership transfer resumes replication to a reachable target by default, while an explicit +/// `Some(false)` leaves the target in backoff. +#[tracing::instrument] +#[test_harness::test(harness = ut_harness)] +async fn transfer_leader_resets_target_backoff_by_default() -> Result<()> { + for reset_backoff_on_transfer_leader in [Some(false), None] { + let config = Arc::new( + Config { + enable_tick: false, + backoff: BACKOFF.to_string(), + reset_backoff_on_transfer_leader, + ..Default::default() + } + .validate()?, + ); + + let mut router = RaftRouter::new(config.clone()); + + tracing::info!(?reset_backoff_on_transfer_leader, "--- bring up a 3-node cluster"); + let mut log_index = router.new_cluster(btreeset! {0,1,2}, btreeset! {}).await?; + + let n0 = router.get_raft_handle(&0)?; + let n2 = router.get_raft_handle(&2)?; + let failures_remaining = Arc::new(AtomicU64::new(0)); + + tracing::info!(log_index, "--- fail one AppendEntries to node-2 to start its backoff"); + { + let hook_failures_remaining = failures_remaining.clone(); + router + .set_rpc_pre_hook(RPCTypes::AppendEntries, move |_router, _req, _from, to| { + let should_fail = to == 2 && hook_failures_remaining.load(Ordering::SeqCst) > 0; + let res = if should_fail { + hook_failures_remaining.fetch_sub(1, Ordering::SeqCst); + Err(RPCError::Unreachable(Unreachable::::from_string( + "injected", + ))) + } else { + Ok(()) + }; + Box::pin(futures::future::ready(res)) + }) + .await; + + failures_remaining.store(1, Ordering::SeqCst); + router.client_request(0, "foo", 1).await?; + log_index += 1; + } + + tracing::info!(log_index, "--- node-2 is reachable but still in backoff"); + { + TypeConfig::sleep(Duration::from_millis(500)).await; + + let failures_remaining = failures_remaining.load(Ordering::SeqCst); + assert_eq!(failures_remaining, 0, "the injected failure is consumed"); + + let n2_metrics = n2.metrics(); + let last_log_index = n2_metrics.borrow_watched().last_log_index; + assert_eq!(last_log_index, Some(log_index - 1), "node-2 is behind during backoff"); + } + + tracing::info!(log_index, ?reset_backoff_on_transfer_leader, "--- transfer to node-2"); + { + n0.trigger().transfer_leader(2).await?; + + if reset_backoff_on_transfer_leader == Some(false) { + TypeConfig::sleep(Duration::from_millis(1_000)).await; + + let n2_metrics = n2.metrics(); + let last_log_index = n2_metrics.borrow_watched().last_log_index; + assert_eq!( + last_log_index, + Some(log_index - 1), + "disabled reset keeps node-2 in backoff" + ); + } else { + router.wait(&2, timeout()).applied_index(Some(log_index), "transfer resumes replication").await?; + } + } + } + + Ok(()) +} + fn timeout() -> Option { Some(Duration::from_millis(1_000)) }