diff --git a/CHANGELOG.md b/CHANGELOG.md index e5ee08bd..31667d9f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +#### Node lifecycle + +- A heartbeat whose send failed no longer counts as one that was delivered. + The peer's "last heartbeat" timestamp was stamped before the send and left + alone whatever came back, so a failure suppressed the next attempt for a + full `node.heartbeat_interval_secs` even though the peer had heard nothing — + on a 10s interval against a 30s `link_dead_timeout_secs`, three failures in + a row were the whole budget. The timestamp now moves only on a send that + returned cleanly, and a separate record of the *attempt* spaces the retries + so a peer that keeps failing is retried in seconds rather than either + hammered every tick or left for a full interval. That retry spacing + applies to the failure path only: gating a healthy peer on it as well + would have floored `node.heartbeat_interval_secs` at two seconds, so a + configured value below that would silently not have been honoured. + #### Data plane - A peer that stops reading can no longer stall the node. TCP, Tor, Nym and diff --git a/src/node/handlers/mmp.rs b/src/node/handlers/mmp.rs index 1f7b3b9d..28651e61 100644 --- a/src/node/handlers/mmp.rs +++ b/src/node/handlers/mmp.rs @@ -20,6 +20,21 @@ use crate::transport::{TransportAddr, TransportId}; use std::time::{Duration, Instant}; use tracing::{debug, info, trace, warn}; +/// How long a peer whose heartbeat send *failed* waits before the next attempt. +/// +/// Applies to the failure path only. Gating a healthy peer on it too would +/// floor `node.heartbeat_interval_secs` at this value without validating or +/// reporting it, which is a configured knob quietly not doing what it says. +/// +/// Short against `heartbeat_interval_secs`, because a failed heartbeat means +/// the peer has heard nothing and the point is to recover well inside +/// `link_dead_timeout_secs` rather than after another full interval. Not +/// shorter still, because the send behind it awaits an unbounded `write_all` +/// on a connection-oriented transport, on the rx loop; retrying that every +/// tick would make a stranded stream a stalled node. Once that write is +/// bounded this can come down to the tick. +const HEARTBEAT_RETRY_INTERVAL: Duration = Duration::from_secs(2); + /// Emit the operator `trace!` point for a processed ReceiverReport outcome. /// /// These log points used to live inside `MmpMetrics::process_receiver_report`; @@ -473,11 +488,42 @@ impl Node { && peer.rekey_msg1_resend_count() < max_resends && peer.rekey_msg1().is_some(); - // Check if heartbeat is due. - let heartbeat_due = match peer.last_heartbeat_sent() { + // Check if heartbeat is due. Two gates, not one. The first is + // the interval since a heartbeat last *landed*; a send that + // failed does not satisfy it, so a peer that has heard nothing + // stays due instead of being suppressed for a full interval by + // an attempt that went nowhere. + // + // The second spaces the retries out. Without it a peer whose + // send keeps failing would be retried on every tick, and the + // send behind this is not always cheap: on a connection-oriented + // transport it awaits an unbounded `write_all` on the rx loop. + // Until that is bounded, retrying a failing peer once a second + // would turn a stranded stream into a stalled node. + let heartbeat_landed_due = match peer.last_heartbeat_sent() { None => true, Some(last) => now.duration_since(last) >= heartbeat_interval, }; + // The retry gate applies only after a *failure*. A successful + // send stamps both timestamps with the same instant, so on a + // healthy peer an unconditional gate would floor the configured + // interval at HEARTBEAT_RETRY_INTERVAL — silently turning a + // configured `heartbeat_interval_secs` of 1 into 2, with no + // validation refusing the value and nothing saying why. An + // attempt strictly newer than the last success is the only + // state that means "the last one did not land". + let last_attempt_failed = + match (peer.last_heartbeat_attempt(), peer.last_heartbeat_sent()) { + (Some(attempt), Some(sent)) => attempt > sent, + (Some(_), None) => true, + _ => false, + }; + let retry_due = !last_attempt_failed + || match peer.last_heartbeat_attempt() { + None => true, + Some(last) => now.duration_since(last) >= HEARTBEAT_RETRY_INTERVAL, + }; + let heartbeat_due = heartbeat_landed_due && retry_due; PeerLivenessSnapshot { peer: *node_addr, @@ -513,14 +559,25 @@ impl Node { self.route_link_dead(peer, now_ms).await; } MmpAction::Heartbeat { peer } => { + // Attempt first, success after: the attempt is recorded + // even if the send below fails or never returns, so the + // retry stays spaced; only a send that came back clean + // moves the interval that says the peer has heard from us. if let Some(p) = self.peers.get_mut(&peer) { - p.mark_heartbeat_sent(now); + p.mark_heartbeat_attempt(now); } - if let Err(e) = self + match self .send_encrypted_link_message(&peer, &heartbeat_msg) .await { - trace!(peer = %self.peer_display_name(&peer), error = %e, "Failed to send heartbeat"); + Ok(()) => { + if let Some(p) = self.peers.get_mut(&peer) { + p.mark_heartbeat_sent(now); + } + } + Err(e) => { + trace!(peer = %self.peer_display_name(&peer), error = %e, "Failed to send heartbeat"); + } } } MmpAction::SendLinkReport { .. } diff --git a/src/node/handlers/netmon.rs b/src/node/handlers/netmon.rs index cddd1326..de6f551d 100644 --- a/src/node/handlers/netmon.rs +++ b/src/node/handlers/netmon.rs @@ -130,7 +130,10 @@ impl Node { /// Send one heartbeat to every peer whose send path cannot block, so each /// learns the node's new source address in one RTT rather than at the next - /// due interval. Returns how many went out. + /// due interval. Returns how many sends actually succeeded, which is what + /// the operator log reports — a count of peers *selected* would read the + /// same whether every frame left or none did, and a medium change is + /// exactly when sends start failing. /// /// The filter was written for a hazard that no longer exists, and it is /// kept deliberately rather than by oversight. It was this: a @@ -161,7 +164,7 @@ impl Node { /// dropping the stale connection /// rather than writing into it, which is a different change with a real /// cost behind it — a Tor peer pays a fresh circuit — and is not this one. - async fn heartbeat_all_peers_after_net_change(&mut self) -> usize { + pub(in crate::node) async fn heartbeat_all_peers_after_net_change(&mut self) -> usize { let now = Instant::now(); let heartbeat = [LinkMessageType::Heartbeat.to_byte()]; let targets: Vec = self @@ -175,17 +178,25 @@ impl Node { .map(|(addr, _)| *addr) .collect(); - let sent = targets.len(); + let mut sent = 0usize; for addr in targets { if let Some(peer) = self.peers.get_mut(&addr) { - peer.mark_heartbeat_sent(now); + peer.mark_heartbeat_attempt(now); } - if let Err(e) = self.send_encrypted_link_message(&addr, &heartbeat).await { - debug!( - peer = %self.peer_display_name(&addr), - error = %e, - "Failed to send post-medium-change heartbeat" - ); + match self.send_encrypted_link_message(&addr, &heartbeat).await { + Ok(()) => { + if let Some(peer) = self.peers.get_mut(&addr) { + peer.mark_heartbeat_sent(now); + } + sent += 1; + } + Err(e) => { + debug!( + peer = %self.peer_display_name(&addr), + error = %e, + "Failed to send post-medium-change heartbeat" + ); + } } } sent diff --git a/src/node/tests/heartbeat.rs b/src/node/tests/heartbeat.rs index b5ff6e3e..5da32b7d 100644 --- a/src/node/tests/heartbeat.rs +++ b/src/node/tests/heartbeat.rs @@ -123,3 +123,163 @@ async fn heartbeat_unaffected_without_rekey() { cleanup_nodes(&mut nodes).await; } + +/// Set `node.heartbeat_interval_secs`, the knob the retry gate must not floor. +fn set_heartbeat_interval(node: &mut crate::node::Node, secs: u64) { + node.replace_context(|ctx| { + let mut cfg = (*ctx.config).clone(); + cfg.node.heartbeat_interval_secs = secs; + ctx.config = std::sync::Arc::new(cfg); + }); +} + +/// Rewind a peer's heartbeat bookkeeping by `age`, as if that long had passed +/// since its last successful send. +/// +/// The sweep reads `std::time::Instant`, which tokio's paused clock does not +/// move, so elapsed time is staged on the peer rather than waited out. Sets +/// both timestamps, which is the state a *healthy* peer is in. +fn age_heartbeat(node: &mut crate::node::Node, addr: &NodeAddr, age: Duration) { + let then = std::time::Instant::now() - age; + node.peers + .get_mut(addr) + .expect("peer present") + .mark_heartbeat_sent(then); +} + +/// **The retry gate must not floor a healthy peer's configured interval.** +/// +/// A successful send stamps `last_heartbeat_sent` and `last_heartbeat_attempt` +/// with the same instant. Gating every peer on the attempt timestamp therefore +/// gates the healthy path too, and the effective interval becomes the larger of +/// the configured value and `HEARTBEAT_RETRY_INTERVAL` — so a configured 1s +/// becomes 2s, silently, with nothing validating the value and nothing saying +/// why. `src/node/tests/tcp.rs` already configures 1s against a 3s dead +/// timeout, which is the margin that would quietly halve. +#[tokio::test] +async fn a_healthy_peer_is_heartbeated_on_its_configured_interval() { + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + set_heartbeat_interval(&mut nodes[0].node, 1); + + // Past the configured interval, short of the failure-retry interval. That + // window is the whole defect: healthy, due, and gated anyway. + age_heartbeat(&mut nodes[0].node, &addr_1, Duration::from_millis(1_200)); + let before = nodes[0] + .node + .get_peer(&addr_1) + .expect("peer 1 is established") + .last_heartbeat_sent() + .expect("staged above"); + + nodes[0].node.check_link_heartbeats().await; + + let after = nodes[0] + .node + .get_peer(&addr_1) + .expect("peer 1 is still established") + .last_heartbeat_sent() + .expect("still sent"); + assert!( + after > before, + "a healthy peer must be heartbeated on its configured interval, not \ + floored at the failure-retry interval" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The other half: after a send that *failed*, the retry is spaced out rather +/// than reattempted on the very next tick. +/// +/// Without that spacing a peer whose send keeps failing is retried every tick, +/// and the send behind it can await an unbounded stream write on the rx loop. +/// The peer is re-pinned onto a UDP transport that was never started, so its +/// send fails with `NotStarted` before touching a socket. +#[tokio::test] +async fn a_failing_peer_is_retried_after_the_gap_and_not_before() { + use crate::transport::{TransportAddr, TransportHandle, TransportId, packet_channel}; + + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + set_heartbeat_interval(&mut nodes[0].node, 1); + + let dead_id = TransportId::new(91); + let (tx, _rx) = packet_channel(64); + nodes[0].node.transports.insert( + dead_id, + TransportHandle::Udp(crate::transport::udp::UdpTransport::new( + dead_id, + None, + crate::config::UdpConfig::default(), + tx, + )), + ); + nodes[0] + .node + .peers + .get_mut(&addr_1) + .expect("peer 1 is established") + .set_current_addr(dead_id, TransportAddr::from_string("10.0.0.2:2121")); + + // Long overdue and healthy-looking, so the sweep will try. + age_heartbeat(&mut nodes[0].node, &addr_1, Duration::from_secs(10)); + nodes[0].node.check_link_heartbeats().await; + + let attempt_1 = nodes[0] + .node + .get_peer(&addr_1) + .expect("peer 1 is established") + .last_heartbeat_attempt() + .expect("a failed send is still an attempt"); + assert!( + nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .last_heartbeat_sent() + .expect("staged") + < attempt_1, + "the failed send must not have stamped a success" + ); + + // Immediately after: still inside the gap, so no second attempt. + nodes[0].node.check_link_heartbeats().await; + assert_eq!( + nodes[0] + .node + .get_peer(&addr_1) + .expect("peer 1 is established") + .last_heartbeat_attempt(), + Some(attempt_1), + "a failing peer must not be retried on the very next tick" + ); + + // Stage the gap as elapsed, keeping the attempt newer than the success so + // the peer still reads as "last one failed". + let past = std::time::Instant::now() - Duration::from_secs(3); + nodes[0] + .node + .peers + .get_mut(&addr_1) + .expect("peer present") + .mark_heartbeat_attempt(past); + + nodes[0].node.check_link_heartbeats().await; + let attempt_2 = nodes[0] + .node + .get_peer(&addr_1) + .expect("peer 1 is established") + .last_heartbeat_attempt() + .expect("still attempted"); + assert!( + attempt_2 > past, + "a failing peer must be retried once the gap has passed" + ); + + cleanup_nodes(&mut nodes).await; +} diff --git a/src/node/tests/netmon.rs b/src/node/tests/netmon.rs index ef6f35a8..72658238 100644 --- a/src/node/tests/netmon.rs +++ b/src/node/tests/netmon.rs @@ -441,3 +441,68 @@ async fn a_peers_connected_socket_publishes_the_source_it_was_pinned_to() { cleanup_nodes(&mut nodes).await; } + +/// A heartbeat that did not go out must not be counted as one that did — in +/// the operator log, or in the peer's own idea of when it was last heard from. +/// +/// A medium change is exactly the condition under which sends start failing, +/// so a count of peers *selected* would read identically whether every frame +/// left or none did, and the peer would then be suppressed for a full +/// `heartbeat_interval_secs` on the strength of a send that never landed. +/// +/// The peer is re-pinned onto a UDP transport that was never started, which is +/// connectionless — so the fan-out selects it — and fails its send with +/// `NotStarted` before touching a socket. +#[tokio::test] +async fn a_heartbeat_that_failed_is_not_counted_and_does_not_suppress_the_next() { + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + let peer_1 = identity_of(&nodes, 1); + configure_auto_peer(&mut nodes[0].node, &peer_1); + + let dead_id = TransportId::new(88); + let (tx, _rx) = packet_channel(64); + nodes[0].node.transports.insert( + dead_id, + TransportHandle::Udp(crate::transport::udp::UdpTransport::new( + dead_id, + None, + crate::config::UdpConfig::default(), + tx, + )), + ); + nodes[0] + .node + .peers + .get_mut(&addr_1) + .expect("peer 1 is established") + .set_current_addr(dead_id, TransportAddr::from_string("10.0.0.2:2121")); + + let before = nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .last_heartbeat_sent(); + + let sent = nodes[0].node.heartbeat_all_peers_after_net_change().await; + + assert_eq!( + sent, 0, + "the count reports sends that succeeded, not peers picked out" + ); + + let peer = nodes[0].node.get_peer(&addr_1).unwrap(); + assert_eq!( + peer.last_heartbeat_sent(), + before, + "a failed heartbeat must not move the interval that says the peer has heard from us" + ); + assert!( + peer.last_heartbeat_attempt().is_some(), + "the attempt is still recorded, or a failing peer would be retried every tick" + ); + + cleanup_nodes(&mut nodes).await; +} diff --git a/src/peer/active.rs b/src/peer/active.rs index 36cafb35..33c7afc3 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -249,8 +249,16 @@ pub struct ActivePeer { remote_epoch: Option<[u8; 8]>, // === Heartbeat === - /// When we last sent a heartbeat to this peer. + /// When a heartbeat to this peer last *succeeded*. A send that failed does + /// not move this: it did not tell the peer anything, and treating it as if + /// it had would leave the peer un-heartbeated for a full interval on the + /// strength of a send that never landed. last_heartbeat_sent: Option, + /// When a heartbeat to this peer was last *attempted*, successfully or + /// not. Paired with the above so a failing peer is retried sooner than the + /// heartbeat interval without being retried on every tick — see + /// `HEARTBEAT_RETRY_INTERVAL`. + last_heartbeat_attempt: Option, // === Handshake Resend === /// Wire-format msg2 for resend on duplicate msg1 (responder only). @@ -311,6 +319,7 @@ impl ActivePeer { authenticated_at, remote_epoch: None, last_heartbeat_sent: None, + last_heartbeat_attempt: None, handshake_msg2: None, session_established_at: now, rekey_jitter_secs: draw_rekey_jitter(), @@ -390,6 +399,7 @@ impl ActivePeer { authenticated_at, remote_epoch, last_heartbeat_sent: None, + last_heartbeat_attempt: None, handshake_msg2: None, session_established_at: now, rekey_jitter_secs: draw_rekey_jitter(), @@ -796,14 +806,32 @@ impl ActivePeer { // === Heartbeat === - /// When we last sent a heartbeat to this peer. + /// When a heartbeat to this peer last succeeded. pub fn last_heartbeat_sent(&self) -> Option { self.last_heartbeat_sent } - /// Record that we sent a heartbeat. + /// Record that a heartbeat reached the transport without error. + /// + /// Call this *after* the send, and only on success. Marking before the + /// send makes a failed heartbeat indistinguishable from a delivered one, + /// which then suppresses the next attempt for a full + /// `heartbeat_interval_secs` even though the peer has heard nothing. pub fn mark_heartbeat_sent(&mut self, now: Instant) { self.last_heartbeat_sent = Some(now); + self.last_heartbeat_attempt = Some(now); + } + + /// When a heartbeat to this peer was last attempted, whatever came of it. + pub fn last_heartbeat_attempt(&self) -> Option { + self.last_heartbeat_attempt + } + + /// Record that a heartbeat send was attempted. Call this before the send, + /// so an attempt that fails — or one that never returns — still spaces the + /// next one out. + pub fn mark_heartbeat_attempt(&mut self, now: Instant) { + self.last_heartbeat_attempt = Some(now); } // === State Updates === diff --git a/src/proto/mmp/core.rs b/src/proto/mmp/core.rs index 4d832bb6..ff882879 100644 --- a/src/proto/mmp/core.rs +++ b/src/proto/mmp/core.rs @@ -39,8 +39,11 @@ pub(crate) struct PeerLivenessSnapshot { /// An FMP rekey handshake is genuinely in flight with retransmission budget /// left; suppresses teardown of an otherwise-silent rekey link. pub rekey_active: bool, - /// A heartbeat is due (`last_heartbeat_sent` is none, or elapsed since it is - /// >= the heartbeat interval). + /// A heartbeat is due. Two conditions, both resolved shell-side: elapsed + /// since the last heartbeat that *landed* is >= the heartbeat interval (or + /// none has), and — only when the last attempt failed — the failure-retry + /// gap has passed since that attempt. The retry gap deliberately does not + /// apply to a healthy peer, or it would floor the configured interval. pub heartbeat_due: bool, } @@ -169,8 +172,15 @@ pub(crate) enum MmpAction { /// Reap a dead peer: the shell runs `remove_active_peer` + /// `schedule_reconnect` (with its wall-clock `now_ms`). ReapPeer { peer: NodeAddr }, - /// Send a heartbeat to `peer`: the shell runs `mark_heartbeat_sent` and the - /// encrypted link send. + /// Send a heartbeat to `peer`: the shell runs `mark_heartbeat_attempt`, + /// then the encrypted link send, and `mark_heartbeat_sent` **only if that + /// send returned cleanly**. + /// + /// The order and the condition are the contract, not an implementation + /// detail. Stamping the success before the send makes a failed heartbeat + /// indistinguishable from a delivered one, which suppresses the next + /// attempt for a full interval even though the peer heard nothing; the + /// separate attempt stamp is what spaces the retries without doing that. Heartbeat { peer: NodeAddr }, /// Build (shell: `proto/mmp/` `build_report` + `encode`) and send the given /// link report over the encrypted link. The interval-advancing