From 2a24942f58b014e4571df4d42ba04cf78a74befb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Sat, 3 Oct 2026 14:27:13 -0700 Subject: [PATCH 1/2] fix raft vote reset, election timer start and startup grace --- CHANGELOG.md | 15 +++ Cargo.lock | 4 +- crates/mqdb-cli/Cargo.toml | 2 +- crates/mqdb-cluster/Cargo.toml | 2 +- .../src/cluster/raft/coordinator/mod.rs | 2 +- .../src/cluster/raft/coordinator/tests.rs | 29 ++++++ crates/mqdb-cluster/src/cluster/raft/node.rs | 92 ++++++++++++++++++- crates/mqdb-cluster/src/cluster/raft/state.rs | 6 +- .../tests/cluster_integration_test.rs | 6 ++ crates/mqdb-cluster/tests/cluster_test.rs | 3 + 10 files changed, 151 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index dbbea9a2..a7f38b8b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,21 @@ All notable changes to this project will be documented in this file. Each entry lists the date and the crate versions that were released. +## 2026-10-03 — mqdb-cluster 0.4.17, mqdb-cli 0.8.41 + +### Fixed + +- **Raft no longer elects two leaders in one term when nodes start with partial peer lists** (#155). Three defects combined: + - A candidate that stepped down to a same-term leader cleared its vote and could vote again in that term. It now keeps its vote. + - The election timer started already expired, so a joining node campaigned within about 100 ms of starting. The timer now starts on the first tick, with a per-node randomized first timeout so nodes started together do not all campaign together. + - The 10-second startup grace for a node with no peers never applied, so that node elected itself at once with a quorum of one. Startup grace now applies. + + Measured on 5 nodes, starting each node with only lower-numbered peers: 16/16 fresh starts settle in 8–9 s with one leader, against 2/16 that never converged and a 64 s median before. On a full peer mesh the first election now waits one election timeout (3–5 s), so a fresh cluster settles in 11–14 s instead of 8–9 s. + +### Notes + +- Quorum is still computed over each node's own peer list, with no committed membership, so starting every node with the full peer list remains required (#157). + ## 2026-10-01 — mqdb-agent 0.8.30, mqdb-cluster 0.4.16, mqdb-vault 0.1.6, mqdb-cli 0.8.40 ### Changed diff --git a/Cargo.lock b/Cargo.lock index d17876a5..b806d14f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1529,7 +1529,7 @@ dependencies = [ [[package]] name = "mqdb-cli" -version = "0.8.40" +version = "0.8.41" dependencies = [ "base64 0.22.1", "bebytes", @@ -1556,7 +1556,7 @@ dependencies = [ [[package]] name = "mqdb-cluster" -version = "0.4.16" +version = "0.4.17" dependencies = [ "arc-swap", "bebytes", diff --git a/crates/mqdb-cli/Cargo.toml b/crates/mqdb-cli/Cargo.toml index a30af071..e964b1be 100644 --- a/crates/mqdb-cli/Cargo.toml +++ b/crates/mqdb-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cli" -version = "0.8.40" +version = "0.8.41" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cluster/Cargo.toml b/crates/mqdb-cluster/Cargo.toml index 9bd645e4..00b696b5 100644 --- a/crates/mqdb-cluster/Cargo.toml +++ b/crates/mqdb-cluster/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cluster" -version = "0.4.16" +version = "0.4.17" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cluster/src/cluster/raft/coordinator/mod.rs b/crates/mqdb-cluster/src/cluster/raft/coordinator/mod.rs index 66f76388..88707d72 100644 --- a/crates/mqdb-cluster/src/cluster/raft/coordinator/mod.rs +++ b/crates/mqdb-cluster/src/cluster/raft/coordinator/mod.rs @@ -80,7 +80,7 @@ impl RaftCoordinator { } fn rebuild_partition_map(&mut self) { - let outputs = self.node.tick(0); + let outputs = self.node.take_committed(); for output in outputs { if let RaftOutput::ApplyCommand(cmd) = output { self.apply_command_local(&cmd); diff --git a/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs b/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs index 1e092fe3..9ebdb42d 100644 --- a/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs +++ b/crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs @@ -110,6 +110,7 @@ async fn coordinator_election_and_propose() { coord1.add_peer(node2); coord2.add_peer(node1); + coord1.tick(0).await; coord1.tick(1000).await; let sent = coord1.transport.sent_messages(); assert_eq!(sent.len(), 1); @@ -140,6 +141,7 @@ async fn coordinator_applies_partition_update() { coord1.add_peer(node2); coord2.add_peer(node1); + coord1.tick(0).await; coord1.tick(1000).await; let request = match &coord1.transport.sent_messages()[0].1 { ClusterMessage::RequestVote(req) => *req, @@ -224,6 +226,7 @@ async fn handle_node_death_reassigns_partitions() { coord2.add_peer(node1); coord2.add_peer(node3); + coord1.tick(0).await; coord1.tick(1000).await; let request = match &coord1.transport.sent_messages()[0].1 { ClusterMessage::RequestVote(req) => *req, @@ -275,3 +278,29 @@ async fn handle_node_death_does_nothing_when_not_leader() { let indices = coord.handle_node_death(node2).await; assert!(indices.is_empty()); } + +#[tokio::test] +async fn coordinator_from_storage_respects_startup_grace_without_peers() { + let node_id = NodeId::validated(1).unwrap(); + let transport = MockTransport::new(node_id); + let backend: Arc = Arc::new(mqdb_core::MemoryBackend::new()); + let config = RaftConfig { + election_timeout_min_ms: 150, + election_timeout_max_ms: 300, + heartbeat_interval_ms: 50, + startup_grace_period_ms: 10_000, + }; + let mut coord = RaftCoordinator::new_with_storage(node_id, transport, config, backend).unwrap(); + let now = 1_790_000_000_000; + + coord.tick(now).await; + coord.tick(now + 5_000).await; + assert!( + !coord.is_leader(), + "a node without peers must wait out the startup grace period before electing itself" + ); + + coord.tick(now + 10_000).await; + coord.tick(now + 10_400).await; + assert!(coord.is_leader()); +} diff --git a/crates/mqdb-cluster/src/cluster/raft/node.rs b/crates/mqdb-cluster/src/cluster/raft/node.rs index cf0c1c61..31bbdb8b 100644 --- a/crates/mqdb-cluster/src/cluster/raft/node.rs +++ b/crates/mqdb-cluster/src/cluster/raft/node.rs @@ -215,6 +215,10 @@ impl RaftNode { pub fn tick(&mut self, now_ms: u64) -> Vec { if self.startup_time.is_none() { self.startup_time = Some(now_ms); + if self.last_heartbeat_time == 0 { + self.last_heartbeat_time = now_ms; + self.reset_election_timeout(); + } } let mut outputs = Vec::new(); @@ -244,7 +248,6 @@ impl RaftNode { self.persist_state(); self.last_election_time = now_ms; self.last_heartbeat_time = now_ms; - self.reset_election_timeout(); if self.state.has_quorum() { self.state.become_leader(); @@ -294,6 +297,10 @@ impl RaftNode { outputs } + pub fn take_committed(&mut self) -> Vec { + self.apply_committed() + } + fn apply_committed(&mut self) -> Vec { let commands: Vec<_> = self .state @@ -341,7 +348,6 @@ impl RaftNode { self.state.grant_vote(request.term, c); self.persist_state(); self.last_heartbeat_time = now_ms; - self.reset_election_timeout(); } RequestVoteResponse::granted(self.state.current_term()) } else { @@ -416,7 +422,6 @@ impl RaftNode { } self.last_heartbeat_time = now_ms; - self.reset_election_timeout(); self.persist_log_entries(&request.entries); @@ -508,6 +513,82 @@ mod tests { RaftNode::create(node_id, test_config()) } + fn heartbeat_from(leader: u16, term: u64, prev_log_index: u64) -> AppendEntriesRequest { + AppendEntriesRequest::create(term, leader, prev_log_index, 0, Vec::new(), 0) + } + + #[test] + fn candidate_stepping_down_in_same_term_keeps_its_vote() { + let mut node = make_node(2); + node.add_peer(NodeId::validated(1).unwrap()); + node.add_peer(NodeId::validated(3).unwrap()); + node.add_peer(NodeId::validated(4).unwrap()); + node.add_peer(NodeId::validated(5).unwrap()); + node.tick(0); + node.tick(1000); + assert_eq!(node.role(), RaftRole::Candidate); + assert_eq!(node.current_term(), 1); + + let _ = node.handle_append_entries( + NodeId::validated(1).unwrap(), + heartbeat_from(1, 1, 257), + 1100, + ); + assert_eq!(node.role(), RaftRole::Follower); + assert_eq!(node.current_term(), 1); + + let request = RequestVoteRequest::create(1, 5, 0, 0); + let (response, _) = node.handle_request_vote(NodeId::validated(5).unwrap(), request, 1150); + assert!( + !response.is_granted(), + "a node that voted for itself in term 1 must not grant a second term-1 vote" + ); + } + + #[test] + fn first_tick_starts_the_election_timer_instead_of_firing_it() { + let mut node = make_node(2); + node.add_peer(NodeId::validated(1).unwrap()); + let now = 1_790_000_000_000; + + let outputs = node.tick(now); + assert!(outputs.is_empty()); + assert_eq!(node.role(), RaftRole::Follower); + + let outputs = node.tick(now + 301); + assert!( + outputs + .iter() + .any(|o| matches!(o, RaftOutput::SendRequestVote { .. })) + ); + } + + fn first_campaign_time(node_id: u16, start_ms: u64) -> Option { + let mut node = make_node(node_id); + node.add_peer(NodeId::validated(if node_id == 1 { 2 } else { 1 }).unwrap()); + node.tick(start_ms); + (start_ms..=start_ms + 300).find(|&now| { + node.tick(now) + .iter() + .any(|o| matches!(o, RaftOutput::SendRequestVote { .. })) + }) + } + + #[test] + fn nodes_started_together_do_not_campaign_together() { + let start = 1_790_000_000_000; + let campaigns: Vec<_> = (1..=5).map(|id| first_campaign_time(id, start)).collect(); + assert!(campaigns.iter().all(Option::is_some)); + let mut distinct = campaigns.clone(); + distinct.sort_unstable(); + distinct.dedup(); + assert_eq!( + distinct.len(), + campaigns.len(), + "first election timeouts must differ per node: {campaigns:?}" + ); + } + #[test] fn starts_as_follower() { let node = make_node(1); @@ -520,6 +601,7 @@ mod tests { node.add_peer(NodeId::validated(2).unwrap()); node.add_peer(NodeId::validated(3).unwrap()); + node.tick(0); let outputs = node.tick(1000); assert!( outputs @@ -538,6 +620,7 @@ mod tests { node.add_peer(peer2); node.add_peer(peer3); + node.tick(0); node.tick(1000); assert_eq!(node.role(), RaftRole::Candidate); @@ -556,6 +639,7 @@ mod tests { fn steps_down_on_higher_term() { let mut node = make_node(1); node.add_peer(NodeId::validated(2).unwrap()); + node.tick(0); node.tick(1000); assert_eq!(node.current_term(), 1); @@ -577,6 +661,7 @@ mod tests { let peer2 = NodeId::validated(2).unwrap(); node.add_peer(peer2); + node.tick(0); node.tick(1000); let response = RequestVoteResponse::granted(1); node.handle_request_vote_response(peer2, response); @@ -600,6 +685,7 @@ mod tests { leader.add_peer(peer2); follower.add_peer(peer1); + leader.tick(0); leader.tick(1000); let response = RequestVoteResponse::granted(1); leader.handle_request_vote_response(peer2, response); diff --git a/crates/mqdb-cluster/src/cluster/raft/state.rs b/crates/mqdb-cluster/src/cluster/raft/state.rs index 778d2ef3..c4724c7e 100644 --- a/crates/mqdb-cluster/src/cluster/raft/state.rs +++ b/crates/mqdb-cluster/src/cluster/raft/state.rs @@ -279,8 +279,10 @@ impl RaftState { pub fn become_follower(&mut self, term: u64, leader: Option) { self.role = RaftRole::Follower; - self.current_term = term; - self.voted_for = None; + if term > self.current_term { + self.current_term = term; + self.voted_for = None; + } self.leader_id = leader; self.votes_received.clear(); } diff --git a/crates/mqdb-cluster/tests/cluster_integration_test.rs b/crates/mqdb-cluster/tests/cluster_integration_test.rs index 19a0a6e6..bdcfd06f 100644 --- a/crates/mqdb-cluster/tests/cluster_integration_test.rs +++ b/crates/mqdb-cluster/tests/cluster_integration_test.rs @@ -302,6 +302,7 @@ async fn cluster_formation_three_nodes() { ); let first = &mut cluster.nodes[0]; + first.raft.tick(0); let outputs = first.raft.tick(1000); assert_eq!(first.raft.role(), RaftRole::Candidate); assert_eq!(first.raft.current_term(), 1); @@ -561,6 +562,7 @@ async fn raft_log_replication_commit() { let n2 = cluster.nodes[1].id; let n3 = cluster.nodes[2].id; + cluster.nodes[0].raft.tick(0); let outputs = cluster.nodes[0].raft.tick(1000); assert_eq!(cluster.nodes[0].raft.role(), RaftRole::Candidate); @@ -1795,6 +1797,7 @@ async fn raft_logs_converge_after_divergence() { let n1 = cluster.nodes[0].id; + cluster.nodes[0].raft.tick(0); let outputs = cluster.nodes[0].raft.tick(1000); for output in outputs { if let RaftOutput::SendRequestVote { to, request } = output { @@ -2272,6 +2275,7 @@ async fn raft_command_application_updates_partition_map() { let n2 = cluster.nodes[1].id; let n3 = cluster.nodes[2].id; + cluster.nodes[0].raft.tick(0); let outputs = cluster.nodes[0].raft.tick(1000); assert_eq!(cluster.nodes[0].raft.role(), RaftRole::Candidate); @@ -2574,6 +2578,7 @@ async fn raft_state_persisted_and_recovered() { node.add_peer(peer_id); + node.tick(0); node.tick(1000); assert_eq!(node.role(), RaftRole::Candidate); assert_eq!(node.current_term(), 1); @@ -2609,6 +2614,7 @@ async fn raft_log_persisted_and_recovered() { node.add_peer(peer_id); + node.tick(0); node.tick(1000); let vote_response = mqdb_cluster::cluster::raft::RequestVoteResponse::granted(1); diff --git a/crates/mqdb-cluster/tests/cluster_test.rs b/crates/mqdb-cluster/tests/cluster_test.rs index 7a74e642..a6e0f873 100644 --- a/crates/mqdb-cluster/tests/cluster_test.rs +++ b/crates/mqdb-cluster/tests/cluster_test.rs @@ -303,6 +303,7 @@ async fn raft_leader_election_three_nodes() { assert_eq!(n2.role(), RaftRole::Follower); assert_eq!(n3.role(), RaftRole::Follower); + n1.tick(0); let outputs = n1.tick(1000); assert_eq!(n1.role(), RaftRole::Candidate); assert_eq!(n1.current_term(), 1); @@ -350,6 +351,7 @@ async fn raft_step_down_on_higher_term() { n2.add_peer(node1); n2.add_peer(node3); + n1.tick(0); n1.tick(1000); assert_eq!(n1.role(), RaftRole::Candidate); assert_eq!(n1.current_term(), 1); @@ -386,6 +388,7 @@ async fn raft_partition_map_updates() { follower.add_peer(node1); follower.add_peer(node3); + leader.tick(0); let outputs = leader.tick(1000); assert_eq!(leader.role(), RaftRole::Candidate); From f4b5fd2ef7b543b892b53163222f5c074a1d298b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Sat, 3 Oct 2026 16:39:23 -0700 Subject: [PATCH 2/2] restore election timeout re-randomization --- crates/mqdb-cluster/src/cluster/raft/node.rs | 26 +++++++++++++++++++- 1 file changed, 25 insertions(+), 1 deletion(-) diff --git a/crates/mqdb-cluster/src/cluster/raft/node.rs b/crates/mqdb-cluster/src/cluster/raft/node.rs index 31bbdb8b..496d65d8 100644 --- a/crates/mqdb-cluster/src/cluster/raft/node.rs +++ b/crates/mqdb-cluster/src/cluster/raft/node.rs @@ -217,8 +217,8 @@ impl RaftNode { self.startup_time = Some(now_ms); if self.last_heartbeat_time == 0 { self.last_heartbeat_time = now_ms; - self.reset_election_timeout(); } + self.reset_election_timeout(); } let mut outputs = Vec::new(); @@ -248,6 +248,7 @@ impl RaftNode { self.persist_state(); self.last_election_time = now_ms; self.last_heartbeat_time = now_ms; + self.reset_election_timeout(); if self.state.has_quorum() { self.state.become_leader(); @@ -348,6 +349,7 @@ impl RaftNode { self.state.grant_vote(request.term, c); self.persist_state(); self.last_heartbeat_time = now_ms; + self.reset_election_timeout(); } RequestVoteResponse::granted(self.state.current_term()) } else { @@ -422,6 +424,7 @@ impl RaftNode { } self.last_heartbeat_time = now_ms; + self.reset_election_timeout(); self.persist_log_entries(&request.entries); @@ -589,6 +592,27 @@ mod tests { ); } + #[test] + fn election_timeout_is_redrawn_for_every_election() { + let mut node = make_node(3); + node.add_peer(NodeId::validated(1).unwrap()); + node.add_peer(NodeId::validated(2).unwrap()); + node.tick(0); + let campaigns: Vec = (1..=2000) + .filter(|&now| { + node.tick(now) + .iter() + .any(|o| matches!(o, RaftOutput::SendRequestVote { .. })) + }) + .collect(); + let intervals: Vec = campaigns.windows(2).map(|w| w[1] - w[0]).collect(); + assert!(intervals.len() >= 4, "campaigns: {campaigns:?}"); + assert!( + intervals.windows(2).any(|w| w[0] != w[1]), + "every election used the same timeout: {intervals:?}" + ); + } + #[test] fn starts_as_follower() { let node = make_node(1);