diff --git a/src/node/handlers/handshake.rs b/src/node/handlers/handshake.rs index b74171b5..9e32fdd3 100644 --- a/src/node/handlers/handshake.rs +++ b/src/node/handlers/handshake.rs @@ -69,7 +69,7 @@ pub(in crate::node) enum Msg1Waiver { } impl EstablishView for Node { - fn establish_snapshot(&self, peer_addr: &NodeAddr) -> EstablishSnapshot { + fn establish_snapshot(&self, peer_addr: &NodeAddr, msg1: &Msg1Digest) -> EstablishSnapshot { let existing = self.peers.get(peer_addr); let max_peers = self.max_peers(); EstablishSnapshot { @@ -85,6 +85,7 @@ impl EstablishView for Node { .unwrap_or(false), rekey_in_progress: existing.map(|p| p.rekey_in_progress()).unwrap_or(false), held_answer: existing.and_then(|p| p.rekey_answer().cloned()), + msg1_answered_before: existing.is_some_and(|p| p.answered_before(msg1)), existing_msg2: existing.and_then(|p| p.handshake_msg2().map(|m| m.to_vec())), at_max_peers: max_peers > 0 && self.peers.len() >= max_peers, has_pending_outbound_to_peer: self.connections().any(|(_, machine)| { @@ -312,6 +313,19 @@ impl Node { .is_none_or(|t| t.accept_connections()) } + /// The transport and address of `peer`'s established link, where a rekey + /// msg2 is sent whatever address its msg1 arrived from. + fn established_link( + &self, + peer: &NodeAddr, + ) -> Option<( + crate::transport::TransportId, + crate::transport::TransportAddr, + )> { + let p = self.peers.get(peer)?; + Some((p.transport_id()?, p.current_addr()?.clone())) + } + /// Handle handshake message 1 (phase 0x1). /// /// This creates a new inbound connection. Rate limiting is applied @@ -563,7 +577,7 @@ impl Node { // session age resolved here, the max-peers cap, our own address for the // tie-break). Taken before this connection is inserted into the // registry, matching the pre-refactor read points. - let est = self.establish_snapshot(&peer_node_addr); + let est = self.establish_snapshot(&peer_node_addr, &wire.msg1_digest); // === PHASE C: structured classification === // Evaluate the inbound decision once on a local establish leg and route @@ -601,7 +615,10 @@ impl Node { .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); } InboundDecision::Reject { - reason: reason @ (InboundReject::PendingSession | InboundReject::DualRekeyWon), + reason: + reason @ (InboundReject::PendingSession + | InboundReject::DualRekeyWon + | InboundReject::AnsweredBefore), } => { // Existing-peer rekey rejects: the classification took the // fresh-context fail path (no actions) and the local machine is @@ -616,6 +633,11 @@ impl Node { peer = %self.peer_display_name(&peer_node_addr), "Dual rekey initiation: we win (smaller addr), dropping their msg1" ), + InboundReject::AnsweredBefore => debug!( + peer = %self.peer_display_name(&peer_node_addr), + remote_addr = %packet.remote_addr, + "Rekey msg1 answered in an ended cycle, dropping the copy" + ), InboundReject::AtMaxPeers => unreachable!(), } // `conn`/`link_id` were never inserted into the registry, so the @@ -626,12 +648,16 @@ impl Node { InboundDecision::ResendMsg2 { msg2 } => { // Duplicate msg1 at the same epoch: the decision carries the // stored msg2 bytes and the inline resend below owns the send; - // the classification touched no state. + // the classification touched no state. It goes on the peer's + // established link, as a rekey msg2 does: a genuine duplicate + // comes from the address the peering was just formed with, + // while a copy can come from anywhere. debug_assert!(actions.is_empty()); if let Some(msg2) = msg2.as_deref() - && let Some(transport) = self.transports.get(&packet.transport_id) + && let Some((tid, addr)) = self.established_link(&peer_node_addr) + && let Some(transport) = self.transports.get(&tid) { - match transport.send(&packet.remote_addr, msg2).await { + match transport.send(&addr, msg2).await { Ok(_) => debug!( peer = %self.peer_display_name(&peer_node_addr), "Resent msg2 for duplicate msg1 (same epoch)" @@ -646,16 +672,10 @@ impl Node { } InboundDecision::ResendRekeyMsg2 { peer, msg2 } => { // A resend of the msg1 that armed the pending we hold: our - // msg2 was lost, so give the same answer again. It goes to the - // peer's established address, not the msg1's source: a - // captured msg1 replayed from elsewhere authenticates the same - // as a resend, and answering its source would reflect. + // msg2 was lost, so give the same answer again, on the peer's + // established link as the first answer went. debug_assert!(actions.is_empty()); - let target = self - .peers - .get(&peer) - .and_then(|p| Some((p.transport_id()?, p.current_addr()?.clone()))); - if let Some((tid, addr)) = target + if let Some((tid, addr)) = self.established_link(&peer) && let Some(transport) = self.transports.get(&tid) { match transport.send(&addr, &msg2).await { @@ -726,30 +746,43 @@ impl Node { } }; - // Send msg2 response using the new handshake. + // Send msg2 response using the new handshake, on the peer's + // established link rather than to the msg1's source. A copy + // of a msg1 authenticates as the peer from any address, so + // answering its source would reflect to an address the + // sender chose. A peer whose address changed is answered at + // the old one until a frame from the new address moves it. let wire_msg2 = build_msg2(our_new_index, wire.their_index, &wire.msg2_payload); - if let Some(transport) = self.transports.get(&packet.transport_id) { - match transport.send(&packet.remote_addr, &wire_msg2).await { - Ok(_) => { - debug!( - peer = %self.peer_display_name(&peer), - new_our_index = %our_new_index, - "Sent rekey msg2 response" - ); - } - Err(e) => { - warn!( - peer = %self.peer_display_name(&peer), - error = %e, - "Failed to send rekey msg2" - ); - let _ = self.index_allocator.free(our_new_index); - self.stats_mut() - .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); - return; - } + let sent = match self.established_link(&peer) { + Some((tid, addr)) => match self.transports.get(&tid) { + Some(transport) => transport + .send(&addr, &wire_msg2) + .await + .map(|_| tid) + .map_err(|e| e.to_string()), + None => Err("no transport for the peer's link".to_string()), + }, + None => Err("the peer has no established link".to_string()), + }; + let link_transport = match sent { + Ok(tid) => tid, + Err(e) => { + warn!( + peer = %self.peer_display_name(&peer), + error = %e, + "Failed to send rekey msg2" + ); + let _ = self.index_allocator.free(our_new_index); + self.stats_mut() + .record_reject(RejectReason::Handshake(HandshakeReject::BadState)); + return; } - } + }; + debug!( + peer = %self.peer_display_name(&peer), + new_our_index = %our_new_index, + "Sent rekey msg2 response" + ); // Store the new session as the responder's pending session. It // is promoted by the initiator's first new-epoch frame, not by @@ -765,9 +798,12 @@ impl Node { existing.record_peer_rekey(); } - // Register new index in peers_by_index. + // Register new index in peers_by_index, under the transport + // the msg2 went out on: the peer's frames on the new session + // arrive there, and retirement removes the entry by the + // peer's transport, not the one the msg1 came in on. self.peers_by_index - .insert((packet.transport_id, our_new_index.as_u32()), peer); + .insert((link_transport, our_new_index.as_u32()), peer); // Do NOT touch addr_to_link — the entry must keep pointing at the // original link so future msg1s from this address are recognized diff --git a/src/node/tests/session.rs b/src/node/tests/session.rs index d6222327..e484be83 100644 --- a/src/node/tests/session.rs +++ b/src/node/tests/session.rs @@ -1222,9 +1222,10 @@ async fn rekey_cutover_preserves_data_plane() { cleanup_nodes(&mut nodes).await; } -/// A two-node pair caught mid FMP rekey, with node 1's msg2 held back from -/// node 0. Built by [`rekey_pair_with_held_msg2`]. -struct HeldMsg2Pair { +/// A two-node pair with an FSP session over an FMP link whose link sessions +/// have been aged, with a TUN receiver on each node. Built by +/// [`aged_link_pair`]. +struct AgedLinkPair { nodes: Vec, node0_addr: NodeAddr, node1_addr: NodeAddr, @@ -1232,36 +1233,16 @@ struct HeldMsg2Pair { fips1: crate::FipsAddress, tun0_rx: std::sync::mpsc::Receiver>, tun1_rx: std::sync::mpsc::Receiver>, - node0_idx_before: Option, - node1_idx_before: Option, - rekey_idx: crate::utils::index::SessionIndex, - held_msg2: crate::transport::ReceivedPacket, } -/// Build a two-node pair with an FSP session, age both link sessions past -/// both rekey gates, start node 0's FMP rekey, deliver its msg1 to node 1 -/// only, and pull node 1's real msg2 out of node 0's queue. -/// -/// node 0 rekeys on time and node 1 only ever responds, so node 1 holds the -/// new session it committed at msg1 and node 0 is mid-cycle when this -/// returns. Both directions are shown to decode before the rekey, so a later -/// delivery failure is the rekey's and not the harness's. -async fn rekey_pair_with_held_msg2() -> HeldMsg2Pair { - use crate::proto::fmp::wire::{CommonPrefix, PHASE_MSG2}; - use crate::transport::ReceivedPacket; - - const REKEY_AFTER_SECS: u64 = 60; - - // node 0 rekeys on time; node 1 only ever responds. - let mut cfg0 = crate::config::Config::new(); - cfg0.node.rekey.enabled = true; - cfg0.node.rekey.after_secs = REKEY_AFTER_SECS; - cfg0.node.rekey.after_messages = u64::MAX; - let mut cfg1 = crate::config::Config::new(); - cfg1.node.rekey.enabled = true; - cfg1.node.rekey.after_secs = u64::MAX; - cfg1.node.rekey.after_messages = u64::MAX; - +/// Build a two-node pair from `cfg0` and `cfg1`, peer node 0 to node 1 over +/// FMP, open an FSP session from node 0, show that both directions decode, +/// then backdate both link sessions by `age`. +async fn aged_link_pair( + cfg0: crate::config::Config, + cfg1: crate::config::Config, + age: Duration, +) -> AgedLinkPair { let mut nodes = vec![ make_test_node_with_config(cfg0, 1280).await, make_test_node_with_config(cfg1, 1280).await, @@ -1304,8 +1285,8 @@ async fn rekey_pair_with_held_msg2() -> HeldMsg2Pair { let fips0 = crate::FipsAddress::from_node_addr(&node0_addr); let fips1 = crate::FipsAddress::from_node_addr(&node1_addr); - // Baseline: both directions decode before the rekey, so a failure below - // is the rekey's and not the harness's. + // Baseline: both directions decode before the caller's scenario, so a + // later failure is the scenario's and not the harness's. let pre_fwd = build_ipv6_packet(&fips0, &fips1, b"pre-rekey 0 to 1"); let pre_rev = build_ipv6_packet(&fips1, &fips0, b"pre-rekey 1 to 0"); nodes[0].node.handle_tun_outbound(pre_fwd.clone()).await; @@ -1316,9 +1297,7 @@ async fn rekey_pair_with_held_msg2() -> HeldMsg2Pair { let got: Vec> = std::iter::from_fn(|| tun0_rx.try_recv().ok()).collect(); assert_eq!(got, vec![pre_rev], "baseline node 1 to node 0 must decode"); - // Age both sessions past both rekey gates: node 0's jittered time trigger, - // and node 1's 30 s floor below which a msg1 is a duplicate, not a rekey. - let age = Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + // Age both link sessions. nodes[0] .node .get_peer_mut(&node1_addr) @@ -1329,6 +1308,70 @@ async fn rekey_pair_with_held_msg2() -> HeldMsg2Pair { .get_peer_mut(&node0_addr) .unwrap() .test_backdate_session_established(age); + + AgedLinkPair { + nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun0_rx, + tun1_rx, + } +} + +/// A two-node pair caught mid FMP rekey, with node 1's msg2 held back from +/// node 0. Built by [`rekey_pair_with_held_msg2`]. +struct HeldMsg2Pair { + nodes: Vec, + node0_addr: NodeAddr, + node1_addr: NodeAddr, + fips0: crate::FipsAddress, + fips1: crate::FipsAddress, + tun0_rx: std::sync::mpsc::Receiver>, + tun1_rx: std::sync::mpsc::Receiver>, + node0_idx_before: Option, + node1_idx_before: Option, + rekey_idx: crate::utils::index::SessionIndex, + held_msg2: crate::transport::ReceivedPacket, +} + +/// Build a two-node pair with an FSP session, age both link sessions past +/// both rekey gates, start node 0's FMP rekey, deliver its msg1 to node 1 +/// only, and pull node 1's real msg2 out of node 0's queue. +/// +/// node 0 rekeys on time and node 1 only ever responds, so node 1 holds the +/// new session it committed at msg1 and node 0 is mid-cycle when this +/// returns. Both directions are shown to decode before the rekey, so a later +/// delivery failure is the rekey's and not the harness's. +async fn rekey_pair_with_held_msg2() -> HeldMsg2Pair { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_MSG2}; + use crate::transport::ReceivedPacket; + + const REKEY_AFTER_SECS: u64 = 60; + + // node 0 rekeys on time; node 1 only ever responds. + let mut cfg0 = crate::config::Config::new(); + cfg0.node.rekey.enabled = true; + cfg0.node.rekey.after_secs = REKEY_AFTER_SECS; + cfg0.node.rekey.after_messages = u64::MAX; + let mut cfg1 = crate::config::Config::new(); + cfg1.node.rekey.enabled = true; + cfg1.node.rekey.after_secs = u64::MAX; + cfg1.node.rekey.after_messages = u64::MAX; + + // Age both sessions past both rekey gates: node 0's jittered time trigger, + // and node 1's 30 s floor below which a msg1 is a duplicate, not a rekey. + let age = Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + let AgedLinkPair { + mut nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun0_rx, + tun1_rx, + } = aged_link_pair(cfg0, cfg1, age).await; let node0_idx_before = nodes[0].node.get_peer(&node1_addr).unwrap().our_index(); let node1_idx_before = nodes[1].node.get_peer(&node0_addr).unwrap().our_index(); @@ -2235,6 +2278,559 @@ async fn a_worker_decrypted_frame_on_the_previous_link_session_does_not_feed_mmp cleanup_nodes(&mut nodes).await; } +/// Where a captured link rekey msg1 is replayed from, once the cycle it +/// started has completed and the new link session is past the 30 s rekey +/// floor. +#[derive(Clone, Copy, Debug)] +enum ReplaySource { + /// No replay: the control the other two are measured against. + Nothing, + /// An address the node has no link with. + ThirdAddress, + /// The peer's own link address. + PeerAddress, +} + +/// The rekey that is attempted after the replay. +#[derive(Clone, Copy, Debug)] +enum ReplayProbe { + /// node 0, the peer, starts its next rekey. + PeerRekey, + /// node 1's own trigger, past its time threshold and its dampening. + OwnTrigger, +} + +/// What a replayed link msg1 left behind. +#[derive(Debug, PartialEq, Eq)] +struct ReplayOutcome { + /// node 1 holds a pending responder session after the replay. + armed_pending: bool, + /// Rekey msg2s node 1 sent to the replay's source address. + msg2_to_source: usize, + /// The probed rekey went ahead: node 0 completed on node 1's answer, or + /// node 1 started one of its own. + probe_proceeded: bool, +} + +/// The outcome when a replay changes nothing. +const REPLAY_HARMLESS: ReplayOutcome = ReplayOutcome { + armed_pending: false, + msg2_to_source: 0, + probe_proceeded: true, +}; + +/// Deliver `packets` to `node` as the network would have. +async fn dispatch_link_packets( + node: &mut TestNode, + packets: Vec, +) { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_ESTABLISHED, PHASE_MSG1, PHASE_MSG2}; + + for packet in packets { + match CommonPrefix::parse(&packet.data).map(|p| p.phase) { + Some(PHASE_MSG1) => node.node.handle_msg1(packet).await, + Some(PHASE_MSG2) => node.node.handle_msg2(packet).await, + Some(PHASE_ESTABLISHED) => node.node.handle_encrypted_frame(packet).await, + _ => {} + } + } +} + +/// Run `cycles` genuine FMP rekeys from node 0 to completion while keeping a +/// copy of the first one's msg1, age node 1's new link session past the 30 s +/// rekey floor, replay the copy to node 1 from `source`, then attempt the +/// rekey `probe` names. +/// +/// Both nodes rekey on time, so node 1 has a trigger of its own to probe; its +/// tick is never run before the probe, so it never initiates earlier. A third +/// node with no link to either stands in for the third address, so a msg2 sent +/// there is observable. +async fn replay_link_msg1_after_its_cycle( + source: ReplaySource, + probe: ReplayProbe, + cycles: usize, +) -> ReplayOutcome { + replay_link_msg1(source, probe, cycles, true).await +} + +/// [`replay_link_msg1_after_its_cycle`], with the replay sent either past the +/// 30 s floor (`past_floor`) or straight after the last cutover, while a +/// same-epoch msg1 is still taken as a duplicate of the link setup. +async fn replay_link_msg1( + source: ReplaySource, + probe: ReplayProbe, + cycles: usize, + past_floor: bool, +) -> ReplayOutcome { + use crate::proto::fmp::wire::{CommonPrefix, PHASE_MSG1, PHASE_MSG2}; + use crate::transport::ReceivedPacket; + + const REKEY_AFTER_SECS: u64 = 60; + let trigger_age = + Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + let config = || { + let mut config = crate::config::Config::new(); + config.node.rekey.enabled = true; + config.node.rekey.after_secs = REKEY_AFTER_SECS; + config.node.rekey.after_messages = u64::MAX; + config + }; + let AgedLinkPair { + mut nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun1_rx, + .. + } = aged_link_pair(config(), config(), trigger_age).await; + nodes.push(make_test_node_with_config(crate::config::Config::new(), 1280).await); + + // The genuine cycles, with a copy of the first one's msg1 kept on the way. + let mut captured = None; + for cycle in 0..cycles { + if cycle > 0 { + nodes[0] + .node + .get_peer_mut(&node1_addr) + .unwrap() + .test_backdate_session_established(trigger_age); + nodes[1] + .node + .get_peer_mut(&node0_addr) + .unwrap() + .test_backdate_session_established(Duration::from_secs(31)); + } + let node1_idx_before = nodes[1].node.get_peer(&node0_addr).unwrap().our_index(); + nodes[0].node.check_rekey().await; + let msg1 = nodes[1] + .packet_rx + .try_recv() + .expect("node 0's rekey msg1 must be queued at node 1"); + assert_eq!( + CommonPrefix::parse(&msg1.data).map(|p| p.phase), + Some(PHASE_MSG1), + "the captured packet must be node 0's rekey msg1" + ); + captured.get_or_insert_with(|| msg1.data.clone()); + nodes[1].node.handle_msg1(msg1).await; + pump_until_quiet(&mut nodes).await; + nodes[0].node.check_rekey().await; + let first = build_ipv6_packet(&fips0, &fips1, b"first frame on the new link epoch"); + nodes[0].node.handle_tun_outbound(first.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![first], + "node 0's first new-epoch frame must decode" + ); + let peer1 = nodes[1].node.get_peer(&node0_addr).unwrap(); + assert!( + peer1.pending_new_session().is_none() && peer1.our_index() != node1_idx_before, + "node 1 must have promoted the genuine cycle's session" + ); + } + let captured = captured.expect("at least one cycle must run"); + + // Past the floor below which a same-epoch msg1 is a duplicate. + if past_floor { + nodes[1] + .node + .get_peer_mut(&node0_addr) + .unwrap() + .test_backdate_session_established(Duration::from_secs(31)); + } + + let from = match source { + ReplaySource::Nothing => None, + ReplaySource::ThirdAddress => Some(2), + ReplaySource::PeerAddress => Some(0), + }; + let mut msg2_to_source = 0; + if let Some(from) = from { + let replay = ReceivedPacket::new(nodes[1].transport_id, nodes[from].addr.clone(), captured); + nodes[1].node.handle_msg1(replay).await; + tokio::time::sleep(Duration::from_millis(10)).await; + let reached: Vec = + std::iter::from_fn(|| nodes[from].packet_rx.try_recv().ok()).collect(); + msg2_to_source = reached + .iter() + .filter(|p| CommonPrefix::parse(&p.data).map(|p| p.phase) == Some(PHASE_MSG2)) + .count(); + dispatch_link_packets(&mut nodes[from], reached).await; + pump_until_quiet(&mut nodes).await; + } + let armed_pending = nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(); + + let probe_proceeded = match probe { + ReplayProbe::PeerRekey => { + nodes[0] + .node + .get_peer_mut(&node1_addr) + .unwrap() + .test_backdate_session_established(trigger_age); + nodes[0].node.check_rekey().await; + assert!( + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .rekey_in_progress(), + "node 0 must start its next rekey" + ); + pump_until_quiet(&mut nodes).await; + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .pending_new_session() + .is_some() + } + ReplayProbe::OwnTrigger => { + let peer = nodes[1].node.get_peer_mut(&node0_addr).unwrap(); + peer.test_backdate_session_established(trigger_age); + peer.backdate_dampener(Duration::from_secs(31)); + nodes[1].node.check_rekey().await; + nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .rekey_in_progress() + } + }; + + cleanup_nodes(&mut nodes).await; + ReplayOutcome { + armed_pending, + msg2_to_source, + probe_proceeded, + } +} + +/// The control for the peer's next rekey: with nothing replayed, the +/// procedure that measures the replays sees node 1 answer node 0's next rekey. +#[tokio::test] +async fn a_peers_next_link_rekey_is_answered_when_nothing_is_replayed() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::Nothing, ReplayProbe::PeerRekey, 1).await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// The control for the node's own trigger: with nothing replayed, node 1 +/// starts its own rekey once past its threshold and its dampening. +#[tokio::test] +async fn a_nodes_own_link_rekey_trigger_fires_when_nothing_is_replayed() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::Nothing, ReplayProbe::OwnTrigger, 1).await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 from a completed cycle, replayed from an address the node has +/// no link with, must not hold a pending that refuses the peer's next rekey. +#[tokio::test] +async fn a_link_msg1_replayed_from_a_third_address_does_not_block_the_peers_next_rekey() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::ThirdAddress, ReplayProbe::PeerRekey, 1) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 from a completed cycle, replayed from an address the node has +/// no link with, must not hold a pending that suppresses the node's own rekey. +#[tokio::test] +async fn a_link_msg1_replayed_from_a_third_address_does_not_suppress_the_nodes_own_rekey() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::ThirdAddress, ReplayProbe::OwnTrigger, 1) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 from a completed cycle, replayed from the peer's own address, +/// must not hold a pending that refuses the peer's next rekey. +#[tokio::test] +async fn a_link_msg1_replayed_from_the_peers_address_does_not_block_the_peers_next_rekey() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::PeerAddress, ReplayProbe::PeerRekey, 1) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 from a completed cycle, replayed from the peer's own address, +/// must not hold a pending that suppresses the node's own rekey. +#[tokio::test] +async fn a_link_msg1_replayed_from_the_peers_address_does_not_suppress_the_nodes_own_rekey() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::PeerAddress, ReplayProbe::OwnTrigger, 1) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 from an earlier cycle than the last one, replayed from an +/// address the node has no link with, is refused as well: the node keeps +/// the msg1s of its ended cycles, not only the latest. +#[tokio::test] +async fn a_link_msg1_replayed_from_an_earlier_cycle_does_not_block_the_peers_next_rekey() { + let outcome = + replay_link_msg1_after_its_cycle(ReplaySource::ThirdAddress, ReplayProbe::PeerRekey, 3) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A link msg1 replayed from an address the node has no link with, inside the +/// 30 s after a cutover, is taken as a duplicate of the link setup. The stored +/// setup msg2 it draws goes to the peer's established link, never to the +/// address the copy came from, and nothing else changes. +#[tokio::test] +async fn a_link_msg1_replayed_inside_the_rekey_floor_draws_nothing_to_its_source() { + let outcome = replay_link_msg1( + ReplaySource::ThirdAddress, + ReplayProbe::OwnTrigger, + 1, + false, + ) + .await; + assert_eq!(outcome, REPLAY_HARMLESS); +} + +/// A fresh rekey msg1 that arrives from an address other than the peer's is +/// answered on the peer's established link, not at the address it came from, +/// and the rekey completes there. +#[tokio::test] +async fn a_rekey_msg2_answers_on_the_peers_established_link_whatever_the_msg1_source() { + use crate::node::tests::spanning_tree::make_test_node; + use crate::proto::fmp::wire::{CommonPrefix, PHASE_MSG2}; + + const REKEY_AFTER_SECS: u64 = 60; + let trigger_age = + Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + let mut cfg0 = crate::config::Config::new(); + cfg0.node.rekey.after_secs = REKEY_AFTER_SECS; + cfg0.node.rekey.after_messages = u64::MAX; + let mut cfg1 = crate::config::Config::new(); + cfg1.node.rekey.after_secs = u64::MAX; + cfg1.node.rekey.after_messages = u64::MAX; + let AgedLinkPair { + mut nodes, + node0_addr, + node1_addr, + .. + } = aged_link_pair(cfg0, cfg1, trigger_age).await; + let mut third = vec![make_test_node().await]; + + nodes[0].node.check_rekey().await; + let mut msg1 = nodes[1] + .packet_rx + .try_recv() + .expect("node 0's rekey msg1 must be queued at node 1"); + msg1.remote_addr = third[0].addr.clone(); + nodes[1].node.handle_msg1(msg1).await; + + assert!( + third[0].packet_rx.try_recv().is_err(), + "nothing may be sent to the address the msg1 came from" + ); + let answered: Vec<_> = std::iter::from_fn(|| nodes[0].packet_rx.try_recv().ok()).collect(); + assert_eq!( + answered + .iter() + .map(|p| CommonPrefix::parse(&p.data).map(|p| p.phase)) + .collect::>(), + vec![Some(PHASE_MSG2)], + "the msg2 must go to node 0 on its established link" + ); + assert!( + nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 1 must hold the session it answered with" + ); + nodes[0] + .node + .handle_msg2(answered.into_iter().next().unwrap()) + .await; + assert!( + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .pending_new_session() + .is_some(), + "node 0 must complete its rekey on the answer" + ); + + cleanup_nodes(&mut nodes).await; + cleanup_nodes(&mut third).await; +} + +/// A rekey msg1 that arrives on a transport other than the peer's link is +/// answered on the link, and the pending session's index is registered under +/// the link's transport: the peer's frames on the new session arrive there, +/// and retirement removes the entry by the peer's transport. +#[tokio::test] +async fn a_rekey_answered_on_the_established_link_registers_its_index_on_that_transport() { + use crate::transport::TransportId; + + const REKEY_AFTER_SECS: u64 = 60; + let trigger_age = + Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + let mut cfg0 = crate::config::Config::new(); + cfg0.node.rekey.after_secs = REKEY_AFTER_SECS; + cfg0.node.rekey.after_messages = u64::MAX; + let mut cfg1 = crate::config::Config::new(); + cfg1.node.rekey.after_secs = u64::MAX; + cfg1.node.rekey.after_messages = u64::MAX; + let AgedLinkPair { + mut nodes, + node0_addr, + node1_addr, + fips0, + fips1, + tun1_rx, + .. + } = aged_link_pair(cfg0, cfg1, trigger_age).await; + let link_transport = nodes[1].transport_id; + let other_transport = TransportId::new(link_transport.as_u32() + 1); + + nodes[0].node.check_rekey().await; + let mut msg1 = nodes[1] + .packet_rx + .try_recv() + .expect("node 0's rekey msg1 must be queued at node 1"); + msg1.transport_id = other_transport; + nodes[1].node.handle_msg1(msg1).await; + let pending_idx = nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_our_index() + .expect("node 1 must answer the msg1 and hold its new session"); + assert!( + nodes[1] + .node + .peers_by_index + .contains_key(&(link_transport, pending_idx.as_u32())), + "the pending index must be registered under the link's transport" + ); + assert!( + !nodes[1] + .node + .peers_by_index + .contains_key(&(other_transport, pending_idx.as_u32())), + "the pending index must not be registered under the msg1's transport" + ); + + // node 0 completes and cuts over; its first new-epoch frame, on the link, + // must find node 1's pending and promote it. + pump_until_quiet(&mut nodes).await; + nodes[0].node.check_rekey().await; + let first = build_ipv6_packet(&fips0, &fips1, b"first frame on the new link epoch"); + nodes[0].node.handle_tun_outbound(first.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![first], + "node 0's first new-epoch frame must decode" + ); + assert_eq!( + nodes[1].node.get_peer(&node0_addr).unwrap().our_index(), + Some(pending_idx), + "node 1 must promote the pending on node 0's first new-epoch frame" + ); + assert!( + nodes[0] + .node + .get_peer(&node1_addr) + .unwrap() + .pending_new_session() + .is_none() + ); + + cleanup_nodes(&mut nodes).await; +} + +/// The record of answered msg1s lives with the peer, so a msg1 captured before +/// the peering was formed again, with the peer's epoch unchanged, is not +/// recognized. This node restarting reaches the same state, as does a link +/// torn down and re-formed. Kept as the record of that residual: it stays red +/// until a msg1 carries something that ties it to one cycle, which is a wire +/// change. +#[tokio::test] +#[ignore = "residual: a link msg1 captured before the peering was re-formed in the same peer epoch still arms a responder pending; the answered-msg1 record does not survive the peering"] +async fn a_link_msg1_captured_before_the_peering_was_re_formed_does_not_arm_a_pending() { + use crate::transport::ReceivedPacket; + + const REKEY_AFTER_SECS: u64 = 60; + let trigger_age = + Duration::from_secs(REKEY_AFTER_SECS + crate::node::REKEY_JITTER_SECS as u64 + 1); + let mut cfg0 = crate::config::Config::new(); + cfg0.node.rekey.after_secs = REKEY_AFTER_SECS; + cfg0.node.rekey.after_messages = u64::MAX; + let mut cfg1 = crate::config::Config::new(); + cfg1.node.rekey.after_secs = u64::MAX; + cfg1.node.rekey.after_messages = u64::MAX; + let AgedLinkPair { + mut nodes, + node0_addr, + node1_addr, + .. + } = aged_link_pair(cfg0, cfg1, trigger_age).await; + + // node 1 answers node 0's rekey; a copy of the msg1 is kept. + nodes[0].node.check_rekey().await; + let msg1 = nodes[1] + .packet_rx + .try_recv() + .expect("node 0's rekey msg1 must be queued at node 1"); + let captured = msg1.data.clone(); + nodes[1].node.handle_msg1(msg1).await; + assert!( + nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_new_session() + .is_some(), + "setup: node 1 must answer the genuine msg1" + ); + + // The peering is torn down and formed again; neither node restarts, so + // node 0's epoch, which the captured msg1 carries, is unchanged. + nodes[0].node.remove_active_peer(&node1_addr); + nodes[1].node.remove_active_peer(&node0_addr); + pump_until_quiet(&mut nodes).await; + initiate_handshake(&mut nodes, 0, 1).await; + drain_all_packets(&mut nodes, false).await; + nodes[1] + .node + .get_peer_mut(&node0_addr) + .expect("setup: the peering must be formed again") + .test_backdate_session_established(Duration::from_secs(31)); + + let replay = ReceivedPacket::new(nodes[1].transport_id, nodes[0].addr.clone(), captured); + nodes[1].node.handle_msg1(replay).await; + assert!( + nodes[1] + .node + .get_peer(&node0_addr) + .unwrap() + .pending_new_session() + .is_none(), + "a copy of a msg1 from before the peering was re-formed must not arm a pending" + ); + + 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/peer/active.rs b/src/peer/active.rs index d9ed1a1b..ae81d974 100644 --- a/src/peer/active.rs +++ b/src/peer/active.rs @@ -7,7 +7,7 @@ use crate::config::MmpConfig; use crate::node::REKEY_JITTER_SECS; use crate::noise::{HandshakeState as NoiseHandshakeState, NoiseError, NoiseSession}; use crate::proto::bloom::BloomFilter; -use crate::proto::fmp::{RekeyAnswer, RekeyRole}; +use crate::proto::fmp::{AnsweredMsg1s, Msg1Digest, RekeyAnswer, RekeyRole}; use crate::proto::mmp::MmpPeerState; use crate::proto::stp::{ParentDeclaration, TreeCoordinate}; use crate::transport::{LinkId, LinkStats, TransportAddr, TransportId}; @@ -293,11 +293,11 @@ pub struct ActivePeer { pending_role: Option, /// When the pending session was installed, for the responder hold. pending_since: Option, - /// The answer that armed a pending this node holds as the responder: the - /// msg1 it answered, by digest, and the msg2 it sent. Set with a responder - /// pending and cleared with it, so a resend of that msg1 can be answered - /// again while the pending is held. - rekey_answer: Option, + /// The rekey msg1s this node answered for this peer as the responder: the + /// whole answer that armed a pending it holds, so a resend of that msg1 + /// can be answered again, and the digests of ended cycles, so a copy of + /// one is refused. Not cycle state: it outlives every pending. + answered: AnsweredMsg1s, // === Published active-send-state (two-tier boundary) === /// The send-critical subset read (and, on roam/responder-cutover, written) @@ -342,7 +342,7 @@ impl ActivePeer { rekey_msg1_resend_count: 0, pending_role: None, pending_since: None, - rekey_answer: None, + answered: AnsweredMsg1s::default(), send: PeerSendState::new(link_id, now, authenticated_at), } } @@ -425,7 +425,7 @@ impl ActivePeer { rekey_msg1_resend_count: 0, pending_role: None, pending_since: None, - rekey_answer: None, + answered: AnsweredMsg1s::default(), send, } } @@ -1004,6 +1004,16 @@ impl ActivePeer { .map(|t| t.checked_sub(age).unwrap_or_else(Instant::now)); } + /// Test-only seam: backdate the last peer-initiated rekey so a test can + /// move past the rekey dampening window without waiting it out. Shifts + /// only the private timestamp; compiled out of release builds. + #[cfg(test)] + pub(crate) fn backdate_dampener(&mut self, age: Duration) { + self.last_peer_rekey = self + .last_peer_rekey + .map(|t| t.checked_sub(age).unwrap_or_else(Instant::now)); + } + /// Test-only seam: install link-layer MMP state with a chosen operating /// mode on a peer that was constructed without a Noise session (the bare /// `new` constructor leaves `mmp` as `None`). This only attaches the same @@ -1120,7 +1130,12 @@ impl ActivePeer { /// The answer that armed the pending session, when this node holds it as /// the rekey responder; `None` otherwise. pub(crate) fn rekey_answer(&self) -> Option<&RekeyAnswer> { - self.rekey_answer.as_ref() + self.answered.held() + } + + /// Whether `msg1` armed a responder cycle with this peer that has ended. + pub(crate) fn answered_before(&self, msg1: &Msg1Digest) -> bool { + self.answered.ended(msg1) } /// Check whether the pending session has been held for at least `hold` @@ -1142,7 +1157,7 @@ impl ActivePeer { their_index: SessionIndex, ) { self.install_pending(session, our_index, their_index, RekeyRole::Initiator); - self.rekey_answer = None; + self.answered.end(); } /// Store the session this node produced by answering the peer's rekey @@ -1160,16 +1175,16 @@ impl ActivePeer { answer: RekeyAnswer, ) { self.install_pending(session, our_index, their_index, RekeyRole::Responder); - self.rekey_answer = Some(answer); + self.answered.arm(answer); } - /// Clear what is recorded beside the pending slot: its role, its install - /// time, and the answer that armed it. Every path that empties the slot - /// calls this, so none of the three outlives the session it describes. + /// Clear what is recorded beside the pending slot: its role and install + /// time, and the answer that armed it, of which only the msg1 digest is + /// kept, as an ended cycle. Every path that empties the slot calls this. fn release_pending(&mut self) { self.pending_role = None; self.pending_since = None; - self.rekey_answer = None; + self.answered.end(); } /// Store a pending session with the role that produced it and the time it @@ -1966,8 +1981,9 @@ mod tests { } /// The answer that armed a responder pending is held exactly as long as - /// the pending: it leaves on retirement, on promotion, and on abandon, and - /// a pending this node initiated carries none. + /// the pending: it leaves on retirement, on promotion, and on abandon, + /// leaving its msg1 recorded as an ended cycle, and a pending this node + /// initiated carries none. #[test] fn the_answer_that_armed_a_pending_leaves_with_it() { type Exit = fn(&mut ActivePeer) -> Option; @@ -1997,6 +2013,10 @@ mod tests { peer.rekey_answer().is_none(), "{name}: the answer must leave with the pending" ); + assert!( + peer.answered_before(&answer().msg1), + "{name}: the answered msg1 must be remembered as an ended cycle" + ); } let (_cur_send, cur_recv) = ik_session_pair(); diff --git a/src/peer/machine.rs b/src/peer/machine.rs index 5bddce0e..a9b28ec9 100644 --- a/src/peer/machine.rs +++ b/src/peer/machine.rs @@ -1909,6 +1909,7 @@ mod tests { pending_new_session: false, rekey_in_progress: false, held_answer: None, + msg1_answered_before: false, existing_msg2: None, at_max_peers: false, has_pending_outbound_to_peer: false, diff --git a/src/proto/fmp/core.rs b/src/proto/fmp/core.rs index 738061af..fabcf780 100644 --- a/src/proto/fmp/core.rs +++ b/src/proto/fmp/core.rs @@ -21,6 +21,7 @@ use super::state::Fmp; use crate::transport::LinkId; use crate::utils::index::SessionIndex; use crate::{NodeAddr, PeerIdentity}; +use std::collections::VecDeque; /// Determine winner of cross-connection tie-breaker. /// @@ -172,6 +173,67 @@ pub(crate) struct RekeyAnswer { pub msg2: Vec, } +/// How many msg1s of ended cycles one peer's [`AnsweredMsg1s`] remembers. +/// +/// A msg1 is answered as a rekey only on a session at least +/// [`REKEY_MIN_SESSION_AGE_SECS`] old, so this covers at least two hours of +/// the peer's cycles, and eight and a half at the default 120 s interval when +/// the message-count trigger does not fire first. 32 bytes each, so 8 KiB per +/// peer at most. +pub(crate) const ENDED_MSG1_RECORD: usize = 256; + +/// The rekey msg1s this node answered as the link-rekey responder for one +/// peer, by digest. +/// +/// A link msg1 carries nothing that ties it to one cycle, so a copy taken off +/// the wire still authenticates as the peer, in the peer's current epoch, +/// after the cycle it started has ended. This record is how a copy is told +/// from a fresh msg1. The answer that armed the pending session this node +/// holds is kept whole, so a resend of that msg1 draws the same msg2. When +/// the pending leaves (adopted, retired or abandoned) its digest moves to the +/// ended list, and a msg1 matching an ended cycle is refused instead of arming +/// a new pending. +/// +/// Retention is bounded: the ended list keeps the last +/// [`ENDED_MSG1_RECORD`] digests, so a msg1 from an older cycle of the same +/// epoch is not recognized. The record lives with the peer, so it starts empty +/// whenever the peering is established again while the peer's epoch stays the +/// same, as after this node restarts or the link is torn down and re-formed. +#[derive(Debug, Default)] +pub(crate) struct AnsweredMsg1s { + held: Option, + ended: VecDeque, +} + +impl AnsweredMsg1s { + /// Record the answer that armed a new responder pending. + pub(crate) fn arm(&mut self, answer: RekeyAnswer) { + self.end(); + self.held = Some(answer); + } + + /// The pending the held answer armed has left: keep only its digest. + pub(crate) fn end(&mut self) { + if let Some(answer) = self.held.take() { + if self.ended.len() == ENDED_MSG1_RECORD { + self.ended.pop_front(); + } + self.ended.push_back(answer.msg1); + } + } + + /// The answer that armed the pending this node holds, if it holds one it + /// answered. + pub(crate) fn held(&self) -> Option<&RekeyAnswer> { + self.held.as_ref() + } + + /// Whether `msg1` armed a cycle that has since ended. + pub(crate) fn ended(&self, msg1: &Msg1Digest) -> bool { + self.ended.contains(msg1) + } +} + /// A snapshot of one active peer's rekey-relevant state, taken by the shell. /// /// Every clock read is resolved shell-side into a plain `u64`/`bool` before the @@ -287,6 +349,9 @@ pub(crate) struct EstablishSnapshot { /// responder, the answer that armed it. `None` with no pending, or with a /// pending this node initiated. pub held_answer: Option, + /// This msg1 armed a responder cycle of this peer's that has since ended + /// (pre-evaluated shell-side against the peer's [`AnsweredMsg1s`]). + pub msg1_answered_before: bool, /// The existing peer's stored msg2 wire bytes (an opaque blob), resent on a /// same-epoch duplicate msg1. `None` when there is no existing peer or it /// has no stored msg2. @@ -447,7 +512,7 @@ pub(crate) enum InboundDecision { } /// Why an inbound msg1 was rejected. Distinguishes only the diagnostic log -/// message; all three reject identically (BadState stat, rate-limiter complete, +/// message; all four reject identically (BadState stat, rate-limiter complete, /// the local not-yet-registered connection dropped). #[derive(Debug)] pub(crate) enum InboundReject { @@ -461,6 +526,9 @@ pub(crate) enum InboundReject { /// Dual rekey initiation and we are the tie-break *winner* (smaller /// NodeAddr): drop the peer's msg1 and keep driving our own rekey. DualRekeyWon, + /// The msg1 armed a rekey cycle with this peer that has already ended: a + /// copy, not a fresh request, and it must not arm a pending. + AnsweredBefore, } /// The classification outcome for one outbound `handle_msg2` completion, decided @@ -507,7 +575,9 @@ pub(crate) trait EstablishView { /// `peer_addr`: the existing peer's epoch/session/rekey state (with the /// session age resolved shell-side), the max-peers cap, and this node's own /// address for the tie-break. - fn establish_snapshot(&self, peer_addr: &NodeAddr) -> EstablishSnapshot; + /// `msg1` is the digest of the msg1 being classified, checked against the + /// peer's record of answered msg1s. + fn establish_snapshot(&self, peer_addr: &NodeAddr, msg1: &Msg1Digest) -> EstablishSnapshot; /// Snapshot the registry state relevant to classifying an outbound msg2 /// completion for `peer_addr`: whether the identity is already an active @@ -724,6 +794,14 @@ impl Fmp { reason: InboundReject::PendingSession, }; } + if snap.msg1_answered_before { + // A copy of a msg1 whose cycle has ended: refuse it + // before it can arm a pending, or, on a tie-break we + // lose, abandon our own rekey. + return InboundDecision::Reject { + reason: InboundReject::AnsweredBefore, + }; + } if snap.rekey_in_progress { // Dual initiation — smaller NodeAddr wins as initiator. // Our own rekey is the outbound/initiator side, so reuse diff --git a/src/proto/fmp/mod.rs b/src/proto/fmp/mod.rs index 85d485fd..9e3f8953 100644 --- a/src/proto/fmp/mod.rs +++ b/src/proto/fmp/mod.rs @@ -35,9 +35,9 @@ pub(crate) mod wire; mod tests; pub(crate) use core::{ - ConnAction, ConnSnapshot, EstablishSnapshot, EstablishView, InboundDecision, InboundReject, - LifecycleView, Msg1Digest, OutboundDecision, OutboundSnapshot, PeerSnapshot, RekeyAnswer, - RekeyCfg, RekeyResendSnapshot, RekeyRole, WireOutcome, + AnsweredMsg1s, ConnAction, ConnSnapshot, EstablishSnapshot, EstablishView, InboundDecision, + InboundReject, LifecycleView, Msg1Digest, OutboundDecision, OutboundSnapshot, PeerSnapshot, + RekeyAnswer, RekeyCfg, RekeyResendSnapshot, RekeyRole, WireOutcome, }; pub use core::{PromotionResult, cross_connection_winner}; pub(crate) use limits::backoff_ms; diff --git a/src/proto/fmp/tests/core.rs b/src/proto/fmp/tests/core.rs index 9de48c5b..435df23a 100644 --- a/src/proto/fmp/tests/core.rs +++ b/src/proto/fmp/tests/core.rs @@ -5,8 +5,9 @@ use super::util::{ wire_outcome, }; use crate::NodeAddr; +use crate::proto::fmp::core::ENDED_MSG1_RECORD; use crate::proto::fmp::{ - ConnAction, Fmp, InboundDecision, InboundReject, Msg1Digest, OutboundDecision, + AnsweredMsg1s, ConnAction, Fmp, InboundDecision, InboundReject, Msg1Digest, OutboundDecision, OutboundSnapshot, RekeyAnswer, RekeyCfg, RekeyRole, cross_connection_winner, }; use crate::testutil::make_node_addr; @@ -568,6 +569,91 @@ fn a_pending_this_node_initiated_answers_no_msg1() { )); } +#[test] +fn a_msg1_from_an_ended_cycle_is_refused_and_arms_nothing() { + // A copy of a msg1 whose cycle has ended must not arm a new pending, even + // with no pending held and the session aged past the rekey floor. + let fmp = Fmp::new(); + let mut snap = snapshot_holding_an_answer(b"msg1"); + snap.pending_new_session = false; + snap.held_answer = None; + snap.msg1_answered_before = true; + let wire = wire_outcome(Some([7u8; 8])); + assert!(matches!( + fmp.establish_inbound(&snap, &wire), + InboundDecision::Reject { + reason: InboundReject::AnsweredBefore + } + )); +} + +#[test] +fn a_msg1_from_an_ended_cycle_does_not_make_us_abandon_our_own_rekey() { + // Mid-rekey and on the losing side of the tie-break, a fresh msg1 makes + // us abandon ours and respond; a copy of an ended cycle's msg1 must not. + let fmp = Fmp::new(); + let mut snap = snapshot_holding_an_answer(b"msg1"); + snap.pending_new_session = false; + snap.held_answer = None; + snap.rekey_in_progress = true; + snap.our_node_addr = max_node_addr(); + let wire = wire_outcome(Some([7u8; 8])); + assert!(matches!( + fmp.establish_inbound(&snap, &wire), + InboundDecision::RekeyRespond { + abandon_first: true, + .. + } + )); + snap.msg1_answered_before = true; + assert!(matches!( + fmp.establish_inbound(&snap, &wire), + InboundDecision::Reject { + reason: InboundReject::AnsweredBefore + } + )); +} + +#[test] +fn the_answered_record_holds_the_armed_answer_then_remembers_its_msg1_once_ended() { + let mut record = AnsweredMsg1s::default(); + let first = Msg1Digest::of(b"first"); + record.arm(RekeyAnswer { + msg1: first, + msg2: vec![0x02; 4], + }); + assert_eq!(record.held().map(|a| a.msg1), Some(first)); + assert!(!record.ended(&first), "a held cycle has not ended"); + + record.end(); + assert!(record.held().is_none()); + assert!(record.ended(&first)); + assert!(!record.ended(&Msg1Digest::of(b"never answered"))); + + // Ending with nothing held records nothing. + record.end(); + assert!(record.ended(&first)); +} + +#[test] +fn the_answered_record_keeps_the_most_recent_ended_cycles_up_to_its_bound() { + let mut record = AnsweredMsg1s::default(); + let digest = |i: usize| Msg1Digest::of(&i.to_le_bytes()); + for i in 0..=ENDED_MSG1_RECORD { + record.arm(RekeyAnswer { + msg1: digest(i), + msg2: Vec::new(), + }); + record.end(); + } + assert!( + !record.ended(&digest(0)), + "the oldest beyond the bound is forgotten" + ); + assert!(record.ended(&digest(1))); + assert!(record.ended(&digest(ENDED_MSG1_RECORD))); +} + #[test] fn establish_inbound_dual_init_we_win_rejects() { // rekey in progress + our addr < peer addr (our = 0x10, peer = pubkey-derived diff --git a/src/proto/fmp/tests/util.rs b/src/proto/fmp/tests/util.rs index 9e50e57e..c791b34f 100644 --- a/src/proto/fmp/tests/util.rs +++ b/src/proto/fmp/tests/util.rs @@ -86,6 +86,7 @@ pub(super) fn establish_snapshot() -> EstablishSnapshot { pending_new_session: false, rekey_in_progress: false, held_answer: None, + msg1_answered_before: false, existing_msg2: None, at_max_peers: false, has_pending_outbound_to_peer: false,