diff --git a/src/node/handlers/mmp.rs b/src/node/handlers/mmp.rs index 28651e61..00d9a112 100644 --- a/src/node/handlers/mmp.rs +++ b/src/node/handlers/mmp.rs @@ -35,6 +35,42 @@ use tracing::{debug, info, trace, warn}; /// bounded this can come down to the tick. const HEARTBEAT_RETRY_INTERVAL: Duration = Duration::from_secs(2); +/// Decide whether a peer is due a heartbeat, from the two timestamps it keeps. +/// +/// Two gates rather than one. `sent` is when a heartbeat last *landed*, and it +/// alone paces a healthy peer. `attempt` is when one was last *tried*, and it +/// gates only a peer whose last try failed, holding the retry off for +/// [`HEARTBEAT_RETRY_INTERVAL`] so a peer whose send keeps failing is not +/// retried on every tick. +/// +/// **The retry gate is deliberately not consulted on the healthy path.** On a +/// peer whose last send succeeded the two timestamps are equal, so gating there +/// would clamp a configured `heartbeat_interval_secs` up to the retry interval, +/// and that setting has no validation floor. +fn heartbeat_due( + sent: Option, + attempt: Option, + now: Instant, + interval: Duration, +) -> bool { + let landed_due = match sent { + None => true, + Some(last) => now.duration_since(last) >= interval, + }; + + // An attempt later than the last success is one that failed, and an attempt + // with no success behind it is the same thing on a peer never reached. + let retry_due = match (attempt, sent) { + (Some(last), Some(landed)) if last > landed => { + now.duration_since(last) >= HEARTBEAT_RETRY_INTERVAL + } + (Some(last), None) => now.duration_since(last) >= HEARTBEAT_RETRY_INTERVAL, + _ => true, + }; + + landed_due && retry_due +} + /// Emit the operator `trace!` point for a processed ReceiverReport outcome. /// /// These log points used to live inside `MmpMetrics::process_receiver_report`; @@ -488,42 +524,19 @@ impl Node { && peer.rekey_msg1_resend_count() < max_resends && peer.rekey_msg1().is_some(); - // 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; + // Check if heartbeat is due. Two gates, not one: a send that + // failed does not satisfy the interval, so a peer that has + // heard nothing stays due instead of being suppressed by an + // attempt that went nowhere, and the retry gap keeps a peer + // whose send keeps failing from being tried on every tick. + // Both are decided by `heartbeat_due`, which is a pure + // function so it can be tested without driving a send. + let heartbeat_due = heartbeat_due( + peer.last_heartbeat_sent(), + peer.last_heartbeat_attempt(), + now, + heartbeat_interval, + ); PeerLivenessSnapshot { peer: *node_addr, @@ -645,3 +658,97 @@ impl Node { } } } + +#[cfg(test)] +mod tests { + use super::{HEARTBEAT_RETRY_INTERVAL, heartbeat_due}; + use std::time::{Duration, Instant}; + + const INTERVAL: Duration = Duration::from_secs(10); + + #[test] + fn a_peer_never_heartbeated_is_due_immediately() { + let now = Instant::now(); + assert!(heartbeat_due(None, None, now, INTERVAL)); + } + + #[test] + fn a_peer_whose_heartbeat_landed_waits_the_configured_interval() { + let landed = Instant::now(); + assert!(!heartbeat_due( + Some(landed), + Some(landed), + landed + INTERVAL - Duration::from_millis(1), + INTERVAL + )); + assert!(heartbeat_due( + Some(landed), + Some(landed), + landed + INTERVAL, + INTERVAL + )); + } + + #[test] + fn a_healthy_peer_is_paced_by_the_configured_interval_and_not_by_the_retry_floor() { + // The interval a peer configures can be shorter than the retry floor. + // Consulting the retry gate on the healthy path would clamp it, and + // `heartbeat_interval_secs` has no validation floor to prevent that. + let short = Duration::from_secs(1); + assert!(short < HEARTBEAT_RETRY_INTERVAL); + let landed = Instant::now(); + assert!(heartbeat_due( + Some(landed), + Some(landed), + landed + short, + short + )); + } + + #[test] + fn a_failed_attempt_does_not_suppress_the_next_heartbeat_for_a_full_interval() { + // A heartbeat landed at t0 and the next attempt, at t0 + INTERVAL, + // failed. Once the retry interval has passed the peer is due again, + // rather than waiting another whole interval on a send that never + // reached it. + let landed = Instant::now(); + let failed = landed + INTERVAL; + assert!(heartbeat_due( + Some(landed), + Some(failed), + failed + HEARTBEAT_RETRY_INTERVAL, + INTERVAL + )); + } + + #[test] + fn a_failed_attempt_is_not_retried_before_the_retry_interval() { + let landed = Instant::now(); + let failed = landed + INTERVAL; + assert!(!heartbeat_due( + Some(landed), + Some(failed), + failed + HEARTBEAT_RETRY_INTERVAL - Duration::from_millis(1), + INTERVAL + )); + } + + #[test] + fn a_peer_never_reached_is_retried_on_the_retry_interval_not_the_heartbeat_interval() { + // No heartbeat has ever landed, so there is no interval to pace by. + // The attempt alone spaces the retries. + let failed = Instant::now(); + assert!(!heartbeat_due( + None, + Some(failed), + failed + Duration::from_millis(1), + INTERVAL + )); + assert!(heartbeat_due( + None, + Some(failed), + failed + HEARTBEAT_RETRY_INTERVAL, + INTERVAL + )); + } +} diff --git a/src/node/tests/heartbeat.rs b/src/node/tests/heartbeat.rs index 5da32b7d..8d66e425 100644 --- a/src/node/tests/heartbeat.rs +++ b/src/node/tests/heartbeat.rs @@ -35,6 +35,70 @@ fn set_link_dead_timeout(node: &mut crate::node::Node, secs: u64) { }); } +/// Set `node.heartbeat_interval_secs` on an already-constructed node, the same +/// way `set_link_dead_timeout` does. This is 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); + }); +} + +/// A heartbeat whose send failed is not recorded as having landed, and the +/// failed attempt is not retried on the very next tick. +/// +/// The failure is forced by taking the node's transport handles away, so the +/// encrypted send fails before any I/O with `TransportNotFound`. Marking the +/// send before it happens, which is what this replaced, would record the peer +/// as heartbeated and suppress the next attempt for a whole interval although +/// the peer heard nothing. +#[tokio::test] +async fn a_failed_heartbeat_send_is_not_recorded_as_landed() { + let mut nodes = run_tree_test(2, &[(0, 1)], false).await; + verify_tree_convergence(&nodes); + + let addr_1 = *nodes[1].node.node_addr(); + assert!(nodes[0].node.get_peer(&addr_1).is_some()); + + // Whatever landed during convergence is the baseline this asserts against. + let landed_before = nodes[0] + .node + .get_peer(&addr_1) + .unwrap() + .last_heartbeat_sent(); + + // Due on every tick, so the only variable is what the send does. + set_heartbeat_interval(&mut nodes[0].node, 0); + nodes[0].node.transports.clear(); + + nodes[0].node.check_link_heartbeats().await; + + let peer = nodes[0].node.get_peer(&addr_1).expect("peer present"); + let failed_at = peer + .last_heartbeat_attempt() + .expect("the attempt is recorded even though the send failed"); + assert_eq!( + peer.last_heartbeat_sent(), + landed_before, + "a heartbeat whose send failed was recorded as having landed" + ); + + // The retry gate spaces the next attempt out rather than letting a failing + // peer be retried on every tick. + nodes[0].node.check_link_heartbeats().await; + + let peer = nodes[0].node.get_peer(&addr_1).expect("peer present"); + assert_eq!( + peer.last_heartbeat_attempt(), + Some(failed_at), + "a peer whose send failed was retried inside the retry interval" + ); + + cleanup_nodes(&mut nodes).await; +} + /// A peer past the link-dead timeout is NOT reaped while an FMP rekey is in /// progress with its msg1 budget unexhausted. #[tokio::test] @@ -124,15 +188,6 @@ 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. /// diff --git a/src/peer/active.rs b/src/peer/active.rs index 33c7afc3..6562cced 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -250,13 +250,13 @@ pub struct ActivePeer { // === Heartbeat === /// 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. + /// not move this: it told the peer nothing, 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 + /// When a heartbeat to this peer was last *attempted*, whatever came of it. + /// Paired with the above so a peer whose send failed is retried sooner than + /// the heartbeat interval without being retried on every tick — see /// `HEARTBEAT_RETRY_INTERVAL`. last_heartbeat_attempt: Option, @@ -813,24 +813,25 @@ impl ActivePeer { /// 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. + /// Call this *after* the send, and only when it returned cleanly. 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` although 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 { + pub(crate) 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) { + /// 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(crate) fn mark_heartbeat_attempt(&mut self, now: Instant) { self.last_heartbeat_attempt = Some(now); } diff --git a/testing/.gitignore b/testing/.gitignore index 0a184bb6..d6174a08 100644 --- a/testing/.gitignore +++ b/testing/.gitignore @@ -4,8 +4,9 @@ fipsctl fipstop fips-gateway -# Generated test configs -generated-configs/ +# Generated test configs. The suffixed form is what a run with +# FIPS_CI_NAME_SUFFIX set writes, so the glob has to cover both. +generated-configs*/ # Simulation results sim-results/