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
42 changes: 42 additions & 0 deletions crates/mqdb-cluster/src/cluster/heartbeat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,17 @@ impl HeartbeatManager {
.collect()
}

#[must_use]
pub fn alive_or_suspected_nodes(&self) -> Vec<NodeId> {
self.nodes
.iter()
.filter(|(_, state)| {
state.status == NodeStatus::Alive || state.status == NodeStatus::Suspected
})
.filter_map(|(&id, _)| NodeId::validated(id))
.collect()
}

#[must_use]
pub fn has_alive_peers(&self) -> bool {
self.nodes
Expand Down Expand Up @@ -351,6 +362,37 @@ mod tests {
assert_eq!(mgr.node_status(node2), NodeStatus::Dead);
}

#[test]
fn alive_or_suspected_excludes_dead_and_unknown() {
let local = NodeId::validated(1).unwrap();
let mut mgr = HeartbeatManager::new(local, config());

let n2 = NodeId::validated(2).unwrap();
let n3 = NodeId::validated(3).unwrap();
let n4 = NodeId::validated(4).unwrap();
let n5 = NodeId::validated(5).unwrap();
for n in [n2, n3, n4, n5] {
mgr.register_node(n);
}

mgr.receive_heartbeat(n2, &Heartbeat::create(n2, 1900), 1900);
mgr.receive_heartbeat(n3, &Heartbeat::create(n3, 1600), 1600);
mgr.receive_heartbeat(n4, &Heartbeat::create(n4, 1300), 1300);

mgr.check_timeouts(2000);
assert_eq!(mgr.node_status(n2), NodeStatus::Alive);
assert_eq!(mgr.node_status(n3), NodeStatus::Suspected);
assert_eq!(mgr.node_status(n4), NodeStatus::Dead);
assert_eq!(mgr.node_status(n5), NodeStatus::Unknown);

let set: std::collections::BTreeSet<NodeId> =
mgr.alive_or_suspected_nodes().into_iter().collect();
assert!(set.contains(&n2), "alive node must be included");
assert!(set.contains(&n3), "suspected node must be included");
assert!(!set.contains(&n4), "dead node must be excluded");
assert!(!set.contains(&n5), "unknown node must be excluded");
}

#[test]
fn should_send_respects_interval() {
use crate::cluster::{Epoch, PartitionAssignment, PartitionId};
Expand Down
4 changes: 4 additions & 0 deletions crates/mqdb-cluster/src/cluster/node_controller/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -649,6 +649,10 @@ impl<T: ClusterTransport> NodeController<T> {
self.heartbeat.alive_nodes()
}

pub fn alive_or_suspected_nodes(&self) -> Vec<NodeId> {
self.heartbeat.alive_or_suspected_nodes()
}

/// Cluster members this node knows about, excluding itself: every node that owns a partition
/// as primary or replica, plus any node named in voter gossip. A node can appear here without
/// being reachable, which is exactly the case `unlinked_nodes` reports.
Expand Down
27 changes: 20 additions & 7 deletions crates/mqdb-cluster/src/cluster/quic_transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,14 @@ use std::net::SocketAddr;
use std::path::Path;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::{Notify, RwLock};
use tracing::{debug, error, info, trace, warn};

const INBOX_CHANNEL_CAPACITY: usize = 16384;
const PEER_CONTROL_QUEUE_CAPACITY: usize = 256;
const PEER_BULK_QUEUE_CAPACITY: usize = 1024;
const REDIAL_DIAL_TIMEOUT: Duration = Duration::from_secs(5);

struct PeerConnection {
_connection: Connection,
Expand Down Expand Up @@ -294,23 +296,34 @@ impl QuicDirectTransport {
/// Driven by the heartbeat liveness view rather than the transport peer
/// map, so a stale entry that has not yet been removed by a send failure
/// does not stop reconnection. `connect_to_peer` replaces any stale entry.
pub async fn redial_unlinked(&self, alive: &[NodeId]) {
pub async fn redial_unlinked(&self, linked: &[NodeId]) {
let targets: Vec<(NodeId, SocketAddr)> = {
let addrs = self.peer_addrs.read().await;
addrs
.iter()
.filter(|&(node, _)| *node != self.node_id && !alive.contains(node))
.filter(|&(node, _)| *node != self.node_id && !linked.contains(node))
.map(|(node, addr)| (*node, *addr))
.collect()
};

let mut dials = tokio::task::JoinSet::new();
for (node, addr) in targets {
if let Err(e) = self.connect_to_peer(node, addr).await {
debug!(peer = node.get(), addr = %addr, error = %e, "redial to peer failed, will retry");
} else {
info!(peer = node.get(), addr = %addr, "redialled peer");
}
let transport = self.clone();
dials.spawn(async move {
match tokio::time::timeout(REDIAL_DIAL_TIMEOUT, transport.connect_to_peer(node, addr))
.await
{
Ok(Ok(())) => info!(peer = node.get(), addr = %addr, "redialled peer"),
Ok(Err(e)) => {
debug!(peer = node.get(), addr = %addr, error = %e, "redial to peer failed, will retry");
}
Err(_) => {
debug!(peer = node.get(), addr = %addr, "redial to peer timed out, will retry");
}
}
});
}
while dials.join_next().await.is_some() {}
}

#[must_use]
Expand Down
13 changes: 10 additions & 3 deletions crates/mqdb-cluster/src/cluster_agent/event_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,10 @@ impl ClusteredAgent {
}
}
}

if !batch.dead_nodes.is_empty() {
self.redial_disconnected_peers().await;
}
}

fn try_resolve_constraint_response(&self, msg: &crate::cluster::InboundMessage) -> bool {
Expand Down Expand Up @@ -559,13 +563,16 @@ impl ClusteredAgent {
}

async fn redial_disconnected_peers(&self) {
let (quic, alive) = {
let (quic, linked) = {
let ctrl = self.controller.read().await;
(ctrl.transport().as_quic().cloned(), ctrl.alive_nodes())
(
ctrl.transport().as_quic().cloned(),
ctrl.alive_or_suspected_nodes(),
)
};
if let Some(quic) = quic {
tokio::spawn(async move {
quic.redial_unlinked(&alive).await;
quic.redial_unlinked(&linked).await;
});
}
}
Expand Down
Loading