Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions openraft/src/config/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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))]
Expand Down Expand Up @@ -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<bool>,

/// 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.
Expand Down Expand Up @@ -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,
}
Expand Down Expand Up @@ -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
Expand Down
23 changes: 23 additions & 0 deletions openraft/src/config/config_clap_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"])?;
Expand Down
14 changes: 14 additions & 0 deletions openraft/src/config/config_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
32 changes: 30 additions & 2 deletions openraft/src/core/raft_core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<C, C::NodeId>) {
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.
Expand Down Expand Up @@ -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(),
Expand All @@ -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)),
);

Expand All @@ -1149,12 +1168,13 @@ where
leader_vote: CommittedVoteOf<C>,
prog: &TargetProgress<C>,
replicate_tx: WatchSenderOf<C, Replicate<C>>,
backoff_reset_tx: WatchSenderOf<C, ()>,
) -> (ReplicationHandle<C>, ReplicationContext<C>) {
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)
}
Expand Down Expand Up @@ -1854,8 +1874,15 @@ where
self.engine.trigger_purge_log(upto);
}
ExternalCommand::TriggerTransferLeader { to } => {
if self.config.get_reset_backoff_on_transfer_leader() {
let targets = BatchOf::<C, _>::of([to.clone()]);
self.trigger_reset_backoff(targets);
}
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() {
Expand Down Expand Up @@ -2305,6 +2332,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)]
Expand Down
14 changes: 14 additions & 0 deletions openraft/src/core/raft_msg/external_command.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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<C, C::NodeId> },

/// Allow or not the next revert of the replication to the specified node.
AllowNextRevert {
to: C::NodeId,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions openraft/src/core/raft_msg/raft_msg_name.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::")]
Expand All @@ -11,6 +12,7 @@ pub enum ExternalCommandName {
Snapshot,
PurgeLog,
TriggerTransferLeader,
ResetBackoff,
AllowNextRevert,
SetMetricsRecorder,
RefreshServerState,
Expand Down
20 changes: 20 additions & 0 deletions openraft/src/raft/trigger.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -102,11 +104,29 @@ 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<C>> {
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<Item = C::NodeId>) -> Result<(), Fatal<C>> {
let to = BatchOf::<C, _>::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.
///
Expand Down
Loading
Loading