Files
fips/src/node/handlers/session.rs
T
Johnathan Corgan 0e4bcac835 Merge branch 'maint' into master
# Conflicts:
#	CHANGELOG.md
#	src/transport/tcp/mod.rs
2026-10-06 01:14:50 +00:00

3422 lines
150 KiB
Rust

//! End-to-end session message handlers.
//!
//! Handles locally-delivered session payloads from SessionDatagram envelopes.
//! Dispatches based on FSP common prefix phase to specific handlers for
//! SessionSetup (Noise XK msg1), SessionAck (msg2), SessionMsg3 (msg3),
//! encrypted data, and error signals (CoordsRequired, PathBroken).
use crate::NodeAddr;
use crate::ipv6tun::icmp::IcmpContext;
use crate::ipv6tun::outbound::{Mesh, Outcome, Route};
use crate::node::handlers::mmp::format_throughput;
use crate::node::rate_limit::Msg1Class;
use crate::node::reject::{RejectReason, SessionReject};
use crate::node::session::{EndToEndState, EpochSlot, SessionEntry};
use crate::node::{Node, NodeError};
use crate::noise::{
HandshakeState, XK_HANDSHAKE_MSG1_SIZE, XK_HANDSHAKE_MSG2_SIZE, XK_HANDSHAKE_MSG3_SIZE,
};
#[cfg(unix)]
use crate::proto::fmp::wire::{
ESTABLISHED_HEADER_SIZE, FLAG_KEY_EPOCH, FLAG_SP, build_established_header,
};
use crate::proto::fsp::quorum::QuorumVerdict;
use crate::proto::fsp::wire::{
FSP_COMMON_PREFIX_SIZE, FSP_FLAG_CP, FSP_FLAG_K, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED,
FSP_PHASE_MSG1, FSP_PHASE_MSG2, FSP_PHASE_MSG3, FSP_PORT_HEADER_SIZE, FSP_PORT_IPV6_SHIM,
FspCommonPrefix, FspEncryptedHeader, build_fsp_header, fsp_prepend_inner_header,
fsp_strip_inner_header, parse_encrypted_coords,
};
use crate::proto::fsp::{
DecryptSlot, EpochReaction, FspAction, FspInnerFlags, SessionAck, SessionMessageType,
SessionMsg3, SessionSetup, mark_ipv6_ecn_ce,
};
#[cfg(unix)]
use crate::proto::link::LinkMessageType;
#[cfg(unix)]
use crate::proto::link::SESSION_DATAGRAM_HEADER_SIZE;
use crate::proto::link::SessionDatagram;
use crate::proto::mmp::{
BackoffUpdate, MmpAction, MmpSessionState, PathMtuNotification, ReceiverReport, SendResult,
SessionReceiverReport, SessionReportKind, SessionReportSnapshot, SessionSenderReport,
};
use crate::proto::mmp::{MAX_SESSION_REPORT_INTERVAL_MS, MIN_SESSION_REPORT_INTERVAL_MS};
use crate::proto::routing::{CoordsRequired, MtuExceeded, PathBroken, RoutingSignalType};
use crate::proto::stp::{coords_wire_size, encode_coords};
#[cfg(unix)]
use crate::transport::TransportHandle;
use crate::upper::icmp::FIPS_OVERHEAD;
use secp256k1::PublicKey;
use tracing::{debug, info, trace, warn};
/// Minimum interval between path-MTU releases driven by `PathBroken` for one
/// destination.
///
/// `PathBroken` is unauthenticated, so a release is a remote party's claim
/// that the path a tightened MTU described is gone. Without an interval the
/// claim can be repeated at line rate, discarding a genuinely learned
/// bottleneck as fast as it is relearned. Raising it defers a legitimate
/// release after a second real break, which costs throughput on the new path
/// but never a blackhole, since the deferred value is the tighter one.
pub(in crate::node) const PATH_MTU_RELEASE_MIN_INTERVAL: std::time::Duration =
std::time::Duration::from_millis(1000);
/// Bytes the link layer adds to an encoded `SessionDatagram` on its way to the
/// wire: the established FMP header, the 4-byte session-relative timestamp and
/// the AEAD tag. Mirrors the buffer `send_encrypted_link_message_with_ce`
/// builds.
///
/// Spelled out in full rather than through the `crate::proto::fmp::wire`
/// import above, which is `#[cfg(unix)]`. This constant feeds `link_wire_len`,
/// whose caller `send_session_datagram` is compiled on every platform, so
/// taking the name from that import fails to build on Windows.
const LINK_FRAME_OVERHEAD: usize =
crate::proto::fmp::wire::ESTABLISHED_HEADER_SIZE + 4 + crate::noise::TAG_SIZE;
/// Wire size of an encoded `SessionDatagram` of `encoded_len` bytes.
fn link_wire_len(encoded_len: usize) -> usize {
encoded_len + LINK_FRAME_OVERHEAD
}
/// Divisor giving the share of the session table that unauthenticated
/// half-open entries may hold, as `max_sessions / DIVISOR`.
///
/// Two means a reconnect storm, where every peer that had a session
/// initiates at once after a restart or a healed partition, still fits in
/// half the table; a tighter share bites four times sooner and is felt by
/// a hub before it is felt by an attacker. Half-open entries are reaped
/// after `handshake_timeout_secs` while established ones survive
/// `idle_timeout_secs`, so they turn over faster than the share suggests.
/// Lowering the divisor raises the share, which lets a handshake flood
/// crowd out peers that complete; raising it refuses legitimate initiators
/// sooner in a storm.
const HALF_OPEN_SHARE_DIVISOR: usize = 2;
/// Inputs to `try_send_session_data_pipelined` — the FSP+FMP pipelined
/// fast path that hands both AEAD operations to the encrypt worker
/// in a single dispatch.
#[cfg(unix)]
struct PipelinedSend<'a> {
dest_addr: &'a NodeAddr,
payload: &'a [u8],
now_ms: u64,
timestamp: u32,
fsp_flags: u8,
inner_plaintext: &'a [u8],
my_coords: Option<&'a crate::proto::stp::TreeCoordinate>,
dest_coords: Option<&'a crate::proto::stp::TreeCoordinate>,
}
/// Outcome of the routing-signal admission test.
///
/// `Unbound` and `Forged` are both refusals, kept apart because they mean
/// different things to an operator. `Unbound` is consistent with a benign
/// race — a signal arriving just after a local session teardown. `Forged`
/// is not consistent with any honest emitter, so it is the sharper
/// indicator and is counted separately.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum SignalVerdict {
/// The named destination is an address this node bound itself.
Admit,
/// The src/dest pairing is structurally impossible for a legitimate
/// emitter.
Forged,
/// No qualifying session entry exists for the named destination.
Unbound,
}
impl SignalVerdict {
/// Short stable label for the `verdict` log field.
fn label(self) -> &'static str {
match self {
Self::Admit => "admit",
Self::Forged => "forged",
Self::Unbound => "unbound",
}
}
}
impl Node {
/// Handle a locally-delivered session datagram payload.
///
/// Called from `handle_session_datagram()` when `dest_addr == self.node_addr()`.
/// Dispatches based on the 4-byte FSP common prefix:
///
/// - Phase 0x1 → SessionSetup (handshake msg1)
/// - Phase 0x2 → SessionAck (handshake msg2)
/// - Phase 0x3 → SessionMsg3 (XK handshake msg3)
/// - Phase 0x0 + U flag → plaintext error signal (CoordsRequired/PathBroken)
/// - Phase 0x0 + !U → encrypted session message (data, reports, etc.)
///
/// `src_addr` is the datagram's claimed source, an envelope field the
/// sender chooses. `link_peer` is the authenticated FMP peer the datagram
/// arrived over, and is the only identity on this path worth keying a
/// limiter on.
pub(in crate::node) async fn handle_session_payload(
&mut self,
src_addr: &NodeAddr,
link_peer: &NodeAddr,
payload: &[u8],
path_mtu: u16,
ce_flag: bool,
) {
let prefix = match FspCommonPrefix::parse(payload) {
Some(p) => p,
None => {
debug!(
len = payload.len(),
"Session payload too short for FSP prefix"
);
return;
}
};
let inner = &payload[FSP_COMMON_PREFIX_SIZE..];
match prefix.phase {
FSP_PHASE_MSG1 => {
self.handle_session_setup(src_addr, link_peer, inner).await;
}
FSP_PHASE_MSG2 => {
self.handle_session_ack(src_addr, inner).await;
}
FSP_PHASE_MSG3 => {
self.handle_session_msg3(src_addr, inner).await;
}
FSP_PHASE_ESTABLISHED if prefix.is_unencrypted() => {
// Plaintext error signals: read msg_type from first byte after prefix
if inner.is_empty() {
debug!("Empty plaintext error signal");
return;
}
let error_type = inner[0];
let error_body = &inner[1..];
match RoutingSignalType::from_byte(error_type) {
Some(RoutingSignalType::CoordsRequired) => {
self.handle_coords_required(src_addr, error_body).await;
}
Some(RoutingSignalType::PathBroken) => {
self.handle_path_broken(src_addr, link_peer, error_body)
.await;
}
Some(RoutingSignalType::MtuExceeded) => {
self.handle_mtu_exceeded(src_addr, error_body).await;
}
_ => {
debug!(error_type, "Unknown plaintext error signal type");
}
}
}
FSP_PHASE_ESTABLISHED => {
self.handle_encrypted_session_msg(src_addr, payload, path_mtu, ce_flag)
.await;
}
_ => {
debug!(phase = prefix.phase, "Unknown FSP phase");
}
}
}
/// Handle an encrypted session message (phase 0x0, U flag clear).
///
/// Full FSP receive pipeline:
/// 1. Parse FspEncryptedHeader (12 bytes) → counter, flags, header_bytes
/// 2. If CP flag: parse cleartext coords, cache them
/// 3. Session lookup (must be Established)
/// 4. AEAD decrypt with AAD = header_bytes
/// 5. Strip FSP inner header → timestamp, msg_type, inner_flags
/// 6. Dispatch by msg_type
async fn handle_encrypted_session_msg(
&mut self,
src_addr: &NodeAddr,
payload: &[u8],
path_mtu: u16,
ce_flag: bool,
) {
// Parse the 12-byte encrypted header (includes the 4-byte prefix)
let header = match FspEncryptedHeader::parse(payload) {
Some(h) => h,
None => {
debug!(
len = payload.len(),
"Encrypted session message too short for FSP header"
);
return;
}
};
// Determine where ciphertext starts (after header, optionally after coords)
let mut ciphertext_offset = FSP_HEADER_SIZE;
// If CP flag set, parse cleartext coords between header and ciphertext
if header.has_coords() {
let coord_data = &payload[FSP_HEADER_SIZE..];
match parse_encrypted_coords(coord_data) {
Ok((src_coords, dest_coords, bytes_consumed)) => {
let now_ms = Self::now_ms();
let my_addr = *self.node_addr();
for action in
self.fsp
.plan_cache_coords(*src_addr, my_addr, src_coords, dest_coords)
{
if let FspAction::CacheCoords { addr, coords } = action {
self.insert_coord_hint(addr, coords, now_ms);
}
}
ciphertext_offset += bytes_consumed;
}
Err(e) => {
debug!(error = %e, "Failed to parse coords from encrypted session message");
return;
}
}
}
let ciphertext = &payload[ciphertext_offset..];
// Look up session entry — must be Established to decrypt
{
let entry = match self.sessions.get(src_addr) {
Some(e) => e,
None => {
debug!(src = %self.peer_display_name(src_addr), "Encrypted session message for unknown session");
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
};
// Drop encrypted data if session is not yet established.
// With XK, the responder must wait for msg3 before it can decrypt.
if !entry.is_established() {
debug!(
src = %self.peer_display_name(src_addr),
"Encrypted message but session not established (awaiting handshake completion)"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::BadState));
return;
}
}
// The received K-bit is only an ordering hint for the
// trial-decrypt cascade — it picks which key epoch to try first.
// Correctness never depends on it; promotion is driven by which
// slot actually authenticates the frame, not by the header bit.
let received_k_bit = header.flags & FSP_FLAG_K != 0;
let mut entry = match self.sessions.remove(src_addr) {
Some(e) => e,
None => return,
};
let now_ms = Self::now_ms();
// Overlapping-epoch trial-decrypt: try current, pending and
// previous so any epoch the peer might have sealed this frame in
// can be decrypted. This makes rekey correctness independent of
// cutover timing — no ordering and no reordering can cause a
// decrypt failure. A successful `previous`-slot decrypt also
// refreshes the drain deadline so the old epoch is retained as
// long as the peer keeps using it.
let (plaintext, slot) = match entry.fsp_trial_decrypt(
ciphertext,
header.counter,
&header.header_bytes,
received_k_bit,
now_ms,
) {
Some(result) => result,
None => {
// Every live slot failed — a genuine drop. The upper
// layer retransmits.
debug!(
src = %self.peer_display_name(src_addr),
counter = header.counter,
"Session AEAD decryption failed (all epochs)"
);
self.sessions.insert(*src_addr, entry);
return;
}
};
// 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
// `SessionEntry` mutation.
let decrypt_slot = match slot {
EpochSlot::Current => DecryptSlot::Current,
EpochSlot::Pending => DecryptSlot::Pending,
EpochSlot::Previous => DecryptSlot::Previous,
};
match self.fsp.classify_epoch(
decrypt_slot,
entry.rekey_msg3_payload().is_some(),
entry.pending_new_session().is_some(),
) {
EpochReaction::PromoteConfirming => {
// A frame that authenticates against `pending` is itself the
// cutover signal — proof the peer derived the new session and
// moved to it. The peer received msg3, so confirm it on the new
// epoch (stop retransmitting) before `handle_peer_kbit_flip`
// consumes the pending session, then promote.
info!(
peer = %self.peer_display_name(src_addr),
"Peer FSP new-epoch frame authenticated, FSP rekey cutover complete, promoting new session"
);
entry.confirm_peer_new_epoch();
entry.handle_peer_kbit_flip(now_ms);
}
EpochReaction::Promote => {
// Promote now: current → previous, pending → current, flip the
// K-bit. The header K-bit is only a hint; the authenticated
// decrypt is the gating event.
info!(
peer = %self.peer_display_name(src_addr),
"Peer FSP new-epoch frame authenticated, FSP rekey cutover complete, promoting new session"
);
entry.handle_peer_kbit_flip(now_ms);
}
EpochReaction::ConfirmResponder => {
// We are the rekey initiator that already cut over on its own
// timer: `current` is now the new epoch, so a frame decrypting
// against it confirms the responder reached it. Stop
// retransmitting msg3.
entry.confirm_peer_new_epoch();
}
EpochReaction::None => {
// Steady-state `current`, or an old-epoch `previous` straggler:
// `fsp_trial_decrypt` already refreshed the drain deadline so
// the `previous` slot is not retired while the peer keeps using
// it — no further state change, just deliver.
}
}
self.sessions.insert(*src_addr, entry);
// Strip FSP inner header (6 bytes)
let (timestamp, msg_type, inner_flags_byte, rest) = match fsp_strip_inner_header(&plaintext)
{
Some(parts) => parts,
None => {
debug!(src = %self.peer_display_name(src_addr), "Decrypted payload too short for FSP inner header");
return;
}
};
// MMP per-message recording on RX path
if let Some(entry) = self.sessions.get_mut(src_addr)
&& let Some(mmp) = entry.mmp_mut()
{
let now_ms = crate::time::mono_ms();
mmp.receiver
.record_recv(header.counter, timestamp, plaintext.len(), ce_flag, now_ms);
// Spin bit: advance state machine for correct TX reflection.
// RTT samples not fed into SRTT — timestamp-echo provides
// accurate RTT; spin bit includes variable inter-frame delays.
let inner_flags = FspInnerFlags::from_byte(inner_flags_byte);
let _spin_rtt = mmp
.spin_bit
.rx_observe(inner_flags.spin_bit, header.counter, now_ms);
}
// Feed path_mtu from datagram envelope to MMP path MTU tracking.
// Done for ALL session messages, not just DataPackets, so the
// destination learns the path MTU even when only reports flow.
if let Some(entry) = self.sessions.get_mut(src_addr)
&& let Some(mmp) = entry.mmp_mut()
{
mmp.path_mtu.observe_incoming_mtu(path_mtu);
}
// Dispatch by msg_type
match SessionMessageType::from_byte(msg_type) {
Some(SessionMessageType::DataPacket) => {
// msg_type 0x10: port-multiplexed service dispatch
if rest.len() < FSP_PORT_HEADER_SIZE {
debug!(len = rest.len(), "DataPacket too short for port header");
return;
}
let src_port = u16::from_le_bytes([rest[0], rest[1]]);
let dst_port = u16::from_le_bytes([rest[2], rest[3]]);
let service_payload = &rest[FSP_PORT_HEADER_SIZE..];
match dst_port {
FSP_PORT_IPV6_SHIM => {
use crate::FipsAddress;
let src_ipv6 = FipsAddress::from_node_addr(src_addr).to_ipv6().octets();
let dst_ipv6 = FipsAddress::from_node_addr(self.node_addr())
.to_ipv6()
.octets();
match crate::upper::ipv6_shim::decompress_ipv6(
service_payload,
src_ipv6,
dst_ipv6,
) {
Some(mut packet) => {
if ce_flag {
mark_ipv6_ecn_ce(&mut packet);
self.metrics().congestion.ce_received.inc();
}
if let Some(tun_tx) = &self.supervisor.ipv6tun.tun_tx {
if let Err(e) = tun_tx.send(packet) {
debug!(error = %e, "Failed to deliver decompressed IPv6 packet to TUN");
}
} else {
trace!(
src = %self.peer_display_name(src_addr),
"IPv6 shim packet decompressed (no TUN interface)"
);
}
}
None => {
debug!(
src = %self.peer_display_name(src_addr),
len = service_payload.len(),
"IPv6 shim decompression failed"
);
}
}
}
_ => {
// Every other port belongs to the native datagram API,
// which decides between an established flow, a listener
// and a drop. The same decision serves the debug
// arrival command, so the rule was exercised before
// this call site existed.
// The peer's key comes from the session entry rather
// than from the wire: the handler above refuses
// anything but an Established session, and an
// Established entry holds the key its handshake
// attested. Nothing has to invert the node address.
let payload = service_payload.to_vec();
// The entry was removed for the trial-decrypt cascade
// and re-inserted above, and this arm is reached only
// for an Established session, so the lookup finds one.
// It is written as a lookup rather than an unwrap
// because a future path that skipped the re-insert
// would otherwise panic on a peer's datagram. Skipping
// is scoped to the native dispatch alone: the idle
// timer and the pending-outbound flush at the end of
// this function are this message's bookkeeping and are
// owed whether or not it had anywhere to go.
if let Some(peer_key) = self
.sessions
.get(src_addr)
.map(|session| session.remote_pubkey().x_only_public_key().0)
{
let outcome = self
.native_deliver(*src_addr, peer_key, src_port, dst_port, payload);
if let crate::native::link::Outcome::Dropped(why) = outcome {
debug!(
src = %self.peer_display_name(src_addr),
dst_port,
why = why.as_str(),
"Unknown FSP service port, dropping DataPacket"
);
}
}
}
}
}
Some(SessionMessageType::SenderReport) => {
self.handle_session_sender_report(src_addr, rest);
}
Some(SessionMessageType::ReceiverReport) => {
self.handle_session_receiver_report(src_addr, rest);
}
Some(SessionMessageType::PathMtuNotification) => {
self.handle_session_path_mtu_notification(src_addr, rest);
}
Some(SessionMessageType::CoordsWarmup) => {
// Standalone coordinate warming — coords already extracted
// from CP flag by transit nodes. No action needed at endpoint.
trace!(src = %self.peer_display_name(src_addr), "CoordsWarmup received");
}
_ => {
debug!(src = %self.peer_display_name(src_addr), msg_type, "Unknown session message type, dropping");
}
}
// Only application data resets the idle timer and traffic counters —
// MMP reports (SenderReport, ReceiverReport, PathMtuNotification) do not.
if msg_type == SessionMessageType::DataPacket.to_byte()
&& let Some(entry) = self.sessions.get_mut(src_addr)
{
entry.record_recv(rest.len());
entry.touch(Self::now_ms());
}
// Flush any pending outbound packets (e.g., simultaneous initiation
// where responder also had queued outbound packets)
self.flush_pending_packets(src_addr).await;
self.flush_pending_native(src_addr).await;
}
/// Handle an incoming SessionSetup (Noise XK msg1).
///
/// The remote node wants to establish an end-to-end session with us.
/// We create an XK responder handshake, process msg1, send SessionAck with msg2,
/// and transition to AwaitingMsg3.
async fn handle_session_setup(
&mut self,
src_addr: &NodeAddr,
link_peer: &NodeAddr,
inner: &[u8],
) {
let setup = match SessionSetup::decode(inner) {
Ok(s) => s,
Err(e) => {
debug!(error = %e, "Malformed SessionSetup");
return;
}
};
if setup.handshake_payload.len() != XK_HANDSHAKE_MSG1_SIZE {
debug!(
len = setup.handshake_payload.len(),
expected = XK_HANDSHAKE_MSG1_SIZE,
"Invalid handshake payload size in SessionSetup"
);
return;
}
// Meter the setup before anything is spent on it. This sits ahead of
// every `send_session_datagram` call in the handler — the duplicate
// ack resend, the rekey ack and the fresh-setup ack alike — so a
// refused msg1 emits nothing at all, which is what bounds the ack
// amplification. It also precedes both responder handshake
// constructions, so a refusal costs no crypto. Moving it below the
// existing-entry lookup would leave the duplicate resend unmetered.
//
// The class is read from the session table but the *key* is the link
// peer: `src_addr` is chosen by the sender, so keying on it would be
// no limit at all. A setup naming an established peer cannot grow the
// table and is metered separately, so that a stranger flood over a
// shared link cannot stop that peer's rekey from arming.
// Population cap, ahead of the limiter so a full table costs no
// token, no responder handshake and no ack. The predicate is "would
// admitting this grow the table", not "is this a stranger": `class`
// is Stranger for an existing Initiating or AwaitingMsg3 entry too,
// and refusing those would break in-flight legitimate handshakes and
// the duplicate-ack resend. Same shape as the pending-destination cap
// in `queue_pending_packet`. Refuse rather than evict: msg1 is
// unauthenticated here, so evicting would hand a stranger a teardown
// primitive it does not have.
if !self.admit_new_session(src_addr) {
return;
}
let class = if self
.sessions
.get(src_addr)
.is_some_and(|e| e.is_established())
{
Msg1Class::EstablishedLink
} else {
Msg1Class::Stranger
};
if !self.setup_rate_limiter.try_admit(link_peer, class) {
debug!(
link_peer = %self.peer_display_name(link_peer),
src = %self.peer_display_name(src_addr),
?class,
"SessionSetup rate limited"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::SetupRateLimited));
return;
}
// Check for existing session with this remote
if let Some(existing) = self.sessions.get(src_addr) {
if existing.is_initiating() {
// Simultaneous initiation: smaller NodeAddr wins as initiator
if crate::proto::fsp::initiation_winner(self.identity().node_addr(), src_addr) {
// We win — drop their setup, they'll process ours
debug!(
src = %self.peer_display_name(src_addr),
"Simultaneous session initiation: we win (smaller addr), dropping their setup"
);
return;
}
// We lose — discard our pending handshake, become responder below
debug!(
src = %self.peer_display_name(src_addr),
"Simultaneous session initiation: we lose, becoming responder"
);
} else if existing.is_awaiting_msg3() {
// Duplicate setup while we already sent msg2 — resend stored ack
if let Some(payload) = existing.handshake_payload() {
debug!(src = %self.peer_display_name(src_addr), "Duplicate SessionSetup, resending SessionAck");
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *src_addr, payload.to_vec())
.with_ttl(self.config().node.session.default_ttl);
if let Err(e) = self.send_session_datagram(&mut datagram).await {
debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to resend SessionAck");
}
} else {
debug!(src = %self.peer_display_name(src_addr), "Duplicate SessionSetup, no stored ack to resend");
}
return;
} else if existing.is_established() {
// A SessionSetup naming an already-established peer is
// unauthenticated: msg1 is a bare ephemeral and the source
// address is an envelope field, so anyone able to reach us
// can claim it. It may therefore only arm a handshake
// alongside the running session, never replace it. The new
// keys are adopted in `handle_session_msg3` and only when the
// authenticated static key matches the key this session was
// opened with, and the cut-over waits for a frame that
// authenticates against the pending epoch. Every exit below
// returns, so an established entry never reaches the
// re-establishment path that replaces it.
let rekey_in_progress = existing.has_rekey_in_progress();
// A completed rekey outranks a fresh setup message while the
// cut-over it is waiting for can still arrive. Once it has
// waited a full idle timeout for a peer that never appeared
// on the new epoch, it stops vetoing: a peer that restarted,
// or one whose own cycle lapsed, would otherwise be refused
// for as long as our own sends kept the session from idling
// out. The pending keys survive both outcomes of this test:
// the veto returns without touching them, and the
// fall-through only arms a handshake beside them. Adopting
// them is still an authenticated msg3's job alone. No arm
// below drops one either.
let pending_outranks = existing.pending_new_session().is_some()
&& !existing.pending_stale(
Self::now_ms(),
self.config().node.session.idle_timeout_secs * 1000,
);
// A handshake the peer armed is held until its msg3 or its
// timeout. Once the peer has read our SessionAck it holds the
// new keys, so this handshake is the only one its msg3 can
// complete, and discarding it for an unauthenticated setup
// splits the session's epochs. A genuine retry that meets a
// stale handshake here completes one handshake timeout
// later, once the handshake has expired.
if rekey_in_progress && !existing.is_rekey_initiator() {
debug!(
src = %self.peer_display_name(src_addr),
"FSP rekey msg1 received while the peer's handshake awaits msg3, dropping"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyHeld));
return;
}
// Dual-initiation detection: both sides sent SessionSetup
// simultaneously. Apply tie-breaker — smaller NodeAddr
// wins as initiator (same as initial session setup).
if rekey_in_progress {
if crate::proto::fsp::initiation_winner(self.identity().node_addr(), src_addr) {
// We win as initiator — drop their msg1.
debug!(
src = %self.peer_display_name(src_addr),
"Dual FSP rekey initiation: we win (smaller addr), dropping their msg1"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyTiebreak));
return;
}
// We lose — abandon our armed handshake, become
// responder below.
//
// `abandon_handshake`, not `abandon_rekey`, although the
// arm above leaves only handshakes this node initiated,
// and an initiator handshake never sits beside a pending
// session (see the SessionAck handler). Only the handshake
// is ours to discard, and discarding it costs nothing:
// the peer has not answered it, so it holds no key
// material either endpoint is using.
debug!(
src = %self.peer_display_name(src_addr),
"Dual FSP rekey initiation: we lose (larger addr), abandoning ours"
);
let entry = self.sessions.get_mut(src_addr).unwrap();
entry.abandon_handshake();
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyYielded));
} else if pending_outranks {
// Guard: already have a pending session waiting for K-bit cutover
debug!(
src = %self.peer_display_name(src_addr),
"FSP rekey msg1 received but already have pending session, dropping"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyPending));
return;
}
// This frame's own copy of the node's long-term private key;
// the handshake state keeps its own and clears that on drop.
let mut our_keypair = self.identity().keypair();
let mut handshake = HandshakeState::new_xk_responder(our_keypair);
our_keypair.non_secure_erase();
handshake.set_local_epoch(self.startup_epoch());
if let Err(e) = handshake.read_xk_message_1(&setup.handshake_payload) {
debug!(
src = %self.peer_display_name(src_addr),
error = %e,
"Failed to process rekey XK msg1"
);
return;
}
// Generate msg2
let msg2 = match handshake.write_xk_message_2() {
Ok(m) => m,
Err(e) => {
debug!(
src = %self.peer_display_name(src_addr),
error = %e,
"Failed to generate rekey XK msg2"
);
return;
}
};
// Build and send SessionAck
let our_coords = self.tree_state.my_coords().clone();
let ack = SessionAck::new(our_coords, setup.src_coords).with_handshake(msg2);
let ack_payload = ack.encode();
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *src_addr, ack_payload)
.with_ttl(self.config().node.session.default_ttl);
if let Err(e) = self.send_session_datagram(&mut datagram).await {
debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send rekey SessionAck");
return;
}
// Store rekey state on the existing entry
let now_ms = Self::now_ms();
let entry = self.sessions.get_mut(src_addr).unwrap();
entry.set_rekey_state(handshake, false);
entry.record_peer_rekey(now_ms);
self.stats_mut().session.rekey_armed += 1;
debug!(
src = %self.peer_display_name(src_addr),
"FSP rekey: processed peer's msg1, sent msg2, awaiting msg3"
);
return;
}
}
// Create XK responder handshake and process msg1
// This frame's own copy of the node's long-term private key; the
// handshake state keeps its own and clears that on drop.
let mut our_keypair = self.identity().keypair();
let mut handshake = HandshakeState::new_xk_responder(our_keypair);
our_keypair.non_secure_erase();
handshake.set_local_epoch(self.startup_epoch());
if let Err(e) = handshake.read_xk_message_1(&setup.handshake_payload) {
debug!(error = %e, "Failed to process Noise XK msg1 in SessionSetup");
return;
}
// XK: responder does NOT learn initiator's identity until msg3
// Use a placeholder pubkey from src_addr for the session entry.
// The real pubkey will be registered when msg3 arrives.
// Generate msg2
let msg2 = match handshake.write_xk_message_2() {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Failed to generate Noise XK msg2 for SessionAck");
return;
}
};
// Build and send SessionAck (include initiator's coords for return-path warming)
let our_coords = self.tree_state.my_coords().clone();
let ack = SessionAck::new(our_coords, setup.src_coords).with_handshake(msg2);
let ack_payload = ack.encode();
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *src_addr, ack_payload.clone())
.with_ttl(self.config().node.session.default_ttl);
// Route the ack back to the initiator
if let Err(e) = self.send_session_datagram(&mut datagram).await {
debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send SessionAck");
return;
}
// Store session entry in AwaitingMsg3 state with ack payload for potential resend.
// Use a dummy pubkey since we don't know the initiator's identity yet.
// We use our own pubkey as placeholder; it will be replaced in handle_session_msg3.
// `keypair()` hands back a copy of the long-term private key, so the
// temporary is bound and erased rather than left to the statement end.
let mut our_keypair = self.identity().keypair();
let placeholder_pubkey = our_keypair.public_key();
our_keypair.non_secure_erase();
let now_ms = Self::now_ms();
let resend_interval = self.config().node.rate_limit.handshake_resend_interval_ms;
let mut entry = SessionEntry::new(
*src_addr,
placeholder_pubkey,
EndToEndState::AwaitingMsg3(handshake),
now_ms,
false,
);
entry.set_handshake_payload(ack_payload, now_ms + resend_interval);
self.sessions.insert(*src_addr, entry);
debug!(src = %self.peer_display_name(src_addr), "SessionSetup processed (XK), SessionAck sent, awaiting msg3");
}
/// Handle an incoming SessionAck (Noise XK msg2).
///
/// Processes msg2, generates and sends msg3, then transitions to Established.
async fn handle_session_ack(&mut self, src_addr: &NodeAddr, inner: &[u8]) {
let ack = match SessionAck::decode(inner) {
Ok(a) => a,
Err(e) => {
debug!(error = %e, "Malformed SessionAck");
return;
}
};
if ack.handshake_payload.len() != XK_HANDSHAKE_MSG2_SIZE {
debug!(
len = ack.handshake_payload.len(),
expected = XK_HANDSHAKE_MSG2_SIZE,
"Invalid handshake payload size in SessionAck"
);
return;
}
// Remove the entry to take ownership of the handshake state
let mut entry = match self.sessions.remove(src_addr) {
Some(e) => e,
None => {
debug!(src = %self.peer_display_name(src_addr), "SessionAck for unknown session");
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
};
// Rekey path: entry is Established with rekey_state
//
// `abandon_rekey` below, not `abandon_handshake` as in the responder
// arm of `handle_session_msg3`, and the difference rests on an
// invariant rather than on a different judgement: an entry with
// `rekey_initiator` set holds no pending session, so the two calls
// are the same action here. Arming as initiator has one caller,
// `initiate_session_rekey` through `begin_rekey`, and
// `set_rekey_state(_, true)` remains only at the restore below;
// `check_session_rekey` never reaches `initiate_session_rekey` for
// an entry holding a pending session; and `set_pending_session`
// clears `rekey_state`, so a completed initiator cycle leaves at
// most one of the two set. If that ever stops holding, these three
// sites become instances of the epoch discard the responder arm was
// fixed for.
if entry.is_established() && entry.has_rekey_in_progress() && entry.is_rekey_initiator() {
let mut handshake = match entry.take_rekey_state() {
Some(hs) => hs,
None => {
self.sessions.insert(*src_addr, entry);
return;
}
};
// Process XK msg2, for the same reason and in the same way as the
// primary arm below. Nothing here has been authenticated: the
// only tie to our rekey is the datagram's source address, which
// the sender chooses. Abandoning would let anyone able to name
// the session end the cycle, so the handshake goes back, rolled
// back to its pre-read state so it can still read the genuine
// ack, and the refusal is counted. The rollback matters because
// `read_xk_message_2` mixes the sender's ephemeral in before it
// authenticates. The restore does not restamp the deadline, which
// runs from the setup this node sent, so an unreadable ack cannot
// hold the rekey open.
if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) {
debug!(error = %e, "Failed to process rekey XK msg2, keeping the rekey");
entry.set_rekey_state(handshake, true);
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::AckHandshakeFailed));
return;
}
// The three abandons below stay abandons. Each follows a msg2
// that already authenticated, so they are local failures rather
// than possible forgeries.
// Generate XK msg3
let msg3 = match handshake.write_xk_message_3() {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Failed to generate rekey XK msg3");
entry.abandon_rekey();
self.sessions.insert(*src_addr, entry);
return;
}
};
// Send SessionMsg3
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.clone())
.with_ttl(self.config().node.session.default_ttl);
if let Err(e) = self.send_session_datagram(&mut datagram).await {
debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send rekey SessionMsg3");
entry.abandon_rekey();
self.sessions.insert(*src_addr, entry);
return;
}
// Complete handshake → store as pending new session
let session = match handshake.into_session() {
Ok(s) => s,
Err(e) => {
debug!(error = %e, "Failed to create session from rekey XK");
entry.abandon_rekey();
self.sessions.insert(*src_addr, entry);
return;
}
};
// Retain msg3 for retransmission (liveness): a single msg3
// loss must not leave the responder without the new session.
// Retransmission runs until the responder is confirmed on the
// new epoch — an authenticated peer frame against `pending` or
// post-cutover `current` — decoupled from this initiator's own
// cutover. The initiator may cut over on its liveness timer
// before the responder receives msg3; overlapping-epoch
// decrypt keeps both directions safe meanwhile.
let now_ms = Self::now_ms();
let resend_interval = self.config().node.rate_limit.handshake_resend_interval_ms;
entry.set_pending_session(session);
entry.set_rekey_completed_ms(now_ms);
entry.set_rekey_msg3_payload(msg3_payload, now_ms + resend_interval);
self.sessions.insert(*src_addr, entry);
debug!(
src = %self.peer_display_name(src_addr),
"FSP rekey: completed XK as initiator, msg3 sent, pending cutover"
);
return;
}
// Must be in Initiating state — check before take to avoid poisoning
if !entry.is_initiating() {
debug!(src = %self.peer_display_name(src_addr), "SessionAck but session not in Initiating state");
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::BadState));
return;
}
let mut handshake = match entry.take_state() {
Some(EndToEndState::Initiating(hs)) => hs,
_ => unreachable!("checked is_initiating above"),
};
// Process XK msg2: read_xk_message_2 (extracts responder's epoch)
//
// Nothing here has been authenticated: the only thing tying this
// message to the initiation is the datagram's source address, an
// envelope field the sender chooses. Dropping the entry would let
// anyone able to reach us cancel any initiation in flight, so the
// entry goes back with the handshake rolled back to its pre-read
// state and the stored msg1 still scheduled for resend. The rollback
// is what makes the reinsert worth anything: `read_xk_message_2`
// mixes the sender's ephemeral in before it authenticates, so a
// handshake put back as it was left could never read the genuine
// msg2. A real peer's corrupt ack is covered by the same path — the
// responder resends its stored msg2, and the handshake sweep reaps
// the entry on its original deadline if none arrives. `touch()` is
// deliberately not called, so a spray cannot push that deadline out.
if let Err(e) = handshake.try_read_xk_message_2(&ack.handshake_payload) {
debug!(error = %e, "Failed to process Noise XK msg2 in SessionAck");
entry.set_state(EndToEndState::Initiating(handshake));
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::AckHandshakeFailed));
return;
}
// The three drops below stay drops. Each is downstream of a msg2 that
// already authenticated, so they are local failures rather than
// possible forgeries, and a handshake left at `Message2Done` cannot
// be re-driven from a resent msg1 anyway.
// Generate XK msg3: write_xk_message_3 (sends encrypted static + epoch)
let msg3 = match handshake.write_xk_message_3() {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Failed to generate Noise XK msg3");
return;
}
};
// Send SessionMsg3 (phase 0x3)
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.clone())
.with_ttl(self.config().node.session.default_ttl);
if let Err(e) = self.send_session_datagram(&mut datagram).await {
debug!(error = %e, dest = %self.peer_display_name(src_addr), "Failed to send SessionMsg3");
return;
}
// Complete the handshake: into_session()
let session = match handshake.into_session() {
Ok(s) => s,
Err(e) => {
debug!(error = %e, "Failed to create session after XK msg3");
return;
}
};
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);
// 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);
// Flush any queued outbound packets for this destination
self.flush_pending_packets(src_addr).await;
self.flush_pending_native(src_addr).await;
info!(src = %self.peer_display_name(src_addr), "Session established (initiator, XK)");
}
/// Handle an incoming SessionMsg3 (Noise XK msg3).
///
/// The initiator reveals their encrypted static key. The responder
/// processes msg3, learns the initiator's identity, and transitions
/// to Established.
async fn handle_session_msg3(&mut self, src_addr: &NodeAddr, inner: &[u8]) {
let msg3 = match SessionMsg3::decode(inner) {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Malformed SessionMsg3");
return;
}
};
if msg3.handshake_payload.len() != XK_HANDSHAKE_MSG3_SIZE {
debug!(
len = msg3.handshake_payload.len(),
expected = XK_HANDSHAKE_MSG3_SIZE,
"Invalid handshake payload size in SessionMsg3"
);
return;
}
// Remove the entry to take ownership of the handshake state
let mut entry = match self.sessions.remove(src_addr) {
Some(e) => e,
None => {
debug!(src = %self.peer_display_name(src_addr), "SessionMsg3 for unknown session");
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
};
// Rekey path: entry is Established with rekey_state (responder side)
//
// Nothing in a msg3 is authenticated until the read has both
// succeeded and produced a static key matching this session's peer,
// so a failure here proves nothing about the sender and must not cost
// the entry anything it would miss. An unreadable msg3 therefore
// costs nothing at all: the handshake goes back, rolled back, for the
// genuine msg3. Every later failure follows a read that
// authenticated its sender and abandons only the handshake. A
// `pending_new_session` beside the handshake is the epoch the real
// peer may already have cut over to, and dropping it kills the
// reverse direction on two unauthenticated messages: a forged msg1 to
// arm the handshake, then any garbage msg3. `abandon_handshake` keeps
// it; `abandon_rekey` does not.
//
// What `abandon_handshake` leaves behind, and why each is safe here:
// `rekey_completed_ms` must survive, since `pending_stale` reads it
// and a zeroed stamp reads as freshly completed. A stranded
// `rekey_msg3_payload` belongs to an initiator cycle and clears
// itself once `resend_pending_session_msg3` exhausts its budget.
// `peer_new_epoch_confirmed` only stops that retransmission, and
// `rekey_initiator` is false throughout this arm by its own gate.
if entry.is_established() && entry.has_rekey_in_progress() && !entry.is_rekey_initiator() {
let mut handshake = match entry.take_rekey_state() {
Some(hs) => hs,
None => {
self.sessions.insert(*src_addr, entry);
return;
}
};
// Process XK msg3. The only tie between this msg3 and the
// handshake is the datagram's source address, which the sender
// chooses, and the peer that read our SessionAck may already hold
// the new keys. Abandoning here would let anyone able to name the
// session discard the handshake the peer's genuine msg3 needs,
// splitting the session's epochs. So the handshake goes back,
// rolled back to its pre-read state, because the read advances
// the cipher nonce before it authenticates. The restore leaves
// the deadline alone, which runs from the peer's setup, so a
// spray cannot hold the handshake open.
if let Err(e) = handshake.try_read_xk_message_3(&msg3.handshake_payload) {
debug!(
src = %self.peer_display_name(src_addr),
error = %e,
"Failed to process rekey XK msg3, keeping the handshake"
);
entry.set_rekey_state(handshake, false);
self.sessions.insert(*src_addr, entry);
return;
}
// The rekey must come from the peer the session was established
// with. Compare x-only keys: the stored key's parity may be a
// synthesized even parity (npubs carry no parity), while the
// handshake learns the true point.
let rekey_pubkey = match handshake.remote_static() {
Some(pk) => *pk,
None => {
// Not independently exercised by any test: a successful
// `read_xk_message_3` always sets the remote static, so
// reaching this needs fault injection into
// `HandshakeState`. Changed with its three siblings so
// the arm has one rule rather than three plus an
// exception.
debug!("No remote static key after processing rekey XK msg3");
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
return;
}
};
if rekey_pubkey.x_only_public_key().0 != entry.remote_pubkey().x_only_public_key().0 {
warn!(
src = %self.peer_display_name(src_addr),
"FSP rekey: initiator static key differs from the established peer key"
);
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::RekeyKeyMismatch));
return;
}
// Complete the handshake → store as pending new session
let session = match handshake.into_session() {
Ok(s) => s,
Err(e) => {
// Also not independently exercised, for the same reason
// as the missing-static arm above: a handshake that read
// msg3 successfully always converts.
debug!(error = %e, "Failed to create session from rekey XK msg3");
entry.abandon_handshake();
self.sessions.insert(*src_addr, entry);
return;
}
};
// A pending session already held for this peer is superseded
// here rather than by any timer: only a msg3 carrying the
// session's own peer key can replace the epoch that peer moved
// to. The keys it displaces may still be in use, so the event is
// counted and logged rather than silent.
let superseded = entry.pending_new_session().is_some();
entry.set_pending_session(session);
entry.set_rekey_completed_ms(Self::now_ms());
self.sessions.insert(*src_addr, entry);
if superseded {
self.stats_mut().session.pending_replaced += 1;
warn!(
src = %self.peer_display_name(src_addr),
"FSP rekey: newly completed session replaced one still awaiting cutover"
);
}
debug!(
src = %self.peer_display_name(src_addr),
"FSP rekey: completed XK as responder, pending cutover"
);
return;
}
// Must be in AwaitingMsg3 state
if !entry.is_awaiting_msg3() {
debug!(src = %self.peer_display_name(src_addr), "SessionMsg3 but session not in AwaitingMsg3 state");
self.sessions.insert(*src_addr, entry);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::BadState));
return;
}
let mut handshake = match entry.take_state() {
Some(EndToEndState::AwaitingMsg3(hs)) => hs,
_ => unreachable!("checked is_awaiting_msg3 above"),
};
// Process XK msg3 (extracts the initiator's static key and epoch).
//
// Nothing here has been authenticated: the only tie to this half-open
// entry is the datagram's source address, which the sender chooses,
// and the initiator considers the session established once it has
// sent msg3. Dropping the entry would let anyone able to name the
// initiator discard the handshake its genuine msg3 and resends need,
// so the entry goes back with the handshake rolled back to its
// pre-read state, as the SessionAck arm does. `touch()` is
// deliberately not called, so a spray cannot push the handshake
// sweep's deadline out. The drops below stay drops: each follows a
// read that authenticated the sender.
if let Err(e) = handshake.try_read_xk_message_3(&msg3.handshake_payload) {
debug!(error = %e, "Failed to process Noise XK msg3, keeping the handshake");
entry.set_state(EndToEndState::AwaitingMsg3(handshake));
self.sessions.insert(*src_addr, entry);
return;
}
// Extract the initiator's static public key (now available after msg3)
let remote_pubkey = match handshake.remote_static() {
Some(pk) => *pk,
None => {
debug!("No remote static key after processing XK msg3");
return;
}
};
// The claimed source address must be derivable from the key we just
// authenticated, or the peer is opening a session under another
// node's address.
let derived_addr = NodeAddr::from_pubkey(&remote_pubkey.x_only_public_key().0);
if derived_addr != *src_addr {
warn!(
src = %self.peer_display_name(src_addr),
derived = %derived_addr,
"SessionMsg3 source address does not match the authenticated static key"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::AddrMismatch));
return; // Entry was already removed
}
// Register the initiator's identity for future TUN → session routing
self.register_identity(*src_addr, remote_pubkey);
// Complete the handshake
let session = match handshake.into_session() {
Ok(s) => s,
Err(e) => {
debug!(error = %e, "Failed to create session from XK handshake");
return;
}
};
let now_ms = Self::now_ms();
// Replace the placeholder pubkey with the real one
let mut new_entry = SessionEntry::new(
*src_addr,
remote_pubkey,
EndToEndState::Established(session),
now_ms,
false,
);
new_entry.set_coords_warmup_remaining(self.config().node.session.coords_warmup_packets);
new_entry.mark_established(now_ms);
new_entry.init_mmp(&self.config().node.session_mmp);
new_entry.touch(now_ms);
self.sessions.insert(*src_addr, new_entry);
// Flush any pending packets
self.flush_pending_packets(src_addr).await;
self.flush_pending_native(src_addr).await;
info!(src = %self.peer_display_name(src_addr), "Session established (responder, XK)");
}
// === Session-layer MMP report handlers ===
/// Check all sessions for pending MMP reports and send them.
///
/// Called from the tick handler. Also emits periodic session MMP logs.
/// Uses the collect-then-send pattern to avoid borrowing conflicts.
pub(in crate::node) async fn check_session_mmp_reports(&mut self) {
let now_ms = crate::time::mono_ms();
// Build one report-gating snapshot per session, resolving every timing
// read shell-side into a `bool`. The snapshots own only
// `NodeAddr`/`MmpMode`/`bool`, so the session-iteration borrow is released
// before the pure decision runs and the driving loop mutates the
// reporting state / performs the sends.
let snapshots: Vec<SessionReportSnapshot> = self
.sessions
.iter()
.filter_map(|(dest_addr, entry)| {
let mmp = entry.mmp()?;
Some(SessionReportSnapshot {
dest: *dest_addr,
mode: mmp.mode(),
sr_due: mmp.sender.should_send_report(now_ms),
rr_due: mmp.receiver.should_send_report(now_ms),
mtu_due: mmp.path_mtu.should_send_notification(now_ms),
log_due: mmp.should_log(now_ms),
})
})
.collect();
let actions = self.mmp.plan_session_reports(&snapshots);
// Drive the planned actions in phase-grouped order (all logs, then the
// sends in per-session SR/RR/MTU order). Logs run first because the
// session operator log reads cumulative_packets_sent, which each send
// advances (send_session_msg -> sender.record_sent); the pre-refactor
// handler logged during its collect pass, before any send. Each build
// (`build_report`/`build_notification`, which advance interval/
// notification state) runs only on its SendSessionReport action, exactly
// as the pre-refactor collect pass did. Per-destination success/failure
// is collected for the backoff dedup + failure-log suppression.
let mut send_results: Vec<SendResult> = Vec::new();
for action in actions {
match action {
MmpAction::LogSession { dest } => {
// Resolve the display name exactly as the pre-refactor loop
// did (alias, else short_npub from the session's remote key).
let session_name = self.peer_aliases.get(&dest).cloned().unwrap_or_else(|| {
self.sessions
.get(&dest)
.map(|entry| {
let (xonly, _) = entry.remote_pubkey().x_only_public_key();
crate::PeerIdentity::from_pubkey(xonly).short_npub()
})
.unwrap_or_default()
});
if let Some(mmp) = self.sessions.get_mut(&dest).and_then(|e| e.mmp_mut()) {
Self::log_session_mmp_metrics(&session_name, mmp);
mmp.mark_logged(now_ms);
}
}
MmpAction::SendSessionReport { dest, kind } => {
let built = self
.sessions
.get_mut(&dest)
.and_then(|entry| entry.mmp_mut())
.and_then(|mmp| match kind {
SessionReportKind::Sender => {
mmp.sender.build_report(now_ms).map(|sr| {
(
SessionMessageType::SenderReport.to_byte(),
SessionSenderReport::from(&sr).encode(),
)
})
}
SessionReportKind::Receiver => {
mmp.receiver.build_report(now_ms).map(|rr| {
(
SessionMessageType::ReceiverReport.to_byte(),
SessionReceiverReport::from(&rr).encode(),
)
})
}
SessionReportKind::PathMtu => {
mmp.path_mtu.build_notification(now_ms).map(|mtu_value| {
(
SessionMessageType::PathMtuNotification.to_byte(),
PathMtuNotification::new(mtu_value).encode(),
)
})
}
});
let Some((msg_type, body)) = built else {
continue;
};
match self.send_session_msg(&dest, msg_type, &body).await {
Ok(()) => send_results.push(SendResult { dest, ok: true }),
Err(e) => {
// Peek at current failure count for log suppression
// (unchanged by the backoff apply, which runs later).
let failures = self
.sessions
.get(&dest)
.and_then(|entry| entry.mmp())
.map(|mmp| mmp.sender.consecutive_send_failures())
.unwrap_or(0);
if failures < 3 {
debug!(
dest = %self.peer_display_name(&dest),
msg_type,
error = %e,
"Failed to send session MMP report"
);
} else if failures == 3 {
debug!(
dest = %self.peer_display_name(&dest),
"Suppressing further session MMP send failure logs"
);
}
// failures > 3: silently suppressed
send_results.push(SendResult { dest, ok: false });
}
}
}
MmpAction::ReapPeer { .. }
| MmpAction::Heartbeat { .. }
| MmpAction::SendLinkReport { .. }
| MmpAction::LogLink { .. } => {}
}
}
// Deduplicate send results per destination (any-ok -> success, all-fail
// -> failure) and apply the backoff state transition for each dest.
for update in self.mmp.plan_backoff(&send_results) {
match update {
BackoffUpdate::Success { dest } => {
if let Some(mmp) = self.sessions.get_mut(&dest).and_then(|e| e.mmp_mut()) {
let prev = mmp.sender.record_send_success();
if prev > 3 {
debug!(
dest = %self.peer_display_name(&dest),
consecutive_failures = prev,
"Resumed session MMP reporting"
);
}
}
}
BackoffUpdate::Failure { dest } => {
if let Some(mmp) = self.sessions.get_mut(&dest).and_then(|e| e.mmp_mut()) {
mmp.sender.record_send_failure();
}
}
}
}
}
/// Emit periodic session MMP metrics.
fn log_session_mmp_metrics(session_name: &str, mmp: &MmpSessionState) {
let m = &mmp.metrics;
let rtt_str = if m.rtt_trend.initialized() {
format!("{:.1}ms", m.rtt_trend.long() / 1000.0)
} else {
"n/a".to_string()
};
let loss_str = if m.loss_trend.initialized() {
format!("{:.1}%", m.loss_trend.long() * 100.0)
} else {
"n/a".to_string()
};
let jitter_ms = mmp.receiver.jitter_us() as f64 / 1000.0;
debug!(
session = %session_name,
rtt = %rtt_str,
loss = %loss_str,
jitter = format_args!("{:.1}ms", jitter_ms),
goodput = %format_throughput(m.goodput_bps()),
mtu = mmp.path_mtu.last_observed_mtu(),
tx_pkts = mmp.sender.cumulative_packets_sent(),
rx_pkts = mmp.receiver.cumulative_packets_recv(),
"MMP session metrics"
);
}
/// Emit a teardown log summarizing lifetime session MMP metrics.
pub(in crate::node) fn log_session_mmp_teardown(session_name: &str, mmp: &MmpSessionState) {
let m = &mmp.metrics;
let jitter_ms = mmp.receiver.jitter_us() as f64 / 1000.0;
let rtt_str = match m.srtt_ms() {
Some(rtt) => format!("{:.1}ms", rtt),
None => "n/a".to_string(),
};
let loss_str = format!("{:.1}%", m.loss_rate() * 100.0);
debug!(
session = %session_name,
rtt = %rtt_str,
loss = %loss_str,
jitter = format_args!("{:.1}ms", jitter_ms),
etx = format_args!("{:.2}", m.etx),
goodput = %format_throughput(m.goodput_bps()),
send_mtu = mmp.path_mtu.current_mtu(),
observed_mtu = mmp.path_mtu.last_observed_mtu(),
tx_pkts = mmp.sender.cumulative_packets_sent(),
tx_bytes = mmp.sender.cumulative_bytes_sent(),
rx_pkts = mmp.receiver.cumulative_packets_recv(),
rx_bytes = mmp.receiver.cumulative_bytes_recv(),
"MMP session teardown"
);
}
/// Handle an incoming session-layer SenderReport (msg_type 0x11).
///
/// Informational only — the peer is telling us about what they sent.
/// Logged but not used for metrics (same pattern as link-layer).
fn handle_session_sender_report(&mut self, src_addr: &NodeAddr, body: &[u8]) {
let sr = match SessionSenderReport::decode(body) {
Ok(sr) => sr,
Err(e) => {
debug!(src = %self.peer_display_name(src_addr), error = %e, "Malformed SessionSenderReport");
return;
}
};
trace!(
src = %self.peer_display_name(src_addr),
cum_pkts = sr.cumulative_packets_sent,
interval_bytes = sr.interval_bytes_sent,
"Received SessionSenderReport"
);
}
/// Handle an incoming session-layer ReceiverReport (msg_type 0x12).
///
/// The peer is telling us about what they received from us. We feed
/// this to our metrics to compute RTT, loss rate, and trend indicators.
fn handle_session_receiver_report(&mut self, src_addr: &NodeAddr, body: &[u8]) {
let session_rr = match SessionReceiverReport::decode(body) {
Ok(rr) => rr,
Err(e) => {
debug!(src = %self.peer_display_name(src_addr), error = %e, "Malformed SessionReceiverReport");
return;
}
};
// Convert to link-layer ReceiverReport for MmpMetrics processing
let rr: ReceiverReport = ReceiverReport::from(&session_rr);
let now_ms = Self::now_ms();
let peer_name = self.peer_display_name(src_addr);
let entry = match self.sessions.get_mut(src_addr) {
Some(e) => e,
None => {
debug!(src = %peer_name, "SessionReceiverReport for unknown session");
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
};
let our_timestamp_ms = entry.session_timestamp(now_ms);
let Some(mmp) = entry.mmp_mut() else {
return;
};
let (_first_rtt, rr_log) =
mmp.metrics
.process_receiver_report(&rr, our_timestamp_ms, crate::time::mono_ms());
// Re-emit the operator trace the core used to log mid-decision.
super::mmp::log_rr_outcome(&rr, our_timestamp_ms, rr_log);
// Feed SRTT back to sender/receiver report interval tuning (session-layer bounds)
if let Some(srtt_ms) = mmp.metrics.srtt_ms() {
let srtt_us = (srtt_ms * 1000.0) as i64;
mmp.sender.update_report_interval_with_bounds(
srtt_us,
MIN_SESSION_REPORT_INTERVAL_MS,
MAX_SESSION_REPORT_INTERVAL_MS,
);
mmp.receiver.update_report_interval_with_bounds(
srtt_us,
MIN_SESSION_REPORT_INTERVAL_MS,
MAX_SESSION_REPORT_INTERVAL_MS,
);
// Also update PathMtu notification interval from SRTT
mmp.path_mtu.update_interval_from_srtt(srtt_ms);
}
// Update reverse delivery ratio from our own receiver state, using per-interval deltas.
let our_recv_packets = mmp.receiver.cumulative_packets_recv();
let peer_highest = mmp.receiver.highest_counter();
mmp.metrics
.update_reverse_delivery(our_recv_packets, peer_highest);
trace!(
src = %peer_name,
rtt_ms = ?mmp.metrics.srtt_ms(),
loss = format_args!("{:.1}%", mmp.metrics.loss_rate() * 100.0),
"Processed SessionReceiverReport"
);
}
/// Handle an incoming PathMtuNotification (msg_type 0x13).
///
/// The destination is telling us the path MTU has changed.
/// Apply source-side rules (decrease immediate, increase validated).
pub(in crate::node) fn handle_session_path_mtu_notification(
&mut self,
src_addr: &NodeAddr,
body: &[u8],
) {
let notif = match PathMtuNotification::decode(body) {
Ok(n) => n,
Err(e) => {
debug!(src = %self.peer_display_name(src_addr), error = %e, "Malformed PathMtuNotification");
return;
}
};
let peer_name = self.peer_display_name(src_addr);
let entry = match self.sessions.get_mut(src_addr) {
Some(e) => e,
None => {
debug!(src = %peer_name, "PathMtuNotification for unknown session");
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
};
let Some(mmp) = entry.mmp_mut() else {
return;
};
// `apply_notification` refuses a sub-floor value, but it returns the
// same `false` it returns for the ordinary "no change" case, which is
// the common one. Test the floor here so the refusal is visible: this
// arrives on the decrypted service-payload path, so a value this low
// means an authenticated peer we hold a session with is sending
// something unusable.
if notif.path_mtu < crate::upper::icmp::MIN_ACTIONABLE_PATH_MTU {
warn!(
src = %peer_name,
reported_mtu = notif.path_mtu,
floor = crate::upper::icmp::MIN_ACTIONABLE_PATH_MTU,
"PathMtuNotification reports a path MTU below the actionable floor; ignoring"
);
self.metrics.errors.path_mtu_notif_below_floor.inc();
return;
}
let old_mtu = mmp.path_mtu.current_mtu();
let changed = mmp
.path_mtu
.apply_notification(notif.path_mtu, crate::time::mono_ms());
let new_mtu = mmp.path_mtu.current_mtu();
if !changed {
return;
}
debug!(
src = %peer_name,
old_mtu,
new_mtu,
"Path MTU changed via notification"
);
// Mirror the new effective MTU into the FipsAddress-keyed lookup used
// by the TUN reader/writer at TCP MSS clamp time. Without this, new
// TCP flows opened on a path the proactive end-to-end echo has
// already tightened keep getting clamped by the staler discovery-
// time value until a reactive MtuExceeded happens to fire. Keep the
// tighter of existing-or-new — never loosen the clamp.
let fips_addr = crate::FipsAddress::from_node_addr(src_addr);
match self.path_mtu_lookup.write() {
Ok(mut map) => {
// Read existing, decide, and apply the write under one guard so
// the keep-tighter update stays atomic.
let prior = map.get(&fips_addr).copied();
let actions =
self.fsp
.plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), new_mtu);
if actions.is_empty() {
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
new_mtu,
existing = prior.map(|e| e.mtu).unwrap_or(new_mtu),
"PathMtuNotification: keeping tighter existing path_mtu_lookup value"
);
}
for action in actions {
if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action {
// Held, not expiring. This value arrives inside a
// session, and a session's teardown or a PathBroken
// naming it already releases the entry. A deadline here
// would instead recreate the gap this mirror exists to
// close: a peer repeating an identical value on a stable
// path takes the unchanged early-return above and never
// rewrites the entry, so an expiring one would vanish
// and stay gone.
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu));
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
new_mtu,
prior = ?prior,
map_len = map.len(),
"PathMtuNotification: tightened path_mtu_lookup"
);
}
}
}
Err(e) => {
warn!(
dest = %peer_name,
fips_addr = %fips_addr,
new_mtu,
error = %e,
"path_mtu_lookup write lock poisoned; PathMtuNotification not reflected"
);
}
}
}
/// Whether a routing signal naming `dest`, arriving in a datagram
/// claiming source `src`, may be acted on, and if not, which kind of
/// refusal it is.
///
/// `src` is the SessionDatagram's `src_addr`: a plain wire field,
/// authenticated only hop-by-hop by FMP Noise and never end to end. It
/// is therefore logged, not trusted. What is enforced here is that this
/// node has bound `dest` itself, either by initiating toward it or by
/// completing Noise XK, which binds the address to the peer's static
/// key (see the address-mismatch check in `handle_session_msg3`). A
/// responder entry that is still awaiting msg3 does NOT qualify: it is
/// keyed on an address the sender merely claimed, so admitting it would
/// let one forged SessionSetup unlock a signal about any address.
///
/// The first two clauses reject nothing legitimate, and so return
/// `Forged` rather than `Unbound`. A datagram whose destination is this
/// node takes the deliver-local branch before any forwarding, so no node
/// ever emits a signal naming us as `dest`; and the emitter is by
/// construction a transit node for the datagram it is reporting on, so
/// it is never itself that datagram's destination.
///
/// This narrows who can be targeted; it does not authenticate the
/// sender, which nothing short of a wire format change can do.
fn signal_verdict(&self, src: &NodeAddr, dest: &NodeAddr) -> SignalVerdict {
if dest == self.node_addr() || src == dest {
return SignalVerdict::Forged;
}
if self
.sessions
.get(dest)
.is_some_and(|e| e.is_established() || e.is_initiator())
{
SignalVerdict::Admit
} else {
SignalVerdict::Unbound
}
}
/// Handle a CoordsRequired error signal from a transit router.
///
/// The router couldn't route our packet because it lacks cached
/// coordinates for the destination. Send a standalone CoordsWarmup
/// immediately (rate-limited), trigger discovery, and reset the
/// warmup counter for subsequent data packets.
///
/// `src_addr` is the datagram's claimed source and is not
/// end-to-end authenticated; see `signal_verdict`.
pub(in crate::node) async fn handle_coords_required(
&mut self,
src_addr: &NodeAddr,
inner: &[u8],
) {
self.metrics().errors.coords_required.inc();
let msg = match CoordsRequired::decode(inner) {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Malformed CoordsRequired");
return;
}
};
// The premise: this signal carries no end-to-end authentication, so
// the body's `dest_addr` is attacker-chosen. Everything below acts on
// it — warmup send, discovery, warmup-counter reset — so the gate has
// to run before any of that, and ahead of the rate limiter, whose
// state would otherwise be keyed on an attacker-chosen address.
let verdict = self.signal_verdict(src_addr, &msg.dest_addr);
if verdict != SignalVerdict::Admit {
debug!(src = %src_addr, dest = %msg.dest_addr, reporter = %msg.reporter,
signal = "CoordsRequired", verdict = verdict.label(),
"Routing signal names an address this node has not bound; dropping");
self.metrics().errors.unbound.coords.inc();
if verdict == SignalVerdict::Forged {
self.metrics().errors.unbound.forged.inc();
}
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
debug!(
dest = %msg.dest_addr,
reporter = %msg.reporter,
"CoordsRequired: transit router needs coordinates"
);
// Send standalone CoordsWarmup immediately (rate-limited)
if self
.coords_response_rate_limiter
.should_send(&msg.dest_addr, Self::now_ms())
{
if let Some(entry) = self.sessions.get(&msg.dest_addr)
&& entry.is_established()
&& let Err(e) = self.send_coords_warmup(&msg.dest_addr).await
{
debug!(dest = %msg.dest_addr, error = %e,
"Failed to send CoordsWarmup in response to CoordsRequired");
}
} else {
trace!(dest = %msg.dest_addr,
"CoordsRequired response rate-limited, skipping standalone CoordsWarmup");
}
// Only trigger discovery if we have the target's identity cached —
// otherwise we can't verify the LookupResponse proof.
let has_cached_identity = self.has_cached_identity(&msg.dest_addr);
let actions = self
.fsp
.plan_coords_required_lookup(msg.dest_addr, has_cached_identity);
if actions.is_empty() {
debug!(dest = %msg.dest_addr,
"Skipping discovery after CoordsRequired: no cached identity for target");
}
for action in actions {
if let FspAction::InitiateLookup { dest } = action {
self.maybe_initiate_lookup(&dest).await;
}
}
// Reset coords warmup counter so the next N packets also include
// COORDS_PRESENT, re-warming transit caches along the path.
let n = self.config().node.session.coords_warmup_packets;
if let Some(entry) = self.sessions.get_mut(&msg.dest_addr) {
entry.set_coords_warmup_remaining(n);
debug!(
dest = %msg.dest_addr,
warmup_packets = n,
"Reset coords warmup counter after CoordsRequired"
);
}
}
/// Handle a PathBroken error signal from a transit router.
///
/// The router has coordinates but still can't route to the destination.
/// Send a standalone CoordsWarmup immediately (rate-limited), re-validate
/// the destination's coordinates by lookup, release its path MTU, and
/// reset the warmup counter.
///
/// Cached coordinates that are only a hint are removed. Coordinates a
/// lookup verified are kept while the lookup re-validates them, and are
/// demoted to a hint only once reports arriving over distinct links reach
/// the quorum: the signal is unauthenticated, and deleting a verified
/// entry on one report is what let the next forged warm replace it.
///
/// `src_addr` is the datagram's claimed source and is not
/// end-to-end authenticated; see `signal_verdict`. `link_peer` is the
/// authenticated peer the datagram arrived over. It is the report's vote
/// in the quorum, because the body's reporter is whatever the sender
/// wrote, and it is compared with the forward path to count mismatches;
/// it never refuses the signal.
pub(in crate::node) async fn handle_path_broken(
&mut self,
src_addr: &NodeAddr,
link_peer: &NodeAddr,
inner: &[u8],
) {
self.metrics().errors.path_broken.inc();
let msg = match PathBroken::decode(inner) {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Malformed PathBroken");
return;
}
};
// The premise: this signal carries no end-to-end authentication, so
// the body's `dest_addr` is attacker-chosen. The coord-cache action,
// the lookup and the path-MTU release below all act on whatever
// address the body names unless the gate refuses it here, in the
// shell, which is the only layer that knows who sent the datagram.
let verdict = self.signal_verdict(src_addr, &msg.dest_addr);
if verdict != SignalVerdict::Admit {
debug!(src = %src_addr, dest = %msg.dest_addr, reporter = %msg.reporter,
signal = "PathBroken", verdict = verdict.label(),
"Routing signal names an address this node has not bound; dropping");
self.metrics().errors.unbound.broken.inc();
if verdict == SignalVerdict::Forged {
self.metrics().errors.unbound.forged.inc();
}
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
debug!(
dest = %msg.dest_addr,
reporter = %msg.reporter,
"PathBroken: transit router reports routing failure"
);
self.count_path_broken_mismatches(link_peer, &msg);
// Send standalone CoordsWarmup immediately (rate-limited)
if self
.coords_response_rate_limiter
.should_send(&msg.dest_addr, Self::now_ms())
{
if let Some(entry) = self.sessions.get(&msg.dest_addr)
&& entry.is_established()
&& let Err(e) = self.send_coords_warmup(&msg.dest_addr).await
{
debug!(dest = %msg.dest_addr, error = %e,
"Failed to send CoordsWarmup in response to PathBroken");
}
} else {
trace!(dest = %msg.dest_addr,
"PathBroken response rate-limited, skipping standalone CoordsWarmup");
}
// Only a live verified entry has anything for the quorum to protect,
// so only a report against one is recorded; that also bounds the
// quorum's keys by the destinations this node has looked up. The vote
// is the link the report arrived over, so a sender on one link counts
// once however many reporters it names.
let now = Self::now_ms();
let verified = self
.coord_cache
.get_entry(&msg.dest_addr)
.is_some_and(|e| e.is_verified(now));
let quorum = if verified {
self.broken_quorum.record(msg.dest_addr, *link_peer, now)
} else {
QuorumVerdict::Below { distinct: 0 }
};
let actions = self.fsp.plan_path_broken(msg.dest_addr, verified, quorum);
for action in actions {
match action {
FspAction::InvalidateCoords { addr } => {
self.coord_cache.remove(&addr);
}
FspAction::DemoteCoords { addr } => {
self.coord_cache.demote(&addr);
self.metrics().errors.broken_demoted.inc();
debug!(dest = %addr,
"PathBroken quorum reached; demoted verified coordinates to a hint");
}
FspAction::InitiateLookup { dest } => {
self.cache_session_identity(&dest);
self.maybe_initiate_lookup(&dest).await;
}
_ => {}
}
}
if let QuorumVerdict::Below { distinct } = quorum
&& verified
{
self.metrics().errors.broken_below_quorum.inc();
debug!(dest = %msg.dest_addr, distinct,
"PathBroken below quorum; keeping verified coordinates pending the lookup");
}
// The path this destination's stored MTU described is gone, so release
// it rather than carrying it onto whatever path replaces it. Rate
// limited per destination on its own budget: PathBroken is
// unauthenticated, and an unlimited release discards a genuinely
// learned bottleneck as fast as it is relearned. The budget is not
// shared with any other signal, so nothing else can spend it.
//
// A cache entry kept above still carries the MTU the lookup stored
// with it, which describes the same path; clear it with the map so
// the two do not disagree.
if self
.path_mtu_release_limiter
.should_send(&msg.dest_addr, Self::now_ms())
{
self.path_mtu_lookup_release(&msg.dest_addr);
self.coord_cache.clear_path_mtu(&msg.dest_addr);
} else {
trace!(dest = %msg.dest_addr,
"PathBroken path MTU release rate-limited, keeping the stored value");
}
// Reset coords warmup counter so the next N packets include
// COORDS_PRESENT, re-warming transit caches along the new path.
let n = self.config().node.session.coords_warmup_packets;
if let Some(entry) = self.sessions.get_mut(&msg.dest_addr) {
entry.set_coords_warmup_remaining(n);
debug!(
dest = %msg.dest_addr,
warmup_packets = n,
"Reset coords warmup counter after PathBroken"
);
}
}
/// Cache the identity of `dest` from its session entry if the identity
/// cache has none.
///
/// A lookup's answer is verified against the target's cached key, so a
/// lookup for a destination whose identity has been evicted runs to its
/// timeout and drops the packets queued for it as unreachable. A session
/// already holds the key: the one this node initiated to, or the one the
/// responder handshake authenticated.
fn cache_session_identity(&mut self, dest: &NodeAddr) {
if self.has_cached_identity(dest) {
return;
}
if let Some(pubkey) = self.sessions.get(dest).map(|e| *e.remote_pubkey()) {
self.register_identity(*dest, pubkey);
}
}
/// Count, without acting on, the two ways an admitted PathBroken can
/// disagree with this node's own view of the path.
///
/// Advisory only: both checks read state an attacker can influence and
/// both have a non-zero healthy floor, so they size the problem rather
/// than refuse anything.
///
/// - Link: the signal arrived over a link other than the one this node
/// would forward to the destination on. A genuine report can also do
/// this when the reverse path differs from the forward one.
/// - Reporter: the reporter is this node or the destination, or its known
/// coordinates are not strictly closer to the destination than this
/// node's. Forwarding makes strict progress under each forwarder's own
/// view of the destination, and a PathBroken arises exactly where that
/// view may differ from this node's, so a genuine report can count here
/// too. A reporter's coordinates are rarely known unless it is a direct
/// peer, so this mostly reads as unknown and is not counted.
fn count_path_broken_mismatches(&self, link_peer: &NodeAddr, msg: &PathBroken) {
let now = Self::now_ms();
let dest = &msg.dest_addr;
let reporter = &msg.reporter;
if let (Some(hop), _) = self.preview_next_hop(dest, now)
&& hop.node_addr != *link_peer
{
self.metrics().errors.broken_link_mismatch.inc();
debug!(dest = %dest, link_peer = %link_peer, next_hop = %hop.node_addr,
"PathBroken arrived off the forward link");
}
let implausible = if reporter == self.node_addr() || reporter == dest {
true
} else {
let dest_coords = self.coord_cache.get(dest, now);
let reporter_coords = self
.tree_state
.peer_coords(reporter)
.or_else(|| self.coord_cache.get(reporter, now));
match (dest_coords, reporter_coords) {
(Some(d), Some(r)) => {
r.distance_to(d) >= self.tree_state.my_coords().distance_to(d)
}
_ => false,
}
};
if implausible {
self.metrics().errors.broken_reporter_mismatch.inc();
debug!(dest = %dest, reporter = %reporter,
"PathBroken reporter is not closer to the destination than this node");
}
}
/// Handle an MtuExceeded error signal from a transit router.
///
/// A transit router couldn't forward our packet because it exceeded the
/// next-hop transport MTU. Apply the reported bottleneck MTU to our
/// PathMtuState for the affected session, causing an immediate decrease.
///
/// `src_addr` is the datagram's claimed source and is not
/// end-to-end authenticated; see `signal_verdict`.
pub(in crate::node) async fn handle_mtu_exceeded(&mut self, src_addr: &NodeAddr, inner: &[u8]) {
self.metrics().errors.mtu_exceeded.inc();
let msg = match MtuExceeded::decode(inner) {
Ok(m) => m,
Err(e) => {
debug!(error = %e, "Malformed MtuExceeded");
return;
}
};
// The premise: this signal carries no end-to-end authentication, so
// the body's `dest_addr` is attacker-chosen, and the `path_mtu_lookup`
// write further down needs no session, no peer relationship and no
// prior state to reach. This gate is about WHICH address may be
// written; the floor guard below is about WHAT value may be written.
// They are independent refusals — a bound destination can still carry
// an unusable value — so neither subsumes the other and each keeps its
// own counter.
let verdict = self.signal_verdict(src_addr, &msg.dest_addr);
if verdict != SignalVerdict::Admit {
debug!(src = %src_addr, dest = %msg.dest_addr, reporter = %msg.reporter,
signal = "MtuExceeded", verdict = verdict.label(),
"Routing signal names an address this node has not bound; dropping");
self.metrics().errors.unbound.mtu.inc();
if verdict == SignalVerdict::Forged {
self.metrics().errors.unbound.forged.inc();
}
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::UnknownSession));
return;
}
let peer_name = self.peer_display_name(&msg.dest_addr);
debug!(
dest = %peer_name,
reporter = %msg.reporter,
bottleneck_mtu = msg.mtu,
"MtuExceeded: transit router reports oversized packet"
);
// Both effects below — the session's own path MTU and the
// FipsAddress-keyed lookup the TUN MSS clamp reads — are refused from
// here, so one return covers both. The guards sit ahead of the apply
// rather than between the two effects, which is what makes the floor
// govern `current_mtu` and not only the lookup table.
// Refuse a bottleneck too small to describe a usable path; a stored
// value that low drives the SYN-time MSS clamp into single digits or
// zero. The reactive carrier is unauthenticated, so it has its own
// floor constant, currently equal to the actionable one.
if msg.mtu < crate::upper::icmp::MIN_REACTIVE_PATH_MTU {
warn!(
dest = %peer_name,
reporter = %msg.reporter,
bottleneck_mtu = msg.mtu,
floor = crate::upper::icmp::MIN_REACTIVE_PATH_MTU,
"MtuExceeded reports a path MTU below the actionable floor; ignoring"
);
self.metrics().errors.mtu_exceeded_below_floor.inc();
return;
}
// Corroboration. The admission gate narrows which destination may be
// named; it cannot authenticate the reporter, so a legal value is a
// legal value from anyone and the floor alone only sets the outcome of
// a forgery rather than preventing it. An honest report exists only
// because a frame this node emitted did not fit some hop, so require
// that this node has actually sent something larger than the value
// being claimed since the last accepted decrease. Honest path-MTU
// discovery satisfies this by construction; a forgery has to wait for
// us to emit a frame bigger than the value it wants to claim, which
// bounds every accepted claim from below by our own traffic.
let sent_wire_len = self
.sessions
.get(&msg.dest_addr)
.map(|e| e.max_sent_wire_len())
.unwrap_or(0);
if msg.mtu >= sent_wire_len {
debug!(
dest = %peer_name,
reporter = %msg.reporter,
bottleneck_mtu = msg.mtu,
max_sent_wire_len = sent_wire_len,
"MtuExceeded reports a bottleneck no smaller than anything this node has sent; ignoring"
);
self.metrics().errors.mtu_exceeded_uncorroborated.inc();
return;
}
// Apply to PathMtuState: immediate decrease via apply_notification()
if let Some(entry) = self.sessions.get_mut(&msg.dest_addr)
&& let Some(mmp) = entry.mmp_mut()
{
let old_mtu = mmp.path_mtu.current_mtu();
if mmp
.path_mtu
.apply_notification(msg.mtu, crate::time::mono_ms())
{
let new_mtu = mmp.path_mtu.current_mtu();
info!(
dest = %peer_name,
old_mtu,
new_mtu,
reporter = %msg.reporter,
"Path MTU decreased via reactive MtuExceeded signal"
);
}
}
// Spent: the evidence vouched for this decrease and does not vouch for
// the next one. An initiating session has no `mmp` and so reaches this
// with the apply above skipped; the reset belongs to the acceptance,
// not to the apply.
if let Some(entry) = self.sessions.get_mut(&msg.dest_addr) {
entry.clear_sent_wire_len();
}
// Mirror the bottleneck into the FipsAddress-keyed lookup used by
// the TUN reader/writer at TCP MSS clamp time. Discovery's reverse-
// path response can carry a value too generous for the actual
// forward path; the reactive signal from a forwarder that actually
// dropped a packet is authoritative for "what fits". Keep the
// tighter of existing-or-new — never loosen the clamp.
//
// The admission gate above, not this block, is what restricts which
// addresses can be written here.
let fips_addr = crate::FipsAddress::from_node_addr(&msg.dest_addr);
match self.path_mtu_lookup.write() {
Ok(mut map) => {
// Read existing, decide, and apply the write under one guard so
// the keep-tighter update stays atomic.
let prior = map.get(&fips_addr).copied();
let actions =
self.fsp
.plan_path_mtu_tighten(fips_addr, prior.map(|e| e.mtu), msg.mtu);
if actions.is_empty() {
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
bottleneck_mtu = msg.mtu,
existing = prior.map(|e| e.mtu).unwrap_or(msg.mtu),
"Reactive MtuExceeded: keeping tighter existing path_mtu_lookup value"
);
}
for action in actions {
if let FspAction::TightenPathMtuLookup { fips_addr, mtu } = action {
// Held, not expiring. The admission gate above requires
// a session for the named destination, and that
// session's teardown releases this entry. Nothing
// re-sends the signal once traffic is sized to fit, so a
// deadline would drop a genuine persistent bottleneck
// and start the next flow at the conservative ceiling.
map.insert(fips_addr, crate::upper::tun::PathMtuEntry::held(mtu));
debug!(
dest = %peer_name,
fips_addr = %fips_addr,
bottleneck_mtu = msg.mtu,
prior = ?prior,
map_len = map.len(),
"Reactive MtuExceeded: tightened path_mtu_lookup"
);
}
}
}
Err(e) => {
warn!(
dest = %peer_name,
fips_addr = %fips_addr,
bottleneck_mtu = msg.mtu,
error = %e,
"path_mtu_lookup write lock poisoned; reactive MtuExceeded not reflected"
);
}
}
}
// === Session Initiation (Send Path) ===
/// Initiate an end-to-end session with a remote node.
///
/// Creates a Noise XK handshake as initiator, wraps msg1 in a
/// SessionSetup, encapsulates in a SessionDatagram, and routes
/// toward the destination.
/// Whether a session for `addr` may be created, given the table cap.
///
/// Returns true when an entry already exists, since admitting it cannot
/// grow the table. Counts its own refusals, so the two reasons are
/// distinguishable without turning on debug logging.
pub(in crate::node) fn admit_new_session(&mut self, addr: &NodeAddr) -> bool {
let max_sessions = self.config().node.limits.max_sessions;
if max_sessions == 0 || self.sessions.contains_key(addr) {
return true;
}
if self.sessions.len() >= max_sessions {
debug!(
src = %self.peer_display_name(addr),
sessions = self.sessions.len(),
max_sessions = max_sessions,
"Session table full, refusing to create a session"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::TableFull));
return false;
}
// Half-open entries are unauthenticated and are reaped after
// `handshake_timeout_secs`, so they are the cheap half of the table
// to fill. Holding them to a share keeps room for peers that
// complete. The outer length test makes the scan unreachable below
// the share, and the table is itself bounded by the cap above.
// At least one, or a table capped at one would admit no inbound
// session at all rather than one.
let half_open_share = (max_sessions / HALF_OPEN_SHARE_DIVISOR).max(1);
if self.sessions.len() >= half_open_share {
let half_open = self
.sessions
.values()
.filter(|e| e.is_awaiting_msg3())
.count();
if half_open >= half_open_share {
debug!(
src = %self.peer_display_name(addr),
half_open = half_open,
half_open_share = half_open_share,
"Half-open session share exhausted, refusing to create a session"
);
self.stats_mut()
.record_reject(RejectReason::Session(SessionReject::HalfOpenFull));
return false;
}
}
true
}
pub(in crate::node) async fn initiate_session(
&mut self,
dest_addr: NodeAddr,
dest_pubkey: PublicKey,
) -> Result<(), NodeError> {
// Check for existing session
if let Some(existing) = self.sessions.get(&dest_addr)
&& (existing.is_established() || existing.is_initiating())
{
return Ok(());
}
// Create Noise XK initiator handshake
// This frame's own copy of the node's long-term private key; the
// handshake state keeps its own and clears that on drop.
let mut our_keypair = self.identity().keypair();
let mut handshake = HandshakeState::new_xk_initiator(our_keypair, dest_pubkey);
our_keypair.non_secure_erase();
handshake.set_local_epoch(self.startup_epoch());
let msg1 = handshake
.write_xk_message_1()
.map_err(|e| NodeError::SendFailed {
node_addr: dest_addr,
reason: format!("Noise XK msg1 generation failed: {}", e),
})?;
// Build SessionSetup with coordinates
let our_coords = self.tree_state.my_coords().clone();
let dest_coords = self.get_dest_coords(&dest_addr);
let setup = SessionSetup::new(our_coords, dest_coords).with_handshake(msg1);
let setup_payload = setup.encode();
// Wrap in SessionDatagram
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, dest_addr, setup_payload.clone())
.with_ttl(self.config().node.session.default_ttl);
// Route toward destination
self.send_session_datagram(&mut datagram).await?;
// Register destination identity for TUN → session routing
self.register_identity(dest_addr, dest_pubkey);
// Store session entry with handshake payload for potential resend
let now_ms = Self::now_ms();
let resend_interval = self.config().node.rate_limit.handshake_resend_interval_ms;
let mut entry = SessionEntry::new(
dest_addr,
dest_pubkey,
EndToEndState::Initiating(handshake),
now_ms,
true,
);
entry.set_handshake_payload(setup_payload, now_ms + resend_interval);
self.sessions.insert(dest_addr, entry);
debug!(dest = %self.peer_display_name(&dest_addr), "Session initiation started");
Ok(())
}
/// Send application data over an established session.
///
/// Uses the FSP pipeline: builds a 12-byte cleartext header (used as AAD),
/// prepends the 6-byte inner header to the plaintext, encrypts with AAD,
/// optionally inserts cleartext coords, and wraps in a SessionDatagram.
///
/// The `src_port` and `dst_port` identify the service. A 4-byte port header
/// `[src_port:2 LE][dst_port:2 LE]` is prepended to `payload` inside the
/// AEAD envelope. The receiver dispatches by `dst_port`.
pub(in crate::node) async fn send_session_data(
&mut self,
dest_addr: &NodeAddr,
src_port: u16,
dst_port: u16,
payload: &[u8],
) -> Result<(), NodeError> {
let now_ms = Self::now_ms();
// First borrow: read session metadata (NLL releases before coord decision)
let entry = self
.sessions
.get(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
let wants_coords = entry.coords_warmup_remaining() > 0;
let timestamp = entry.session_timestamp(now_ms);
let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit());
if !entry.is_established() {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "session not established".into(),
});
}
// Build port-prefixed plaintext: [src_port:2 LE][dst_port:2 LE][payload...]
let mut port_payload = Vec::with_capacity(FSP_PORT_HEADER_SIZE + payload.len());
port_payload.extend_from_slice(&src_port.to_le_bytes());
port_payload.extend_from_slice(&dst_port.to_le_bytes());
port_payload.extend_from_slice(payload);
// Build inner plaintext (doesn't depend on counter)
let msg_type = SessionMessageType::DataPacket.to_byte(); // 0x10
let inner_flags = FspInnerFlags { spin_bit }.to_byte();
let inner_plaintext =
fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &port_payload);
// With no coordinates cached for the destination, send without CP
// and leave the warmup budget unspent: our own coordinates in the
// destination's slot would be filed under its address by every
// receiver, and the budget is better spent on the first frames after
// discovery refills the cache.
let cached_dst = if wants_coords {
self.cached_dest_coords(dest_addr)
} else {
None
};
let warming = cached_dst.is_some();
// Determine whether coords fit within transport MTU.
// If not, send standalone CoordsWarmup before the data packet.
let (include_coords, my_coords, dest_coords) = if let Some(dst) = cached_dst {
let src = self.tree_state.my_coords().clone();
let coords_size = coords_wire_size(&src) + coords_wire_size(&dst);
let total_wire =
FIPS_OVERHEAD as usize + FSP_PORT_HEADER_SIZE + coords_size + payload.len();
if total_wire <= self.transport_mtu() as usize {
(true, Some(src), Some(dst))
} else {
// Coords don't fit piggybacked — send standalone CoordsWarmup first
if let Err(e) = self.send_coords_warmup(dest_addr).await {
debug!(dest = %self.peer_display_name(dest_addr), error = %e,
"Failed to send standalone CoordsWarmup before data packet");
}
(false, None, None)
}
} else {
(false, None, None)
};
// Decrement warmup counter if we sent coords (piggybacked or standalone)
if warming && let Some(entry) = self.sessions.get_mut(dest_addr) {
entry.set_coords_warmup_remaining(entry.coords_warmup_remaining() - 1);
}
// Build FSP flags (CP flag if coords, K-bit for key epoch)
let mut flags = if include_coords { FSP_FLAG_CP } else { 0 };
if let Some(entry) = self.sessions.get(dest_addr)
&& entry.current_k_bit()
{
flags |= FSP_FLAG_K;
}
// ── Pipelined FSP+FMP fast path (unix + UDP) ──
// Both AEAD layers + the sendmmsg syscall run on the encrypt
// worker. The rx_loop only builds the wire buffer + reserves
// counters on both sessions. Falls through to the legacy
// sync FSP encrypt + send_session_datagram path on any
// prereq miss (non-UDP, no worker, no cipher, …).
#[cfg(unix)]
if self
.try_send_session_data_pipelined(PipelinedSend {
dest_addr,
payload,
now_ms,
timestamp,
fsp_flags: flags,
inner_plaintext: &inner_plaintext,
my_coords: my_coords.as_ref(),
dest_coords: dest_coords.as_ref(),
})
.await?
{
return Ok(());
}
// Borrow session for counter + encryption (after potential standalone send)
let entry = self
.sessions
.get_mut(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
let session = match entry.state_mut() {
EndToEndState::Established(s) => s,
_ => {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "session not established".into(),
});
}
};
let counter = session.current_send_counter();
// Build 12-byte FSP header (used as AAD for AEAD)
let payload_len = inner_plaintext.len() as u16;
let header = build_fsp_header(counter, flags, payload_len);
// Encrypt with AAD binding to the FSP header
let ciphertext = session
.encrypt_with_aad(&inner_plaintext, &header)
.map_err(|e| NodeError::SendFailed {
node_addr: *dest_addr,
reason: format!("session encrypt failed: {}", e),
})?;
// Assemble: header(12) + [coords] + ciphertext
let mut fsp_payload = Vec::with_capacity(FSP_HEADER_SIZE + ciphertext.len() + 200);
fsp_payload.extend_from_slice(&header);
if let (Some(src), Some(dst)) = (&my_coords, &dest_coords) {
encode_coords(src, &mut fsp_payload);
encode_coords(dst, &mut fsp_payload);
}
fsp_payload.extend_from_slice(&ciphertext);
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *dest_addr, fsp_payload)
.with_ttl(self.config().node.session.default_ttl);
self.send_session_datagram(&mut datagram).await?;
// Re-borrow after send (which borrowed &mut self)
if let Some(entry) = self.sessions.get_mut(dest_addr) {
entry.record_sent(payload.len());
if let Some(mmp) = entry.mmp_mut() {
mmp.sender.record_sent(counter, timestamp, ciphertext.len());
}
entry.touch(now_ms);
}
Ok(())
}
/// Pipelined send: FSP+FMP AEAD + sendmsg on the encrypt worker.
///
/// Returns `Ok(true)` when the worker accepted the job (caller is
/// done), `Ok(false)` on prereq miss (non-UDP, no worker, no
/// cipher) so the caller falls back to the legacy synchronous
/// path. `Err` on routing / state errors that prevent any send.
///
/// Wire layout built directly (no intermediate `inner_plaintext`
/// or `fsp_payload` Vecs):
/// ```text
/// [16 FMP header][8 link-ts][1 SessionDatagram][1 ttl][2 mtu]
/// [16 src_addr][16 dest_addr][12 FSP header][coords?]
/// [FSP plaintext] <- worker seals here, appends FSP tag,
/// then seals the full FMP plaintext and
/// appends the FMP tag.
/// ```
#[cfg(unix)]
async fn try_send_session_data_pipelined(
&mut self,
send: PipelinedSend<'_>,
) -> Result<bool, NodeError> {
let dest_addr = send.dest_addr;
let Some(workers) = self.supervisor.encrypt_workers.as_ref().cloned() else {
return Ok(false);
};
let Some(next_hop_addr) = self.find_next_hop(dest_addr).map(|peer| *peer.node_addr())
else {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no route to destination".into(),
});
};
// Read the next hop's transport / link MTU for the
// SessionDatagram's path_mtu field. Saves a path_mtu round-
// trip vs the legacy path where we'd seed in
// send_session_datagram.
let mut path_mtu = u16::MAX;
if let Some(peer) = self.peers.get(&next_hop_addr)
&& let Some(tid) = peer.transport_id()
&& let Some(transport) = self.transports.get(&tid)
{
if let Some(addr) = peer.current_addr() {
path_mtu = path_mtu.min(transport.link_mtu(addr));
} else {
path_mtu = path_mtu.min(transport.mtu());
}
}
// Extract next-hop info + FMP cipher in one peers/sessions borrow scope.
let (their_index, transport_id, remote_addr, timestamp_ms, fmp_flags, fmp_cipher) = {
let peer = self
.peers
.get_mut(&next_hop_addr)
.ok_or(NodeError::PeerNotFound(next_hop_addr))?;
let their_index = peer.their_index().ok_or_else(|| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: "no their_index".into(),
})?;
let transport_id = peer.transport_id().ok_or_else(|| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: "no transport_id".into(),
})?;
let remote_addr =
peer.current_addr()
.cloned()
.ok_or_else(|| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: "no current_addr".into(),
})?;
let timestamp_ms = peer.session_elapsed_ms();
let sp_flag = peer.mmp().map(|m| m.spin_bit.tx_bit()).unwrap_or(false);
let mut fmp_flags = if sp_flag { FLAG_SP } else { 0 };
if peer.current_k_bit() {
fmp_flags |= FLAG_KEY_EPOCH;
}
let session = peer
.noise_session_mut()
.ok_or_else(|| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: "no noise session".into(),
})?;
let Some(fmp_cipher) = session.send_cipher_clone() else {
return Ok(false);
};
(
their_index,
transport_id,
remote_addr,
timestamp_ms,
fmp_flags,
fmp_cipher,
)
};
#[cfg(any(target_os = "linux", target_os = "macos"))]
let connected_socket = self
.peers
.get(&next_hop_addr)
.and_then(|peer| peer.connected_udp());
// Need a UDP transport (this whole path is sendmsg-on-raw-fd).
let transport = self
.transports
.get(&transport_id)
.ok_or(NodeError::TransportNotFound(transport_id))?;
let TransportHandle::Udp(udp) = transport else {
return Ok(false);
};
let socket_addr_opt = {
#[cfg(any(target_os = "linux", target_os = "macos"))]
{
match connected_socket.as_ref() {
Some(s) => Some(s.peer_addr()),
None => udp.resolve_for_off_task(&remote_addr).await.ok(),
}
}
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
{
udp.resolve_for_off_task(&remote_addr).await.ok()
}
};
let Some(socket_addr) = socket_addr_opt else {
return Ok(false);
};
let Some(socket) = udp.async_socket() else {
return Ok(false);
};
let stats = udp.stats().clone();
// FSP cipher + counter — separate session from next-hop FMP session.
let (fsp_counter, fsp_cipher) = {
let entry = self
.sessions
.get_mut(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
if let Some(mmp) = entry.mmp_mut() {
mmp.path_mtu.seed_source_mtu(path_mtu);
}
let session = match entry.state_mut() {
EndToEndState::Established(s) => s,
_ => {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "session not established".into(),
});
}
};
let Some(fsp_cipher) = session.send_cipher_clone() else {
return Ok(false);
};
let counter = session
.take_send_counter()
.map_err(|e| NodeError::SendFailed {
node_addr: *dest_addr,
reason: format!("session counter reservation failed: {}", e),
})?;
(counter, fsp_cipher)
};
// FSP header (the AAD for the inner AEAD seal).
let fsp_header = build_fsp_header(
fsp_counter,
send.fsp_flags,
send.inner_plaintext.len() as u16,
);
let coords_size = match (send.my_coords, send.dest_coords) {
(Some(src), Some(dst)) => coords_wire_size(src) + coords_wire_size(dst),
_ => 0,
};
// FMP inner = [link-ts u64][SessionDatagram-encoded after msg_type] + [FSP-encrypted blob]
// SessionDatagram-encoded includes: [0x00 type][ttl][mtu][src][dest][fsp_payload]
// fsp_payload = [fsp_header][coords][fsp_plaintext][16-byte FSP tag]
let link_plaintext_len = SESSION_DATAGRAM_HEADER_SIZE
+ FSP_HEADER_SIZE
+ coords_size
+ send.inner_plaintext.len();
let fmp_inner_len = 4 + link_plaintext_len + crate::noise::TAG_SIZE;
// FMP counter + header (the AAD for the outer AEAD seal).
let fmp_counter = {
let peer = self
.peers
.get_mut(&next_hop_addr)
.ok_or(NodeError::PeerNotFound(next_hop_addr))?;
let session = peer
.noise_session_mut()
.ok_or_else(|| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: "no noise session".into(),
})?;
session
.take_send_counter()
.map_err(|e| NodeError::SendFailed {
node_addr: next_hop_addr,
reason: format!("counter reservation failed: {}", e),
})?
};
let fmp_header =
build_established_header(their_index, fmp_counter, fmp_flags, fmp_inner_len as u16);
// Build the wire buffer once, in place. The two FSP/FMP tags
// are appended by the worker after the seals — reserve their
// capacity here so the worker doesn't have to re-grow.
let wire_capacity = ESTABLISHED_HEADER_SIZE + fmp_inner_len + crate::noise::TAG_SIZE;
let mut wire_buf = Vec::with_capacity(wire_capacity);
wire_buf.extend_from_slice(&fmp_header);
wire_buf.extend_from_slice(&timestamp_ms.to_le_bytes());
// SessionDatagram-encoded layout (matches `SessionDatagram::encode`):
wire_buf.push(LinkMessageType::SessionDatagram.to_byte());
wire_buf.push(self.config().node.session.default_ttl);
wire_buf.extend_from_slice(&path_mtu.to_le_bytes());
wire_buf.extend_from_slice(self.node_addr().as_bytes());
wire_buf.extend_from_slice(dest_addr.as_bytes());
// FSP layer (worker seals on `inner_plaintext` portion):
let fsp_aad_offset = wire_buf.len();
wire_buf.extend_from_slice(&fsp_header);
if let (Some(src), Some(dst)) = (send.my_coords, send.dest_coords) {
encode_coords(src, &mut wire_buf);
encode_coords(dst, &mut wire_buf);
}
let fsp_plaintext_offset = wire_buf.len();
wire_buf.extend_from_slice(send.inner_plaintext);
// Stats update — predict size exactly (ChaCha20-Poly1305 tag
// is constant 16 bytes).
let predicted_bytes = wire_capacity;
if let Some(peer) = self.peers.get_mut(&next_hop_addr) {
peer.link_stats_mut().record_sent(predicted_bytes);
if let Some(mmp) = peer.mmp_mut() {
mmp.sender
.record_sent(fmp_counter, timestamp_ms, predicted_bytes);
}
}
self.metrics()
.forwarding
.record_originated(link_plaintext_len + crate::noise::TAG_SIZE);
if let Some(entry) = self.sessions.get_mut(dest_addr) {
entry.record_sent(send.payload.len());
entry.record_sent_wire_len(wire_capacity);
if let Some(mmp) = entry.mmp_mut() {
mmp.sender.record_sent(
fsp_counter,
send.timestamp,
send.inner_plaintext.len() + crate::noise::TAG_SIZE,
);
}
entry.touch(send.now_ms);
}
let dispatched = workers.dispatch(crate::node::encrypt_worker::FmpSendJob {
cipher: fmp_cipher,
counter: fmp_counter,
wire_buf,
fsp_seal: Some(crate::node::encrypt_worker::FspSealJob {
cipher: fsp_cipher,
counter: fsp_counter,
aad_offset: fsp_aad_offset,
plaintext_offset: fsp_plaintext_offset,
}),
socket,
dest_addr: socket_addr,
#[cfg(any(target_os = "linux", target_os = "macos"))]
connected_socket,
stats,
// Bulk endpoint data: drop on UDP backpressure so the
// worker queue keeps moving instead of stranding under
// sustained congestion.
drop_on_backpressure: true,
queued_at: None,
});
if let Err(job) = dispatched {
self.send_refused_job_inline(*job, transport_id, &remote_addr, next_hop_addr)
.await;
}
Ok(true)
}
/// Seal and send a job the encrypt worker for its next hop refused
/// because that worker has exited, using the FSP and FMP counters the
/// job already reserved so neither counter is skipped. Stats were
/// recorded before dispatch and now describe this packet.
///
/// A failure is logged and swallowed, as the worker does with its own:
/// the caller sees the same result whichever of the two sent the packet.
#[cfg(unix)]
async fn send_refused_job_inline(
&self,
job: crate::node::encrypt_worker::FmpSendJob,
transport_id: crate::transport::TransportId,
remote_addr: &crate::transport::TransportAddr,
next_hop_addr: NodeAddr,
) {
let wire = match job.seal_inline() {
Ok(wire) => wire,
Err(error) => {
debug!(next_hop = %next_hop_addr, %error, "Inline seal of session data failed");
return;
}
};
let Some(transport) = self.transports.get(&transport_id) else {
debug!(next_hop = %next_hop_addr, "Transport gone before inline send of session data");
return;
};
if let Err(error) = transport.send_existing(remote_addr, &wire).await {
debug!(next_hop = %next_hop_addr, %error, "Inline send of session data failed");
}
}
/// Send an IPv6 packet through the IPv6 shim (port 256) with header compression.
///
/// Compresses the IPv6 header (format 0x00), then sends via `send_session_data`
/// with `src_port=256, dst_port=256`.
pub(in crate::node) async fn send_ipv6_packet(
&mut self,
dest_addr: &NodeAddr,
ipv6_packet: &[u8],
) -> Result<(), NodeError> {
let compressed = crate::upper::ipv6_shim::compress_ipv6(ipv6_packet).ok_or_else(|| {
NodeError::SendFailed {
node_addr: *dest_addr,
reason: "IPv6 header compression failed".into(),
}
})?;
self.send_session_data(
dest_addr,
FSP_PORT_IPV6_SHIM,
FSP_PORT_IPV6_SHIM,
&compressed,
)
.await
}
/// Send a non-data session message (reports, notifications) over an established session.
///
/// Similar to `send_session_data()` but:
/// - Takes an explicit `msg_type` byte (0x11, 0x12, 0x13, etc.)
/// - Never includes COORDS_PRESENT (reports are lightweight)
/// - Reads spin bit from MMP state for the inner header
/// - Records the send in MMP sender state
pub(in crate::node) async fn send_session_msg(
&mut self,
dest_addr: &NodeAddr,
msg_type: u8,
payload: &[u8],
) -> Result<(), NodeError> {
let now_ms = Self::now_ms();
// Read spin bit and session timestamp from entry
let entry = self
.sessions
.get(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
let timestamp = entry.session_timestamp(now_ms);
let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit());
// Build inner flags with spin bit
let inner_flags = FspInnerFlags { spin_bit }.to_byte();
// Get mutable access for encryption
let entry = self
.sessions
.get_mut(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
// Read K-bit before mutable borrow of session state
let k_flags = if entry.current_k_bit() { FSP_FLAG_K } else { 0 };
let session = match entry.state_mut() {
EndToEndState::Established(s) => s,
_ => {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "session not established".into(),
});
}
};
let counter = session.current_send_counter();
// FSP inner header + plaintext
let inner_plaintext = fsp_prepend_inner_header(timestamp, msg_type, inner_flags, payload);
// Build 12-byte FSP header (K-bit for key epoch, no CP for reports)
let payload_len = inner_plaintext.len() as u16;
let header = build_fsp_header(counter, k_flags, payload_len);
// Encrypt with AAD
let ciphertext = session
.encrypt_with_aad(&inner_plaintext, &header)
.map_err(|e| NodeError::SendFailed {
node_addr: *dest_addr,
reason: format!("session encrypt failed: {}", e),
})?;
// Assemble: header(12) + ciphertext (no coords)
let mut fsp_payload = Vec::with_capacity(FSP_HEADER_SIZE + ciphertext.len());
fsp_payload.extend_from_slice(&header);
fsp_payload.extend_from_slice(&ciphertext);
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *dest_addr, fsp_payload)
.with_ttl(self.config().node.session.default_ttl);
self.send_session_datagram(&mut datagram).await?;
// Record in MMP sender state (no touch — MMP reports don't reset idle timer)
if let Some(entry) = self.sessions.get_mut(dest_addr)
&& let Some(mmp) = entry.mmp_mut()
{
mmp.sender.record_sent(counter, timestamp, ciphertext.len());
}
Ok(())
}
/// Send a standalone CoordsWarmup message to warm transit node caches.
///
/// Constructs an encrypted FSP message with CP flag set and
/// msg_type=CoordsWarmup. Transit nodes extract the cleartext
/// coordinates via `try_warm_coord_cache()` (same as CP-flagged data
/// packets). The encrypted inner payload is the 6-byte inner header
/// with no application data.
pub(in crate::node) async fn send_coords_warmup(
&mut self,
dest_addr: &NodeAddr,
) -> Result<(), NodeError> {
let now_ms = Self::now_ms();
// A warmup's only content is the two coordinates; with none cached
// for the destination, ours would stand in for its own.
let Some(dest_coords) = self.cached_dest_coords(dest_addr) else {
trace!(
dest = %self.peer_display_name(dest_addr),
"No cached coordinates for destination, skipping CoordsWarmup"
);
return Ok(());
};
let my_coords = self.tree_state.my_coords().clone();
// Read session metadata
let entry = self
.sessions
.get(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
let timestamp = entry.session_timestamp(now_ms);
let spin_bit = entry.mmp().is_some_and(|m| m.spin_bit.tx_bit());
// Get mutable access for encryption
let entry = self
.sessions
.get_mut(dest_addr)
.ok_or_else(|| NodeError::SendFailed {
node_addr: *dest_addr,
reason: "no session".into(),
})?;
let session = match entry.state_mut() {
EndToEndState::Established(s) => s,
_ => {
return Err(NodeError::SendFailed {
node_addr: *dest_addr,
reason: "session not established".into(),
});
}
};
let counter = session.current_send_counter();
// FSP inner header only, no body payload
let msg_type = SessionMessageType::CoordsWarmup.to_byte();
let inner_flags = FspInnerFlags { spin_bit }.to_byte();
let inner_plaintext = fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &[]);
// Build FSP header with CP flag
let payload_len = inner_plaintext.len() as u16;
let header = build_fsp_header(counter, FSP_FLAG_CP, payload_len);
// Encrypt with AAD
let ciphertext = session
.encrypt_with_aad(&inner_plaintext, &header)
.map_err(|e| NodeError::SendFailed {
node_addr: *dest_addr,
reason: format!("session encrypt failed: {}", e),
})?;
// Assemble: header(12) + coords + ciphertext
let coords_size = coords_wire_size(&my_coords) + coords_wire_size(&dest_coords);
let mut fsp_payload = Vec::with_capacity(FSP_HEADER_SIZE + coords_size + ciphertext.len());
fsp_payload.extend_from_slice(&header);
encode_coords(&my_coords, &mut fsp_payload);
encode_coords(&dest_coords, &mut fsp_payload);
fsp_payload.extend_from_slice(&ciphertext);
let my_addr = *self.node_addr();
let mut datagram = SessionDatagram::new(my_addr, *dest_addr, fsp_payload)
.with_ttl(self.config().node.session.default_ttl);
self.send_session_datagram(&mut datagram).await?;
// Record in MMP (infrastructure traffic — no idle timer touch)
if let Some(entry) = self.sessions.get_mut(dest_addr)
&& let Some(mmp) = entry.mmp_mut()
{
mmp.sender.record_sent(counter, timestamp, ciphertext.len());
}
debug!(dest = %self.peer_display_name(dest_addr), "Sent standalone CoordsWarmup");
Ok(())
}
/// Route and send a SessionDatagram through the mesh.
///
/// Finds the next hop for the destination, seeds path_mtu from the
/// first-hop transport MTU, and sends as an encrypted link message.
pub(in crate::node) async fn send_session_datagram(
&mut self,
datagram: &mut SessionDatagram,
) -> Result<(), NodeError> {
let next_hop_addr = match self.find_next_hop(&datagram.dest_addr) {
Some(peer) => *peer.node_addr(),
None => {
return Err(NodeError::SendFailed {
node_addr: datagram.dest_addr,
reason: "no route to destination".into(),
});
}
};
// Seed path_mtu from the first-hop transport MTU (same as forwarding path)
if let Some(peer) = self.peers.get(&next_hop_addr)
&& let Some(tid) = peer.transport_id()
&& let Some(transport) = self.transports.get(&tid)
{
if let Some(addr) = peer.current_addr() {
datagram.path_mtu = datagram.path_mtu.min(transport.link_mtu(addr));
} else {
datagram.path_mtu = datagram.path_mtu.min(transport.mtu());
}
}
// Source-side: seed our PathMtuState.current_mtu from the outbound
// transport MTU so it doesn't stay at u16::MAX until the destination
// sends a PathMtuNotification back.
if let Some(entry) = self.sessions.get_mut(&datagram.dest_addr)
&& let Some(mmp) = entry.mmp_mut()
{
mmp.path_mtu.seed_source_mtu(datagram.path_mtu);
}
let encoded = datagram.encode();
self.send_encrypted_link_message(&next_hop_addr, &encoded)
.await?;
self.metrics().forwarding.record_originated(encoded.len());
// Evidence for the reactive path-MTU carrier. A transit hop
// re-encapsulates what it forwards, so the frame that overflows a
// downstream link is the size this frame is here.
if let Some(entry) = self.sessions.get_mut(&datagram.dest_addr) {
entry.record_sent_wire_len(link_wire_len(encoded.len()));
}
Ok(())
}
/// Look up destination coordinates in the coordinate cache, with no
/// fallback.
///
/// Use this wherever the coordinates go on the wire as the destination's
/// own: a miss must not be filled with ours, which every receiver would
/// file under the destination's address.
pub(in crate::node) fn cached_dest_coords(
&self,
dest: &NodeAddr,
) -> Option<crate::proto::stp::TreeCoordinate> {
self.coord_cache.get(dest, Self::now_ms()).cloned()
}
/// Look up destination coordinates from available caches.
///
/// Returns our own coordinates as a fallback (the SessionSetup will
/// carry src_coords for return path routing; empty dest_coords
/// would fail wire encoding since TreeCoordinate requires ≥1 entry).
pub(in crate::node) fn get_dest_coords(
&self,
dest: &NodeAddr,
) -> crate::proto::stp::TreeCoordinate {
if let Some(coords) = self.cached_dest_coords(dest) {
return coords;
}
// Fallback: use our own coordinates. The SessionSetup dest_coords
// field cannot be empty (wire format requires ≥1 entry). Using our
// own coords is safe — transit routers will still cache them, and
// the destination will return its actual coords in the SessionAck.
self.tree_state.my_coords().clone()
}
/// Current Unix time in milliseconds.
pub(in crate::node) fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
// === TUN Outbound (Data Plane) ===
/// Handle an outbound IPv6 packet from the TUN reader.
///
/// The host-side checks and ICMPv6 replies are in
/// [`crate::ipv6tun::outbound::forward`], which reaches back into the
/// node through its [`Mesh`] impl: the destination prefix is resolved
/// in the identity cache, and the packet is sent on an established
/// session or queued while one is set up.
///
/// Also performs MTU checking: if the packet (plus FIPS overhead) exceeds
/// the transport MTU, an ICMP Packet Too Big message is sent back to the
/// source and the packet is dropped.
pub(in crate::node) async fn handle_tun_outbound(&mut self, ipv6_packet: Vec<u8>) {
crate::ipv6tun::outbound::forward(self, ipv6_packet).await;
}
/// Send a TUN packet to a resolved destination, or queue it while the
/// destination's session is set up.
///
/// Sends through an established session, queues behind one still being
/// set up, and otherwise initiates a session and queues. With no route
/// for the initiation it starts discovery and still queues. Refuses the
/// packet, handing it back, when a new session would exceed the session
/// table.
async fn send_outbound(
&mut self,
dest_addr: NodeAddr,
dest_pubkey: PublicKey,
ipv6_packet: Vec<u8>,
) -> Outcome {
// Check for established session
if let Some(entry) = self.sessions.get(&dest_addr) {
if entry.is_established() {
if let Err(e) = self.send_ipv6_packet(&dest_addr, &ipv6_packet).await {
debug!(dest = %self.peer_display_name(&dest_addr), error = %e, "Failed to send TUN packet via session");
}
return Outcome::Sent;
}
// Session exists but not yet established — queue the packet
self.queue_pending_packet(dest_addr, ipv6_packet);
return Outcome::Queued;
}
// No session, so this one would grow the table. Answer the local
// application the way an unroutable destination is answered rather
// than returning an error from `initiate_session`: the caller reads
// an error as "no route" and responds with a discovery lookup and a
// queued packet, which is outbound traffic on a node already at its
// limit.
if !self.admit_new_session(&dest_addr) {
return Outcome::TableFull(ipv6_packet);
}
// No session: initiate one and queue the packet.
// If session initiation fails (no route), trigger discovery and
// queue the packet for retry when discovery completes.
if let Err(e) = self.initiate_session(dest_addr, dest_pubkey).await {
debug!(dest = %self.peer_display_name(&dest_addr), error = %e, "Failed to initiate session, trying discovery");
self.maybe_initiate_lookup(&dest_addr).await;
self.queue_pending_packet(dest_addr, ipv6_packet);
return Outcome::Queued;
}
self.queue_pending_packet(dest_addr, ipv6_packet);
Outcome::Queued
}
/// Borrow the host-facing ICMPv6 sender: the TUN channel, our address
/// and the Packet Too Big rate limiter.
pub(in crate::node) fn host_icmp(&mut self) -> IcmpContext<'_> {
let our_ipv6 = crate::FipsAddress::from_node_addr(self.node_addr()).to_ipv6();
IcmpContext::new(
self.supervisor.ipv6tun.tun_tx.as_ref(),
our_ipv6,
&mut self.icmp_rate_limiter,
)
}
/// Queue a packet while waiting for session establishment.
fn queue_pending_packet(&mut self, dest_addr: NodeAddr, packet: Vec<u8>) {
// Reject if we already have too many pending destinations
let max_dests = self.config().node.session.pending_max_destinations;
if !self.pending_tun_packets.contains_key(&dest_addr)
&& self.pending_tun_packets.len() >= max_dests
{
return;
}
let per_dest = self.config().node.session.pending_packets_per_dest;
let queue = self.pending_tun_packets.entry(dest_addr).or_default();
crate::proto::fsp::push_bounded_pending(queue, packet, per_dest);
}
/// Flush pending packets for a destination whose session just reached Established.
async fn flush_pending_packets(&mut self, dest_addr: &NodeAddr) {
let packets = match self.pending_tun_packets.remove(dest_addr) {
Some(q) => q,
None => return,
};
for packet in packets {
if let Err(e) = self.send_ipv6_packet(dest_addr, &packet).await {
debug!(dest = %self.peer_display_name(dest_addr), error = %e, "Failed to send queued TUN packet");
break;
}
}
}
/// Retry session initiation after discovery provided coordinates.
///
/// Called when a LookupResponse arrives and we have pending TUN packets
/// or native datagrams for the discovered target. The coord_cache now
/// has coords, so `find_next_hop()` should succeed and the SessionSetup
/// can be sent.
pub(in crate::node) async fn retry_session_after_discovery(&mut self, dest_addr: NodeAddr) {
// Look up the destination's public key from the identity cache
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&dest_addr.as_bytes()[0..15]);
let dest_pubkey = match self.lookup_by_fips_prefix(&prefix) {
Some((_, pk)) => pk,
None => {
debug!(dest = %self.peer_display_name(&dest_addr), "Discovery complete but no identity for session retry");
return;
}
};
// Skip if a session already exists
if let Some(existing) = self.sessions.get(&dest_addr)
&& (existing.is_established() || existing.is_initiating())
{
return;
}
if !self.admit_new_session(&dest_addr) {
return;
}
match self.initiate_session(dest_addr, dest_pubkey).await {
Ok(()) => {
debug!(dest = %self.peer_display_name(&dest_addr), "Session initiated after discovery");
}
Err(e) => {
debug!(dest = %self.peer_display_name(&dest_addr), error = %e, "Session retry after discovery failed");
}
}
}
}
/// The mesh side of TUN outbound forwarding.
impl Mesh for Node {
type Dest = (NodeAddr, PublicKey);
fn ipv6_mtu(&self) -> u16 {
self.effective_ipv6_mtu()
}
fn resolve(&mut self, prefix: &[u8; 15]) -> Option<Route<Self::Dest>> {
// Look up in identity cache
let (dest_addr, dest_pubkey) = self.lookup_by_fips_prefix(prefix)?;
let path_mtu = self
.sessions
.get(&dest_addr)
.filter(|entry| entry.is_established())
.and_then(|entry| entry.mmp())
.map(|mmp| mmp.path_mtu.current_mtu());
Some(Route {
dest: (dest_addr, dest_pubkey),
path_mtu,
})
}
fn send(&mut self, dest: Self::Dest, packet: Vec<u8>) -> impl Future<Output = Outcome> + Send {
let (dest_addr, dest_pubkey) = dest;
self.send_outbound(dest_addr, dest_pubkey, packet)
}
fn icmp(&mut self) -> IcmpContext<'_> {
self.host_icmp()
}
}