Skip to content

Commit 12944e7

Browse files
committed
feat(nodedb-cluster): add QUIC peer pre-warming on startup
Add `NexarTransport::warm_peer(node_id)` which dials a registered peer and inserts the `quinn::Connection` into the connection cache. First replicated request after boot reuses the cached connection rather than paying a cold dial. Add `control::cluster::warm_peers` module with `warm_known_peers()` which fans out warm attempts in parallel across all topology peers (excluding self) with a configurable per-peer timeout, and returns a `PeerWarmReport` summarising successes and failures.
1 parent 4d2e90d commit 12944e7

5 files changed

Lines changed: 323 additions & 0 deletions

File tree

nodedb-cluster/src/transport/client.rs

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,23 @@ impl NexarTransport {
152152
debug!(node_id, %addr, "peer registered");
153153
}
154154

155+
/// Pre-warm the QUIC connection cache for a peer by
156+
/// performing the full dial + handshake and inserting
157+
/// the connection into the peer cache. On success, the
158+
/// next `send_rpc(target, ...)` skips the dial entirely
159+
/// and reuses the cached `quinn::Connection`.
160+
///
161+
/// Caller MUST have called [`register_peer`] first — this
162+
/// function resolves the peer address from the
163+
/// `peer_addrs` map. Used by the startup `warm_peers`
164+
/// phase so the first replicated request after boot
165+
/// doesn't pay a cold-connect penalty.
166+
///
167+
/// [`register_peer`]: Self::register_peer
168+
pub async fn warm_peer(&self, node_id: u64) -> Result<()> {
169+
self.get_or_connect(node_id).await.map(|_| ())
170+
}
171+
155172
/// Run the inbound RPC accept loop until shutdown.
156173
///
157174
/// For each incoming connection, spawns a task that accepts

nodedb/src/control/cluster/mod.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,10 +18,12 @@ pub mod init;
1818
pub mod metadata_applier;
1919
pub mod spsc_applier;
2020
pub mod start_raft;
21+
pub mod warm_peers;
2122

2223
pub use applied_index_watcher::AppliedIndexWatcher;
2324
pub use handle::ClusterHandle;
2425
pub use init::{init_cluster, init_cluster_with_transport};
2526
pub use metadata_applier::MetadataCommitApplier;
2627
pub use spsc_applier::SpscCommitApplier;
2728
pub use start_raft::start_raft;
29+
pub use warm_peers::{PeerWarmReport, warm_known_peers};
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
//! Pre-warm the QUIC peer cache after `TransportBind` so the
2+
//! first replicated request after boot doesn't pay a cold
3+
//! connect. Slots into the startup sequencer between
4+
//! `TransportBind` and `WarmPeers` phases.
5+
6+
pub mod report;
7+
pub mod warm;
8+
9+
pub use report::PeerWarmReport;
10+
pub use warm::warm_known_peers;
Lines changed: 118 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,118 @@
1+
//! Report returned by [`super::warm_known_peers`].
2+
//!
3+
//! Logged at INFO regardless of outcome so operators see
4+
//! `peer_warm: 3 succeeded, 1 failed in 412ms` in the
5+
//! startup log without grepping for individual lines.
6+
7+
use std::fmt;
8+
use std::time::Duration;
9+
10+
#[derive(Debug, Clone)]
11+
pub struct PeerWarmReport {
12+
/// Number of peers attempted (everything in the topology
13+
/// except this node).
14+
pub attempted: usize,
15+
/// Node ids that successfully warmed.
16+
pub succeeded: Vec<u64>,
17+
/// Per-peer failure with the underlying transport error
18+
/// stringified — kept as String because the cluster
19+
/// transport returns its own typed error which we don't
20+
/// want to leak into this struct's API surface.
21+
pub failed: Vec<(u64, String)>,
22+
/// Total wall-clock for the whole warm phase. Includes
23+
/// the parallel-dial time, NOT the sum of individual
24+
/// per-peer dials.
25+
pub elapsed: Duration,
26+
}
27+
28+
impl PeerWarmReport {
29+
/// Whether at least one peer was warmed. Useful for
30+
/// tests; the production path logs the full report
31+
/// regardless.
32+
pub fn is_success(&self) -> bool {
33+
!self.succeeded.is_empty()
34+
}
35+
36+
/// Whether every attempted peer was reached.
37+
pub fn is_complete(&self) -> bool {
38+
self.failed.is_empty() && self.succeeded.len() == self.attempted
39+
}
40+
41+
/// Empty report — used when the topology has no peers
42+
/// (single-node mode).
43+
pub fn empty() -> Self {
44+
Self {
45+
attempted: 0,
46+
succeeded: Vec::new(),
47+
failed: Vec::new(),
48+
elapsed: Duration::ZERO,
49+
}
50+
}
51+
}
52+
53+
impl fmt::Display for PeerWarmReport {
54+
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
55+
write!(
56+
f,
57+
"peer_warm: {}/{} succeeded ({} failed) in {:?}",
58+
self.succeeded.len(),
59+
self.attempted,
60+
self.failed.len(),
61+
self.elapsed
62+
)?;
63+
for (id, err) in &self.failed {
64+
write!(f, "\n node {id}: {err}")?;
65+
}
66+
Ok(())
67+
}
68+
}
69+
70+
#[cfg(test)]
71+
mod tests {
72+
use super::*;
73+
74+
#[test]
75+
fn empty_report_is_complete_and_unsuccessful() {
76+
let r = PeerWarmReport::empty();
77+
assert!(r.is_complete());
78+
assert!(!r.is_success());
79+
}
80+
81+
#[test]
82+
fn complete_report() {
83+
let r = PeerWarmReport {
84+
attempted: 2,
85+
succeeded: vec![1, 2],
86+
failed: vec![],
87+
elapsed: Duration::from_millis(50),
88+
};
89+
assert!(r.is_complete());
90+
assert!(r.is_success());
91+
}
92+
93+
#[test]
94+
fn partial_failure_not_complete() {
95+
let r = PeerWarmReport {
96+
attempted: 3,
97+
succeeded: vec![1],
98+
failed: vec![(2, "timeout".into()), (3, "refused".into())],
99+
elapsed: Duration::from_millis(2000),
100+
};
101+
assert!(!r.is_complete());
102+
assert!(r.is_success());
103+
}
104+
105+
#[test]
106+
fn display_includes_failures() {
107+
let r = PeerWarmReport {
108+
attempted: 2,
109+
succeeded: vec![1],
110+
failed: vec![(2, "connection refused".into())],
111+
elapsed: Duration::from_millis(100),
112+
};
113+
let s = r.to_string();
114+
assert!(s.contains("1/2 succeeded"));
115+
assert!(s.contains("node 2"));
116+
assert!(s.contains("connection refused"));
117+
}
118+
}
Lines changed: 176 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,176 @@
1+
//! Parallel peer-cache warm-up.
2+
//!
3+
//! Dials every node in the live topology except `self_id`,
4+
//! caps each dial at half the total deadline, and aggregates
5+
//! the outcome into a [`PeerWarmReport`]. Failures are
6+
//! non-fatal — operators see the per-peer failure in the
7+
//! returned report and the startup phase advances regardless.
8+
9+
use std::collections::HashSet;
10+
use std::time::{Duration, Instant};
11+
12+
use futures::future::join_all;
13+
use nodedb_cluster::{ClusterTopology, NexarTransport};
14+
15+
use super::report::PeerWarmReport;
16+
17+
/// Pre-warm every peer in `topology` except `self_id`.
18+
///
19+
/// `total_deadline` bounds the entire operation. Each
20+
/// individual peer dial is capped at `total_deadline / 2`
21+
/// so a single slow peer cannot eat the whole budget. All
22+
/// dials run in parallel, so the total wall-clock is
23+
/// roughly `max(per-peer-latency, total_deadline)`.
24+
///
25+
/// Empty topology → returns [`PeerWarmReport::empty`].
26+
/// Caller logs the returned report once; this function
27+
/// only logs per-peer failures if verbose tracing is
28+
/// enabled on the transport layer.
29+
pub async fn warm_known_peers(
30+
transport: &NexarTransport,
31+
topology: &ClusterTopology,
32+
self_id: u64,
33+
total_deadline: Duration,
34+
) -> PeerWarmReport {
35+
let start = Instant::now();
36+
let per_peer_deadline = total_deadline.checked_div(2).unwrap_or(total_deadline);
37+
38+
let peers: Vec<(u64, std::net::SocketAddr)> = topology
39+
.all_nodes()
40+
.filter(|n| n.node_id != self_id)
41+
.filter_map(|n| n.socket_addr().map(|a| (n.node_id, a)))
42+
.collect();
43+
44+
if peers.is_empty() {
45+
return PeerWarmReport::empty();
46+
}
47+
48+
let attempted = peers.len();
49+
50+
// Ensure the transport's address map has every peer.
51+
// `register_peer` is idempotent — safe to call again.
52+
for (id, addr) in &peers {
53+
transport.register_peer(*id, *addr);
54+
}
55+
56+
let futs = peers.iter().map(|(id, _addr)| {
57+
let id = *id;
58+
async move {
59+
match tokio::time::timeout(per_peer_deadline, transport.warm_peer(id)).await {
60+
Ok(Ok(())) => (id, Ok(())),
61+
Ok(Err(e)) => (id, Err(format!("{e}"))),
62+
Err(_) => (id, Err(format!("dial timeout after {per_peer_deadline:?}"))),
63+
}
64+
}
65+
});
66+
67+
// Outer deadline bounds the entire parallel batch. On
68+
// outer timeout, every still-pending future is dropped
69+
// (cancelling the in-flight dials) and the missing peers
70+
// are reported as deadline exceedances.
71+
let outcome = tokio::time::timeout(total_deadline, join_all(futs)).await;
72+
73+
let mut succeeded: Vec<u64> = Vec::with_capacity(attempted);
74+
let mut failed: Vec<(u64, String)> = Vec::new();
75+
76+
match outcome {
77+
Ok(results) => {
78+
for (id, outcome) in results {
79+
match outcome {
80+
Ok(()) => succeeded.push(id),
81+
Err(msg) => failed.push((id, msg)),
82+
}
83+
}
84+
}
85+
Err(_) => {
86+
// Outer deadline expired — we don't know which
87+
// futures completed vs which were in flight, so
88+
// mark every peer as a deadline exceedance. The
89+
// report is accurate about the worst case.
90+
for (id, _) in &peers {
91+
failed.push((
92+
*id,
93+
format!("outer warm deadline exceeded after {total_deadline:?}"),
94+
));
95+
}
96+
}
97+
}
98+
99+
// Sanity: account every attempted peer. If a future was
100+
// dropped between matching on the outer result and
101+
// collecting (shouldn't happen with join_all but cheap
102+
// to double-check), any missing ids become failures.
103+
let reported: HashSet<u64> = succeeded
104+
.iter()
105+
.copied()
106+
.chain(failed.iter().map(|(id, _)| *id))
107+
.collect();
108+
for (id, _) in &peers {
109+
if !reported.contains(id) {
110+
failed.push((*id, "unaccounted peer — internal bookkeeping bug".into()));
111+
}
112+
}
113+
114+
PeerWarmReport {
115+
attempted,
116+
succeeded,
117+
failed,
118+
elapsed: start.elapsed(),
119+
}
120+
}
121+
122+
#[cfg(test)]
123+
mod tests {
124+
use super::*;
125+
use nodedb_cluster::ClusterTopology;
126+
127+
#[tokio::test]
128+
async fn empty_topology_returns_empty_report() {
129+
let transport = test_transport();
130+
let topo = ClusterTopology::new();
131+
let report = warm_known_peers(&transport, &topo, 1, Duration::from_secs(1)).await;
132+
assert_eq!(report.attempted, 0);
133+
assert!(report.is_complete());
134+
assert_eq!(report.elapsed, Duration::ZERO);
135+
}
136+
137+
#[tokio::test]
138+
async fn self_only_topology_skips_self() {
139+
let transport = test_transport();
140+
let mut topo = ClusterTopology::new();
141+
topo.add_node(nodedb_cluster::NodeInfo::new(
142+
1,
143+
"127.0.0.1:65530".parse().unwrap(),
144+
nodedb_cluster::NodeState::Active,
145+
));
146+
let report = warm_known_peers(&transport, &topo, 1, Duration::from_secs(1)).await;
147+
assert_eq!(report.attempted, 0);
148+
}
149+
150+
#[tokio::test]
151+
async fn dead_peer_is_reported_as_failure() {
152+
let transport = test_transport();
153+
let mut topo = ClusterTopology::new();
154+
topo.add_node(nodedb_cluster::NodeInfo::new(
155+
1,
156+
"127.0.0.1:65530".parse().unwrap(),
157+
nodedb_cluster::NodeState::Active,
158+
));
159+
// 127.0.0.1:1 — reserved port, guaranteed to fail.
160+
topo.add_node(nodedb_cluster::NodeInfo::new(
161+
2,
162+
"127.0.0.1:1".parse().unwrap(),
163+
nodedb_cluster::NodeState::Active,
164+
));
165+
let report = warm_known_peers(&transport, &topo, 1, Duration::from_millis(500)).await;
166+
assert_eq!(report.attempted, 1);
167+
assert_eq!(report.succeeded.len(), 0);
168+
assert_eq!(report.failed.len(), 1);
169+
assert_eq!(report.failed[0].0, 2);
170+
assert!(report.elapsed <= Duration::from_millis(700));
171+
}
172+
173+
fn test_transport() -> NexarTransport {
174+
NexarTransport::new(999, "127.0.0.1:0".parse().unwrap()).expect("test transport bind")
175+
}
176+
}

0 commit comments

Comments
 (0)