Files
fips/src/node/tests/forwarding.rs
T
Johnathan Corgan 5974aeee01 Give path-MTU entries a way back, and charge the malformed counter the whole frame
Three changes to the path-MTU release machinery and the forwarding counters.

The release helper now resets the session own current_mtu alongside the
address-keyed map. It previously reseeded the link MTU and left the
tightened value in place, so a path declared dead recovered only through
the increase ladder: three consecutive identical higher values spanning
two notification intervals.

Entries written by the discovery lookup carrier now carry a deadline and
age out. The sender-binding gate made the path-broken route require a
session, and the two timeout routes fire only on session removal, so an
entry for a destination this node never opens a session with had no
release path at all. One response carrying a floor value pinned that
destination clamp until restart, and an unknown request id still
classifies as originator, so a captured response replayed indefinitely.
Only the lookup-carrier write stamps a deadline; the notification mirror
deliberately does not, which is why the uniform version was dropped.

The malformed-frame byte counter is charged the outer frame like every
sibling counter in the same struct, rather than the inner payload. The
gap was fixed per frame and therefore proportionally largest for the
minimal attack frame the counter exists to make visible.

Green: fmt, build, clippy and test --lib, 1530 passed.
2026-08-15 07:21:45 +00:00

1114 lines
38 KiB
Rust

//! SessionDatagram forwarding tests.
//!
//! Tests for the handle_session_datagram handler including decode errors,
//! TTL enforcement, local delivery, coordinate cache warming, and
//! multi-hop forwarding through live node topologies.
use super::*;
use crate::node::session_wire::{FSP_FLAG_CP, build_fsp_header};
use crate::protocol::{SessionAck, SessionDatagram, SessionSetup, encode_coords};
use crate::tree::TreeCoordinate;
use spanning_tree::{
TestNode, cleanup_nodes, process_available_packets, run_tree_test, verify_tree_convergence,
};
// ============================================================================
// Unit Tests
// ============================================================================
// --- Decode errors ---
#[tokio::test]
async fn test_forwarding_decode_error() {
let mut node = make_node();
let from = make_node_addr(0xAA);
// Too-short payload: should log error and return without panic
node.handle_session_datagram(&from, &[0x00; 5], false).await;
}
// --- TTL ---
#[tokio::test]
async fn test_forwarding_hop_limit_exhausted() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src = make_node_addr(0x01);
let dest = make_node_addr(0x02);
let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(0);
let encoded = dg.encode();
// Dispatch with payload after msg_type byte
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
// No panic, no send (node has no peers)
let fwd = &node.metrics().forwarding;
assert_eq!(
fwd.ttl_exhausted_packets.get(),
1,
"transit ttl=0 should be charged to TtlExhausted"
);
assert_eq!(
fwd.drop_no_route_packets.get(),
0,
"transit ttl=0 should never reach the routing step"
);
}
#[tokio::test]
async fn test_forwarding_ttl_one_local_delivery_is_not_gated() {
// dest == self, so this is local delivery, not transit: the TTL gate
// does not apply and the datagram is handed to the session layer.
let mut node = make_node();
let from = make_node_addr(0xAA);
let my_addr = *node.node_addr();
let src = make_node_addr(0x01);
let dg = SessionDatagram::new(src, my_addr, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(1);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
let fwd = &node.metrics().forwarding;
assert_eq!(fwd.delivered_packets.get(), 1, "ttl=1 should be delivered");
assert_eq!(fwd.ttl_exhausted_packets.get(), 0);
}
/// Acceptance: a datagram addressed to this node with ttl=0 is delivered
/// locally. The TTL governs forwarding, not delivery to the addressed host,
/// so the gate must sit after the local-delivery test — and the
/// `TtlExhausted` reject must not be charged for a delivered datagram.
#[tokio::test]
async fn test_forwarding_ttl_zero_local_delivery_is_not_gated() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let my_addr = *node.node_addr();
let src = make_node_addr(0x01);
let dg = SessionDatagram::new(src, my_addr, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(0);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
let fwd = &node.metrics().forwarding;
assert_eq!(
fwd.delivered_packets.get(),
1,
"ttl=0 addressed to this node must still be delivered locally"
);
assert_eq!(
fwd.ttl_exhausted_packets.get(),
0,
"local delivery must not be charged to the TtlExhausted reject"
);
assert_eq!(fwd.drop_no_route_packets.get(), 0);
}
/// Acceptance: a transit datagram arriving with ttl=1 would leave with ttl=0,
/// so it is dropped here rather than transmitted. Reaching the routing step at
/// all (`drop_no_route`) would mean it had been handed to the forwarder.
#[tokio::test]
async fn test_forwarding_ttl_one_transit_dropped_before_routing() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src = make_node_addr(0x01);
let dest = make_node_addr(0x02);
let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(1);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
let fwd = &node.metrics().forwarding;
assert_eq!(
fwd.ttl_exhausted_packets.get(),
1,
"transit ttl=1 must be dropped as TTL-exhausted, not forwarded"
);
assert_eq!(
fwd.drop_no_route_packets.get(),
0,
"transit ttl=1 must not reach the routing step"
);
assert_eq!(fwd.forwarded_packets.get(), 0);
assert_eq!(fwd.delivered_packets.get(), 0);
}
/// The other side of the same boundary: ttl=2 clears the gate. This node has
/// no peers, so it fails at the routing step instead — which is the evidence
/// that the TTL gate passed it through.
#[tokio::test]
async fn test_forwarding_ttl_two_transit_clears_the_gate() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src = make_node_addr(0x01);
let dest = make_node_addr(0x02);
let dg = SessionDatagram::new(src, dest, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(2);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
let fwd = &node.metrics().forwarding;
assert_eq!(
fwd.ttl_exhausted_packets.get(),
0,
"transit ttl=2 must clear the TTL gate"
);
assert_eq!(
fwd.drop_no_route_packets.get(),
1,
"transit ttl=2 should have reached the routing step and found no route"
);
}
// --- Local delivery ---
#[tokio::test]
async fn test_forwarding_local_delivery() {
let mut node = make_node();
let my_addr = *node.node_addr();
let from = make_node_addr(0xAA);
let dg = SessionDatagram::new(from, my_addr, vec![0x10, 0x00, 0x00, 0x00]);
let encoded = dg.encode();
// Should detect local delivery and return without forwarding
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
}
// --- Direct peer forwarding ---
#[tokio::test]
async fn test_forwarding_direct_peer() {
// Set up a node with one peer. Send a datagram destined for that peer.
// The handler should forward it directly.
let edges = vec![(0, 1)];
let mut nodes = run_tree_test(2, &edges, false).await;
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
// Build a datagram from some external source destined for node 1
let external_src = make_node_addr(0xEE);
let dg = SessionDatagram::new(external_src, node1_addr, vec![0x10, 0x00, 0x00, 0x00]);
let encoded = dg.encode();
// Handle on node 0: should forward to node 1 (direct peer)
nodes[0]
.node
.handle_session_datagram(&node0_addr, &encoded[1..], false)
.await;
// Process packets — node 1 should receive the forwarded datagram
tokio::time::sleep(Duration::from_millis(50)).await;
let count = process_available_packets(&mut nodes).await;
assert!(count > 0, "Expected forwarded packet to arrive at node 1");
cleanup_nodes(&mut nodes).await;
}
// ============================================================================
// Coordinate Cache Warming Tests
// ============================================================================
#[tokio::test]
async fn test_coord_cache_warming_session_setup() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src_addr = make_node_addr(0x01);
let dest_addr = make_node_addr(0x02);
let root_addr = make_node_addr(0xF0);
let src_coords = TreeCoordinate::from_addrs(vec![src_addr, root_addr]).unwrap();
let dest_coords = TreeCoordinate::from_addrs(vec![dest_addr, root_addr]).unwrap();
let setup = SessionSetup::new(src_coords.clone(), dest_coords.clone());
let setup_payload = setup.encode();
let dg = SessionDatagram::new(src_addr, dest_addr, setup_payload);
let encoded = dg.encode();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
// Before: cache is empty
assert!(node.coord_cache().get(&src_addr, now_ms).is_none());
assert!(node.coord_cache().get(&dest_addr, now_ms).is_none());
// Handle the datagram (will be local delivery or no-route, but cache warming
// happens before routing decision)
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
// After: both src and dest coords should be cached
let cached_src = node.coord_cache().get(&src_addr, now_ms);
let cached_dest = node.coord_cache().get(&dest_addr, now_ms);
assert!(cached_src.is_some(), "src_addr coords not cached");
assert!(cached_dest.is_some(), "dest_addr coords not cached");
// Verify the cached coords have the right root
let cached_src = cached_src.unwrap();
let cached_dest = cached_dest.unwrap();
assert_eq!(cached_src.root_id(), &root_addr);
assert_eq!(cached_dest.root_id(), &root_addr);
}
#[tokio::test]
async fn test_coord_cache_warming_session_ack() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src_addr = make_node_addr(0x01);
let dest_addr = make_node_addr(0x02);
let root_addr = make_node_addr(0xF0);
let src_coords = TreeCoordinate::from_addrs(vec![src_addr, root_addr]).unwrap();
let dest_coords = TreeCoordinate::from_addrs(vec![dest_addr, root_addr]).unwrap();
let ack = SessionAck::new(src_coords.clone(), dest_coords.clone());
let ack_payload = ack.encode();
let dg = SessionDatagram::new(src_addr, dest_addr, ack_payload);
let encoded = dg.encode();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
assert!(node.coord_cache().get(&src_addr, now_ms).is_none());
assert!(node.coord_cache().get(&dest_addr, now_ms).is_none());
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
// SessionAck caches both src_coords and dest_coords
let cached_src = node.coord_cache().get(&src_addr, now_ms);
assert!(
cached_src.is_some(),
"src_addr coords not cached from SessionAck"
);
assert_eq!(cached_src.unwrap().root_id(), &root_addr);
let cached_dest = node.coord_cache().get(&dest_addr, now_ms);
assert!(
cached_dest.is_some(),
"dest_addr coords not cached from SessionAck"
);
assert_eq!(cached_dest.unwrap().root_id(), &root_addr);
}
#[tokio::test]
async fn test_coord_cache_warming_encrypted_msg_with_coords() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src_addr = make_node_addr(0x01);
let dest_addr = make_node_addr(0x02);
let root_addr = make_node_addr(0xF0);
let src_coords = TreeCoordinate::from_addrs(vec![src_addr, root_addr]).unwrap();
let dest_coords = TreeCoordinate::from_addrs(vec![dest_addr, root_addr]).unwrap();
// Build FSP encrypted message with CP flag: header(12) + coords + fake_ciphertext
let header = build_fsp_header(0, FSP_FLAG_CP, 20);
let mut data_payload = Vec::new();
data_payload.extend_from_slice(&header);
encode_coords(&src_coords, &mut data_payload);
encode_coords(&dest_coords, &mut data_payload);
data_payload.extend_from_slice(&[0xCC; 36]); // fake ciphertext (20 payload + 16 tag)
let dg = SessionDatagram::new(src_addr, dest_addr, data_payload);
let encoded = dg.encode();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
assert!(node.coord_cache().get(&src_addr, now_ms).is_none());
assert!(node.coord_cache().get(&dest_addr, now_ms).is_none());
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
assert!(
node.coord_cache().get(&src_addr, now_ms).is_some(),
"src coords not cached from encrypted message"
);
assert!(
node.coord_cache().get(&dest_addr, now_ms).is_some(),
"dest coords not cached from encrypted message"
);
// Changing what the malformed counter charges is close enough to changing
// when it fires that the well-formed case is pinned in the same place.
assert_eq!(
node.metrics().forwarding.warm_malformed_packets.get(),
0,
"a well-formed CP datagram must not be counted as an abandoned warm attempt"
);
}
#[tokio::test]
async fn test_coord_cache_warming_encrypted_msg_no_coords() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src_addr = make_node_addr(0x01);
let dest_addr = make_node_addr(0x02);
// Build FSP encrypted message without CP flag: header(12) + fake_ciphertext
let header = build_fsp_header(0, 0, 20);
let mut data_payload = Vec::new();
data_payload.extend_from_slice(&header);
data_payload.extend_from_slice(&[0xCC; 36]); // fake ciphertext (20 payload + 16 tag)
let dg = SessionDatagram::new(src_addr, dest_addr, data_payload);
let encoded = dg.encode();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
assert!(
node.coord_cache().get(&src_addr, now_ms).is_none(),
"Should not cache coords from message without CP flag"
);
assert!(
node.coord_cache().get(&dest_addr, now_ms).is_none(),
"Should not cache coords from message without CP flag"
);
}
/// Acceptance: an inner FSP payload of 4 to 11 bytes with phase 0x0 and the
/// CP flag set is dropped rather than panicking the forwarding path. That
/// window sits between the common prefix parser's 4-byte floor and the
/// 12-byte header slice the warm path takes, so before the fix the first
/// iteration panicked with a range start index out of range.
#[tokio::test]
async fn test_coord_cache_warming_short_inner_payload_is_dropped_not_panic() {
let mut node = make_node();
let from = make_node_addr(0xAA);
let src_addr = make_node_addr(0x01);
let dest_addr = make_node_addr(0x02);
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
for extra in 0..=7 {
let mut data_payload = vec![0x00, FSP_FLAG_CP, 0x00, 0x00];
data_payload.resize(4 + extra, 0x00);
let dg = SessionDatagram::new(src_addr, dest_addr, data_payload).with_ttl(1);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
}
assert!(
node.coord_cache().get(&src_addr, now_ms).is_none(),
"Short inner payload must not warm src coords"
);
assert!(
node.coord_cache().get(&dest_addr, now_ms).is_none(),
"Short inner payload must not warm dest coords"
);
// Anti-vacuity: only a datagram that ran past the warm call reaches the
// TTL gate. `received_packets` is charged before decode and so would
// count a datagram rejected earlier.
assert_eq!(
node.metrics().forwarding.ttl_exhausted_packets.get(),
8,
"each short-inner-payload datagram must run past the warm call to the TTL gate"
);
// Discriminating: separates "the guard fired" from "coords parsed and
// yielded nothing", which the cache assertions above cannot tell apart.
assert_eq!(
node.metrics().forwarding.warm_malformed_packets.get(),
8,
"each short-inner-payload datagram must be counted as an abandoned warm attempt"
);
// Inner lengths 12 to 27 document the new 28-byte floor: they do not
// panic today either, so this half is not discriminating.
for len in 12..=27 {
let mut data_payload = vec![0x00, FSP_FLAG_CP, 0x00, 0x00];
data_payload.resize(len, 0x00);
let dg = SessionDatagram::new(src_addr, dest_addr, data_payload).with_ttl(1);
let encoded = dg.encode();
node.handle_session_datagram(&from, &encoded[1..], false)
.await;
}
assert!(
node.coord_cache().get(&src_addr, now_ms).is_none(),
"Payload below the encrypted minimum must not warm src coords"
);
assert!(
node.coord_cache().get(&dest_addr, now_ms).is_none(),
"Payload below the encrypted minimum must not warm dest coords"
);
assert_eq!(
node.metrics().forwarding.ttl_exhausted_packets.get(),
24,
"every datagram in both loops must reach the TTL gate"
);
assert_eq!(
node.metrics().forwarding.warm_malformed_packets.get(),
24,
"every datagram in both loops must be counted as an abandoned warm attempt"
);
// The byte counter shares a fipstop row with `received_bytes` and
// `decode_error_bytes`, so it has to measure the same population: the
// outer SessionDatagram payload, not the inner FSP one. Two assertions
// produced two different ways, because a single one cannot tell "the
// basis matches" from "two counters are wrong in the same direction".
//
// Self-derived: every one of the 24 datagrams reaches the warm guard, as
// the two packet counts above already pin, and `record_received` charges
// the identical outer slice.
assert_eq!(
node.metrics().forwarding.warm_malformed_bytes.get(),
node.metrics().forwarding.received_bytes.get(),
"the byte counter must be charged the same outer payload as its \
siblings on the same row"
);
// Literal cross-check. The first loop sends inner lengths 4..=11, so
// outer 39..=46, summing to 340; the second sends inner 12..=27, so
// outer 47..=62, summing to 872. Charging the inner payload instead
// reads 60 + 312 = 372, about 15% of the wire volume that arrived.
assert_eq!(
node.metrics().forwarding.warm_malformed_bytes.get(),
1212,
"24 frames of 39..=46 and 47..=62 outer bytes sum to 1212"
);
}
// ============================================================================
// Integration Tests
// ============================================================================
/// Helper: populate all coordinate caches across a set of test nodes.
fn populate_all_coord_caches(nodes: &mut [TestNode]) {
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
// Collect all coords first to avoid borrow conflicts
let all_coords: Vec<(NodeAddr, TreeCoordinate)> = nodes
.iter()
.map(|tn| {
(
*tn.node.node_addr(),
tn.node.tree_state().my_coords().clone(),
)
})
.collect();
for tn in nodes.iter_mut() {
for (addr, coords) in &all_coords {
if addr != tn.node.node_addr() {
tn.node
.coord_cache_mut()
.insert(*addr, coords.clone(), now_ms);
}
}
}
}
#[tokio::test]
async fn test_forwarding_single_hop() {
// 3-node chain: 0 -- 1 -- 2
// Send datagram from node 0 destined for node 2.
// Node 1 should forward it.
let edges = vec![(0, 1), (1, 2)];
let mut nodes = run_tree_test(3, &edges, false).await;
verify_tree_convergence(&nodes);
populate_all_coord_caches(&mut nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let node2_addr = *nodes[2].node.node_addr();
// Build a SessionDatagram from node 0 to node 2
let dg = SessionDatagram::new(
node0_addr,
node2_addr,
vec![0x10, 0x00, 0x04, 0x00, 1, 2, 3, 4],
);
let encoded = dg.encode();
// Send from node 0 to node 1 (the first hop)
nodes[0]
.node
.send_encrypted_link_message(&node1_addr, &encoded)
.await
.unwrap();
// Process: node 1 receives, decrypts, dispatches to handler, forwards to node 2
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
// Give time for the forwarded packet to arrive at node 2
tokio::time::sleep(Duration::from_millis(50)).await;
let count = process_available_packets(&mut nodes).await;
// Node 2 should have received the forwarded datagram
// (it sees dest_addr == self, treats as local delivery)
// We verify the chain completed by checking packets were processed.
assert!(count > 0, "Expected forwarded packet at node 2");
cleanup_nodes(&mut nodes).await;
}
#[tokio::test]
async fn test_forwarding_multi_hop() {
// 5-node chain: 0 -- 1 -- 2 -- 3 -- 4
// Send datagram from node 0 destined for node 4.
let edges = vec![(0, 1), (1, 2), (2, 3), (3, 4)];
let mut nodes = run_tree_test(5, &edges, false).await;
verify_tree_convergence(&nodes);
populate_all_coord_caches(&mut nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let node4_addr = *nodes[4].node.node_addr();
// Build a SessionDatagram with enough TTL for 4 hops
let dg = SessionDatagram::new(
node0_addr,
node4_addr,
vec![0x10, 0x00, 0x04, 0x00, 1, 2, 3, 4],
);
let encoded = dg.encode();
// Inject at node 0 → node 1
nodes[0]
.node
.send_encrypted_link_message(&node1_addr, &encoded)
.await
.unwrap();
// Process multiple rounds to let the datagram traverse the chain
for _ in 0..5 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
// Verify no crashes — the datagram should have traversed 1→2→3→4
// and been delivered locally at node 4.
cleanup_nodes(&mut nodes).await;
}
#[tokio::test]
async fn test_forwarding_hop_limit_prevents_infinite_loops() {
// 3-node chain: 0 -- 1 -- 2
// Send a datagram with ttl=2. Node 1 forwards it as transit (2 -> 1) and
// node 2 delivers it locally, which is not TTL-gated. Had node 2 been
// transit instead, the arriving ttl=1 would have stopped it there.
let edges = vec![(0, 1), (1, 2)];
let mut nodes = run_tree_test(3, &edges, false).await;
verify_tree_convergence(&nodes);
populate_all_coord_caches(&mut nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let node2_addr = *nodes[2].node.node_addr();
let dg = SessionDatagram::new(
node0_addr,
node2_addr,
vec![0x10, 0x00, 0x04, 0x00, 1, 2, 3, 4],
)
.with_ttl(2); // Node 1 forwards with ttl=1; node 2 is the destination
let encoded = dg.encode();
nodes[0]
.node
.send_encrypted_link_message(&node1_addr, &encoded)
.await
.unwrap();
for _ in 0..3 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
// No panic, no infinite loop
cleanup_nodes(&mut nodes).await;
}
/// Acceptance: a transit datagram arriving with ttl=2 leaves with ttl=1.
///
/// Pinned on a live 3-node chain (0 -- 1 -- 2) by where the datagram stops,
/// since the TTL that leaves node 0 is only observable through what the next
/// hop does with it. Both injections are transit at node 0 (external source,
/// destined for node 2), so node 1 is a forwarder in both.
///
/// - ttl=2 in: node 0 must emit ttl=1, which node 1 (transit) drops. If node 0
/// emitted ttl=2 unchanged, node 1 would forward and node 2 would deliver.
/// - ttl=3 in: node 0 emits 2, node 1 emits 1, node 2 delivers (delivery is
/// not TTL-gated). If either hop decremented by more than one, the datagram
/// would have died at node 1 instead.
///
/// Together the two pin the decrement at exactly one per hop and the drop at
/// would-leave-zero.
#[tokio::test]
async fn test_forwarding_ttl_decrement_is_one_per_hop() {
let edges = vec![(0, 1), (1, 2)];
let mut nodes = run_tree_test(3, &edges, false).await;
verify_tree_convergence(&nodes);
populate_all_coord_caches(&mut nodes);
let node0_addr = *nodes[0].node.node_addr();
let node2_addr = *nodes[2].node.node_addr();
let external_src = make_node_addr(0xEE);
// --- ttl=2: must die at node 1, one hop short of the destination ---
let dg =
SessionDatagram::new(external_src, node2_addr, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(2);
let encoded = dg.encode();
nodes[0]
.node
.handle_session_datagram(&node0_addr, &encoded[1..], false)
.await;
for _ in 0..3 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
assert_eq!(
nodes[0].node.metrics().forwarding.forwarded_packets.get(),
1,
"node 0 should have forwarded the ttl=2 datagram"
);
assert_eq!(
nodes[1]
.node
.metrics()
.forwarding
.ttl_exhausted_packets
.get(),
1,
"node 1 should have received ttl=1 and dropped it as TTL-exhausted"
);
assert_eq!(
nodes[1].node.metrics().forwarding.forwarded_packets.get(),
0,
"node 1 must not forward a datagram that would leave with ttl=0"
);
assert_eq!(
nodes[2].node.metrics().forwarding.delivered_packets.get(),
0,
"node 2 must never see the ttl=2 datagram"
);
// --- ttl=3: must survive both transit hops and be delivered at node 2 ---
let dg =
SessionDatagram::new(external_src, node2_addr, vec![0x10, 0x00, 0x00, 0x00]).with_ttl(3);
let encoded = dg.encode();
nodes[0]
.node
.handle_session_datagram(&node0_addr, &encoded[1..], false)
.await;
for _ in 0..3 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
assert_eq!(
nodes[0].node.metrics().forwarding.forwarded_packets.get(),
2,
"node 0 should have forwarded the ttl=3 datagram too"
);
assert_eq!(
nodes[1].node.metrics().forwarding.forwarded_packets.get(),
1,
"node 1 should have forwarded the ttl=2 it received"
);
assert_eq!(
nodes[1]
.node
.metrics()
.forwarding
.ttl_exhausted_packets
.get(),
1,
"node 1 should not have dropped the second datagram"
);
assert_eq!(
nodes[2].node.metrics().forwarding.delivered_packets.get(),
1,
"node 2 should have delivered the datagram that arrived with ttl=1"
);
cleanup_nodes(&mut nodes).await;
}
#[tokio::test]
async fn test_forwarding_no_route_generates_error() {
// 2-node network: 0 -- 1
// Node 0 receives a datagram from node 1 destined for unknown node.
// Node 0 should generate CoordsRequired back to node 1.
let edges = vec![(0, 1)];
let mut nodes = run_tree_test(2, &edges, false).await;
verify_tree_convergence(&nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let unknown_dest = make_node_addr(0xFF);
// Node 1 sends a datagram to unknown dest via node 0
let dg = SessionDatagram::new(node1_addr, unknown_dest, vec![0x10, 0x00, 0x00, 0x00]);
let encoded = dg.encode();
// Inject at node 1 → node 0
nodes[1]
.node
.send_encrypted_link_message(&node0_addr, &encoded)
.await
.unwrap();
// Process: node 0 receives, can't route to unknown_dest, sends error back to node 1
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
// Process the error signal arriving at node 1
tokio::time::sleep(Duration::from_millis(50)).await;
let count = process_available_packets(&mut nodes).await;
assert!(count > 0, "Expected error signal to arrive at node 1");
cleanup_nodes(&mut nodes).await;
}
#[tokio::test]
async fn test_forwarding_with_cache_warming_enables_routing() {
// 4-node chain: 0 -- 1 -- 2 -- 3
// Initially, only populate coord caches at node 0.
// Send a SessionSetup from node 0 to node 3.
// As it traverses 1 and 2, those nodes should cache coordinates from the
// SessionSetup. Then verify the caches were warmed.
let edges = vec![(0, 1), (1, 2), (2, 3)];
let mut nodes = run_tree_test(4, &edges, false).await;
verify_tree_convergence(&nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let _node2_addr = *nodes[2].node.node_addr();
let node3_addr = *nodes[3].node.node_addr();
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
// Only populate node 0's cache with all coords (the source knows where to send)
let all_coords: Vec<(NodeAddr, TreeCoordinate)> = nodes
.iter()
.map(|tn| {
(
*tn.node.node_addr(),
tn.node.tree_state().my_coords().clone(),
)
})
.collect();
// Node 0 gets full cache
for (addr, coords) in &all_coords {
if addr != nodes[0].node.node_addr() {
nodes[0]
.node
.coord_cache_mut()
.insert(*addr, coords.clone(), now_ms);
}
}
// Nodes 1 and 2 only get their direct peers' coords (from tree state)
// but NOT node 0 or node 3's coords (the endpoints)
// Actually, they need bloom filter hits to route, so let's also ensure
// bloom filters are converged (which they should be from run_tree_test).
// But nodes 1 and 2 need cached coords to make loop-free forwarding
// decisions. Without coords, find_next_hop returns None.
// This is exactly what the SessionSetup cache warming solves!
// Populate enough so nodes can route to their adjacent peers,
// but NOT the distant endpoint coords.
for i in 0..4 {
for j in 0..4 {
if i != j {
// Give each node coords for its direct peers only
let j_addr = *nodes[j].node.node_addr();
if nodes[i].node.get_peer(&j_addr).is_some() {
let coords = all_coords
.iter()
.find(|(a, _)| a == &j_addr)
.unwrap()
.1
.clone();
nodes[i]
.node
.coord_cache_mut()
.insert(j_addr, coords, now_ms);
}
}
}
}
// Build SessionSetup with real coordinates
let src_coords = nodes[0].node.tree_state().my_coords().clone();
let dest_coords = nodes[3].node.tree_state().my_coords().clone();
let setup = SessionSetup::new(src_coords, dest_coords);
let setup_payload = setup.encode();
let dg = SessionDatagram::new(node0_addr, node3_addr, setup_payload);
let encoded = dg.encode();
// Inject: node 0 → node 1
nodes[0]
.node
.send_encrypted_link_message(&node1_addr, &encoded)
.await
.unwrap();
// Process multiple rounds for the datagram to traverse 1→2→3
for _ in 0..5 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
// Verify cache warming: nodes 1 and 2 should now have cached coords
// for both node 0 and node 3 (from the SessionSetup)
let cached_0_at_1 = nodes[1].node.coord_cache().get(&node0_addr, now_ms);
let cached_3_at_1 = nodes[1].node.coord_cache().get(&node3_addr, now_ms);
assert!(
cached_0_at_1.is_some(),
"Node 1 should have cached node 0's coords from SessionSetup"
);
assert!(
cached_3_at_1.is_some(),
"Node 1 should have cached node 3's coords from SessionSetup"
);
let cached_0_at_2 = nodes[2].node.coord_cache().get(&node0_addr, now_ms);
let cached_3_at_2 = nodes[2].node.coord_cache().get(&node3_addr, now_ms);
assert!(
cached_0_at_2.is_some(),
"Node 2 should have cached node 0's coords from SessionSetup"
);
assert!(
cached_3_at_2.is_some(),
"Node 2 should have cached node 3's coords from SessionSetup"
);
cleanup_nodes(&mut nodes).await;
}
// ============================================================================
// ECN Tests
// ============================================================================
use crate::node::TransportDropState;
use crate::node::handlers::session::mark_ipv6_ecn_ce;
use crate::transport::TransportId;
/// Build a minimal IPv6 header (40 bytes) with specified ECN bits.
fn make_ipv6_packet_with_ecn(ecn: u8) -> Vec<u8> {
let mut pkt = vec![0u8; 40];
let tc = ecn; // DSCP=0, ECN=ecn
pkt[0] = 0x60 | (tc >> 4);
pkt[1] = tc << 4;
pkt
}
/// Extract ECN bits from an IPv6 packet.
fn read_ecn(pkt: &[u8]) -> u8 {
let tc = ((pkt[0] & 0x0F) << 4) | (pkt[1] >> 4);
tc & 0x03
}
#[test]
fn test_mark_ecn_ce_on_ect0() {
let mut pkt = make_ipv6_packet_with_ecn(0b10);
assert_eq!(read_ecn(&pkt), 0b10);
mark_ipv6_ecn_ce(&mut pkt);
assert_eq!(read_ecn(&pkt), 0b11);
}
#[test]
fn test_mark_ecn_ce_on_ect1() {
let mut pkt = make_ipv6_packet_with_ecn(0b01);
assert_eq!(read_ecn(&pkt), 0b01);
mark_ipv6_ecn_ce(&mut pkt);
assert_eq!(read_ecn(&pkt), 0b11);
}
#[test]
fn test_mark_ecn_ce_on_not_ect() {
let mut pkt = make_ipv6_packet_with_ecn(0b00);
mark_ipv6_ecn_ce(&mut pkt);
assert_eq!(read_ecn(&pkt), 0b00);
}
#[test]
fn test_mark_ecn_ce_already_ce() {
let mut pkt = make_ipv6_packet_with_ecn(0b11);
mark_ipv6_ecn_ce(&mut pkt);
assert_eq!(read_ecn(&pkt), 0b11);
}
#[test]
fn test_mark_ecn_ce_preserves_dscp_and_flow_label() {
let mut pkt = vec![0u8; 40];
// DSCP=0b101100 (46=EF), ECN=ECT(0)=0b10 → TC=0xB2
let tc: u8 = 0xB2;
pkt[0] = 0x60 | (tc >> 4); // 0x6B
pkt[1] = (tc << 4) | 0x0A; // 0x2A (flow label high nibble = 0xA)
pkt[2] = 0xBC;
pkt[3] = 0xDE;
mark_ipv6_ecn_ce(&mut pkt);
let new_tc = ((pkt[0] & 0x0F) << 4) | (pkt[1] >> 4);
assert_eq!(new_tc, 0xB3, "TC should be 0xB3 (DSCP preserved, ECN=CE)");
assert_eq!(pkt[0] >> 4, 6, "Version nibble preserved");
assert_eq!(pkt[1] & 0x0F, 0x0A, "Flow label high nibble preserved");
assert_eq!(pkt[2], 0xBC, "Flow label byte 2 preserved");
assert_eq!(pkt[3], 0xDE, "Flow label byte 3 preserved");
}
#[test]
fn test_mark_ecn_ce_short_packet() {
let mut pkt = vec![0x60];
mark_ipv6_ecn_ce(&mut pkt);
assert_eq!(pkt, vec![0x60]);
let mut empty: Vec<u8> = vec![];
mark_ipv6_ecn_ce(&mut empty);
assert!(empty.is_empty());
}
#[tokio::test]
async fn test_ce_relay_through_forwarding() {
// 3-node chain: 0 -- 1 -- 2
// Send a datagram with CE set from node 0 to node 1.
// Node 1 should relay CE to node 2.
let edges = vec![(0, 1), (1, 2)];
let mut nodes = run_tree_test(3, &edges, false).await;
verify_tree_convergence(&nodes);
populate_all_coord_caches(&mut nodes);
let node0_addr = *nodes[0].node.node_addr();
let node1_addr = *nodes[1].node.node_addr();
let node2_addr = *nodes[2].node.node_addr();
// Record ecn_ce_count at node 2 before
let ce_before = nodes[2]
.node
.get_peer(&node1_addr)
.and_then(|p| p.mmp())
.map(|m| m.receiver.ecn_ce_count())
.unwrap_or(0);
// Build a SessionDatagram from node 0 to node 2
let dg = SessionDatagram::new(
node0_addr,
node2_addr,
vec![0x10, 0x00, 0x04, 0x00, 1, 2, 3, 4],
);
let encoded = dg.encode();
// Send from node 0 to node 1 with CE flag set
nodes[0]
.node
.send_encrypted_link_message_with_ce(&node1_addr, &encoded, true)
.await
.unwrap();
// Process: node 1 receives (CE set), forwards to node 2 (CE relayed)
for _ in 0..3 {
tokio::time::sleep(Duration::from_millis(50)).await;
process_available_packets(&mut nodes).await;
}
// Node 2's link-layer MMP should have received a CE-flagged frame from node 1
let ce_after = nodes[2]
.node
.get_peer(&node1_addr)
.and_then(|p| p.mmp())
.map(|m| m.receiver.ecn_ce_count())
.unwrap_or(0);
assert!(
ce_after > ce_before,
"Node 2 should see CE flag relayed from node 1 (before={ce_before}, after={ce_after})"
);
cleanup_nodes(&mut nodes).await;
}
#[test]
fn test_detect_congestion_with_transport_drops() {
let mut node = make_node();
// No drops — detect_congestion should return false for any address
let fake_addr = NodeAddr::from_bytes([1; 16]);
assert!(!node.detect_congestion(&fake_addr));
// Simulate transport kernel drops
let tid = TransportId::new(1);
node.transport_drops.insert(
tid,
TransportDropState {
prev_drops: 100,
dropping: true,
},
);
// Now detect_congestion should return true (local transport congestion)
assert!(node.detect_congestion(&fake_addr));
// Clear the dropping flag — should return false again
node.transport_drops.get_mut(&tid).unwrap().dropping = false;
assert!(!node.detect_congestion(&fake_addr));
}
#[test]
fn test_detect_congestion_disabled_ecn() {
let mut config = Config::new();
config.node.ecn.enabled = false;
let mut node = Node::new(config).unwrap();
// Even with transport drops, disabled ECN should return false
let tid = TransportId::new(1);
node.transport_drops.insert(
tid,
TransportDropState {
prev_drops: 50,
dropping: true,
},
);
let fake_addr = NodeAddr::from_bytes([1; 16]);
assert!(!node.detect_congestion(&fake_addr));
}
#[test]
fn test_sample_transport_congestion() {
let mut node = make_node();
// Insert a transport drop state with a baseline
let tid = TransportId::new(1);
node.transport_drops.insert(
tid,
TransportDropState {
prev_drops: 0,
dropping: false,
},
);
// No transports registered — sample_transport_congestion is a no-op
// (transport_drops entry stays unchanged)
node.sample_transport_congestion();
assert!(!node.transport_drops[&tid].dropping);
}