diff --git a/crates/mqdb-cluster/src/cluster/heartbeat.rs b/crates/mqdb-cluster/src/cluster/heartbeat.rs index 888451e..141cdd7 100644 --- a/crates/mqdb-cluster/src/cluster/heartbeat.rs +++ b/crates/mqdb-cluster/src/cluster/heartbeat.rs @@ -272,6 +272,17 @@ impl HeartbeatManager { .collect() } + #[must_use] + pub fn alive_or_suspected_nodes(&self) -> Vec { + 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 @@ -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 = + 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}; diff --git a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs index 29c03be..2704cde 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/mod.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/mod.rs @@ -649,6 +649,10 @@ impl NodeController { self.heartbeat.alive_nodes() } + pub fn alive_or_suspected_nodes(&self) -> Vec { + 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. diff --git a/crates/mqdb-cluster/src/cluster/quic_transport.rs b/crates/mqdb-cluster/src/cluster/quic_transport.rs index c9c19df..64da969 100644 --- a/crates/mqdb-cluster/src/cluster/quic_transport.rs +++ b/crates/mqdb-cluster/src/cluster/quic_transport.rs @@ -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, @@ -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] diff --git a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs index d9bf1cb..17734ce 100644 --- a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs +++ b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs @@ -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 { @@ -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; }); } }