diff --git a/CHANGELOG.md b/CHANGELOG.md index e53fa254..f169df8e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -101,6 +101,20 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 interval instead. That retry interval gates only a peer whose last attempt failed, so it cannot clamp a `heartbeat_interval_secs` configured below it. +#### Session setup + +- A session whose last handshake message is lost no longer stays one-sided. + The initiator sent msg3 once and treated the session as established at once; + when that one datagram was lost, the responder kept waiting for it and + dropped every frame the initiator sent, and nothing sent msg3 again, because + the responder's repeated SessionAck was refused as arriving in the wrong + state. The session stayed that way until the next session rekey, or with + periodic rekey switched off, indefinitely. The initiator now keeps its msg3 + and resends it on the handshake resend interval, with backoff, until a frame + from the responder authenticates or `handshake_max_resends` resends have gone + out. The wire format is unchanged: the resend carries the same msg3, and a + responder that already completed the session refuses the duplicate as before. + #### Link and session rekey - A forged rekey msg2 no longer takes the link down. The rekey initiator gave diff --git a/src/node/handlers/session.rs b/src/node/handlers/session.rs index 25ac5b4d..4ec11836 100644 --- a/src/node/handlers/session.rs +++ b/src/node/handlers/session.rs @@ -335,6 +335,12 @@ impl Node { } }; + // A frame that authenticates on this session, in any epoch slot, + // proves the peer completed the handshake, so a msg3 still held for + // resend has arrived. Only an established initiator holds one; for + // every other entry this is a no-op. + entry.clear_handshake_payload(); + // React to the epoch the frame decrypted against. The shell opened // the frame; the core classifies the post-decrypt reaction over the // plain-data slot + session flags, and the shell applies the @@ -1044,7 +1050,7 @@ impl Node { let msg3_wire = SessionMsg3::new(msg3); let msg3_payload = msg3_wire.encode(); let my_addr = *self.node_addr(); - let mut datagram = SessionDatagram::new(my_addr, *src_addr, msg3_payload) + let mut datagram = SessionDatagram::new(my_addr, *src_addr, msg3_payload.clone()) .with_ttl(self.config().node.session.default_ttl); if let Err(e) = self.send_session_datagram(&mut datagram).await { @@ -1062,11 +1068,19 @@ impl Node { }; let now_ms = Self::now_ms(); + let resend_interval = self.config().node.rate_limit.handshake_resend_interval_ms; entry.set_state(EndToEndState::Established(session)); entry.set_coords_warmup_remaining(self.config().node.session.coords_warmup_packets); entry.mark_established(now_ms); entry.init_mmp(&self.config().node.session_mmp); - entry.clear_handshake_payload(); + // Keep msg3 for resend. This end is established once msg3 leaves, the + // responder only once it arrives, and nothing else repairs a lost + // msg3: the responder's resent SessionAck lands on the not-initiating + // arm above and is refused. `resend_pending_session_handshakes` + // resends it until a frame from the peer authenticates on this + // session or the resend budget is spent. The rekey arm keeps its + // msg3 for the same reason. + entry.set_handshake_payload(msg3_payload, now_ms + resend_interval); entry.touch(now_ms); self.sessions.insert(*src_addr, entry); self.insert_coord_hint(*src_addr, ack.src_coords.clone(), now_ms); diff --git a/src/node/handlers/timeout.rs b/src/node/handlers/timeout.rs index 7f62f9ac..aa833b2f 100644 --- a/src/node/handlers/timeout.rs +++ b/src/node/handlers/timeout.rs @@ -6,6 +6,7 @@ use crate::peer::machine::TimerKind; use crate::proto::fmp::{ ConnAction, ConnSnapshot, LifecycleView, PeerSnapshot, RekeyResendSnapshot, }; +use crate::proto::fsp::{FspAction, InitialMsg3ResendSnapshot}; use crate::transport::LinkId; use tracing::{debug, info, warn}; @@ -403,6 +404,10 @@ impl Node { /// - If the handshake has exceeded the timeout window, remove the session. /// - If a resend is due and under max resends, resend the stored payload /// wrapped in a fresh SessionDatagram (so routing can adapt). + /// + /// For an established initiator still holding its msg3, resend it until + /// the peer is heard from or the budget is spent (see + /// `resend_initial_msg3`). Established sessions are never removed here. pub(in crate::node) async fn resend_pending_session_handshakes(&mut self, now_ms: u64) { if self.sessions.is_empty() { return; @@ -474,6 +479,96 @@ impl Node { ); } } + + self.resend_initial_msg3(now_ms).await; + } + + /// Resend an established initiator's retained msg3 until the responder is + /// heard from, and stop retaining it once the resend budget is spent. + /// + /// The initiator is established the moment msg3 leaves; the responder only + /// once it arrives. Nothing else repairs a lost msg3: the responder's + /// resent SessionAck reaches an entry that is no longer initiating and is + /// refused. The payload is released by the first inbound frame that + /// authenticates on the session (`handle_encrypted_session_msg`) or, here, + /// when the budget is spent. Runs whether or not periodic rekey is enabled. + async fn resend_initial_msg3(&mut self, now_ms: u64) { + use crate::proto::link::SessionDatagram; + + let candidates = self.initial_msg3_resend_snapshots(now_ms); + if candidates.is_empty() { + return; + } + let max_resends = self.config().node.rate_limit.handshake_max_resends; + let interval_ms = self.config().node.rate_limit.handshake_resend_interval_ms; + let backoff = self.config().node.rate_limit.handshake_resend_backoff; + let ttl = self.config().node.session.default_ttl; + let my_addr = *self.node_addr(); + + for action in self.fsp.poll_initial_msg3_resends(candidates, max_resends) { + match action { + FspAction::ReleaseInitialMsg3 { addr } => { + if let Some(entry) = self.sessions.get_mut(&addr) { + entry.clear_handshake_payload(); + } + info!( + dest = %self.peer_display_name(&addr), + "Session msg3 unconfirmed after max resends, no longer resending" + ); + } + FspAction::ResendInitialMsg3 { addr } => { + let payload = match self.sessions.get(&addr).and_then(|e| e.handshake_payload()) + { + Some(p) => p.to_vec(), + None => continue, + }; + let mut datagram = SessionDatagram::new(my_addr, addr, payload).with_ttl(ttl); + let sent = match self.send_session_datagram(&mut datagram).await { + Ok(_) => true, + Err(e) => { + debug!( + dest = %self.peer_display_name(&addr), + error = %e, + "Session msg3 resend failed" + ); + false + } + }; + if sent && let Some(entry) = self.sessions.get_mut(&addr) { + let count = entry.resend_count() + 1; + let next = + now_ms + (interval_ms as f64 * backoff.powi(count as i32)) as u64; + entry.record_resend(next); + debug!( + dest = %self.peer_display_name(&addr), + resend = count, + "Resent session msg3" + ); + } + } + #[allow(unreachable_patterns)] + _ => {} + } + } + } + + /// Snapshot every established session still retaining its initial msg3, + /// pre-evaluating the resend-due predicate against `now_ms` so the core + /// reads no clock. + /// + /// `is_established()` partitions the shared handshake resend slot: a + /// non-established entry's SessionSetup or SessionAck belongs to the passes + /// above, an established entry's msg3 to this one. + fn initial_msg3_resend_snapshots(&self, now_ms: u64) -> Vec { + self.sessions + .iter() + .filter(|(_, entry)| entry.is_established() && entry.handshake_payload().is_some()) + .map(|(addr, entry)| InitialMsg3ResendSnapshot { + addr: *addr, + resend_count: entry.resend_count(), + resend_due: entry.next_resend_at_ms() != 0 && now_ms >= entry.next_resend_at_ms(), + }) + .collect() } /// Remove established sessions that have been idle too long. diff --git a/src/node/session/mod.rs b/src/node/session/mod.rs index 635c81f9..f06a3bf7 100644 --- a/src/node/session/mod.rs +++ b/src/node/session/mod.rs @@ -123,8 +123,11 @@ pub(crate) struct SessionEntry { bytes_recv: u64, // === Handshake Resend === - /// Encoded session-layer payload for resend (SessionSetup or SessionAck). - /// Cleared on Established transition. + /// The last initial-handshake message this side sent and has no proof the + /// peer received: SessionSetup while initiating, SessionAck while awaiting + /// msg3, and on an established initiator its SessionMsg3 until an inbound + /// frame authenticates on the session or the resend budget is spent. The + /// state says which one it is, and which sweep resends it. handshake_payload: Option>, /// Number of resends performed. resend_count: u32, @@ -408,6 +411,8 @@ impl SessionEntry { /// /// For initiators, this is the SessionSetup payload bytes. /// For responders, this is the SessionAck payload bytes. + /// For an initiator that has just become established, this is its + /// SessionMsg3 payload bytes, held until the responder is heard from. /// The payload is re-wrapped in a fresh SessionDatagram on each resend /// so routing can adapt to topology changes. pub(crate) fn set_handshake_payload(&mut self, payload: Vec, next_resend_at_ms: u64) { @@ -421,7 +426,8 @@ impl SessionEntry { self.handshake_payload.as_deref() } - /// Clear the stored handshake payload (called on Established transition). + /// Clear the stored handshake payload (an inbound frame authenticated on + /// the session, or the msg3 resend budget is spent). pub(crate) fn clear_handshake_payload(&mut self) { self.handshake_payload = None; self.next_resend_at_ms = 0; diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index 78524156..65a262ae 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -5090,6 +5090,302 @@ async fn test_forged_session_ack_leaves_the_initiation_able_to_complete_on_the_g cleanup_nodes(&mut nodes).await; } +// ============================================================================ +// Integration tests: a lost initial msg3 +// ============================================================================ + +/// Build a two-node pair where node 0's initial msg3 was sent and dropped, so +/// node 0 is established and node 1 is still waiting for msg3. +/// +/// Periodic rekey is off on both nodes: that is the configuration with no +/// other recovery, and it keeps the rekey drivers out of the picture. Every +/// step asserts its packet count, so a harness surprise fails loudly instead +/// of being read as the defect. +async fn pair_with_lost_initial_msg3() -> Vec { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_ESTABLISHED}; + + let mut nodes = make_rekey_disabled_pair().await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let node1_pubkey = nodes[1].node.identity().pubkey_full(); + + nodes[0] + .node + .initiate_session(node1_addr, node1_pubkey) + .await + .expect("initiate_session failed"); + + assert_eq!( + process_available_packets(&mut nodes[1..]).await, + 1, + "node 1 must have exactly node 0's SessionSetup queued" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_awaiting_msg3(), + "node 1 must be awaiting msg3 after answering the SessionSetup" + ); + + assert_eq!( + process_available_packets(&mut nodes[..1]).await, + 1, + "node 0 must have exactly node 1's SessionAck queued" + ); + assert!( + nodes[0] + .node + .get_session(&node1_addr) + .expect("initiator entry present") + .is_established(), + "node 0 must be established once it has sent msg3" + ); + + // Drop msg3: take it out of node 1's queue instead of processing it. + let dropped: Vec<_> = std::iter::from_fn(|| nodes[1].packet_rx.try_recv().ok()).collect(); + assert_eq!( + dropped.len(), + 1, + "node 1 must have only node 0's msg3 queued" + ); + assert_eq!( + CommonPrefix::parse(&dropped[0].data).map(|p| p.phase), + Some(PHASE_ESTABLISHED), + "the dropped packet must be a link data frame carrying the msg3" + ); + + pump_until_quiet(&mut nodes).await; + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_awaiting_msg3(), + "node 1 must still be awaiting msg3 after the drop" + ); + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "nothing may be left queued at node 1" + ); + + nodes +} + +/// One lost initial msg3 must cost one resend, not the session, and the +/// recovery must not depend on periodic rekey. +#[tokio::test] +async fn a_lost_initial_msg3_is_resent_and_the_responder_completes_the_session() { + let mut nodes = pair_with_lost_initial_msg3().await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); + nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + nodes[0] + .node + .resend_pending_session_handshakes(Node::now_ms() + interval_ms + 1) + .await; + pump_until_quiet(&mut nodes).await; + + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder entry present") + .is_established(), + "node 1 must complete the session once node 0 resends its msg3" + ); + + let fwd = build_ipv6_packet(&fips0, &fips1, b"after msg3 resend 0 to 1"); + let rev = build_ipv6_packet(&fips1, &fips0, b"after msg3 resend 1 to 0"); + nodes[0].node.handle_tun_outbound(fwd.clone()).await; + nodes[1].node.handle_tun_outbound(rev.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![fwd], "node 0 to node 1 must decode"); + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![rev], "node 1 to node 0 must decode"); + + cleanup_nodes(&mut nodes).await; +} + +/// Drain and count the packets queued at `node` without processing them. +fn drain_queued(node: &mut TestNode) -> usize { + std::iter::from_fn(|| node.packet_rx.try_recv().ok()).count() +} + +/// On a healthy session the retained msg3 is released by the responder's +/// first frame, so the resend window is one round trip wide and a later tick +/// sends nothing. +#[tokio::test] +async fn an_initiator_stops_resending_msg3_once_a_responder_frame_authenticates() { + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + + // Control: before the responder has sent anything, the sweep does resend, + // so a silent sweep at the end is the release and not a dead driver. + let t1 = Node::now_ms() + interval_ms + 1; + nodes[0].node.resend_pending_session_handshakes(t1).await; + assert_eq!( + drain_queued(&mut nodes[1]), + 1, + "control: the sweep must resend msg3 while the responder is unheard" + ); + + let rev = build_ipv6_packet(&fips1, &fips0, b"responder's first frame"); + nodes[1].node.handle_tun_outbound(rev.clone()).await; + pump_until_quiet(&mut nodes).await; + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!( + got, + vec![rev], + "the responder's frame must authenticate at node 0" + ); + assert_eq!(nodes[1].packet_rx.len(), 0, "nothing may be left queued"); + + nodes[0] + .node + .resend_pending_session_handshakes(t1 + 64_000) + .await; + assert_eq!( + nodes[1].packet_rx.len(), + 0, + "a msg3 the responder has answered must not be resent" + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The resend is harmless to a responder that already completed: it is +/// refused as a bad-state reject and both directions keep decoding. This is +/// the wire-neutrality claim, observed rather than argued. +#[tokio::test] +async fn a_resent_msg3_reaching_an_established_responder_is_refused_and_the_session_keeps_working() +{ + let mut nodes = make_rekey_disabled_pair().await; + establish_pair_session(&mut nodes).await; + let node0_addr = *nodes[0].node.node_addr(); + let node1_addr = *nodes[1].node.node_addr(); + let (tun0_tx, tun0_rx) = std::sync::mpsc::channel(); + nodes[0].node.supervisor.tun_tx = Some(tun0_tx); + let (tun1_tx, tun1_rx) = std::sync::mpsc::channel(); + nodes[1].node.supervisor.tun_tx = Some(tun1_tx); + let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); + let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); + let interval_ms = nodes[0] + .node + .config() + .node + .rate_limit + .handshake_resend_interval_ms; + let before = nodes[1].node.stats().session.bad_state; + + nodes[0] + .node + .resend_pending_session_handshakes(Node::now_ms() + interval_ms + 1) + .await; + assert_eq!( + nodes[1].packet_rx.len(), + 1, + "precondition: node 0 must have resent its msg3 to node 1" + ); + pump_until_quiet(&mut nodes).await; + + assert_eq!( + nodes[1].node.stats().session.bad_state, + before + 1, + "the duplicate msg3 must be refused as a bad-state reject" + ); + assert!( + nodes[1] + .node + .get_session(&node0_addr) + .expect("responder session present") + .is_established(), + "the duplicate msg3 must not disturb the responder's session" + ); + + let fwd = build_ipv6_packet(&fips0, &fips1, b"after duplicate msg3 0 to 1"); + let rev = build_ipv6_packet(&fips1, &fips0, b"after duplicate msg3 1 to 0"); + nodes[0].node.handle_tun_outbound(fwd.clone()).await; + nodes[1].node.handle_tun_outbound(rev.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![fwd], "node 0 to node 1 must decode"); + let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); + assert_eq!(got, vec![rev], "node 1 to node 0 must decode"); + + cleanup_nodes(&mut nodes).await; +} + +/// The resend ladder is bounded in count and in time: after +/// `handshake_max_resends` resends nothing more is sent and the payload is no +/// longer held, while the session itself is kept. +#[tokio::test] +async fn initial_msg3_resends_stop_at_the_budget_and_release_the_payload() { + let mut nodes = pair_with_lost_initial_msg3().await; + let node1_addr = *nodes[1].node.node_addr(); + let max_resends = nodes[0].node.config().node.rate_limit.handshake_max_resends; + + // 64 s steps pass any single backoff interval at stock settings, so each + // step is due. Node 0 holds only its established entry, so the sweep's + // timeout pass cannot remove anything on it. + let mut now = Node::now_ms(); + let mut sent = Vec::new(); + for _ in 0..(max_resends + 2) { + now += 64_000; + nodes[0].node.resend_pending_session_handshakes(now).await; + sent.push(drain_queued(&mut nodes[1])); + } + let mut expected = vec![1; max_resends as usize]; + expected.extend([0, 0]); + assert_eq!( + sent, expected, + "one resend per due tick up to the budget, then none" + ); + + let entry = nodes[0] + .node + .get_session(&node1_addr) + .expect("the release must not tear the session down"); + assert!( + entry.handshake_payload().is_none(), + "the msg3 must no longer be held once the budget is spent" + ); + assert!( + entry.is_established(), + "node 0's session must stay established after the release" + ); + + cleanup_nodes(&mut nodes).await; +} + /// A SessionAck that fails to read must not end an FSP rekey the node /// initiated. /// diff --git a/src/proto/fsp/core.rs b/src/proto/fsp/core.rs index dc061cb0..05c5a5d5 100644 --- a/src/proto/fsp/core.rs +++ b/src/proto/fsp/core.rs @@ -2,8 +2,8 @@ //! //! Pure, runtime-agnostic decisions for the FSP end-to-end session lifecycle: //! the per-tick rekey choreography (initiator cutover, drain completion, rekey -//! trigger), msg3 retransmission classification, and the post-decrypt epoch -//! reaction. The async I/O adapters in `node::handlers::{rekey,session}` build +//! trigger), rekey and initial-handshake msg3 resend classification, and the +//! post-decrypt epoch reaction. The async I/O adapters in `node::handlers::{rekey,session}` build //! the plain-data snapshots (pre-computing every clock read into `u64`/`bool`), //! call these decisions, and drive the returned effects — the sends, the //! `SessionEntry` mutations, metrics, and logging. No I/O, no clock, no crypto, @@ -68,6 +68,14 @@ pub(crate) enum FspAction { /// 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 }, + /// Resend `addr`'s retained initial-handshake msg3 (the shell re-reads the + /// payload from the entry's handshake resend slot, sends it, then records + /// the resend on success). + ResendInitialMsg3 { addr: NodeAddr }, + /// Stop retaining `addr`'s initial-handshake msg3: the resend budget is + /// spent without an inbound frame showing the peer received it. The session + /// itself is kept; only the retained payload is dropped. + ReleaseInitialMsg3 { addr: NodeAddr }, /// Cache `coords` for `addr` in the shared coordinate cache /// (`coord_cache.insert`). CacheCoords { @@ -159,6 +167,18 @@ pub(crate) struct RekeyMsg3ResendSnapshot { pub resend_due: bool, } +/// A snapshot of one established session whose initiator still retains its +/// initial-handshake msg3, taken by the shell for the resend decision. +pub(crate) struct InitialMsg3ResendSnapshot { + /// The session's remote node address (release/resend target). + pub addr: NodeAddr, + /// How many msg3 resends have already happened. + pub resend_count: u32, + /// The retained msg3 is due as of the shell's `now_ms` (pre-evaluated: + /// `next_resend_at_ms != 0 && now_ms >= next_resend_at_ms`). + pub resend_due: bool, +} + /// Which key epoch a just-decrypted frame authenticated against — the shell-side /// [`EpochSlot`](crate::node::session::EpochSlot) mapped to a proto-local /// plain-data mirror so the core carries no `node` dependency. @@ -292,6 +312,33 @@ impl Fsp { abandons } + /// Decide the initial-handshake msg3 resends for the established sessions + /// the shell snapshotted as retaining one. A due candidate whose budget is + /// spent is released; an in-budget due candidate is resent; a candidate not + /// yet due gets nothing this tick. Releases come first, as abandons do in + /// [`poll_rekey_msg3_resends`](Self::poll_rekey_msg3_resends); the shell + /// commits a resend's count and reschedule only on a successful send. + pub(crate) fn poll_initial_msg3_resends( + &self, + candidates: Vec, + max_resends: u32, + ) -> Vec { + let mut releases = Vec::new(); + let mut resends = Vec::new(); + for c in candidates { + if !c.resend_due { + continue; + } + if c.resend_count >= max_resends { + releases.push(FspAction::ReleaseInitialMsg3 { addr: c.addr }); + continue; + } + resends.push(FspAction::ResendInitialMsg3 { addr: c.addr }); + } + releases.extend(resends); + releases + } + /// Classify the reaction to a frame that authenticated against `slot`. Pure /// over the slot and the two plain-data session flags; the shell applies the /// resulting `SessionEntry` mutation and observability. diff --git a/src/proto/fsp/mod.rs b/src/proto/fsp/mod.rs index 50c73e73..6ebb179c 100644 --- a/src/proto/fsp/mod.rs +++ b/src/proto/fsp/mod.rs @@ -12,7 +12,8 @@ //! address-only coordinate helpers downward from `crate::proto::stp`. //! //! - `core.rs` — the stateless [`Fsp`] anchor + [`FspAction`]: the pure rekey -//! choreography (`poll_rekey`/`poll_rekey_msg3_resends`), the post-decrypt +//! choreography (`poll_rekey`/`poll_rekey_msg3_resends`), the initial-handshake +//! msg3 resend decision (`poll_initial_msg3_resends`), the post-decrypt //! `classify_epoch`, the initiation tie-break, and the pure MTU-clamp / //! bounded-queue / ECN transforms. No clock/crypto/I/O/tracing. //! - `limits.rs` — the session-rekey timing constants. @@ -27,8 +28,9 @@ pub(crate) mod wire; mod tests; pub(crate) use core::{ - DecryptSlot, EpochReaction, Fsp, FspAction, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, - cutover_timer_elapsed, initiation_winner, mark_ipv6_ecn_ce, push_bounded_pending, + DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, + RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, + mark_ipv6_ecn_ce, push_bounded_pending, }; pub use wire::{ FspInnerFlags, SessionAck, SessionFlags, SessionMessageType, SessionMsg3, SessionSetup, diff --git a/src/proto/fsp/tests/core.rs b/src/proto/fsp/tests/core.rs index f23d8ea7..b98dcee3 100644 --- a/src/proto/fsp/tests/core.rs +++ b/src/proto/fsp/tests/core.rs @@ -2,9 +2,9 @@ use crate::FipsAddress; use crate::proto::fsp::core::{ - DecryptSlot, EpochReaction, Fsp, FspAction, RekeyCfg, RekeyMsg3ResendSnapshot, SessionSnapshot, - cutover_timer_elapsed, initiation_winner, mark_ipv6_ecn_ce, push_bounded_pending, - should_apply_path_mtu, + DecryptSlot, EpochReaction, Fsp, FspAction, InitialMsg3ResendSnapshot, RekeyCfg, + RekeyMsg3ResendSnapshot, SessionSnapshot, cutover_timer_elapsed, initiation_winner, + mark_ipv6_ecn_ce, push_bounded_pending, should_apply_path_mtu, }; use crate::proto::fsp::limits::FSP_CUTOVER_DELAY_MS; use crate::proto::stp::TreeCoordinate; @@ -441,6 +441,85 @@ fn poll_msg3_abandons_first() { ); } +// ===== poll_initial_msg3_resends ===== + +/// Build an initial-handshake msg3 resend snapshot for the decision under test. +fn initial_msg3_snapshot( + addr_byte: u8, + resend_count: u32, + resend_due: bool, +) -> InitialMsg3ResendSnapshot { + InitialMsg3ResendSnapshot { + addr: make_node_addr(addr_byte), + resend_count, + resend_due, + } +} + +/// A candidate that is not yet due gets nothing, whether or not its budget is +/// spent. +#[test] +fn poll_initial_msg3_not_due_is_noop() { + let fsp = Fsp::new(); + assert!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(1, 0, false)], 3) + .is_empty() + ); + assert!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(1, 99, false)], 3) + .is_empty() + ); +} + +/// A due candidate within its budget is resent. +#[test] +fn poll_initial_msg3_resends_when_due_in_budget() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(2, 1, true)], 3), + vec![FspAction::ResendInitialMsg3 { + addr: make_node_addr(2) + }] + ); +} + +/// A due candidate whose budget is spent is released, not resent. +#[test] +fn poll_initial_msg3_releases_when_due_at_budget() { + let fsp = Fsp::new(); + assert_eq!( + fsp.poll_initial_msg3_resends(vec![initial_msg3_snapshot(3, 3, true)], 3), + vec![FspAction::ReleaseInitialMsg3 { + addr: make_node_addr(3) + }] + ); +} + +/// Releases are returned before resends. +#[test] +fn poll_initial_msg3_releases_before_resends() { + let fsp = Fsp::new(); + let actions = fsp.poll_initial_msg3_resends( + vec![ + initial_msg3_snapshot(1, 0, true), + initial_msg3_snapshot(2, 5, true), + ], + 3, + ); + assert_eq!( + actions, + vec![ + FspAction::ReleaseInitialMsg3 { + addr: make_node_addr(2) + }, + FspAction::ResendInitialMsg3 { + addr: make_node_addr(1) + }, + ], + "releases are grouped before resends" + ); +} + // ===== classify_epoch ===== #[test]