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
82 changes: 1 addition & 81 deletions d-engine-core/src/raft_role/leader_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ use crate::alias::MOF;
use crate::alias::ROF;
use crate::alias::SMHOF;
use crate::alias::TROF;
use crate::ensure_safe_join;
use crate::event::ClientCmd;
use crate::network::Transport;
use crate::stream::create_production_snapshot_stream;
Expand Down Expand Up @@ -2737,69 +2736,6 @@ impl<T: TypeConfig> LeaderState<T> {
Ok(())
}

#[allow(dead_code)]
pub async fn batch_promote_learners(
&mut self,
ready_learners_ids: Vec<u32>,
ctx: &RaftContext<T>,
role_tx: &mpsc::UnboundedSender<RoleEvent>,
) -> Result<()> {
// 1. Determine optimal promotion status based on quorum safety
debug!("1. Determine optimal promotion status based on quorum safety");
let membership = ctx.membership();
let current_voters = membership.voters().await.len();
// let syncings = membership.nodes_with_status(NodeStatus::Syncing).len();
let new_active_count = current_voters + ready_learners_ids.len();

// Determine target status based on quorum safety
trace!(
?current_voters,
?ready_learners_ids,
"[Node-{}] new_active_count: {}",
self.node_id(),
new_active_count
);
let target_status = if ensure_safe_join(self.node_id(), new_active_count).is_ok() {
trace!(
"Going to update nodes-{:?} status to Active",
ready_learners_ids
);
NodeStatus::Active
} else {
trace!(
"Not enough quorum to promote learners: {:?}",
ready_learners_ids
);
return Ok(());
};

// 2. Create configuration change payload
debug!("2. Create configuration change payload");
let config_change = Change::BatchPromote(BatchPromote {
node_ids: ready_learners_ids.clone(),
new_status: target_status as i32,
});

info!(?config_change, "Replicating cluster config");

// 3. Submit single config change for all ready learners (fire-and-forget).
// Commit confirmation comes via the normal Raft flow; maintenance re-runs on next tick.
debug!("3. Submit single config change for all ready learners");
self.execute_request_immediately(
RaftRequestWithSignal {
id: { rand::distr::Alphanumeric.sample_string(&mut rand::rng(), 21) },
payloads: vec![EntryPayload::config(config_change)],
senders: vec![],
wait_for_apply_event: false,
},
ctx,
role_tx,
)
.await?;

Ok(())
}

/// Calculate new submission index
#[instrument(skip(self))]
fn calculate_new_commit_index(
Expand Down Expand Up @@ -2831,22 +2767,6 @@ impl<T: TypeConfig> LeaderState<T> {
}
}

#[allow(dead_code)]
fn if_update_commit_index(
&self,
new_commit_index_option: Option<u64>,
) -> (bool, u64) {
let current_commit_index = self.commit_index();
if let Some(new_commit_index) = new_commit_index_option {
debug!("Leader::update_commit_index: {:?}", new_commit_index);
if current_commit_index < new_commit_index {
return (true, new_commit_index);
}
}
debug!("Leader::update_commit_index: false");
(false, current_commit_index)
}

/// Calculate safe read index for linearizable reads.
///
/// Returns max(commitIndex, noopIndex) to ensure:
Expand Down Expand Up @@ -3690,7 +3610,7 @@ impl<T: TypeConfig> LeaderState<T> {
}

/// Remove non-Active zombie nodes that exceed failure threshold
async fn conditionally_purge_zombie_nodes(
pub(super) async fn conditionally_purge_zombie_nodes(
&mut self,
role_tx: &mpsc::UnboundedSender<RoleEvent>,
ctx: &RaftContext<T>,
Expand Down
56 changes: 56 additions & 0 deletions d-engine-core/src/raft_role/leader_state_test/backpressure_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -355,3 +355,59 @@ async fn test_backpressure_all_read_policies() {
"Should reject with ResourceExhausted"
);
}

// ============================================================================
// BackpressureMetrics Unit Tests
// ============================================================================

/// BackpressureMetrics::record_rejection with enabled=true: write and read paths.
/// Verifies both branches complete without panic when metrics recording is active.
#[test]
fn test_backpressure_metrics_record_rejection_enabled() {
use crate::raft_role::leader_state::BackpressureMetrics;
let m = BackpressureMetrics::new(1, true, 1);
// Both calls must not panic
m.record_rejection(true); // write rejection
m.record_rejection(false); // read rejection
}

/// BackpressureMetrics::record_rejection with enabled=false: both paths are no-ops.
#[test]
fn test_backpressure_metrics_record_rejection_disabled() {
use crate::raft_role::leader_state::BackpressureMetrics;
let m = BackpressureMetrics::new(1, false, 1);
m.record_rejection(true);
m.record_rejection(false);
}

/// BackpressureMetrics::record_buffer_utilization with enabled=true: sampling logic.
/// sample_rate=1 means every call records; verifies write and read labels both work.
#[test]
fn test_backpressure_metrics_record_buffer_utilization_enabled() {
use crate::raft_role::leader_state::BackpressureMetrics;
let m = BackpressureMetrics::new(1, true, 1);
m.record_buffer_utilization(0.0, true); // write, empty buffer
m.record_buffer_utilization(0.5, true); // write, half full
m.record_buffer_utilization(1.0, true); // write, full
m.record_buffer_utilization(0.5, false); // read, half full
}

/// BackpressureMetrics::record_buffer_utilization with enabled=false: no-op.
#[test]
fn test_backpressure_metrics_record_buffer_utilization_disabled() {
use crate::raft_role::leader_state::BackpressureMetrics;
let m = BackpressureMetrics::new(1, false, 1);
m.record_buffer_utilization(0.5, true);
m.record_buffer_utilization(0.5, false);
}

/// BackpressureMetrics: sample_rate>1 skips most recordings (counter-based sampling).
/// With sample_rate=3, only every 3rd call records; all calls must not panic.
#[test]
fn test_backpressure_metrics_sampling_rate() {
use crate::raft_role::leader_state::BackpressureMetrics;
let m = BackpressureMetrics::new(1, true, 3);
for i in 0..10 {
m.record_buffer_utilization(i as f64 / 10.0, i % 2 == 0);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -936,3 +936,35 @@ async fn test_drain_read_buffer_clears_pending_reads_on_stepdown() {
"drained reads must get Unavailable error"
);
}

// ============================================================================
// StepDownSelfRemoved Tests
// ============================================================================

/// Leader receives StepDownSelfRemoved → sends BecomeFollower and returns Ok.
///
/// Per Raft protocol, after a leader commits its own removal from the cluster
/// it must immediately step down to Follower. This test verifies:
/// 1. handle_raft_event returns Ok(())
/// 2. A BecomeFollower(None) event is sent on role_tx
#[tokio::test]
async fn test_step_down_self_removed_sends_become_follower() {
let (_graceful_tx, graceful_rx) = watch::channel(());
let raft_context = mock_raft_context("/tmp/test_step_down_self_removed", graceful_rx, None);

let node_config = raft_context.node_config();
let mut leader_state = LeaderState::<MockTypeConfig>::new(1, node_config);

let (role_tx, mut role_rx) = mpsc::unbounded_channel();
let result = leader_state
.handle_raft_event(RaftEvent::StepDownSelfRemoved, &raft_context, role_tx)
.await;

assert!(result.is_ok(), "StepDownSelfRemoved must return Ok");

let event = role_rx.try_recv().expect("BecomeFollower event must be sent");
assert!(
matches!(event, RoleEvent::BecomeFollower(None)),
"Expected BecomeFollower(None), got: {event:?}"
);
}
Loading
Loading