diff --git a/src/node/dataplane/encrypted.rs b/src/node/dataplane/encrypted.rs index 0f58e0a7..aee31c2e 100644 --- a/src/node/dataplane/encrypted.rs +++ b/src/node/dataplane/encrypted.rs @@ -5,12 +5,33 @@ use crate::noise::NoiseError; use crate::proto::fmp::wire::{ EncryptedHeader, FLAG_CE, FLAG_KEY_EPOCH, FLAG_SP, strip_inner_header, }; +use crate::proto::link::LinkMessageType; use crate::transport::ReceivedPacket; use tracing::{debug, trace, warn}; /// Force-remove a peer after this many consecutive decryption failures. const DECRYPT_FAILURE_THRESHOLD: u32 = 20; +/// Which of a peer's link sessions authenticated an inbound frame. +/// +/// After a rekey cutover the previous session still decrypts during the +/// drain, so frames the peer sealed before it moved are delivered. They +/// describe the session the cutover retired, though, and the new session's +/// MMP state was reset at the cutover: a previous-session frame is not +/// counted by the MMP receiver or the spin bit, and a ReceiverReport it +/// carries is not processed. A frame authenticated by a pending session is +/// promoted to current before it is processed, so it is `Current`. +/// +/// The link-dead check reads the MMP receiver's last-received time, which the +/// cutover clears, so previous-session frames no longer hold the link alive: +/// after a cutover the link-dead timer runs from the cutover until a frame on +/// the new session arrives. Peer `touch` and link statistics still see them. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(in crate::node) enum LinkSlot { + Current, + Previous, +} + impl Node { /// Handle an encrypted frame (phase 0x0). /// @@ -137,6 +158,7 @@ impl Node { let sp_flag = header.flags & FLAG_SP != 0; self.process_authentic_fmp_plaintext( &node_addr, + LinkSlot::Current, packet.transport_id, &packet.remote_addr, packet.timestamp_ms, @@ -195,7 +217,7 @@ impl Node { // Decrypt: try current session first, then previous (drain fallback) let ciphertext = &packet.data[header.ciphertext_offset()..]; - let plaintext = { + let (plaintext, slot) = { let peer = self.peers.get_mut(&node_addr).unwrap(); let session = match peer.noise_session_mut() { Some(s) => s, @@ -215,7 +237,7 @@ impl Node { ) { Ok(p) => { peer.reset_decrypt_failures(); - p + (p, LinkSlot::Current) } Err(e) => { // Current session failed — try previous session (drain window) @@ -227,7 +249,7 @@ impl Node { ) { Ok(p) => { peer.reset_decrypt_failures(); - p + (p, LinkSlot::Previous) } Err(_) => { self.log_decrypt_failure(&node_addr, &header, &e); @@ -266,7 +288,9 @@ impl Node { let mut address_changed = false; if let Some(peer) = self.peers.get_mut(&node_addr) { - if let Some(mmp) = peer.mmp_mut() { + if slot == LinkSlot::Current + && let Some(mmp) = peer.mmp_mut() + { mmp.receiver.record_recv( header.counter, timestamp, @@ -299,7 +323,32 @@ impl Node { } // Dispatch to link message handler - self.dispatch_link_message(&node_addr, link_message, ce_flag) + self.dispatch_authentic(&node_addr, slot, link_message, ce_flag) + .await; + } + + /// Dispatch the link message of an authenticated frame. A ReceiverReport + /// on the previous session reports on this node's traffic in the session + /// a cutover retired, and is dropped rather than judged against the new + /// session's metrics; every other message is dispatched whichever session + /// carried it. + async fn dispatch_authentic( + &mut self, + node_addr: &crate::NodeAddr, + slot: LinkSlot, + link_message: &[u8], + ce_flag: bool, + ) { + if slot == LinkSlot::Previous + && link_message.first() == Some(&(LinkMessageType::ReceiverReport as u8)) + { + trace!( + peer = %self.peer_display_name(node_addr), + "Dropping a ReceiverReport carried on the previous link session" + ); + return; + } + self.dispatch_link_message(node_addr, link_message, ce_flag) .await; } @@ -353,6 +402,7 @@ impl Node { pub(in crate::node) async fn process_authentic_fmp_plaintext( &mut self, node_addr: &crate::NodeAddr, + slot: LinkSlot, transport_id: crate::transport::TransportId, remote_addr: &crate::transport::TransportAddr, packet_timestamp_ms: u64, @@ -381,7 +431,9 @@ impl Node { peer.link_stats_mut() .record_recv(packet_len, packet_timestamp_ms); peer.touch(packet_timestamp_ms); - if let Some(mmp) = peer.mmp_mut() { + if slot == LinkSlot::Current + && let Some(mmp) = peer.mmp_mut() + { mmp.receiver .record_recv(fmp_counter, inner_ts, packet_len, ce_flag, now_ms); let _spin_rtt = mmp.spin_bit.rx_observe(sp_flag, fmp_counter, now_ms); @@ -399,7 +451,7 @@ impl Node { let _ = address_changed; } let link_message = &fmp_plaintext[INNER_TIMESTAMP_LEN..]; - self.dispatch_link_message(node_addr, link_message, ce_flag) + self.dispatch_authentic(node_addr, slot, link_message, ce_flag) .await; } @@ -414,8 +466,23 @@ impl Node { let sp_flag = fallback.fmp_flags & FLAG_SP != 0; let plaintext = &fallback.packet_data[fallback.fmp_plaintext_offset ..fallback.fmp_plaintext_offset + fallback.fmp_plaintext_len]; + // The worker decrypted with the session registered under the + // frame's receiver index. Only the current session's index is + // current; any other is the previous session draining, or one + // already gone by the time the bounce is processed. + let current = self + .peers + .get(&fallback.source_node_addr) + .and_then(|p| p.our_index()) + .is_some_and(|idx| idx.as_u32() == fallback.receiver_idx); + let slot = if current { + LinkSlot::Current + } else { + LinkSlot::Previous + }; self.process_authentic_fmp_plaintext( &fallback.source_node_addr, + slot, fallback.transport_id, &fallback.remote_addr, fallback.timestamp_ms, diff --git a/src/node/decrypt_worker.rs b/src/node/decrypt_worker.rs index 10fb7bb6..9566169a 100644 --- a/src/node/decrypt_worker.rs +++ b/src/node/decrypt_worker.rs @@ -150,6 +150,10 @@ pub(crate) struct DecryptFallback { /// MMP's 30-second link-dead timer fires even though packets /// are arriving fine. pub packet_len: usize, + /// The frame's receiver index: the index of the session that decrypted + /// it, so rx_loop can tell a current-session frame from one on the + /// previous session during a rekey drain. + pub receiver_idx: u32, pub fmp_counter: u64, pub fmp_flags: u8, /// Original received wire buffer, mutated in place by the FMP @@ -521,6 +525,7 @@ fn handle_job( remote_addr, timestamp_ms, packet_len, + receiver_idx: cache_key.1, fmp_counter, fmp_flags, packet_data, diff --git a/src/node/tests/bloom.rs b/src/node/tests/bloom.rs index 70104193..0c32848b 100644 --- a/src/node/tests/bloom.rs +++ b/src/node/tests/bloom.rs @@ -1485,10 +1485,12 @@ fn next_counter(node: &Node, peer: &NodeAddr) -> u64 { .current_send_counter() } -/// Around a rekey, reports that describe the previous session reach both -/// ends of the link: the initiator accepts one the responder built before it -/// switched, and the responder's frames from the old session pollute the -/// initiator's receiver. Neither kind of report may trigger a resend. +/// Around a rekey, frames the responder sent on the old session reach the +/// initiator after it cut over: a report about the initiator's old-session +/// traffic, and data frames carrying the responder's old-session counters. +/// Neither may feed the new session's MMP state, so the reports either end +/// holds after the cutover describe the new session, an announce sent after +/// it is confirmed by them, and nothing is resent. #[tokio::test] async fn test_bloom_reports_from_the_previous_session_do_not_trigger_resends() { let mut fx = flip_fixture(true).await; @@ -1567,17 +1569,46 @@ async fn test_bloom_reports_from_the_previous_session_do_not_trigger_resends() { for packet in held { fx.nodes[M].node.handle_encrypted_frame(packet).await; } + assert_eq!( + rr_counters(&fx.nodes[M].node, &p), + None, + "M must not take P's report about M's previous session" + ); mmp_round(&mut fx.nodes).await; - - let m_rr = rr_counters(&fx.nodes[M].node, &p).expect("setup: M holds a report"); - let p_rr = rr_counters(&fx.nodes[P].node, &m).expect("setup: P holds a report"); + let describes_new = |rr: Option<(u64, u64, u32)>, next: u64| rr.is_none_or(|rr| rr.0 < next); assert!( - m_rr.0 >= next_counter(&fx.nodes[M].node, &p), - "setup: M's report must describe M's previous session" + describes_new( + rr_counters(&fx.nodes[M].node, &p), + next_counter(&fx.nodes[M].node, &p) + ), + "M's report must describe M's new session" ); assert!( - p_rr.0 >= next_counter(&fx.nodes[P].node, &m), - "setup: P's report must carry P's previous-session counter" + describes_new( + rr_counters(&fx.nodes[P].node, &m), + next_counter(&fx.nodes[P].node, &m) + ), + "P's report must not carry P's previous-session counter" + ); + + // Both ends come to hold a usable report on the new session, which an + // announce sent next is measured from. Under the old behaviour the + // previous-session reports kept both ends from ever holding one here. + let usable = |fx: &FlipFixture| { + [(M, p), (P, m)].into_iter().all(|(i, remote)| { + rr_counters(&fx.nodes[i].node, &remote) + .is_some_and(|rr| rr.0 < next_counter(&fx.nodes[i].node, &remote)) + }) + }; + for _ in 0..5 { + if usable(&fx) { + break; + } + mmp_round(&mut fx.nodes).await; + } + assert!( + usable(&fx), + "both ends must hold a usable new-session report" ); fx.nodes[M].node.bloom_state.mark_update_needed(p); @@ -1594,49 +1625,31 @@ async fn test_bloom_reports_from_the_previous_session_do_not_trigger_resends() { let sent_m = fx.nodes[M].node.metrics().bloom.sent.get(); let sent_p = fx.nodes[P].node.metrics().bloom.sent.get(); let guard = switch_guard(&fx); - let mut last = rr_counters(&fx.nodes[P].node, &m); - let mut changes = 0; for _ in 0..5 { mmp_round(&mut fx.nodes).await; fx.nodes[M].node.check_bloom_state().await; fx.nodes[P].node.check_bloom_state().await; process_available_packets(&mut fx.nodes).await; - let now = rr_counters(&fx.nodes[P].node, &m); - if now != last { - changes += 1; - } - last = now; } - assert!( - changes >= 2, - "setup: P must accept at least two reports from M, saw {changes}" - ); assert_unswitched(&fx, &guard); - let m_rr = rr_counters(&fx.nodes[M].node, &p).expect("setup: M holds a report"); - let p_rr = rr_counters(&fx.nodes[P].node, &m).expect("setup: P holds a report"); - assert!( - m_rr.0 >= next_counter(&fx.nodes[M].node, &p) - && p_rr.0 >= next_counter(&fx.nodes[P].node, &m), - "setup: both reports must still describe a previous session" - ); assert_eq!( fx.nodes[M].node.metrics().bloom.sent.get(), sent_m, - "M must not resend on reports from the previous session" + "M must not resend a delivered announce after the rekey" ); assert_eq!( fx.nodes[P].node.metrics().bloom.sent.get(), sent_p, - "P must not resend on reports carrying its previous-session counter" + "P must not resend a delivered announce after the rekey" ); assert!( - fx.nodes[M].node.bloom_state.announce_outstanding(&p), - "M must still hold its announce to P" + !fx.nodes[M].node.bloom_state.announce_outstanding(&p), + "M's announce must be confirmed by a new-session report" ); assert!( - fx.nodes[P].node.bloom_state.announce_outstanding(&m), - "P must still hold its announce to M" + !fx.nodes[P].node.bloom_state.announce_outstanding(&m), + "P's announce must be confirmed by a new-session report" ); cleanup_nodes(&mut fx.nodes).await; } diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 01e7aa2c..d6222327 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -2032,6 +2032,209 @@ async fn a_held_rekey_msg2_is_resent_only_on_the_peers_established_link() { cleanup_nodes(&mut third).await; } +/// A ReceiverReport about `highest` frames, as the link-layer message a peer +/// sends: the type byte followed by the body. +fn receiver_report_message(highest: u64) -> Vec { + crate::proto::mmp::ReceiverReport { + highest_counter: highest, + cumulative_packets_recv: highest, + cumulative_bytes_recv: highest * 100, + timestamp_echo: 0, + dwell_time: 0, + max_burst_loss: 0, + mean_burst_loss: 0, + jitter: 0, + ecn_ce_count: 0, + owd_trend: 0, + burst_loss_count: 0, + cumulative_reorder_count: 0, + interval_packets_recv: highest as u32, + interval_bytes_recv: (highest * 100) as u32, + } + .encode() +} + +/// A link rekey initiator that has cut over keeps its new session's MMP state +/// free of frames the responder sealed on the old session. +/// +/// node 0 cuts over and resets its MMP receiver and metrics. node 1 has not +/// yet seen a frame on the new epoch, so it still sends on the old session, +/// and node 0 decrypts those frames against its previous session during the +/// drain. Their payload is still delivered, but they describe the old +/// session: a data frame's counter must not become the new session's highest +/// counter, which would make every new-session frame count as a reorder, and +/// a ReceiverReport about node 0's old-session traffic must not become the +/// baseline later reports are judged against, which would reject every +/// new-session report as regressed. Once node 1 promotes, its frames and +/// reports on the new session are counted and accepted. +#[tokio::test] +async fn frames_on_the_previous_link_session_do_not_feed_the_new_sessions_mmp() { + let HeldMsg2Pair { + mut nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun0_rx, + tun1_rx, + node0_idx_before, + node1_idx_before, + held_msg2, + .. + } = rekey_pair_with_held_msg2().await; + + // node 0 completes its rekey and cuts over on its own tick; nothing is + // delivered to node 1, which stays on the old session. + nodes[0].node.handle_msg2(held_msg2).await; + nodes[0].node.check_rekey().await; + let peer = nodes[0].node.get_peer(&node1_addr).unwrap(); + assert_ne!(peer.our_index(), node0_idx_before, "node 0 must cut over"); + let mmp = peer.mmp().expect("node 0 must run link MMP"); + assert_eq!(mmp.receiver.highest_counter(), 0, "the cutover resets MMP"); + assert_eq!(mmp.metrics.rr_counters(), None, "the cutover resets MMP"); + let reports_before = mmp.metrics.reports_seen(); + + // node 1, still on the old session, sends data and then a report about + // node 0's old-session traffic. Only node 0's queue is delivered. + let old_rev = build_ipv6_packet(&fips1, &fips0, b"old session 1 to 0"); + nodes[1].node.handle_tun_outbound(old_rev.clone()).await; + nodes[1] + .node + .send_encrypted_link_message(&node0_addr, &receiver_report_message(1_000)) + .await + .unwrap(); + assert_eq!( + nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + node1_idx_before, + "node 1 must still be on the old session" + ); + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(10)).await; + process_available_packets(&mut nodes[..1]).await; + } + + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!( + got, + vec![old_rev], + "an old-session frame's payload must still be delivered" + ); + let mmp = nodes[0].node.get_peer(&node1_addr).unwrap().mmp().unwrap(); + assert_eq!( + mmp.metrics.reports_seen(), + reports_before, + "an old-session ReceiverReport must not reach the new session's metrics" + ); + assert_eq!(mmp.metrics.rr_counters(), None); + assert_eq!( + mmp.receiver.highest_counter(), + 0, + "an old-session frame's counter must not become the new session's highest" + ); + + // node 0's first new-epoch frame promotes node 1; node 1's frames and + // reports on the new session then feed node 0's MMP. + let new_fwd = build_ipv6_packet(&fips0, &fips1, b"new session 0 to 1"); + nodes[0].node.handle_tun_outbound(new_fwd.clone()).await; + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun1_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![new_fwd]); + assert_ne!( + nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + node1_idx_before, + "node 1 must promote on node 0's first new-epoch frame" + ); + let new_rev = build_ipv6_packet(&fips1, &fips0, b"new session 1 to 0"); + nodes[1].node.handle_tun_outbound(new_rev.clone()).await; + nodes[1] + .node + .send_encrypted_link_message(&node0_addr, &receiver_report_message(5)) + .await + .unwrap(); + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![new_rev]); + let mmp = nodes[0].node.get_peer(&node1_addr).unwrap().mmp().unwrap(); + assert!( + mmp.receiver.highest_counter() > 0, + "new-session frames must be counted" + ); + assert_eq!( + mmp.metrics.rr_counters().map(|(highest, _, _)| highest), + Some(5), + "a new-session ReceiverReport must be accepted" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The decrypt-worker path tells the sessions apart too. A frame the worker +/// decrypted under the previous session's index, bounced back after the +/// cutover, does not feed the new session's MMP; the same frame under the +/// current session's index does. +#[cfg(unix)] +#[tokio::test] +async fn a_worker_decrypted_frame_on_the_previous_link_session_does_not_feed_mmp() { + use crate::node::decrypt_worker::DecryptFallback; + + let HeldMsg2Pair { + mut nodes, + node1_addr, + node0_idx_before, + held_msg2, + .. + } = rekey_pair_with_held_msg2().await; + nodes[0].node.handle_msg2(held_msg2).await; + nodes[0].node.check_rekey().await; + let peer = nodes[0].node.get_peer(&node1_addr).unwrap(); + let previous_idx = peer.previous_our_index().expect("node 0 must be draining"); + assert_eq!(Some(previous_idx), node0_idx_before); + let current_idx = peer.our_index().expect("node 0 must have cut over"); + + // A bounced worker plaintext: the 4-byte session timestamp, then a + // ReceiverReport about 1000 frames. + let (transport_id, remote_addr) = (nodes[0].transport_id, nodes[1].addr.clone()); + let bounce = |receiver_idx: u32, counter: u64| { + let mut data = 0u32.to_le_bytes().to_vec(); + data.extend(receiver_report_message(1_000)); + DecryptFallback { + source_node_addr: node1_addr, + transport_id, + remote_addr: remote_addr.clone(), + timestamp_ms: Node::now_ms(), + packet_len: data.len() + 32, + receiver_idx, + fmp_counter: counter, + fmp_flags: 0, + fmp_plaintext_len: data.len(), + packet_data: data, + fmp_plaintext_offset: 0, + } + }; + + let previous = bounce(previous_idx.as_u32(), 900); + nodes[0].node.process_decrypt_fallback(previous).await; + let mmp = nodes[0].node.get_peer(&node1_addr).unwrap().mmp().unwrap(); + assert_eq!( + mmp.receiver.highest_counter(), + 0, + "a previous-session frame's counter must not be counted" + ); + assert_eq!( + mmp.metrics.rr_counters(), + None, + "a previous-session ReceiverReport must not be processed" + ); + + let current = bounce(current_idx.as_u32(), 1); + nodes[0].node.process_decrypt_fallback(current).await; + let mmp = nodes[0].node.get_peer(&node1_addr).unwrap().mmp().unwrap(); + assert_eq!(mmp.receiver.highest_counter(), 1); + assert_eq!(mmp.metrics.rr_counters().map(|(h, _, _)| h), Some(1_000)); + + cleanup_nodes(&mut nodes).await; +} + /// The responder hold is the drain ceiling at stock settings, and a raised /// link-dead timeout or heartbeat interval raises it past that ceiling, taking /// the larger of the two rather than their sum. diff --git a/src/proto/mmp/tests/state.rs b/src/proto/mmp/tests/state.rs index 34b959e6..07336644 100644 --- a/src/proto/mmp/tests/state.rs +++ b/src/proto/mmp/tests/state.rs @@ -797,3 +797,37 @@ fn test_gap_tracker_saturates_when_advancing_onto_the_ceiling_counter() { "a saturated expectation must not keep opening bursts" ); } + +/// Why a previous-session ReceiverReport must not reach the metrics after a +/// rekey: the reset leaves no baseline, so a report about the old session's +/// high counters is accepted as the first one, and every report on the new +/// session then reads as regressed against it until the new session's +/// counters pass the old ones. This pins the mechanism the data plane guards +/// against by dropping such reports; it is not a red-first test of that fix, +/// which lives in the node-level rekey tests. +#[test] +fn a_report_about_the_old_session_after_a_rekey_reset_rejects_the_new_sessions_reports() { + let mut m = MmpMetrics::new(); + m.process_receiver_report(&make_rr(1_000, 1_000, 100_000, 1_000, 0, 0), 1_050, 0); + m.reset_for_rekey(); + assert_eq!(m.rr_counters(), None, "the reset clears the baseline"); + + // A report the peer sent about the old session, arriving after the reset. + m.process_receiver_report(&make_rr(1_200, 1_200, 120_000, 1_100, 0, 0), 1_150, 1_000); + assert_eq!(m.rr_counters().map(|(h, _, _)| h), Some(1_200)); + + // The new session's reports start from small counters and are refused. + m.process_receiver_report(&make_rr(10, 10, 1_000, 1_200, 0, 0), 1_250, 2_000); + assert_eq!( + m.rr_counters().map(|(h, _, _)| h), + Some(1_200), + "the new session's report is rejected as regressed" + ); + + // Without the old-session report, the same new-session report is taken. + let mut m = MmpMetrics::new(); + m.process_receiver_report(&make_rr(1_000, 1_000, 100_000, 1_000, 0, 0), 1_050, 0); + m.reset_for_rekey(); + m.process_receiver_report(&make_rr(10, 10, 1_000, 1_200, 0, 0), 1_250, 2_000); + assert_eq!(m.rr_counters().map(|(h, _, _)| h), Some(10)); +}