Files
fips/src/node/dataplane/forwarding.rs
T
Johnathan Corgan 1d8c22a91b Share one hop-limit rule between the forwarding pre-check and the routing core
The forwarding handler resolves a next hop only for datagrams the routing
core can forward, which keeps the coordinate-cache touch in that
resolution scoped to genuine forwards. Its TTL test (ttl > 1) restated
the core's drop (decrement, then drop at zero) in a different form, and
no test could see the two disagree: reverting the pre-check to its older
ttl != 0 form passed every unit test.

ttl_after_hop now holds the rule. The routing core drops on it,
SessionDatagram::can_forward uses it, and a new
SessionDatagramRef::can_forward is what the handler calls. A unit test
pins the function, a routing test checks the core and can_forward agree
for every TTL, and a handler test checks that a last-hop transit
datagram does not refresh the destination's cached coordinates while a
ttl=2 one does.
2026-10-02 18:14:01 +00:00

626 lines
27 KiB
Rust

//! SessionDatagram forwarding handler.
//!
//! Handles incoming SessionDatagram (0x00) link messages: decodes the
//! envelope, performs coordinate cache warming from plaintext session-layer
//! headers, pre-resolves the next hop for forwardable transit datagrams, and
//! drives the routing core's outcome — local delivery when the datagram is
//! addressed to this node, the transit hop-limit drop, the forward, or the
//! error signal generated on routing failure.
use crate::NodeAddr;
use crate::node::reject::ForwardingReject;
use crate::node::{Node, NodeError, NodeRoutingView};
use crate::proto::fsp::wire::{
FSP_COMMON_PREFIX_SIZE, FSP_HEADER_SIZE, FSP_PHASE_ESTABLISHED, FSP_PHASE_MSG1, FSP_PHASE_MSG2,
FspCommonPrefix, FspEncryptedHeader, parse_encrypted_coords,
};
use crate::proto::fsp::{SessionAck, SessionSetup};
use crate::proto::link::{SessionDatagram, SessionDatagramRef};
use crate::proto::routing::{DropReason, LimitVerdict, NextHop, RouteAction, RouteOutcome};
use crate::proto::stp::TreeCoordinate;
use std::time::{Duration, Instant};
use tracing::{debug, trace, warn};
impl Node {
/// Handle an incoming SessionDatagram from a peer.
///
/// Called by `dispatch_link_message` for msg_type 0x00. The payload
/// has already had its msg_type byte stripped by dispatch.
pub(in crate::node) async fn handle_session_datagram(
&mut self,
from: &NodeAddr,
payload: &[u8],
incoming_ce: bool,
) {
self.metrics().forwarding.record_received(payload.len());
let datagram_ref = match SessionDatagramRef::decode(payload) {
Ok(dg) => dg,
Err(e) => {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::DecodeError, payload.len());
debug!(error = %e, "Malformed SessionDatagram");
return;
}
};
let my_addr = *self.node_addr();
// Coordinate cache warming from plaintext session-layer headers. Runs
// ahead of both the delivery and the TTL decisions the core makes: the
// coords a peer put on the wire are equally valid whichever way those
// go, and the only arrivals this newly warms from are those with an
// exhausted TTL, whose every insert is already achievable at TTL 1.
self.try_warm_coord_cache_ref(&datagram_ref, payload.len());
// Pre-resolve the next hop only for datagrams the core can actually
// forward: not locally destined, and passing `can_forward`, which is
// the core's own hop-limit rule. This keeps `find_next_hop`'s
// coord-cache LRU-touch side effect scoped to genuine forwards, as it
// was when the TTL test ran inline ahead of it. Warming above has
// already run, so the resolution observes freshly cached coords.
let next_hop = if datagram_ref.dest_addr != my_addr && datagram_ref.can_forward() {
self.resolve_next_hop(&datagram_ref.dest_addr)
} else {
None
};
// Read local congestion once and reuse it for both the CE decision
// (via the view) and the congestion metric/log below, keeping
// `detect_congestion` the single source of truth.
let congested = next_hop
.as_ref()
.map(|nh| self.detect_congestion(&nh.addr))
.unwrap_or(false);
// Borrow the routing tables disjointly from `&mut self.routing` for
// the pure decision, then release both before driving the outcome.
let outcome = {
let view = NodeRoutingView {
coord_cache: &self.coord_cache,
peers: &self.peers,
tree_state: &self.tree_state,
congested,
};
self.routing
.route(&datagram_ref, &my_addr, incoming_ce, next_hop, &view)
};
match outcome {
RouteOutcome::Drop {
reason: DropReason::TtlExhausted,
} => {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::TtlExhausted, payload.len());
debug!(
src = %datagram_ref.src_addr,
dest = %datagram_ref.dest_addr,
ttl = datagram_ref.ttl,
"SessionDatagram TTL exhausted, dropping"
);
}
RouteOutcome::DeliverLocal => {
// Local delivery: dispatch to session layer handlers without
// materializing an owned SessionDatagram payload Vec.
self.metrics().forwarding.record_delivered(payload.len());
self.handle_session_payload(
&datagram_ref.src_addr,
from,
datagram_ref.payload,
datagram_ref.path_mtu,
incoming_ce,
)
.await;
}
RouteOutcome::NoRoute => {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::NoRoute, payload.len());
let original = datagram_ref.into_owned();
debug!(
src = %self.peer_display_name(&original.src_addr),
dest = %self.peer_display_name(&original.dest_addr),
bytes = payload.len(),
"Dropping transit SessionDatagram: no route to destination"
);
self.send_routing_error(from, &original).await;
}
RouteOutcome::Forward {
next_hop,
bytes,
outgoing_ce,
} => {
let dest = datagram_ref.dest_addr;
// ECN CE relay: congestion was detected locally above; emit the
// metric and rate-limited log at the transit chokepoint.
if congested {
self.metrics().congestion.congestion_detected.inc();
let now = Instant::now();
let should_log = self
.last_congestion_log
.map(|t| now.duration_since(t) >= Duration::from_secs(5))
.unwrap_or(true);
if should_log {
self.last_congestion_log = Some(now);
debug!(next_hop = %next_hop, "Congestion detected, CE flag set on forwarded packet");
}
}
match self
.send_encrypted_link_message_with_ce(&next_hop, &bytes, outgoing_ce)
.await
{
Err(NodeError::MtuExceeded { mtu, .. }) => {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::MtuExceeded, payload.len());
self.send_mtu_exceeded_error(from, dest, datagram_ref.src_addr, mtu)
.await;
}
Err(e) => {
self.metrics()
.forwarding
.record_reject_bytes(ForwardingReject::SendError, payload.len());
debug!(
next_hop = %next_hop,
dest = %dest,
error = %e,
"Failed to forward SessionDatagram"
);
}
Ok(()) => {
self.metrics().forwarding.record_forwarded(bytes.len());
// Classify this transit forward by route class (partition
// of forwarded_packets). Done here, at the data-plane
// chokepoint, so the error-signal routing callers of
// find_next_hop are excluded.
let class = self.classify_forward(&dest, &next_hop);
self.metrics().forwarding.record_route_class(class);
if outgoing_ce {
self.metrics().congestion.ce_forwarded.inc();
}
}
}
}
}
}
/// Resolve the next hop toward `dest` into its address plus the outgoing
/// link's transport MTU. Returns `None` when there is no route.
///
/// The MTU defaults to `u16::MAX` (a no-op min-fold) when the peer's
/// transport is not resolvable, matching the pre-refactor inline behavior
/// where the MTU `if let` chain simply did not fire.
fn resolve_next_hop(&mut self, dest: &NodeAddr) -> Option<NextHop> {
let addr = *self.find_next_hop(dest)?.node_addr();
let link_mtu = if let Some(peer) = self.peers.get(&addr)
&& let Some(tid) = peer.transport_id()
&& let Some(transport) = self.transports.get(&tid)
{
match peer.current_addr() {
Some(link_addr) => transport.link_mtu(link_addr),
None => transport.mtu(),
}
} else {
u16::MAX
};
Some(NextHop { addr, link_mtu })
}
/// Attempt to warm the coordinate cache from session-layer payload headers.
///
/// Transit routers parse the 4-byte FSP common prefix to identify message
/// type, then extract plaintext coordinate fields from:
/// - SessionSetup (phase 0x1): src_coords + dest_coords
/// - SessionAck (phase 0x2): src_coords
/// - Encrypted with CP flag (phase 0x0): cleartext coords between header and ciphertext
///
/// Decode failures are logged and silently ignored — they don't block
/// forwarding.
///
/// `outer_len` is the length of the msg_type-stripped `SessionDatagram`
/// buffer this view was decoded from. It is carried in rather than
/// reconstructed from the header size so the malformed-frame byte counter
/// measures the same population as its siblings — which are charged the
/// outer slice — instead of the inner FSP payload.
/// Warm one coordinate-cache entry from a plaintext session header, after
/// the two write-side sanity checks.
///
/// The key and the value both come off the wire unauthenticated, so this
/// is the only place a warm write can be filtered at all. Two checks, and
/// they are deliberately of different strengths:
///
/// **Foreign root: refused.** A coordinate under a root other than ours
/// can never route. `StpState::find_next_hop` returns `None` outright on a
/// root mismatch, and the bloom fallback compares against a `my_distance`
/// of `usize::MAX`, so no candidate is ever strictly closer. Caching one
/// therefore buys nothing and costs something real: the entry's presence
/// is what `synth_routing_error` reads to choose `PathBroken` over
/// `CoordsRequired`, so a foreign-root plant turns this node into a
/// one-packet reflector aimed at whatever source the datagram claimed.
/// `CoordCache::invalidate_other_roots` already applies this same
/// invariant whenever our own tree position moves; this applies it at
/// write time instead of waiting for the next move.
///
/// **Key mismatch: counted only.** A coordinate whose first element is not
/// the address it is filed under is wrong, but refusing it here would also
/// refuse a write honest nodes make: a sender whose own cache missed puts
/// its *own* coordinates in `SessionSetup.dest_coords`, by way of
/// `get_dest_coords`. What that costs a transit node on first contact is
/// not established, so this counts and does not refuse. It is **not** a
/// security check either way — an attacker satisfies it by naming the
/// victim as its own child, which is the forgery worth making.
fn warm_coord(&mut self, key: NodeAddr, coords: TreeCoordinate, now_ms: u64) {
if coords.root_id() != self.tree_state.my_coords().root_id() {
self.metrics().forwarding.record_warm_foreign_root();
trace!(addr = %key, "Warm write names a foreign root; not caching");
return;
}
if *coords.node_addr() != key {
self.metrics().forwarding.record_warm_key_mismatch();
}
self.insert_coord_hint(key, coords, now_ms);
}
fn try_warm_coord_cache_ref(&mut self, datagram: &SessionDatagramRef<'_>, outer_len: usize) {
let prefix = match FspCommonPrefix::parse(datagram.payload) {
Some(p) => p,
None => return,
};
let inner = &datagram.payload[FSP_COMMON_PREFIX_SIZE..];
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
match prefix.phase {
FSP_PHASE_MSG1 => match SessionSetup::decode(inner) {
Ok(setup) => {
self.warm_coord(datagram.src_addr, setup.src_coords, now_ms);
self.warm_coord(datagram.dest_addr, setup.dest_coords, now_ms);
debug!(
src = %datagram.src_addr,
dest = %datagram.dest_addr,
"Cached coords from SessionSetup"
);
}
Err(e) => {
debug!(error = %e, "Failed to decode SessionSetup for cache warming");
}
},
FSP_PHASE_MSG2 => match SessionAck::decode(inner) {
Ok(ack) => {
self.warm_coord(datagram.src_addr, ack.src_coords, now_ms);
self.warm_coord(datagram.dest_addr, ack.dest_coords, now_ms);
debug!(
src = %datagram.src_addr,
dest = %datagram.dest_addr,
"Cached coords from SessionAck"
);
}
Err(e) => {
debug!(error = %e, "Failed to decode SessionAck for cache warming");
}
},
FSP_PHASE_ESTABLISHED if prefix.has_coords() => {
// CP flag set: coords in cleartext between header and ciphertext.
// Parse coords from the cleartext section after the 12-byte header.
// Re-parse with the encrypted-header parser — the same guard the
// local-delivery path uses — so the slice below is bounded by
// FSP_ENCRYPTED_MIN_SIZE and not by the 4-byte prefix check.
if FspEncryptedHeader::parse(datagram.payload).is_none() {
// Counter is the always-on surface; the debug fields are the
// drill-down that separates a short frame from a bad version
// or a U-flagged one. The level stays at debug: any peer past
// the handshake can drive this at line rate.
self.metrics().forwarding.record_warm_malformed(outer_len);
debug!(
len = datagram.payload.len(),
outer_len,
version = prefix.version,
flags = prefix.flags,
"Not a well-formed encrypted FSP message; not warming coords"
);
return;
}
let coord_data = &datagram.payload[FSP_HEADER_SIZE..];
match parse_encrypted_coords(coord_data) {
Ok((src_coords, dest_coords, _bytes_consumed)) => {
if let Some(coords) = src_coords {
self.warm_coord(datagram.src_addr, coords, now_ms);
}
if let Some(coords) = dest_coords {
self.warm_coord(datagram.dest_addr, coords, now_ms);
}
debug!(
src = %datagram.src_addr,
dest = %datagram.dest_addr,
"Cached coords from encrypted message"
);
}
Err(e) => {
debug!(error = %e, "Failed to parse coords for cache warming");
}
}
}
_ => {
// Phase 0x0 without CP, error signals, unknown: no coords to cache
}
}
}
/// Spend one peer-budget token, after a gate further down has admitted.
///
/// The budget is keyed on the authenticated link peer the frame arrived
/// over, which is the one value at the emission point a sender cannot
/// mint: every field of the datagram itself is chosen by whoever sent it,
/// so a per-destination or per-source gate is escaped by varying the field
/// it keys on.
///
/// Peek and commit are separate because the per-destination interval gate
/// lives inside `routing::synth_routing_error` and runs after this.
/// Charging a suppressed signal would let a single unroutable destination
/// behind a high-fanout peer spend that peer's whole budget on emissions
/// nothing sends, silencing every other destination behind it.
fn commit_error_emission(&mut self, from: &NodeAddr) {
self.peer_error_budget.commit(from, Instant::now());
}
/// Count what the core's per-destination gate decided about one candidate
/// error signal.
///
/// The three verdicts are counted apart because they mean different
/// things to an operator: `Suppress` is the interval doing its job during
/// an outage, while `AdmitAtCapacity` says the destination map is full and
/// the interval is no longer suppressing anything for this destination, so
/// only the per-peer budget is still bounding emission.
fn record_error_verdict(&mut self, verdict: LimitVerdict) {
match verdict {
LimitVerdict::Suppress => self.metrics().errors.emit_over_dest_interval.inc(),
LimitVerdict::AdmitAtCapacity => self.metrics().errors.emit_limiter_at_capacity.inc(),
LimitVerdict::Admit => {}
}
}
/// Generate and send a routing error signal back to the datagram's source.
///
/// If we have cached coords for the destination, send PathBroken (we know
/// where it is but can't reach it). Otherwise send CoordsRequired (we
/// don't know where it is).
///
/// If we can't route the error back to the source either, drop silently.
/// No cascading errors.
/// `from` is the authenticated link peer the original datagram arrived
/// from, and is what the emission is charged against. It is the one value
/// at this point a sender cannot mint: every field of the datagram itself
/// is chosen by whoever sent it.
async fn send_routing_error(&mut self, from: &NodeAddr, original: &SessionDatagram) {
// Peeked, not spent. The destination gate inside the core may still
// suppress this signal, and charging a suppressed emission would let a
// single unroutable destination behind a high-fanout peer burn that
// peer's whole budget on signals nothing sends.
if !self.peer_error_budget.has_token(from, Instant::now()) {
self.metrics().errors.emit_over_peer_budget.inc();
return;
}
let my_addr = *self.node_addr();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let default_ttl = self.config().node.session.default_ttl;
// Pure decision: rate-limit gate + PathBroken/CoordsRequired choice +
// error-PDU encode. Borrow the routing tables disjointly from
// `&mut self.routing`, then release them before the reverse-hop lookup.
let action = {
let view = NodeRoutingView {
coord_cache: &self.coord_cache,
peers: &self.peers,
tree_state: &self.tree_state,
congested: false,
};
self.routing.synth_routing_error(
&original.dest_addr,
&original.src_addr,
&my_addr,
&view,
now_ms,
default_ttl,
)
};
self.record_error_verdict(action.verdict);
let RouteAction::SendError { toward, bytes } = match action.action {
Some(action) => action,
// Rate limited: drop silently. No cascading errors.
None => return,
};
// Both gates have admitted, so the token peeked above is now spent.
// Charged here rather than at the peek so a destination the core
// suppressed costs the link peer nothing; see
// `commit_error_emission`. A later failure to resolve the reverse hop
// still leaves the token spent, which is deliberate: the work the
// budget bounds is the synthesis this node was induced to perform,
// not whether a hop happened to exist for it.
self.commit_error_emission(from);
// Resolve the reverse link hop only now, after the gate passed, so
// `find_next_hop`'s coord-cache touch keeps its pre-refactor scope.
let next_hop_addr = match self.find_next_hop(&toward) {
Some(peer) => *peer.node_addr(),
None => {
debug!(
src = %original.src_addr,
dest = %original.dest_addr,
"Cannot route error signal back to source, dropping"
);
return;
}
};
if let Err(e) = self
.send_encrypted_link_message(&next_hop_addr, &bytes)
.await
{
debug!(
next_hop = %next_hop_addr,
error = %e,
"Failed to send routing error signal"
);
} else {
debug!(
original_dest = %original.dest_addr,
error_dest = %original.src_addr,
"Sent routing error signal"
);
}
}
/// Generate and send an MtuExceeded error signal back to the datagram's source.
///
/// Called when `send_encrypted_link_message()` fails with
/// `NodeError::MtuExceeded` during forwarding. The signal tells the
/// source the bottleneck MTU so it can immediately reduce its path MTU.
///
/// `dest` is the failed datagram's destination (rate-limit key); `toward`
/// is its source, where the signal is routed back.
///
/// `from` is the authenticated link peer the original datagram arrived
/// from, and is what the emission is charged against. MtuExceeded shares
/// the link peer's budget with the routing errors rather than holding its
/// own: a separate bucket would insulate path-MTU discovery from
/// routing-error pressure, at the cost of a second knob and of letting one
/// peer induce twice the total emission.
async fn send_mtu_exceeded_error(
&mut self,
from: &NodeAddr,
dest: NodeAddr,
toward: NodeAddr,
bottleneck_mtu: u16,
) {
// Peeked, not spent, for the same reason as in `send_routing_error`:
// the per-destination gate inside the core runs below and may still
// suppress this signal.
if !self.peer_error_budget.has_token(from, Instant::now()) {
self.metrics().errors.emit_over_peer_budget.inc();
return;
}
let my_addr = *self.node_addr();
let now_ms = Self::now_ms();
let default_ttl = self.config().node.session.default_ttl;
// Pure decision: rate-limit gate + MtuExceeded PDU + encode.
let action = self.routing.synth_mtu_exceeded(
&dest,
&toward,
&my_addr,
bottleneck_mtu,
now_ms,
default_ttl,
);
self.record_error_verdict(action.verdict);
let RouteAction::SendError { toward, bytes } = match action.action {
Some(action) => action,
// Rate limited: drop silently. No cascading errors.
None => return,
};
// Both gates have admitted; spend the token peeked above.
self.commit_error_emission(from);
// Resolve the reverse link hop only now, after the gate passed, so
// `find_next_hop`'s coord-cache touch keeps its pre-refactor scope.
let next_hop_addr = match self.find_next_hop(&toward) {
Some(peer) => *peer.node_addr(),
None => {
debug!(
src = %toward,
dest = %dest,
"Cannot route MtuExceeded signal back to source, dropping"
);
return;
}
};
if let Err(e) = self
.send_encrypted_link_message(&next_hop_addr, &bytes)
.await
{
debug!(
next_hop = %next_hop_addr,
error = %e,
"Failed to send MtuExceeded error signal"
);
} else {
debug!(
original_dest = %dest,
error_dest = %toward,
bottleneck_mtu,
"Sent MtuExceeded error signal"
);
}
}
/// Detect congestion for CE marking on forwarded datagrams.
///
/// Checks two signal sources:
/// 1. Outgoing link MMP metrics (loss rate, ETX) against configured thresholds
/// 2. Local transport congestion (kernel drops on any transport)
///
/// Returns `true` if any signal indicates congestion.
pub(in crate::node) fn detect_congestion(&self, next_hop: &NodeAddr) -> bool {
if !self.config().node.ecn.enabled {
return false;
}
// Outgoing link MMP metrics
if let Some(peer) = self.peers.get(next_hop)
&& let Some(mmp) = peer.mmp()
{
let metrics = &mmp.metrics;
if metrics.loss_rate() >= self.config().node.ecn.loss_threshold
|| metrics.etx >= self.config().node.ecn.etx_threshold
{
return true;
}
}
// Local transport congestion (kernel drops)
self.transport_drops.values().any(|s| s.dropping)
}
/// Sample transport congestion indicators.
///
/// Called from the tick handler (1s interval). For each transport,
/// queries the cumulative kernel drop counter and sets the `dropping`
/// flag if new drops occurred since the previous sample.
pub(in crate::node) fn sample_transport_congestion(&mut self) {
let mut new_drop_events = Vec::new();
for (&tid, transport) in &self.transports {
let congestion = transport.congestion();
let state = self.transport_drops.entry(tid).or_default();
if let Some(current) = congestion.recv_drops {
let new_drops = current > state.prev_drops;
if new_drops && !state.dropping {
new_drop_events.push(tid);
}
state.dropping = new_drops;
state.prev_drops = current;
}
}
for tid in new_drop_events {
self.metrics().congestion.kernel_drop_events.inc();
warn!(
transport_id = tid.as_u32(),
"Kernel recv drops first observed on transport"
);
}
}
}