Merge branch 'refactor-node' into refactor-node-next

Forward-merge the peering homeostatic reconciler onto the XX/v2 handshake
line. The three scattered peering mechanisms (auto-connect and retry,
overlay discovery, and opportunistic transport-neighbor growth) are now
unified in the sans-IO reconciler on both lines.

Hand-resolved on the first-contact surface: the reconciler opportunistic
layer emits a connect for anonymous, identity-unknown transport legs
unconditionally (bypassing the connect budget and per-peer cap, gated only
by the addr_to_link dedup), reproducing the XX dial-before-identity-known
behavior; named legs route through the reconciler with inputs identical to
the master line. handle_msg1/msg3 keep the XX path; the peer-loss reflexes
use the relocated wrappers. Behavior-neutral versus the pre-merge next line.
This commit is contained in:
Johnathan Corgan
2026-07-13 06:40:05 +00:00
14 changed files with 2266 additions and 774 deletions
+1 -1
View File
@@ -93,7 +93,7 @@ impl Node {
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.schedule_reconnect(addr, now_ms);
self.note_link_dead(addr, now_ms);
}
/// Remove an active peer and clean up all associated state.
+1 -1
View File
@@ -543,7 +543,7 @@ impl Node {
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.schedule_reconnect(addr, now_ms);
self.note_link_dead(addr, now_ms);
}
}
}
+9 -3
View File
@@ -1309,7 +1309,7 @@ impl Node {
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.schedule_reconnect(peer, now_ms);
self.note_link_dead(peer, now_ms);
// Fall through to process as new connection.
}
InboundDecision::Promote => {
@@ -1597,7 +1597,10 @@ impl Node {
self.peers.insert(peer_node_addr, new_peer);
self.peers_by_index
.insert((transport_id, our_index.as_u32()), peer_node_addr);
self.retry_pending.remove(&peer_node_addr);
self.peering
.reconciler
.retry_pending
.remove(&peer_node_addr);
self.register_identity(peer_node_addr, verified_identity.pubkey_full());
// Non-routing peers don't send filters; include them as
@@ -1711,7 +1714,10 @@ impl Node {
self.peers.insert(peer_node_addr, new_peer);
self.peers_by_index
.insert((transport_id, our_index.as_u32()), peer_node_addr);
self.retry_pending.remove(&peer_node_addr);
self.peering
.reconciler
.retry_pending
.remove(&peer_node_addr);
self.register_identity(peer_node_addr, verified_identity.pubkey_full());
// Non-routing peers don't send filters; include them as
+1 -1
View File
@@ -506,7 +506,7 @@ impl Node {
"Removing peer: link dead timeout"
);
self.remove_active_peer(&peer);
self.schedule_reconnect(peer, now_ms);
self.note_link_dead(peer, now_ms);
}
MmpAction::Heartbeat { peer } => {
if let Some(p) = self.peers.get_mut(&peer) {
+1 -1
View File
@@ -85,7 +85,7 @@ impl Node {
let stale = self.stale_connections(now_ms, timeout_ms);
for action in self.fmp.poll_timeouts(stale) {
match action {
ConnAction::ScheduleRetry { peer } => self.schedule_retry(peer, now_ms),
ConnAction::ScheduleRetry { peer } => self.note_handshake_timeout(peer, now_ms),
ConnAction::Teardown { link } => {
// Log before cleanup (needs live connection state).
if let Some(conn) = self.connections.get(&link) {
+432 -360
View File
File diff suppressed because it is too large Load Diff
+20 -20
View File
@@ -15,10 +15,10 @@ pub(crate) mod encrypt_worker;
mod handlers;
mod lifecycle;
pub(crate) mod metrics;
mod peering;
mod rate_limit;
pub(crate) mod reject;
mod reloadable;
mod retry;
pub(crate) mod session;
pub(crate) mod stats;
pub(crate) mod stats_history;
@@ -458,19 +458,18 @@ pub struct Node {
/// Rate limiter for source-side CoordsRequired/PathBroken responses.
coords_response_rate_limiter: RoutingErrorRateLimiter,
// === Pending Transport Connects ===
/// Links waiting for transport-level connection establishment before
/// sending handshake msg1. For connection-oriented transports (TCP, Tor),
/// the transport connect runs in the background; the tick handler polls
/// connection_state() and initiates the handshake when connected.
pending_connects: Vec<PendingConnect>,
// === Connection Retry ===
/// Retry state for peers whose outbound connections have failed.
/// Keyed by NodeAddr. Entries are created when a handshake times out
/// or fails, and removed on successful promotion or when max retries
/// are exhausted.
retry_pending: HashMap<NodeAddr, retry::RetryState>,
// === Peering Homeostasis ===
/// Owner of the peering-reconciler state relocated off `Node`: the sans-IO
/// retry-schedule core (`peering.reconciler.retry_pending`) and the pending
/// transport connects (`peering.pending_connects`) awaiting transport-level
/// connection establishment before handshake msg1 (for connection-oriented
/// transports the transport connect runs in the background; the tick handler
/// polls connection_state() and initiates the handshake when connected). The
/// retry entries are created when a handshake times out or fails, and removed
/// on successful promotion or when max retries are exhausted. Still driven by
/// the imperative methods; the reconciler core is unwired until the driver
/// cutover.
peering: peering::reconcile::Peering,
// === Periodic Parent Re-evaluation ===
/// Timestamp of last periodic parent re-evaluation (for pacing).
@@ -664,8 +663,7 @@ impl Node {
LookupBackoff::with_params(backoff_base_secs, backoff_max_secs),
LookupForwardRateLimiter::with_interval_ms(forward_min_interval_secs * 1000),
),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
peering: peering::reconcile::Peering::new(),
last_parent_reeval: None,
last_congestion_log: None,
estimated_mesh_size: None,
@@ -809,8 +807,7 @@ impl Node {
coords_response_interval_ms,
),
lookup: Lookup::new(LookupBackoff::new(), LookupForwardRateLimiter::new()),
pending_connects: Vec::new(),
retry_pending: HashMap::new(),
peering: peering::reconcile::Peering::new(),
last_parent_reeval: None,
last_congestion_log: None,
estimated_mesh_size: None,
@@ -2282,6 +2279,7 @@ impl Node {
.values()
.any(|peer| peer.transport_id() == Some(transport_id))
|| self
.peering
.pending_connects
.iter()
.any(|pending| pending.transport_id == transport_id);
@@ -2535,8 +2533,10 @@ impl Node {
}
/// Iterate over retry state for diagnostics.
pub fn retry_state_iter(&self) -> impl Iterator<Item = (&NodeAddr, &retry::RetryState)> {
self.retry_pending.iter()
pub fn retry_state_iter(
&self,
) -> impl Iterator<Item = (&NodeAddr, &peering::retry::RetryState)> {
self.peering.reconciler.retry_pending.iter()
}
// === Routing ===
+148
View File
@@ -0,0 +1,148 @@
//! Thin async driver for the peering reconciler.
//!
//! These `impl Node` methods are the I/O edge of the sans-IO
//! [`super::reconcile::PeeringReconciler`]: they snapshot the live dataplane
//! maps into the reconciler's plain-data inputs, invoke the pure core, and
//! perform the dial / advert-refetch I/O each [`PeeringAction`] names. They also
//! host the two gate-guarded reflex wrappers every peer-loss call site routes
//! through, so drain suppression and the connected-guard live in one place.
//!
//! The `Policy` / `Observed` / `Budget` builders these methods consume live in
//! [`crate::node::lifecycle`] next to the surviving budget helpers and limit
//! constants they wrap.
use crate::identity::NodeAddr;
use crate::node::{Node, NodeError};
use tracing::warn;
use super::reconcile::{DiscoveryPools, Gate, PeeringAction};
impl Node {
/// Reflex: an outbound handshake timed out (replaces the old
/// `Node::schedule_retry` call sites).
///
/// Replicates `schedule_retry`'s connected-guard (obligation O2) — the pure
/// core cannot observe the peers map, so the driver drops the event when the
/// peer is already connected — then feeds the gate-guarded reconciler reflex
/// with the gate derived from the live published state (obligation O4).
pub(in crate::node) fn note_handshake_timeout(&mut self, node_addr: NodeAddr, now_ms: u64) {
if self.peers.contains_key(&node_addr) {
return;
}
let policy =
self.build_peering_policy(self.config().auto_connect_peers().cloned().collect());
let gate = Gate::from_state(self.supervisor.state);
let _ = self
.peering
.reconciler
.on_handshake_timeout(node_addr, now_ms, &policy, gate);
}
/// Reflex: a link went dead / a peer was lost (replaces the old
/// `Node::schedule_reconnect` call sites).
///
/// No connected-guard — the peer is already gone by the time a link-dead /
/// disconnect event fires (`schedule_reconnect` had none). The gate is
/// derived from the live published state so a drain self-suppresses the
/// reconnect (obligation O4, design §8 correctness trap).
pub(in crate::node) fn note_link_dead(&mut self, node_addr: NodeAddr, now_ms: u64) {
let policy =
self.build_peering_policy(self.config().auto_connect_peers().cloned().collect());
let gate = Gate::from_state(self.supervisor.state);
let _ = self
.peering
.reconciler
.on_link_dead(node_addr, now_ms, &policy, gate);
}
/// Process pending retries whose time has arrived (replaces the old
/// `Node::process_pending_retries` body).
///
/// The pure retry-dial phase owns the decision — drop expired entries, refuse
/// to grow when admission binds, dial the first `retry_per_tick` due entries
/// (bumping their `retry_after_ms` past the handshake window). This driver
/// performs the advert-refetch + dial I/O each emitted `Connect` names, and
/// on an immediate dial error feeds the `on_handshake_timeout` reflex so the
/// optimistic re-fire suppression is overwritten by proper backoff
/// (obligation O3). During a drain the gate is `Suspended`, so the reconcile
/// clears the schedule and emits nothing (the design §8 gate).
pub(in crate::node) async fn process_pending_retries(&mut self, now_ms: u64) {
if self.peering.reconciler.retry_pending.is_empty() {
return;
}
// Retry-dial cadence slot: empty config floor (design §3 D4) and empty
// discovery pools, so only the retry-dial phase acts.
let policy = self.build_peering_policy(Vec::new());
let observed = self.observe_peering();
let budget = self.build_peering_budget();
let gate = Gate::from_state(self.supervisor.state);
let actions = self.peering.reconciler.reconcile(
&policy,
&observed,
&budget,
&DiscoveryPools::default(),
now_ms,
gate,
);
for action in actions {
let PeeringAction::Connect(candidate) = action else {
continue;
};
let Some(identity) = candidate.identity else {
continue;
};
let node_addr = *identity.node_addr();
let Some(peer_config) = self
.peering
.reconciler
.retry_pending
.get(&node_addr)
.map(|state| state.peer_config.clone())
else {
continue;
};
// Refresh the peer's overlay advert before retrying. The cache is
// read-only on hit, so a retry without a refetch dials the same
// cached endpoint — and the most common reason a peer landed in the
// retry schedule is that endpoint just stopped working (NAT rebind,
// port change, peer restart). Cheap (one Filter fetch, bounded by
// the retry backoff cadence).
if let Some(bootstrap) = self.supervisor.nostr_rendezvous.engine_arc() {
let _ = bootstrap
.refetch_advert_for_stale_check(&peer_config.npub)
.await;
}
match self.initiate_peer_connection(&peer_config).await {
// The core already pushed `retry_after_ms` past the handshake
// window; a successful promotion clears the entry, a later
// timeout re-fires the reflex with proper backoff.
Ok(()) => {}
Err(e) => {
warn!(
peer = %self.peer_display_name(&node_addr),
error = %e,
"Retry connection initiation failed"
);
// No-transport failures usually mean the cached overlay
// advert is stale; force a re-fetch so the next tick picks up
// fresh endpoints.
if matches!(e, NodeError::NoTransportForType(_))
&& let Some(bootstrap) = self.supervisor.nostr_rendezvous.engine_arc()
{
let npub = peer_config.npub.clone();
tokio::spawn(async move {
let _ = bootstrap.refetch_advert_for_stale_check(&npub).await;
});
}
// Immediate failure counts as an attempt: overwrite the
// optimistic re-fire suppression with backoff (obligation O3).
self.note_handshake_timeout(node_addr, now_ms);
}
}
}
}
}
+14
View File
@@ -0,0 +1,14 @@
//! Peering homeostasis — the desired-state controller for the node's peer set.
//!
//! This module is the home for the peering-reconciler concept: config defines a
//! desired peer set; the reconciler converges the observed set toward it
//! (auto-connect floor, overlay pool, transport-neighbor growth) under the
//! `node.limits` ceiling. Startup and steady-state are the same loop.
//!
//! The cross-attempt retry schedule (`retry.rs`) lives here because a fresh
//! connection is created per re-dial, so the escalating backoff count must
//! persist in the reconciler, not per-connection.
pub(in crate::node) mod driver;
pub(in crate::node) mod reconcile;
pub(in crate::node) mod retry;
File diff suppressed because it is too large Load Diff
+49
View File
@@ -0,0 +1,49 @@
//! Cross-attempt retry state for auto-connect peers.
//!
//! [`RetryState`] is the durable per-peer schedule entry the peering reconciler
//! owns (it lives in [`crate::node::peering::reconcile::PeeringReconciler`], not
//! on a per-connection object, because a fresh connection is created per re-dial
//! so the escalating backoff count must persist across attempts). The decision
//! logic that reads and mutates it — the retry-dial phase and the
//! `on_handshake_timeout` / `on_link_dead` reflexes — lives in the sans-IO
//! reconciler core; the driver wrappers that feed it (retry-dial I/O, the
//! gate-guarded reflex call sites) live in [`super::driver`].
use crate::config::PeerConfig;
/// Per-tick cap on retry-dial connection attempts (design §6 ceiling).
pub(in crate::node) const MAX_RETRY_CONNECTIONS_PER_TICK: usize = 16;
/// Tracks retry state for a peer across connection attempts.
pub struct RetryState {
/// The peer config to use for initiating retries.
pub peer_config: PeerConfig,
/// Number of retries attempted so far.
pub retry_count: u32,
/// Timestamp (Unix ms) when the next retry should be attempted.
pub retry_after_ms: u64,
/// Whether this is an auto-reconnect (unlimited retries, ignores max_retries).
pub reconnect: bool,
/// Optional absolute expiry for this retry entry (Unix ms).
///
/// When set, retries are dropped after this point even if reconnect logic
/// would otherwise continue.
pub expires_at_ms: Option<u64>,
}
impl RetryState {
/// Create a new retry state for a peer.
pub fn new(peer_config: PeerConfig) -> Self {
Self {
peer_config,
retry_count: 0,
retry_after_ms: 0,
reconnect: false,
expires_at_ms: None,
}
}
}
-335
View File
@@ -1,335 +0,0 @@
//! Connection retry logic for auto-connect peers.
//!
//! When an outbound handshake fails (timeout or send error), the node can
//! automatically retry with exponential backoff. Retry state lives on Node
//! (not PeerConnection) because each retry creates a fresh connection.
use super::{Node, NodeError};
use crate::PeerIdentity;
use crate::config::PeerConfig;
use crate::identity::NodeAddr;
use crate::proto::fmp::backoff_ms;
use tracing::{debug, info, warn};
// MAX_BACKOFF_MS is now derived from config: node.retry.max_backoff_secs * 1000
const MAX_RETRY_CONNECTIONS_PER_TICK: usize = 16;
/// Tracks retry state for a peer across connection attempts.
pub struct RetryState {
/// The peer config to use for initiating retries.
pub peer_config: PeerConfig,
/// Number of retries attempted so far.
pub retry_count: u32,
/// Timestamp (Unix ms) when the next retry should be attempted.
pub retry_after_ms: u64,
/// Whether this is an auto-reconnect (unlimited retries, ignores max_retries).
pub reconnect: bool,
/// Optional absolute expiry for this retry entry (Unix ms).
///
/// When set, retries are dropped after this point even if reconnect logic
/// would otherwise continue.
pub expires_at_ms: Option<u64>,
}
impl RetryState {
/// Create a new retry state for a peer.
pub fn new(peer_config: PeerConfig) -> Self {
Self {
peer_config,
retry_count: 0,
retry_after_ms: 0,
reconnect: false,
expires_at_ms: None,
}
}
}
impl Node {
/// Schedule a retry for a failed outbound connection, if applicable.
///
/// Only schedules if the peer is an auto-connect peer and max retries
/// have not been exhausted (unless `reconnect` is true, which retries
/// indefinitely). Does nothing if the peer is already connected or has
/// a connection in progress.
pub(super) fn schedule_retry(&mut self, node_addr: NodeAddr, now_ms: u64) {
let retry_cfg = &self.config().node.retry;
let max_retries = retry_cfg.max_retries;
if max_retries == 0 {
return;
}
// Don't retry if peer is already connected
if self.peers.contains_key(&node_addr) {
return;
}
let base_interval_ms = retry_cfg.base_interval_secs * 1000;
let max_backoff_ms = retry_cfg.max_backoff_secs * 1000;
let peer_name = self.peer_display_name(&node_addr);
if let Some(state) = self.retry_pending.get_mut(&node_addr) {
// Already tracking — increment
state.retry_count += 1;
if !state.reconnect && state.retry_count > max_retries {
info!(
peer = %peer_name,
attempts = state.retry_count,
"Max retries exhausted, giving up on peer"
);
self.retry_pending.remove(&node_addr);
return;
}
let delay = backoff_ms(state.retry_count, base_interval_ms, max_backoff_ms);
state.retry_after_ms = now_ms + delay;
debug!(
peer = %peer_name,
retry = state.retry_count,
reconnect = state.reconnect,
delay_secs = delay / 1000,
"Scheduling connection retry"
);
} else {
// First failure — find the matching PeerConfig
let peer_config = self
.config()
.auto_connect_peers()
.find(|pc| {
PeerIdentity::from_npub(&pc.npub)
.map(|id| *id.node_addr() == node_addr)
.unwrap_or(false)
})
.cloned();
if let Some(pc) = peer_config {
let mut state = RetryState::new(pc);
state.retry_count = 1;
state.reconnect = true;
let delay = backoff_ms(state.retry_count, base_interval_ms, max_backoff_ms);
state.retry_after_ms = now_ms + delay;
debug!(
peer = %self.peer_display_name(&node_addr),
delay_secs = delay / 1000,
"First connection attempt failed, scheduling retry"
);
self.retry_pending.insert(node_addr, state);
}
// If not found in auto_connect_peers, no retry (one-shot connection)
}
}
/// Schedule auto-reconnect for a peer removed by MMP dead timeout.
///
/// Looks up the peer in auto-connect config and checks `auto_reconnect`.
/// If enabled, feeds the peer into the retry system with unlimited retries.
///
/// If a retry entry already exists (e.g. from a previous failed handshake
/// attempt during an earlier reconnect cycle), the existing retry count is
/// preserved and incremented rather than reset to zero. This ensures
/// exponential backoff accumulates across repeated link-dead events instead
/// of resetting to the base interval on every peer removal.
pub(super) fn schedule_reconnect(&mut self, node_addr: NodeAddr, now_ms: u64) {
// Find peer in auto-connect config
let peer_config = self
.config()
.auto_connect_peers()
.find(|pc| {
PeerIdentity::from_npub(&pc.npub)
.map(|id| *id.node_addr() == node_addr)
.unwrap_or(false)
})
.cloned();
let Some(pc) = peer_config else {
return; // Not an auto-connect peer, no reconnect
};
if !pc.auto_reconnect {
debug!(
peer = %self.peer_display_name(&node_addr),
"Auto-reconnect disabled for peer, skipping"
);
return;
}
let base_interval_ms = self.config().node.retry.base_interval_secs * 1000;
let max_backoff_ms = self.config().node.retry.max_backoff_secs * 1000;
let peer_name = self.peer_display_name(&node_addr);
// If we already have accumulated backoff from previous failed attempts,
// preserve and bump it rather than resetting to zero. This prevents the
// exponential backoff from being discarded on each link-dead cycle.
if let Some(state) = self.retry_pending.get_mut(&node_addr) {
state.reconnect = true;
state.retry_count += 1;
let delay = backoff_ms(state.retry_count, base_interval_ms, max_backoff_ms);
state.retry_after_ms = now_ms + delay;
debug!(
peer = %peer_name,
retry = state.retry_count,
delay_secs = delay / 1000,
"Scheduling auto-reconnect after link-dead removal (backoff preserved)"
);
return;
}
let mut state = RetryState::new(pc);
state.reconnect = true;
let delay = backoff_ms(state.retry_count, base_interval_ms, max_backoff_ms);
state.retry_after_ms = now_ms + delay;
debug!(
peer = %peer_name,
delay_secs = delay / 1000,
"Scheduling auto-reconnect after link-dead removal"
);
self.retry_pending.insert(node_addr, state);
}
/// Process pending retries whose time has arrived.
///
/// For each due retry, initiates a fresh connection attempt. The retry
/// entry stays in `retry_pending` until the connection succeeds (cleared
/// in `promote_connection`) or max retries are exhausted (cleared in
/// `schedule_retry`).
pub(super) async fn process_pending_retries(&mut self, now_ms: u64) {
if self.retry_pending.is_empty() {
return;
}
let expired: Vec<NodeAddr> = self
.retry_pending
.iter()
.filter_map(|(addr, state)| {
state
.expires_at_ms
.filter(|expires_at_ms| now_ms >= *expires_at_ms)
.map(|_| *addr)
})
.collect();
for node_addr in expired {
self.retry_pending.remove(&node_addr);
info!(
peer = %self.peer_display_name(&node_addr),
"Retry window expired, dropping pending retry state"
);
}
if self.retry_pending.is_empty() {
return;
}
if !self.outbound_admission_check() {
debug!(
peers = self.peers.len(),
max_peers = self.max_peers(),
retry_pending = self.retry_pending.len(),
"Suppressing auto-reconnect retries: at capacity"
);
return;
}
// Collect retries that are due
let due: Vec<NodeAddr> = self
.retry_pending
.iter()
.filter(|(_, state)| now_ms >= state.retry_after_ms)
.map(|(addr, _)| *addr)
.collect();
let deferred = due.len().saturating_sub(MAX_RETRY_CONNECTIONS_PER_TICK);
if deferred > 0 {
debug!(
due = due.len(),
processing = MAX_RETRY_CONNECTIONS_PER_TICK,
deferred,
"Retry processing budget exhausted; deferring remaining peers"
);
}
for node_addr in due.into_iter().take(MAX_RETRY_CONNECTIONS_PER_TICK) {
// Peer may have connected inbound while we waited
if self.peers.contains_key(&node_addr) {
self.retry_pending.remove(&node_addr);
continue;
}
let state = match self.retry_pending.get(&node_addr) {
Some(s) => s,
None => continue,
};
debug!(
peer = %self.peer_display_name(&node_addr),
retry = state.retry_count,
"Attempting connection retry"
);
let peer_config = state.peer_config.clone();
// Refresh the peer's overlay advert before retrying. The cache is
// read-only on hit (see fetch_advert), so every retry without a
// refetch dials the same cached endpoint — and the most common
// reason a peer ended up in retry_pending is that the cached
// endpoint just stopped working (NAT rebind, port change, peer
// restart on a different port). Without this refresh the retry
// loop dials the same dead address forever.
//
// refetch_advert_for_stale_check uses the relay's advert as
// ground truth: replaces the cache if there's a newer one,
// evicts if the relay has nothing, otherwise leaves it. Cheap
// (one Filter fetch with 2s timeout) and bounded by the retry
// backoff cadence.
if let Some(bootstrap) = self.supervisor.nostr_rendezvous.engine_arc() {
let _ = bootstrap
.refetch_advert_for_stale_check(&peer_config.npub)
.await;
}
match self.initiate_peer_connection(&peer_config).await {
Ok(()) => {
// Push retry_after_ms past the handshake timeout window so
// we don't re-fire on the next tick. If the handshake
// succeeds, promote_connection() clears retry_pending. If
// it times out, check_timeouts() calls schedule_retry()
// which bumps the counter and applies proper backoff.
let hs_timeout_ms = self.config().node.rate_limit.handshake_timeout_secs * 1000;
if let Some(state) = self.retry_pending.get_mut(&node_addr) {
state.retry_after_ms = now_ms + hs_timeout_ms;
}
debug!(
peer = %self.peer_display_name(&node_addr),
"Retry connection initiated, suppressing re-fire for {}s",
self.config().node.rate_limit.handshake_timeout_secs,
);
}
Err(e) => {
warn!(
peer = %self.peer_display_name(&node_addr),
error = %e,
"Retry connection initiation failed"
);
// No-transport failures usually mean the cached overlay
// advert is stale (peer rebound NAT, switched relay, etc.).
// The advert cache is read-only inside fetch_advert, so
// every retry returns the same dead address until the
// entry expires. Force a re-fetch so the next retry tick
// picks up fresh endpoints.
if matches!(e, NodeError::NoTransportForType(_))
&& let Some(bootstrap) = self.supervisor.nostr_rendezvous.engine_arc()
{
let npub = peer_config.npub.clone();
tokio::spawn(async move {
let _ = bootstrap.refetch_advert_for_stale_check(&npub).await;
});
}
// Immediate failure counts as an attempt — schedule next retry
// (reconnect flag is preserved on existing retry_pending entry)
self.schedule_retry(node_addr, now_ms);
}
}
}
}
}
+28 -7
View File
@@ -1142,32 +1142,53 @@ async fn test_open_discovery_sweep_queues_eligible_skips_filtered() {
bootstrap.insert_advert_for_test(npub.clone(), advert).await;
}
// The sweep now runs through the gate-checked reconciler overlay layer,
// which is inert unless the node is Running/Degraded. In production the
// sweep fires only from the rx_loop tick (which spins after `start()`
// returns Running), so drive the node into `Running` to reflect that.
node.supervisor.state = crate::node::NodeState::Running;
// Run the sweep.
node.run_open_discovery_sweep(&bootstrap, Some(3_600), "test")
.await;
node.run_open_discovery_sweep(&bootstrap, Some(3_600)).await;
// Eligible peer was queued.
assert!(
node.retry_pending.contains_key(&eligible_node_addr),
node.peering
.reconciler
.retry_pending
.contains_key(&eligible_node_addr),
"eligible advert should be queued for retry"
);
let queued = node.retry_pending.get(&eligible_node_addr).unwrap();
let queued = node
.peering
.reconciler
.retry_pending
.get(&eligible_node_addr)
.unwrap();
assert_eq!(queued.peer_config.npub, eligible_npub);
// Connected-peer skip filter held.
assert!(
!node.retry_pending.contains_key(&connected_node_addr),
!node
.peering
.reconciler
.retry_pending
.contains_key(&connected_node_addr),
"advert for already-connected peer must not be queued"
);
// Self skip filter held.
assert!(
!node.retry_pending.contains_key(&self_node_addr),
!node
.peering
.reconciler
.retry_pending
.contains_key(&self_node_addr),
"advert authored by own node must not be queued"
);
// Exactly one queued entry from the three injected adverts.
assert_eq!(node.retry_pending.len(), 1);
assert_eq!(node.peering.reconciler.retry_pending.len(), 1);
}
// ============================================================================
+126 -45
View File
@@ -773,12 +773,17 @@ fn test_schedule_retry_creates_entry() {
let mut node = Node::new(config).unwrap();
assert!(node.retry_pending.is_empty());
assert!(node.peering.reconciler.retry_pending.is_empty());
node.schedule_retry(peer_node_addr, 1000);
node.note_handshake_timeout(peer_node_addr, 1000);
assert_eq!(node.retry_pending.len(), 1);
let state = node.retry_pending.get(&peer_node_addr).unwrap();
assert_eq!(node.peering.reconciler.retry_pending.len(), 1);
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap();
assert_eq!(state.retry_count, 1);
assert!(
state.reconnect,
@@ -806,15 +811,25 @@ fn test_schedule_retry_increments() {
let mut node = Node::new(config).unwrap();
// First failure
node.schedule_retry(peer_node_addr, 1000);
node.note_handshake_timeout(peer_node_addr, 1000);
assert_eq!(
node.retry_pending.get(&peer_node_addr).unwrap().retry_count,
node.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap()
.retry_count,
1
);
// Second failure
node.schedule_retry(peer_node_addr, 11_000);
let state = node.retry_pending.get(&peer_node_addr).unwrap();
node.note_handshake_timeout(peer_node_addr, 11_000);
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap();
assert_eq!(state.retry_count, 2);
// backoff_ms(5000) with retry_count=2 = 5000 * 4 = 20000
assert_eq!(state.retry_after_ms, 11_000 + 20_000);
@@ -832,9 +847,9 @@ async fn test_process_pending_retries_is_budgeted_per_tick() {
let npub = identity.npub();
let peer_identity = PeerIdentity::from_npub(&npub).unwrap();
let node_addr = *peer_identity.node_addr();
node.retry_pending.insert(
node.peering.reconciler.retry_pending.insert(
node_addr,
crate::node::retry::RetryState {
crate::node::peering::retry::RetryState {
peer_config: crate::config::PeerConfig::new(npub, "udp", "10.0.0.2:2121"),
retry_count: 0,
retry_after_ms: 0,
@@ -845,12 +860,18 @@ async fn test_process_pending_retries_is_budgeted_per_tick() {
addrs.push(node_addr);
}
// The retry-dial tick runs only under a `Reconciling` gate (it fires from the
// rx loop, which spins only after start() reaches Running). Put the node in
// Running so the retry-dial budget — the property under test — is exercised.
node.supervisor.state = crate::node::NodeState::Running;
node.process_pending_retries(1).await;
let processed = addrs
.iter()
.filter(|addr| {
node.retry_pending
node.peering
.reconciler
.retry_pending
.get(addr)
.is_some_and(|state| state.retry_count > 0)
})
@@ -859,7 +880,7 @@ async fn test_process_pending_retries_is_budgeted_per_tick() {
assert_eq!(processed, 16);
assert_eq!(deferred, 4);
assert_eq!(node.retry_pending.len(), 20);
assert_eq!(node.peering.reconciler.retry_pending.len(), 20);
}
/// Test that auto-connect peers retry indefinitely (never exhaust).
@@ -880,20 +901,38 @@ fn test_schedule_retry_auto_connect_never_exhausts() {
let mut node = Node::new(config).unwrap();
// All attempts should keep the entry alive despite max_retries=2
node.schedule_retry(peer_node_addr, 1000);
assert!(node.retry_pending.contains_key(&peer_node_addr));
node.note_handshake_timeout(peer_node_addr, 1000);
assert!(
node.peering
.reconciler
.retry_pending
.contains_key(&peer_node_addr)
);
node.schedule_retry(peer_node_addr, 2000);
assert!(node.retry_pending.contains_key(&peer_node_addr));
node.note_handshake_timeout(peer_node_addr, 2000);
assert!(
node.peering
.reconciler
.retry_pending
.contains_key(&peer_node_addr)
);
// Attempt 3 would have exhausted before, but now retries indefinitely
node.schedule_retry(peer_node_addr, 3000);
node.note_handshake_timeout(peer_node_addr, 3000);
assert!(
node.retry_pending.contains_key(&peer_node_addr),
node.peering
.reconciler
.retry_pending
.contains_key(&peer_node_addr),
"Auto-connect peers should never exhaust retries"
);
assert_eq!(
node.retry_pending.get(&peer_node_addr).unwrap().retry_count,
node.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap()
.retry_count,
3
);
}
@@ -915,9 +954,9 @@ fn test_schedule_retry_disabled() {
let mut node = Node::new(config).unwrap();
node.schedule_retry(peer_node_addr, 1000);
node.note_handshake_timeout(peer_node_addr, 1000);
assert!(
node.retry_pending.is_empty(),
node.peering.reconciler.retry_pending.is_empty(),
"No retry should be scheduled when max_retries=0"
);
}
@@ -931,9 +970,9 @@ fn test_schedule_retry_ignores_non_autoconnect() {
// No peers configured at all
let mut node = make_node();
node.schedule_retry(peer_node_addr, 1000);
node.note_handshake_timeout(peer_node_addr, 1000);
assert!(
node.retry_pending.is_empty(),
node.peering.reconciler.retry_pending.is_empty(),
"No retry for unconfigured peer"
);
}
@@ -953,9 +992,9 @@ fn test_schedule_retry_skips_connected_peer() {
assert_eq!(node.peer_count(), 1);
// Scheduling a retry for an already-connected peer should be a no-op
node.schedule_retry(node_addr, 3000);
node.note_handshake_timeout(node_addr, 3000);
assert!(
node.retry_pending.is_empty(),
node.peering.reconciler.retry_pending.is_empty(),
"No retry for already-connected peer"
);
}
@@ -1199,7 +1238,7 @@ async fn test_nostr_traversal_failure_skips_connected_peer() {
"stale failures for connected peers must not affect traversal cooldown"
);
assert!(
node.retry_pending.is_empty(),
node.peering.reconciler.retry_pending.is_empty(),
"stale failures for connected peers must not enqueue reconnect attempts"
);
}
@@ -1247,7 +1286,7 @@ async fn test_nostr_traversal_established_skips_connected_peer() {
"stale established handoff must not start a new handshake"
);
assert!(
node.retry_pending.is_empty(),
node.peering.reconciler.retry_pending.is_empty(),
"stale established handoff must not enqueue a reconnect"
);
}
@@ -1259,7 +1298,7 @@ async fn test_process_pending_retries_drops_expired_entries() {
let peer_npub = peer_identity.npub();
let peer_node_addr = *PeerIdentity::from_npub(&peer_npub).unwrap().node_addr();
let mut state = super::super::retry::RetryState::new(crate::config::PeerConfig::new(
let mut state = super::super::peering::retry::RetryState::new(crate::config::PeerConfig::new(
peer_npub,
"udp",
"127.0.0.1:9",
@@ -1267,12 +1306,22 @@ async fn test_process_pending_retries_drops_expired_entries() {
state.retry_after_ms = 0;
state.expires_at_ms = Some(1_000);
state.reconnect = true;
node.retry_pending.insert(peer_node_addr, state);
node.peering
.reconciler
.retry_pending
.insert(peer_node_addr, state);
// Retry-dial runs only under a `Reconciling` gate; put the node in Running so
// the expired-entry drop (the property under test) is reached.
node.supervisor.state = crate::node::NodeState::Running;
node.process_pending_retries(1_000).await;
assert!(
!node.retry_pending.contains_key(&peer_node_addr),
!node
.peering
.reconciler
.retry_pending
.contains_key(&peer_node_addr),
"expired retry entries should be dropped before retry processing"
);
}
@@ -1300,19 +1349,29 @@ fn test_schedule_reconnect_preserves_backoff() {
let mut node = Node::new(config).unwrap();
// Simulate two stale handshake timeouts incrementing the retry count.
node.schedule_retry(peer_node_addr, 1_000); // count=1, delay=10s
node.schedule_retry(peer_node_addr, 11_000); // count=2, delay=20s
node.note_handshake_timeout(peer_node_addr, 1_000); // count=1, delay=10s
node.note_handshake_timeout(peer_node_addr, 11_000); // count=2, delay=20s
{
let state = node.retry_pending.get(&peer_node_addr).unwrap();
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap();
assert_eq!(state.retry_count, 2, "Two failures should yield count=2");
}
// Now simulate a link-dead removal triggering schedule_reconnect.
// The existing retry entry (count=2) should be preserved and bumped to 3,
// NOT reset to 0 as it was before the fix.
node.schedule_reconnect(peer_node_addr, 31_000);
node.note_link_dead(peer_node_addr, 31_000);
let state = node.retry_pending.get(&peer_node_addr).unwrap();
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap();
assert!(state.reconnect, "Entry should be marked as reconnect");
assert_eq!(
state.retry_count, 3,
@@ -1347,9 +1406,14 @@ fn test_schedule_reconnect_fresh_state() {
let mut node = Node::new(config).unwrap();
// No prior retry entry — first reconnect should use base delay.
node.schedule_reconnect(peer_node_addr, 1_000);
node.note_link_dead(peer_node_addr, 1_000);
let state = node.retry_pending.get(&peer_node_addr).unwrap();
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.unwrap();
assert!(state.reconnect, "Entry should be marked as reconnect");
assert_eq!(
state.retry_count, 0,
@@ -1389,6 +1453,8 @@ fn test_disconnect_schedules_reconnect() {
node.handle_disconnect(&peer_node_addr, &payload);
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.expect("handle_disconnect should schedule reconnect for auto-connect peer");
@@ -1410,17 +1476,21 @@ fn test_promote_clears_retry_pending() {
let node_addr = *identity.node_addr();
// Simulate a retry entry existing for this peer
node.retry_pending.insert(
node.peering.reconciler.retry_pending.insert(
node_addr,
super::super::retry::RetryState::new(crate::config::PeerConfig::default()),
super::super::peering::retry::RetryState::new(crate::config::PeerConfig::default()),
);
assert_eq!(node.retry_pending.len(), 1);
assert_eq!(node.peering.reconciler.retry_pending.len(), 1);
node.add_connection(conn).unwrap();
node.promote_connection(link_id, identity, 2000).unwrap();
assert!(
!node.retry_pending.contains_key(&node_addr),
!node
.peering
.reconciler
.retry_pending
.contains_key(&node_addr),
"retry_pending should be cleared on successful promotion"
);
}
@@ -1446,12 +1516,15 @@ async fn test_initiate_peer_connections_schedules_retry_on_no_transport() {
));
let mut node = Node::new(config).unwrap();
assert!(node.retry_pending.is_empty());
assert!(node.peering.reconciler.retry_pending.is_empty());
node.initiate_peer_connections().await;
assert!(
node.retry_pending.contains_key(&peer_node_addr),
node.peering
.reconciler
.retry_pending
.contains_key(&peer_node_addr),
"startup peer-init failure must enqueue a retry so the peer can recover \
without a daemon restart"
);
@@ -1723,18 +1796,24 @@ async fn process_pending_retries_gated_at_capacity() {
let peer_identity = Identity::generate();
let peer_npub = peer_identity.npub();
let peer_node_addr = *PeerIdentity::from_npub(&peer_npub).unwrap().node_addr();
let mut state = super::super::retry::RetryState::new(crate::config::PeerConfig::new(
let mut state = super::super::peering::retry::RetryState::new(crate::config::PeerConfig::new(
peer_npub,
"udp",
"127.0.0.1:9",
));
state.retry_after_ms = 0;
state.reconnect = true;
node.retry_pending.insert(peer_node_addr, state);
node.peering
.reconciler
.retry_pending
.insert(peer_node_addr, state);
let before_peers = node.peer_count();
let before_connections = node.connection_count();
// Running gate so the admission short-circuit (not the startup gate) is what
// suppresses the dial — the fingerprint this test asserts on.
node.supervisor.state = crate::node::NodeState::Running;
node.process_pending_retries(1_000).await;
// At capacity: gate short-circuits before due-list collection. The
@@ -1744,6 +1823,8 @@ async fn process_pending_retries_gated_at_capacity() {
// (which fails without a registered transport), and the failure
// handler would call `schedule_retry`, bumping `retry_count` to 1.
let state = node
.peering
.reconciler
.retry_pending
.get(&peer_node_addr)
.expect("retry entry must be preserved when suppressed at capacity");