Skip to content

Commit edafc6b

Browse files
authored
Merge pull request #45 from NodeDB-Lab/cluster
feat(cluster): distributed cluster infrastructure
2 parents bc78809 + 5f6abeb commit edafc6b

91 files changed

Lines changed: 6760 additions & 230 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

Cargo.lock

Lines changed: 18 additions & 18 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ members = [
2222
resolver = "2"
2323

2424
[workspace.package]
25-
version = "0.0.3"
25+
version = "0.0.4"
2626
edition = "2024"
2727
rust-version = "1.94"
2828
license = "BUSL-1.1"

nodedb-cluster/src/bootstrap/bootstrap_fn.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,8 @@ pub(super) fn bootstrap(config: &ClusterConfig, catalog: &ClusterCatalog) -> Res
3535
);
3636

3737
// Create MultiRaft with all groups (single-node, no peers).
38-
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone());
38+
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone())
39+
.with_election_timeout(config.election_timeout_min, config.election_timeout_max);
3940
for group_id in routing.group_ids() {
4041
multi_raft.add_group(group_id, vec![])?;
4142
}
@@ -81,6 +82,7 @@ fn generate_cluster_id() -> u64 {
8182
mod tests {
8283
use super::*;
8384
use crate::catalog::ClusterCatalog;
85+
use std::time::Duration;
8486

8587
fn temp_catalog() -> (tempfile::TempDir, ClusterCatalog) {
8688
let dir = tempfile::tempdir().unwrap();
@@ -102,6 +104,8 @@ mod tests {
102104
force_bootstrap: false,
103105
join_retry: Default::default(),
104106
swim_udp_addr: None,
107+
election_timeout_min: Duration::from_millis(150),
108+
election_timeout_max: Duration::from_millis(300),
105109
};
106110

107111
let state = bootstrap(&config, &catalog).unwrap();

nodedb-cluster/src/bootstrap/config.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,10 @@ pub struct ClusterConfig {
9191
/// [`crate::spawn_swim`] after the cluster is up and feed the
9292
/// seed list from `seed_nodes`.
9393
pub swim_udp_addr: Option<SocketAddr>,
94+
/// Raft election timeout range. Controls how long a follower waits
95+
/// before starting an election after losing contact with the leader.
96+
pub election_timeout_min: Duration,
97+
pub election_timeout_max: Duration,
9498
}
9599

96100
/// Result of cluster startup — everything needed to run the Raft loop.

nodedb-cluster/src/bootstrap/join.rs

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -288,7 +288,8 @@ fn apply_join_response(
288288
// learners). A learner-started group boots in the `Learner`
289289
// role and will not run an election until a subsequent
290290
// `PromoteLearner` conf change is applied.
291-
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone());
291+
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone())
292+
.with_election_timeout(config.election_timeout_min, config.election_timeout_max);
292293
for g in &resp.groups {
293294
let is_voter = g.members.contains(&config.node_id);
294295
let is_learner = g.learners.contains(&config.node_id);
@@ -450,6 +451,8 @@ mod tests {
450451
force_bootstrap: false,
451452
join_retry: Default::default(),
452453
swim_udp_addr: None,
454+
election_timeout_min: Duration::from_millis(150),
455+
election_timeout_max: Duration::from_millis(300),
453456
};
454457
let state1 = bootstrap(&config1, &catalog1).unwrap();
455458

@@ -499,6 +502,8 @@ mod tests {
499502
force_bootstrap: false,
500503
join_retry: Default::default(),
501504
swim_udp_addr: None,
505+
election_timeout_min: Duration::from_millis(150),
506+
election_timeout_max: Duration::from_millis(300),
502507
};
503508

504509
let lifecycle = ClusterLifecycleTracker::new();

nodedb-cluster/src/bootstrap/probe.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -223,6 +223,8 @@ mod tests {
223223
force_bootstrap: false,
224224
join_retry: Default::default(),
225225
swim_udp_addr: None,
226+
election_timeout_min: Duration::from_millis(150),
227+
election_timeout_max: Duration::from_millis(300),
226228
}
227229
}
228230

nodedb-cluster/src/bootstrap/restart.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,8 @@ pub(super) fn restart(
3535
// as a learner on restart; dropping the group entirely would
3636
// leave the node permanently without any copy of it and
3737
// silently broken.
38-
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone());
38+
let mut multi_raft = MultiRaft::new(config.node_id, routing.clone(), config.data_dir.clone())
39+
.with_election_timeout(config.election_timeout_min, config.election_timeout_max);
3940
for (group_id, info) in routing.group_members() {
4041
let is_voter = info.members.contains(&config.node_id);
4142
let is_learner = info.learners.contains(&config.node_id);
@@ -91,6 +92,7 @@ mod tests {
9192
use super::super::bootstrap_fn::bootstrap;
9293
use super::*;
9394
use crate::catalog::ClusterCatalog;
95+
use std::time::Duration;
9496

9597
fn temp_catalog() -> (tempfile::TempDir, ClusterCatalog) {
9698
let dir = tempfile::tempdir().unwrap();
@@ -112,6 +114,8 @@ mod tests {
112114
force_bootstrap: false,
113115
join_retry: Default::default(),
114116
swim_udp_addr: None,
117+
election_timeout_min: Duration::from_millis(150),
118+
election_timeout_max: Duration::from_millis(300),
115119
};
116120

117121
// Bootstrap first.

nodedb-cluster/src/circuit_breaker.rs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -145,6 +145,26 @@ impl CircuitBreaker {
145145
.unwrap_or(CircuitState::Closed)
146146
}
147147

148+
/// Return the ids of every peer whose breaker is currently Open.
149+
///
150+
/// Used by the reachability driver to find peers that need an
151+
/// active probe — without a periodic poke these peers never
152+
/// transition back to HalfOpen (no traffic → no `check()` call
153+
/// → no cooldown re-evaluation).
154+
pub fn open_peers(&self) -> Vec<u64> {
155+
let peers = self.peers.read().unwrap_or_else(|p| p.into_inner());
156+
peers
157+
.iter()
158+
.filter_map(|(id, b)| {
159+
if b.state == CircuitState::Open {
160+
Some(*id)
161+
} else {
162+
None
163+
}
164+
})
165+
.collect()
166+
}
167+
148168
/// Get consecutive failure count for a peer.
149169
pub fn failure_count(&self, peer: u64) -> u32 {
150170
let peers = self.peers.read().unwrap_or_else(|p| p.into_inner());

0 commit comments

Comments
 (0)