diff --git a/src/node/handlers/rekey.rs b/src/node/handlers/rekey.rs index 29eccea6..a0d7dc65 100644 --- a/src/node/handlers/rekey.rs +++ b/src/node/handlers/rekey.rs @@ -627,6 +627,9 @@ impl Node { /// out, abandon it (the handshake only — a completed rekey session /// is never discarded on a timer, see /// [`FspAction::AbandonHandshake`]) + /// - If a handshake this node initiated got no SessionAck within the + /// handshake timeout, drop it so the trigger can start a fresh one + /// (see [`FspAction::ExpireInitiation`]) /// - If the rekey timer/counter fires, initiate a new XK handshake /// (this last one only when `node.rekey.enabled`) /// @@ -692,6 +695,24 @@ impl Node { ); } } + FspAction::ExpireInitiation { addr } => { + // The handshake only: the trigger starts a fresh rekey + // on a later tick. + let age_ms = self + .sessions + .get(&addr) + .map(|entry| now_ms.saturating_sub(entry.initiated_ms())) + .unwrap_or(0); + if let Some(entry) = self.sessions.get_mut(&addr) { + entry.abandon_handshake(); + self.stats_mut().session.rekey_unanswered += 1; + info!( + peer = %self.peer_display_name(&addr), + age_ms, + "FSP rekey we initiated got no answer within the handshake timeout, retrying, session retained" + ); + } + } FspAction::InitiateRekey { addr } => { self.initiate_session_rekey(&addr).await; } @@ -711,6 +732,14 @@ impl Node { // finished, anchored on the peer's last accepted setup message, which // is the only stamp that path writes. A *completed* rekey has no such // bound and must not acquire one: see `FspAction::AbandonHandshake`. + // + // A handshake this node initiated has its own deadline, measured from + // the setup this node sent, and neither clock may stand in for the + // other: the peer's stamp says nothing about our setup, and ours + // says nothing about the peer's. Unlike the peer predicate, ours has + // no `!= 0` conjunct, deliberately: an initiator handshake armed + // without the stamp reads as expired, because dropping one costs a + // retry and keeping one stops rotation. let stale_handshake_ms = self.config().node.rate_limit.handshake_timeout_secs * 1000; // Absolute ceiling on `previous`-slot retention, measured from the // cutover. The sliding drain deadline is peer-progress-aware, so an @@ -734,6 +763,8 @@ impl Node { is_dampened: entry.is_rekey_dampened(now_ms, dampening_ms), armed_handshake_expired: entry.last_peer_rekey_ms() != 0 && now_ms.saturating_sub(entry.last_peer_rekey_ms()) > stale_handshake_ms, + initiation_expired: now_ms.saturating_sub(entry.initiated_ms()) + > stale_handshake_ms, elapsed_secs: now_ms.saturating_sub(entry.session_start_ms()) / 1000, counter: entry.send_counter(), jitter_secs: entry.rekey_jitter_secs(), @@ -744,7 +775,8 @@ impl Node { /// Initiate an FSP session rekey. /// /// Creates a new XK handshake as initiator, sends SessionSetup msg1 - /// through the mesh, and stores the handshake state on the existing entry. + /// through the mesh, and stores the handshake state on the existing entry, + /// stamping the deadline by which a SessionAck must complete it. async fn initiate_session_rekey(&mut self, dest_addr: &NodeAddr) { // Check route availability before paying crypto cost if self.find_next_hop(dest_addr).is_none() { @@ -811,9 +843,10 @@ impl Node { return; } - // Store rekey state on the existing session entry + // Store rekey state on the existing session entry. The deadline is + // stamped only now, so it runs from the setup actually on the wire. if let Some(entry) = self.sessions.get_mut(dest_addr) { - entry.set_rekey_state(handshake, true); + entry.begin_rekey(handshake, Self::now_ms()); } debug!( diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 4ec11836..15aad913 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -899,13 +899,15 @@ impl Node { // arm of `handle_session_msg3`, and the difference rests on an // invariant rather than on a different judgement: an entry with // `rekey_initiator` set holds no pending session, so the two calls - // are the same action here. `set_rekey_state(_, true)` has one - // caller, `initiate_session_rekey`, which `check_session_rekey` - // never reaches for an entry holding a pending session; and - // `set_pending_session` clears `rekey_state`, so a completed - // initiator cycle leaves at most one of the two set. If that ever - // stops holding, these three sites become instances of the epoch - // discard the responder arm was fixed for. + // are the same action here. Arming as initiator has one caller, + // `initiate_session_rekey` through `begin_rekey`, and + // `set_rekey_state(_, true)` remains only at the restore below; + // `check_session_rekey` never reaches `initiate_session_rekey` for + // an entry holding a pending session; and `set_pending_session` + // clears `rekey_state`, so a completed initiator cycle leaves at + // most one of the two set. If that ever stops holding, these three + // sites become instances of the epoch discard the responder arm was + // fixed for. if entry.is_established() && entry.has_rekey_in_progress() && entry.is_rekey_initiator() { let mut handshake = match entry.take_rekey_state() { Some(hs) => hs, @@ -923,7 +925,9 @@ impl Node { // back to its pre-read state so it can still read the genuine // ack, and the refusal is counted. The rollback matters because // `read_xk_message_2` mixes the sender's ephemeral in before it - // authenticates. + // authenticates. The restore does not restamp the deadline, which + // runs from the setup this node sent, so an unreadable ack cannot + // hold the rekey open. if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) { debug!(error = %e, "Failed to process rekey XK msg2, keeping the rekey"); entry.set_rekey_state(handshake, true); diff --git a/src/node/session/mod.rs b/src/node/session/mod.rs index f06a3bf7..18c76b4c 100644 --- a/src/node/session/mod.rs +++ b/src/node/session/mod.rs @@ -163,6 +163,14 @@ pub(crate) struct SessionEntry { rekey_initiator: bool, /// Dampening: last time peer sent us a rekey msg1 (Unix ms). last_peer_rekey_ms: u64, + /// When this node armed its current rekey handshake as initiator (Unix + /// ms), which bounds how long that handshake waits for a SessionAck. + /// + /// Written only by `begin_rekey`, so putting the handshake back after + /// an unreadable SessionAck does not move it. Read only while a + /// handshake is armed and `rekey_initiator` is set, which is why it is + /// never cleared. + initiated_ms: u64, /// When this side's FSP rekey handshake completed and produced the /// `pending` session (Unix ms): the initiator sending msg3, or the /// responder accepting it. Cleared on cutover. @@ -232,6 +240,7 @@ impl SessionEntry { pending_new_session: None, rekey_initiator: false, last_peer_rekey_ms: 0, + initiated_ms: 0, rekey_completed_ms: 0, rekey_msg3_payload: None, rekey_msg3_next_resend_ms: 0, @@ -487,6 +496,15 @@ impl SessionEntry { self.last_peer_rekey_ms } + /// When this node armed its current rekey handshake as initiator, or 0 + /// if it never has. + /// + /// Bounds the age of a handshake this node armed, independently of any + /// rekey the peer started. + pub(crate) fn initiated_ms(&self) -> u64 { + self.initiated_ms + } + /// Record that the peer initiated a rekey (for dampening). pub(crate) fn record_peer_rekey(&mut self, now_ms: u64) { self.last_peer_rekey_ms = now_ms; @@ -680,6 +698,18 @@ impl SessionEntry { self.rekey_initiator = is_initiator; } + /// Arm a rekey handshake this node initiated, and stamp its deadline. + /// + /// The only writer of `initiated_ms`. The restore after an unreadable + /// SessionAck goes through `set_rekey_state` instead and must stay + /// there: a restamping restore would let a stream of forged acks hold + /// the rekey open for good. + pub(crate) fn begin_rekey(&mut self, state: HandshakeState, now_ms: u64) { + self.rekey_state = Some(state); + self.rekey_initiator = true; + self.initiated_ms = now_ms; + } + /// Take the rekey state for processing. pub(crate) fn take_rekey_state(&mut self) -> Option { self.rekey_state.take() @@ -831,10 +861,11 @@ impl SessionEntry { /// Abandon an in-progress rekey handshake, keeping any completed /// session already waiting for cut-over. /// - /// Used when a handshake the peer armed times out without its msg3. - /// The handshake holds no key material either endpoint can be using, - /// so dropping it costs nothing; a `pending` session alongside it is - /// the epoch the peer may already have moved to and must survive. + /// Used when an armed handshake times out: one the peer armed without + /// its msg3, or one this node armed without a SessionAck. The handshake + /// holds no key material either endpoint can be using, so dropping it + /// costs nothing; a `pending` session alongside it is the epoch the + /// peer may already have moved to and must survive. /// /// `rekey_initiator` is deliberately left alone: it describes whichever /// rekey artefact the entry still holds, and every reader gates on a diff --git a/src/node/stats.rs b/src/node/stats.rs index 0d82f62f..a638f131 100644 --- a/src/node/stats.rs +++ b/src/node/stats.rs @@ -64,6 +64,12 @@ pub struct SessionStats { /// the handshake timeout without a msg3 and was discarded. The /// established session is retained. pub rekey_expired: u64, + /// A rekey this node initiated got no readable SessionAck within the + /// handshake timeout, and its handshake was discarded so the trigger + /// can retry. The established session is retained. A sustained rate + /// means setups or acks to that peer are being lost, or the peer holds + /// a stuck handshake of its own and wins the tie-break. + pub rekey_unanswered: u64, /// A completed rekey session still waiting for the peer's cut-over was /// replaced by a newer one, completed from a msg3 carrying the same /// authenticated peer key. This drops key material the peer may @@ -101,6 +107,7 @@ impl SessionStats { rekey_yielded: self.rekey_yielded, rekey_pending: self.rekey_pending, rekey_expired: self.rekey_expired, + rekey_unanswered: self.rekey_unanswered, pending_replaced: self.pending_replaced, ack_handshake_failed: self.ack_handshake_failed, setup_rate_limited: self.setup_rate_limited, @@ -406,6 +413,7 @@ pub struct SessionStatsSnapshot { pub rekey_yielded: u64, pub rekey_pending: u64, pub rekey_expired: u64, + pub rekey_unanswered: u64, pub pending_replaced: u64, pub ack_handshake_failed: u64, pub setup_rate_limited: u64, diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 65a262ae..52c4b022 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -6473,6 +6473,493 @@ async fn test_fresh_peer_armed_rekey_is_not_expired() { ); } +// ============================================================================ +// Integration tests: a rekey this node initiated that is never answered +// ============================================================================ + +/// Build an established two-node pair, with the handshake timeout set to +/// `timeout_secs` on both nodes when given and left at its default otherwise. +/// +/// `rekeys[i]` says whether node `i` rekeys after a single sent message; +/// otherwise it never starts a rekey of its own. The session is opened from +/// node 0, which says nothing about which node later initiates a rekey. +async fn rekey_pair(rekeys: [bool; 2], timeout_secs: Option) -> Vec { + let configs = rekeys + .iter() + .map(|&rekeys| { + let mut config = Config::new(); + if let Some(secs) = timeout_secs { + config.node.rate_limit.handshake_timeout_secs = secs; + } + if rekeys { + config.node.rekey.after_messages = 1; + } else { + config.node.rekey.after_messages = u64::MAX; + config.node.rekey.after_secs = u64::MAX; + } + config + }) + .collect(); + let mut nodes = run_tree_test_with_configs(configs, &[(0, 1)]).await; + verify_tree_convergence(&nodes); + populate_all_coord_caches(&mut nodes); + establish_pair_session(&mut nodes).await; + nodes +} + +/// Send one data frame from `nodes[from]` across its rekey trigger, deliver +/// it, then run `nodes[from]`'s tick so it sends a rekey SessionSetup, and +/// assert it now holds an initiated rekey. +async fn start_rekey(nodes: &mut [TestNode], from: usize) { + let peer = *nodes[1 - from].node.node_addr(); + nodes[from] + .node + .send_session_data(&peer, 0, 0, b"crosses the rekey trigger") + .await + .expect("send_session_data failed"); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(nodes).await; + + nodes[from].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[from], &peer), + "node {from} must have initiated a rekey" + ); +} + +/// Whether `node` holds a rekey handshake toward `peer` that it initiated. +fn rekey_initiated(node: &TestNode, peer: &NodeAddr) -> bool { + node.node + .get_session(peer) + .is_some_and(|e| e.has_rekey_in_progress() && e.is_rekey_initiator()) +} + +/// Whether `node`'s session with `peer` holds a completed rekey session. +fn holds_pending(node: &TestNode, peer: &NodeAddr) -> bool { + node.node + .get_session(peer) + .is_some_and(|e| e.pending_new_session().is_some()) +} + +/// Discard every packet queued at `node` without processing it, as a lossy +/// link would, and return how many were dropped. +fn drop_queued(node: &mut TestNode) -> usize { + std::iter::from_fn(|| node.packet_rx.try_recv().ok()).count() +} + +/// Run three delivery passes over every node. +async fn pump_all(nodes: &mut [TestNode]) { + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(nodes).await; + } +} + +/// A rekey whose SessionSetup was lost is retired once the handshake timeout +/// has passed, and the trigger then starts a fresh rekey that completes. +#[tokio::test] +async fn test_a_rekey_whose_session_setup_was_lost_is_retired_after_the_handshake_timeout_and_retried() + { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[1]) > 0, + "node 0's SessionSetup must have been queued at node 1 to be lost" + ); + assert_eq!(nodes[1].node.stats().session.rekey_armed, 0); + assert!( + !nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .has_rekey_in_progress(), + "the setup must really have been dropped before node 1 armed" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[0].node.check_session_rekey().await; + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + nodes[0].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + pump_all(&mut nodes).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 1, + "the retried setup must reach node 1" + ); + assert!( + holds_pending(&nodes[0], &node1_addr) && holds_pending(&nodes[1], &node0_addr), + "the retried rekey must complete on both nodes" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + assert_eq!(nodes[1].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// A rekey whose SessionAck was lost is retired once the handshake timeout +/// has passed, beside the responder's own expiry, and the retry completes. +#[tokio::test] +async fn test_a_rekey_whose_session_ack_was_lost_is_retired_after_the_handshake_timeout_and_retried() + { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[1..]).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 1, + "node 1 must have armed as the rekey responder" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[0]) > 0, + "node 1's SessionAck must have been queued at node 0 to be lost" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[1].node.check_session_rekey().await; + assert_eq!( + nodes[1].node.stats().session.rekey_expired, + 1, + "node 1's own handshake must expire on the existing responder rule" + ); + assert!( + !nodes[1] + .node + .get_session(&node0_addr) + .unwrap() + .has_rekey_in_progress() + ); + + nodes[0].node.check_session_rekey().await; + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + nodes[0].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + pump_all(&mut nodes).await; + assert_eq!( + nodes[1].node.stats().session.rekey_armed, + 2, + "the retried setup must reach node 1" + ); + assert!( + holds_pending(&nodes[0], &node1_addr) && holds_pending(&nodes[1], &node0_addr), + "the retried rekey must complete on both nodes" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + assert_eq!(nodes[1].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// Drive a lost SessionAck, then a retry that meets the responder's own +/// handshake from the first attempt, still armed because nothing has run the +/// responder's expiry yet. +/// +/// Both nodes would rekey after one message; the initiator is picked at run +/// time so that the responder holds the smaller address when +/// `responder_wins`, and the larger otherwise. Only the initiator's tick is +/// run before the responder has armed, after which the responder is +/// dampened and cannot start a rekey of its own inside the test. +/// +/// A responder that wins the tie-break drops the retry before arming +/// anything, so both handshakes expire and the next retry completes one +/// timeout later. A responder that loses yields and answers the retry at +/// once. +async fn lostack_retry(responder_wins: bool) { + let mut nodes = rekey_pair([true, true], Some(1)).await; + let node0_smaller = + crate::proto::fsp::initiation_winner(nodes[0].node.node_addr(), nodes[1].node.node_addr()); + let resp = if responder_wins == node0_smaller { + 0 + } else { + 1 + }; + let init = 1 - resp; + let init_addr = *nodes[init].node.node_addr(); + let resp_addr = *nodes[resp].node.node_addr(); + + // First attempt: the responder arms, and its SessionAck is lost. + start_rekey(&mut nodes, init).await; + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[resp..=resp]).await; + assert_eq!( + nodes[resp].node.stats().session.rekey_armed, + 1, + "the responder must have armed" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[init]) > 0, + "the responder's SessionAck must have been queued to be lost" + ); + + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[init].node.check_session_rekey().await; + assert!( + !nodes[init] + .node + .get_session(&resp_addr) + .unwrap() + .has_rekey_in_progress(), + "an unanswered rekey this node initiated must be retired after the \ + handshake timeout, or the trigger stays vetoed for good" + ); + + // The retry reaches a responder still holding the first handshake. + nodes[init].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[init], &resp_addr), + "the trigger must start a fresh rekey once the old one is retired" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + process_available_packets(&mut nodes[resp..=resp]).await; + let stats = &nodes[resp].node.stats().session; + let (tiebreak, yielded) = (stats.rekey_tiebreak, stats.rekey_yielded); + assert_eq!( + tiebreak + yielded, + 1, + "the retry must have met the responder's stale handshake" + ); + if responder_wins { + assert_eq!(tiebreak, 1, "the smaller responder must win"); + } else { + assert_eq!(yielded, 1, "the larger responder must yield"); + } + pump_all(&mut nodes).await; + + if responder_wins { + assert!( + !holds_pending(&nodes[init], &resp_addr) && !holds_pending(&nodes[resp], &init_addr), + "a retry the responder dropped completes nothing" + ); + tokio::time::sleep(Duration::from_millis(1200)).await; + nodes[resp].node.check_session_rekey().await; + assert_eq!( + nodes[resp].node.stats().session.rekey_expired, + 1, + "the responder's stale handshake must expire on its own rule" + ); + nodes[init].node.check_session_rekey().await; + assert!( + !nodes[init] + .node + .get_session(&resp_addr) + .unwrap() + .has_rekey_in_progress(), + "the retry the responder dropped must itself be retired" + ); + nodes[init].node.check_session_rekey().await; + assert!( + rekey_initiated(&nodes[init], &resp_addr), + "the trigger must start a second retry" + ); + pump_all(&mut nodes).await; + } + + assert!( + holds_pending(&nodes[init], &resp_addr) && holds_pending(&nodes[resp], &init_addr), + "the retried rekey must complete on both nodes" + ); + // The first handshake always; the retry too when the responder dropped it. + assert_eq!( + nodes[init].node.stats().session.rekey_unanswered, + if responder_wins { 2 } else { 1 }, + "every retired handshake must be counted once" + ); + assert_eq!(nodes[resp].node.stats().session.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + +/// A retry dropped on the tie-break by a smaller responder still holding its +/// stale handshake completes once both handshakes have expired. +#[tokio::test] +async fn test_a_retry_dropped_by_a_smaller_responders_stale_handshake_completes_one_timeout_later() +{ + lostack_retry(true).await; +} + +/// A retry that a larger responder, still holding its stale handshake, +/// yields to completes at once. +#[tokio::test] +async fn test_a_retry_that_a_larger_responder_yields_to_completes_at_once() { + lostack_retry(false).await; +} + +/// A forged SessionAck arriving midway through an unanswered rekey must not +/// restart its deadline: the rekey is retired on the timeout measured from +/// the setup this node sent. +#[tokio::test] +async fn test_forged_session_acks_do_not_hold_an_unanswered_rekey_open_past_its_deadline() { + let mut nodes = rekey_pair([true, false], Some(1)).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + tokio::time::sleep(Duration::from_millis(20)).await; + assert!( + drop_queued(&mut nodes[1]) > 0, + "node 0's SessionSetup must have been queued at node 1 to be lost" + ); + + tokio::time::sleep(Duration::from_millis(600)).await; + let forged = forged_session_ack(&nodes[1]); + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + assert_eq!(nodes[0].node.stats().session.ack_handshake_failed, 1); + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the unreadable ack must have put the handshake back" + ); + + let forged_at = std::time::Instant::now(); + tokio::time::sleep(Duration::from_millis(600)).await; + nodes[0].node.check_session_rekey().await; + println!( + "forged ack to check: {} ms", + forged_at.elapsed().as_millis() + ); + assert!( + !nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .has_rekey_in_progress(), + "a forged SessionAck must not restart the deadline of the rekey this \ + node initiated" + ); + assert_eq!( + nodes[0].node.stats().session.rekey_unanswered, + 1, + "the retired handshake must be counted once" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// Putting a handshake back after an unreadable SessionAck leaves the stamp +/// its deadline runs from exactly where arming wrote it. +#[tokio::test] +async fn test_an_unreadable_session_ack_does_not_push_out_the_rekey_deadline() { + let mut nodes = rekey_pair([true, false], None).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + let armed_at = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(); + assert_ne!(armed_at, 0, "arming must stamp the deadline"); + + tokio::time::sleep(Duration::from_millis(5)).await; + let forged = forged_session_ack(&nodes[1]); + nodes[0] + .node + .handle_session_payload(&node1_addr, &node1_addr, &forged, 1280, false) + .await; + assert_eq!(nodes[0].node.stats().session.ack_handshake_failed, 1); + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "the unreadable ack must have put the handshake back" + ); + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(), + armed_at, + "the restore must not restamp the deadline" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// A fresh rekey this node initiated is not retired because the peer's own +/// last rekey is older than the handshake timeout. +#[tokio::test] +async fn test_a_fresh_initiated_rekey_is_not_retired_on_the_peers_older_rekey_stamp() { + let mut nodes = rekey_pair([true, false], None).await; + let node1_addr = *nodes[1].node.node_addr(); + + start_rekey(&mut nodes, 0).await; + // A peer rekey long past both the handshake timeout and the dampening. + nodes[0] + .node + .sessions + .get_mut(&node1_addr) + .unwrap() + .record_peer_rekey(wall_clock_ms() - 60_000); + let armed_at = nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(); + + nodes[0].node.check_session_rekey().await; + + assert!( + rekey_initiated(&nodes[0], &node1_addr), + "a fresh rekey this node initiated must not be retired on the peer's clock" + ); + assert_eq!( + nodes[0] + .node + .get_session(&node1_addr) + .unwrap() + .initiated_ms(), + armed_at, + "the rekey must be the same one, not retired and re-armed" + ); + let stats = &nodes[0].node.stats().session; + assert_eq!(stats.rekey_expired, 0); + assert_eq!(stats.rekey_unanswered, 0); + + cleanup_nodes(&mut nodes).await; +} + // ============================================================================ // Integration tests: a peer's cutover after a long silence // ============================================================================ diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index 05c5a5d5..33fc8b62 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -65,6 +65,16 @@ pub(crate) enum FspAction { /// Discarding it would make every later frame from that peer /// undecryptable, so the handshake alone is dropped. AbandonHandshake { addr: NodeAddr }, + /// Drop only `addr`'s handshake that this node armed by sending a setup + /// message never answered within the handshake timeout + /// (`SessionEntry::abandon_handshake`). The trigger starts a fresh rekey + /// on a later tick. + /// + /// The handshake only, not [`AbandonRekey`](Self::AbandonRekey): an + /// entry whose handshake this node armed holds no pending session, and + /// dropping only the handshake keeps any such session safe even if that + /// ever stops holding. + ExpireInitiation { addr: NodeAddr }, /// Retransmit `addr`'s retained rekey msg3 (the shell re-reads the payload /// from the entry, sends it, then records the retransmission on success). ResendSessionMsg3 { addr: NodeAddr }, @@ -143,9 +153,14 @@ pub(crate) struct SessionSnapshot { /// A handshake armed by the peer's setup message has passed the handshake /// timeout without its msg3 (pre-evaluated: `last_peer_rekey_ms != 0 && /// now - last_peer_rekey_ms > handshake_timeout`). False when the entry - /// carries no peer-rekey stamp, so a handshake this side armed is never - /// aged out on the peer's clock. + /// carries no peer-rekey stamp. A handshake this side armed is never + /// aged out on the peer's clock; it has its own deadline, + /// [`initiation_expired`](Self::initiation_expired). pub armed_handshake_expired: bool, + /// This node's own armed handshake has passed the handshake timeout + /// measured from the setup message it sent (pre-evaluated: `now - + /// initiated_ms > handshake_timeout`). + pub initiation_expired: bool, /// Monotonic session age in seconds (`(now - session_start_ms) / 1000`). pub elapsed_secs: u64, /// Current Noise send counter. @@ -219,18 +234,22 @@ impl Fsp { /// in-flight rekey, and an elapsed liveness timer cuts over and is /// considered for nothing else. /// - Otherwise an expired drain window is completed, and — independently — - /// a handshake the peer armed and never finished is abandoned, which is - /// the last word on that session this tick. + /// a handshake the peer armed and never finished, or a handshake this + /// node armed and never got an answer to, is abandoned, which is the + /// last word on that session this tick. /// - Failing both, the rekey trigger fires when the session is neither /// mid-rekey, holding a pending session, retaining a msg3 payload, nor /// dampened, and its jittered time threshold or send counter is reached. - /// Only this last decision is gated on `cfg.enabled`: the other three - /// maintain state a peer's setup message can create with periodic rekey - /// switched off. + /// Only this last decision is gated on `cfg.enabled`: the cutover, the + /// drain and the abandon of a peer-armed handshake maintain state a + /// peer's setup message can create with periodic rekey switched off. + /// A handshake this node armed exists only when the trigger fired, but + /// its expiry sits above the gate too, so whether an existing + /// handshake is retired does not depend on whether new ones may start. /// /// Actions are returned phase-grouped (all cutovers, then all drains, then - /// all abandoned handshakes, then all rekey initiations) to preserve the - /// pre-refactor execution order. + /// all abandoned or expired handshakes, then all rekey initiations) to + /// preserve the pre-refactor execution order. pub(crate) fn poll_rekey( &self, sessions: Vec, @@ -264,6 +283,15 @@ impl Fsp { abandons.push(FspAction::AbandonHandshake { addr: s.addr }); continue; } + // 3b. Retire a handshake this node armed whose setup or + // SessionAck was lost. Arm 3 cannot: it runs on the peer's + // clock, which says nothing about our own setup. Only the + // handshake goes; the trigger below starts a fresh rekey on a + // later tick. + if s.is_rekey_initiator && s.rekey_in_progress && s.initiation_expired { + abandons.push(FspAction::ExpireInitiation { addr: s.addr }); + continue; + } // 4. Rekey trigger. if !cfg.enabled { continue; diff --git a/src/proto/fsp/tests/core.rs b/src/proto/fsp/tests/core.rs index b98dcee3..0e80f3af 100644 --- a/src/proto/fsp/tests/core.rs +++ b/src/proto/fsp/tests/core.rs @@ -16,7 +16,8 @@ fn coords(byte: u8) -> TreeCoordinate { } /// A quiescent established-session snapshot: no pending cutover, no drain, no -/// dampening, zero ages/counter/jitter. Tests set only the fields they exercise. +/// dampening, no expired initiation, zero ages/counter/jitter. Tests set only +/// the fields they exercise. fn session_snapshot(addr_byte: u8) -> SessionSnapshot { SessionSnapshot { addr: make_node_addr(addr_byte), @@ -29,6 +30,7 @@ fn session_snapshot(addr_byte: u8) -> SessionSnapshot { has_rekey_msg3_payload: false, is_dampened: false, armed_handshake_expired: false, + initiation_expired: false, elapsed_secs: 0, counter: 0, jitter_secs: 0, @@ -252,28 +254,135 @@ fn poll_rekey_abandons_the_handshake_without_touching_a_completed_pending() { ); } +/// A handshake this node armed whose setup or SessionAck was lost, past the +/// handshake timeout on its own clock and carrying no peer stamp, with the +/// send counter over the rekey threshold. +fn unanswered_initiation(addr_byte: u8) -> SessionSnapshot { + let mut s = session_snapshot(addr_byte); + s.rekey_in_progress = true; + s.is_rekey_initiator = true; + s.initiation_expired = true; + s.armed_handshake_expired = false; + s.counter = 5000; + s +} + +/// A rekey this node initiated whose setup or ack was lost is retired on its +/// own deadline, and once the handshake is gone the trigger fires again. #[test] -fn poll_rekey_does_not_abandon_a_fresh_or_locally_initiated_handshake() { +fn poll_rekey_retires_an_expired_handshake_this_node_initiated_so_the_trigger_fires_again() { let fsp = Fsp::new(); - // Still inside the handshake timeout: the peer's msg3 may be in flight. + let addr = make_node_addr(9); + + let mut s = unanswered_initiation(9); + assert_eq!( + fsp.poll_rekey(vec![unanswered_initiation(9)], &cfg(100, 1000)), + vec![FspAction::ExpireInitiation { addr }], + "the unanswered handshake must be retired, and the trigger must not \ + fire in the same tick" + ); + + // What the executor does for ExpireInitiation: `abandon_handshake` drops + // the handshake and leaves `rekey_initiator` as it was. + s.rekey_in_progress = false; + assert_eq!( + fsp.poll_rekey(vec![s], &cfg(100, 1000)), + vec![FspAction::InitiateRekey { addr }], + "with the handshake retired, the trigger must start a fresh rekey" + ); +} + +/// A handshake is retired only on the deadline of the side that armed it, +/// and nothing that is not an armed handshake is retired at all. +#[test] +fn poll_rekey_retires_an_initiated_handshake_only_on_its_own_deadline() { + let fsp = Fsp::new(); + let mut fresh = session_snapshot(9); fresh.rekey_in_progress = true; - fresh.armed_handshake_expired = false; - assert!(fsp.poll_rekey(vec![fresh], &cfg(100, 1000)).is_empty()); + fresh.is_rekey_initiator = true; + assert!( + fsp.poll_rekey(vec![fresh], &cfg(100, 1000)).is_empty(), + "a handshake this node armed inside the timeout may still be answered" + ); - // Expired, but this side is the initiator: the abandon-on-timeout rule is - // anchored on the peer's setup message and does not reach our own cycle, - // which the msg3 retransmission budget bounds instead. - let mut ours = session_snapshot(9); - ours.rekey_in_progress = true; - ours.is_rekey_initiator = true; - ours.armed_handshake_expired = true; - assert!(fsp.poll_rekey(vec![ours], &cfg(100, 1000)).is_empty()); + let mut peer_older = session_snapshot(9); + peer_older.rekey_in_progress = true; + peer_older.is_rekey_initiator = true; + peer_older.armed_handshake_expired = true; + assert!( + fsp.poll_rekey(vec![peer_older], &cfg(100, 1000)).is_empty(), + "the peer's older rekey stamp must not retire a fresh handshake this \ + node armed" + ); + + let mut theirs = session_snapshot(9); + theirs.rekey_in_progress = true; + theirs.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![theirs], &cfg(100, 1000)).is_empty(), + "this node's stale initiation stamp must not retire a handshake the \ + peer armed" + ); + + let mut responder_fresh = session_snapshot(9); + responder_fresh.rekey_in_progress = true; + assert!( + fsp.poll_rekey(vec![responder_fresh], &cfg(100, 1000)) + .is_empty(), + "a handshake the peer armed inside the timeout may still get its msg3" + ); - // Expired stamp but no handshake left to abandon. let mut none = session_snapshot(9); none.armed_handshake_expired = true; - assert!(fsp.poll_rekey(vec![none], &cfg(100, 1000)).is_empty()); + none.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![none], &cfg(100, 1000)).is_empty(), + "expired stamps with no handshake armed leave nothing to retire" + ); + + let mut completed = session_snapshot(9); + completed.is_rekey_initiator = true; + completed.has_pending = true; + completed.initiation_expired = true; + assert!( + fsp.poll_rekey(vec![completed], &cfg(100, 1000)).is_empty(), + "a completed cycle awaiting its cutover is not an expiring handshake" + ); +} + +/// The expiry of a handshake this node armed does not depend on whether +/// periodic rekey is enabled, and it is grouped with the other abandoned +/// handshakes in input order. +#[test] +fn poll_rekey_retires_an_expired_initiation_with_periodic_rekey_off() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_rekey(vec![unanswered_initiation(9)], &cfg_rekey_off(100, 1000)), + vec![FspAction::ExpireInitiation { + addr: make_node_addr(9) + }], + "an existing handshake must be retired whatever the trigger policy" + ); + + let mut peer_armed = session_snapshot(7); + peer_armed.rekey_in_progress = true; + peer_armed.armed_handshake_expired = true; + assert_eq!( + fsp.poll_rekey( + vec![peer_armed, unanswered_initiation(9)], + &cfg_rekey_off(100, 1000) + ), + vec![ + FspAction::AbandonHandshake { + addr: make_node_addr(7) + }, + FspAction::ExpireInitiation { + addr: make_node_addr(9) + }, + ], + "both expiries sit in the abandon group, in input order" + ); } #[test]