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
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions Cargo.lock

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

2 changes: 1 addition & 1 deletion crates/mqdb-cli/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cluster/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cluster/src/cluster/raft/coordinator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ impl<T: ClusterTransport> RaftCoordinator<T> {
}

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);
Expand Down
29 changes: 29 additions & 0 deletions crates/mqdb-cluster/src/cluster/raft/coordinator/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<dyn mqdb_core::StorageBackend> = 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());
}
110 changes: 110 additions & 0 deletions crates/mqdb-cluster/src/cluster/raft/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,10 @@ impl RaftNode {
pub fn tick(&mut self, now_ms: u64) -> Vec<RaftOutput> {
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();
Expand Down Expand Up @@ -294,6 +298,10 @@ impl RaftNode {
outputs
}

pub fn take_committed(&mut self) -> Vec<RaftOutput> {
self.apply_committed()
}

fn apply_committed(&mut self) -> Vec<RaftOutput> {
let commands: Vec<_> = self
.state
Expand Down Expand Up @@ -508,6 +516,103 @@ 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<u64> {
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 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<u64> = (1..=2000)
.filter(|&now| {
node.tick(now)
.iter()
.any(|o| matches!(o, RaftOutput::SendRequestVote { .. }))
})
.collect();
let intervals: Vec<u64> = 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);
Expand All @@ -520,6 +625,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
Expand All @@ -538,6 +644,7 @@ mod tests {
node.add_peer(peer2);
node.add_peer(peer3);

node.tick(0);
node.tick(1000);
assert_eq!(node.role(), RaftRole::Candidate);

Expand All @@ -556,6 +663,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);

Expand All @@ -577,6 +685,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);
Expand All @@ -600,6 +709,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);
Expand Down
6 changes: 4 additions & 2 deletions crates/mqdb-cluster/src/cluster/raft/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -279,8 +279,10 @@ impl RaftState {

pub fn become_follower(&mut self, term: u64, leader: Option<NodeId>) {
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();
}
Expand Down
6 changes: 6 additions & 0 deletions crates/mqdb-cluster/tests/cluster_integration_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
3 changes: 3 additions & 0 deletions crates/mqdb-cluster/tests/cluster_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);

Expand Down
Loading