From 4990525b62afdd56bc56ef2f3a9c4c0024506cb8 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Wed, 9 Sep 2026 00:00:22 +0000 Subject: [PATCH] fix(peering): do not count a failed heartbeat send as a delivered one The heartbeat was recorded as sent before the send was attempted, so a peer whose heartbeat could not go out was treated as heartbeated and was not tried again for a whole heartbeat_interval_secs, although it had heard nothing and its own link-dead timer was running. Record the attempt and the delivery separately. The interval that paces a healthy peer now advances only on a send that returned cleanly, and a peer whose send failed is retried after a shorter fixed interval instead of after a full heartbeat interval. The retry interval gates only a peer whose last attempt failed. On a healthy peer the two timestamps are equal, so consulting it there would clamp a heartbeat_interval_secs configured below the retry interval, and that setting has no validation floor. The due-or-not decision moves into a small pure function so it can be tested without driving a send. The tests cover both halves and both were checked against the defect rather than only against the fix: restoring the old ordering fails the integration test on the recorded-as-landed assertion, and removing the retry gate fails it on the not-retried-every-tick assertion. Also widen the generated-config ignore glob. The rule matched generated-configs/ exactly, while the scripts write generated-configs${FIPS_CI_NAME_SUFFIX}, so every suffixed run left an untracked directory behind. --- CHANGELOG.md | 15 +++- src/node/handlers/mmp.rs | 172 ++++++++++++++++++++++++++++++++++-- src/node/tests/heartbeat.rs | 63 +++++++++++++ src/peer/active.rs | 34 ++++++- src/proto/mmp/core.rs | 11 ++- testing/.gitignore | 5 +- 6 files changed, 281 insertions(+), 19 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 194f5509..879bf7d6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,8 +7,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] -Nothing yet. Everything previously staged here is folded into -`[0.5.1]` below. +### Fixed + +#### Peering + +- A heartbeat whose send failed no longer counts as one that was delivered. The + send was recorded before it was attempted, so a peer whose heartbeat could not + go out was treated as heartbeated and was not tried again for a whole + `heartbeat_interval_secs`, although it had heard nothing and its own link-dead + timer was running. The attempt and the delivery are now recorded separately: + the interval that paces a healthy peer advances only on a send that returned + cleanly, and a peer whose send failed is retried after a shorter fixed + interval instead. That retry interval gates only a peer whose last attempt + failed, so it cannot clamp a `heartbeat_interval_secs` configured below it. ## [0.5.1] - 2026-09-06 diff --git a/src/node/handlers/mmp.rs b/src/node/handlers/mmp.rs index 1f7b3b9d..606a3bde 100644 --- a/src/node/handlers/mmp.rs +++ b/src/node/handlers/mmp.rs @@ -20,6 +20,52 @@ 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. +/// +/// 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. +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`; @@ -473,11 +519,16 @@ 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() { - None => true, - Some(last) => now.duration_since(last) >= heartbeat_interval, - }; + // 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. + let heartbeat_due = heartbeat_due( + peer.last_heartbeat_sent(), + peer.last_heartbeat_attempt(), + now, + heartbeat_interval, + ); PeerLivenessSnapshot { peer: *node_addr, @@ -513,14 +564,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 { .. } @@ -588,3 +650,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 b5ff6e3e..11e684f7 100644 --- a/src/node/tests/heartbeat.rs +++ b/src/node/tests/heartbeat.rs @@ -35,6 +35,69 @@ fn set_link_dead_timeout(node: &mut crate::node::Node, secs: u64) { }); } +/// Set `heartbeat_interval_secs` on an already-constructed node, the same way +/// `set_link_dead_timeout` does. +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] diff --git a/src/peer/active.rs b/src/peer/active.rs index 36cafb35..3216ea1f 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -249,8 +249,15 @@ 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 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*, 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. + last_heartbeat_attempt: Option, // === Handshake Resend === /// Wire-format msg2 for resend on duplicate msg1 (responder only). @@ -311,6 +318,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 +398,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 +805,33 @@ 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 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(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(crate) 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..24894c29 100644 --- a/src/proto/mmp/core.rs +++ b/src/proto/mmp/core.rs @@ -39,8 +39,9 @@ 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: the interval has elapsed since one last *landed*, + /// and, on a peer whose last attempt failed, the shorter retry interval has + /// elapsed since that attempt. Both gates are resolved shell-side. pub heartbeat_due: bool, } @@ -169,8 +170,10 @@ 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, then `mark_heartbeat_sent` only if that send + /// returned cleanly. Marking the send before it happens would let a failed + /// heartbeat suppress the next one for a full interval. Heartbeat { peer: NodeAddr }, /// Build (shell: `proto/mmp/` `build_report` + `encode`) and send the given /// link report over the encrypted link. The interval-advancing 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/