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
20 changes: 7 additions & 13 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 0 additions & 6 deletions config/base/raft.toml
Original file line number Diff line number Diff line change
Expand Up @@ -122,9 +122,6 @@ cluster_healthcheck_probe_service_name = "d_engine.server.cluster.ClusterManagem
# Timeout for leader to keep retrying no-op entries to confirm leadership (seconds).
verify_leadership_persistent_timeout = { secs = 3600 }

# Interval for membership maintenance tasks (seconds).
membership_maintenance_interval = { secs = 30 }

[raft.membership.zombie]
# Consecutive connection failures before marking a node as zombie.
threshold = 3
Expand All @@ -136,9 +133,6 @@ purge_interval = { secs = 30 }
# Duration after which a learner is considered stale if not promoted (seconds).
stale_learner_threshold = { secs = 300 }

# Interval for checking stale learners (seconds).
stale_check_interval = { secs = 30 }

[raft.persistence]
# Persistence strategy: "MemFirst" — entries written to memory first, flushed to disk async.
# For per-write durability (process-crash safe on every write), set:
Expand Down
2 changes: 1 addition & 1 deletion d-engine-core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ tokio-stream = "0.1.16"
# Packing/unpacking TAR archives
astral-tokio-tar = "0.6"
rand = "0.9"
lru = "0.16"
lru = "0.17"
memmap2 = "0.9.5"
crc32fast = "1.4.2"
# used in stream
Expand Down
1 change: 0 additions & 1 deletion d-engine-core/benches/leader_state_bench.rs
Original file line number Diff line number Diff line change
Expand Up @@ -255,7 +255,6 @@ impl BenchFixture {
nodes: vec![],
current_leader_id: None,
});
membership.expect_get_zombie_candidates().returning(Vec::new);
membership.expect_get_peers_id_with_condition().returning(|_| vec![]);
membership.expect_is_single_node_cluster().returning(|| false);
membership.expect_initial_cluster_size().returning(|| 3);
Expand Down
16 changes: 0 additions & 16 deletions d-engine-core/src/config/raft.rs
Original file line number Diff line number Diff line change
Expand Up @@ -370,9 +370,6 @@ pub struct MembershipConfig {
#[serde(default = "default_verify_leadership_persistent_timeout")]
pub verify_leadership_persistent_timeout: Duration,

#[serde(default = "default_membership_maintenance_interval")]
pub membership_maintenance_interval: Duration,

#[serde(default)]
pub zombie: ZombieConfig,

Expand All @@ -385,7 +382,6 @@ impl Default for MembershipConfig {
Self {
cluster_healthcheck_probe_service_name: default_probe_service(),
verify_leadership_persistent_timeout: default_verify_leadership_persistent_timeout(),
membership_maintenance_interval: default_membership_maintenance_interval(),
zombie: ZombieConfig::default(),
promotion: PromotionConfig::default(),
}
Expand All @@ -395,11 +391,6 @@ fn default_probe_service() -> String {
"d_engine.server.cluster.ClusterManagementService".to_string()
}

// 30 seconds
fn default_membership_maintenance_interval() -> Duration {
Duration::from_secs(30)
}

/// Default timeout for leader to keep verifying its leadership.
///
/// In Raft, the leader may retry sending no-op entries to confirm it still holds leadership.
Expand Down Expand Up @@ -731,15 +722,12 @@ fn default_zombie_purge_interval() -> Duration {
pub struct PromotionConfig {
#[serde(default = "default_stale_learner_threshold")]
pub stale_learner_threshold: Duration,
#[serde(default = "default_stale_check_interval")]
pub stale_check_interval: Duration,
}

impl Default for PromotionConfig {
fn default() -> Self {
Self {
stale_learner_threshold: default_stale_learner_threshold(),
stale_check_interval: default_stale_check_interval(),
}
}
}
Expand All @@ -748,10 +736,6 @@ impl Default for PromotionConfig {
fn default_stale_learner_threshold() -> Duration {
Duration::from_secs(300)
}
// 30 seconds
fn default_stale_check_interval() -> Duration {
Duration::from_secs(30)
}
/// Defines how Raft log entries are persisted and accessed.
///
/// All strategies use a configurable [`FlushPolicy`] to control when memory contents
Expand Down
5 changes: 5 additions & 0 deletions d-engine-core/src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,11 @@ pub enum RoleEvent {
peer_id: u32,
},

/// A peer's connection failure count crossed zombie_threshold.
/// Emitted by RaftHealthMonitor (server layer) via an injected Sender<u32>.
/// Leader responds by proposing a BatchRemove config change for that node.
ZombieDetected(u32),

/// Fatal error from SM worker — node must shutdown.
/// Sent via role_tx (P2) so it is not blocked behind external RPCs on event_tx (P4).
FatalError {
Expand Down
2 changes: 0 additions & 2 deletions d-engine-core/src/membership.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,8 +208,6 @@ where
index: u64,
);

async fn get_zombie_candidates(&self) -> Vec<u32>;

/// If new node could rejoin the cluster again
async fn can_rejoin(
&self,
Expand Down
9 changes: 9 additions & 0 deletions d-engine-core/src/raft.rs
Original file line number Diff line number Diff line change
Expand Up @@ -601,6 +601,15 @@ where
self.role.handle_peer_stream_error(peer_id);
}

RoleEvent::ZombieDetected(node_id) => {
debug!(%node_id, "ZombieDetected: forwarding to leader for BatchRemove");
if let Err(e) =
self.role.handle_zombie_detected(node_id, &self.role_tx, &self.ctx).await
{
error!(%node_id, ?e, "handle_zombie_detected failed");
}
}

RoleEvent::SnapshotPushCompleted { peer_id, success } => {
debug!(%peer_id, %success, "SnapshotPushCompleted");
if success {
Expand Down
Loading
Loading