fix(session): resend a lost initial msg3 until the responder is heard from

The initiator of a new session sent its XK msg3 once, marked the session
established and cleared its stored handshake payload. When that one
datagram was lost, the responder kept waiting for msg3 and dropped every
frame the initiator sent, and nothing sent msg3 again: the responder's
resent SessionAck reached an entry that was no longer initiating and was
refused as a bad-state reject, and the handshake resend sweep skips
established sessions. The rekey arm of the same handler already keeps its
msg3 for this reason; the first-contact arm never did.

Keep the encoded msg3 in the entry's handshake resend slot at
establishment instead of clearing it, and resend it from the session
handshake sweep through a new pure decision in the FSP core: resend when
due and within handshake_max_resends, and stop retaining it once the
budget is spent. The first inbound frame that authenticates on the session
releases it too, which on a healthy session is the responder's first
frame. The sweep's existing passes, including its timeout removal, still
skip established sessions, and the resend runs whether or not periodic
rekey is enabled.

Before this change the new two-node test fails because nothing is
retained and nothing resends it, so the responder never leaves its
waiting state; the test does not exercise the SessionAck handler's
bad-state arm, which is unchanged. No wire format change: the resend
carries the same msg3 bytes in a fresh datagram, and a responder that
already completed the session refuses the duplicate as before.
This commit is contained in:
Johnathan Corgan
2026-09-22 18:48:53 +00:00
parent f8f37b9e6c
commit 34f97d27b6
8 changed files with 566 additions and 13 deletions
+14
View File
@@ -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
+16 -2
View File
@@ -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);
+95
View File
@@ -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<InitialMsg3ResendSnapshot> {
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.
+9 -3
View File
@@ -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<Vec<u8>>,
/// 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<u8>, 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;
+296
View File
@@ -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<TestNode> {
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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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<Vec<u8>> = 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.
///
+49 -2
View File
@@ -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<InitialMsg3ResendSnapshot>,
max_resends: u32,
) -> Vec<FspAction> {
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.
+5 -3
View File
@@ -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,
+82 -3
View File
@@ -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]