Implement FSP port multiplexing and IPv6 header compression

Breaking wire format change: DataPacket payloads inside the AEAD envelope
now carry a 4-byte port header [src_port:2 LE][dst_port:2 LE] before the
service payload. The receiver dispatches by destination port.

Port multiplexing:
- send_session_data() takes src_port/dst_port params, prepends port header
- New send_ipv6_packet() compresses IPv6 header and sends on port 256
- Receive path dispatches DataPackets by port: port 256 decompresses IPv6
  header from session context and delivers to TUN, unknown ports dropped
- Port constants: FSP_PORT_HEADER_SIZE (4 bytes), FSP_PORT_IPV6_SHIM (256)

IPv6 header compression:
- New ipv6_shim module with compress_ipv6()/decompress_ipv6() pure functions
- Strips src/dst addresses (32 bytes) and payload length (2 bytes) from each
  packet, preserving traffic class, flow label, next header, and hop limit
  as 6-byte residual fields
- Addresses reconstructed from session context on receive side
- Net savings: 29 bytes per packet (overhead 106 → 77 bytes)
- FIPS_IPV6_OVERHEAD constant (77 bytes), effective_ipv6_mtu() updated
- 16 unit tests for round-trip fidelity, field preservation, error cases

Documentation:
- fips-wire-formats: DataPacket port header, port registry, IPv6 shim
  format tables, updated encapsulation walkthrough and overhead budget
- fips-ipv6-adapter: FIPS_IPV6_OVERHEAD (77 bytes), updated MTU numbers,
  TUN reader/writer flow with compression steps, impl status
- fips-session-layer: port-based service dispatch section, data transfer
  description, impl status
- fips-intro: IPv6 adapter as port 256 service, node architecture updated
- fips-mesh-operation: packet size summary with compressed overhead
- DataPacket doc updated with port header and dispatch model
- session_wire.rs module doc: DataPacket Port Multiplexing section
This commit is contained in:
Johnathan Corgan
2026-03-11 12:53:32 +00:00
parent f37eb4b846
commit 6ab8b35755
12 changed files with 644 additions and 77 deletions
+82 -20
View File
@@ -10,7 +10,7 @@ use crate::node::session_wire::{
build_fsp_header, fsp_prepend_inner_header, fsp_strip_inner_header,
parse_encrypted_coords, FspCommonPrefix, FspEncryptedHeader, 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_PHASE_MSG2, FSP_PHASE_MSG3, FSP_PORT_HEADER_SIZE, FSP_PORT_IPV6_SHIM,
};
use crate::protocol::{coords_wire_size, encode_coords};
use crate::upper::icmp::FIPS_OVERHEAD;
@@ -271,21 +271,53 @@ impl Node {
// Dispatch by msg_type
match SessionMessageType::from_byte(msg_type) {
Some(SessionMessageType::DataPacket) => {
// msg_type 0x10: deliver rest (IPv6 payload) to TUN
let mut packet = rest.to_vec();
if ce_flag {
mark_ipv6_ecn_ce(&mut packet);
self.stats_mut().congestion.record_ce_received();
// 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;
}
if let Some(tun_tx) = &self.tun_tx {
if let Err(e) = tun_tx.send(packet) {
debug!(error = %e, "Failed to deliver decrypted packet to TUN");
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.stats_mut().congestion.record_ce_received();
}
if let Some(tun_tx) = &self.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"
);
}
}
}
_ => {
debug!(
src = %self.peer_display_name(src_addr),
dst_port,
"Unknown FSP service port, dropping DataPacket"
);
}
} else {
trace!(
src = %self.peer_display_name(src_addr),
"DataPacket decrypted (no TUN interface, plaintext dropped)"
);
}
}
Some(SessionMessageType::SenderReport) => {
@@ -1118,10 +1150,16 @@ impl Node {
/// 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,
plaintext: &[u8],
src_port: u16,
dst_port: u16,
payload: &[u8],
) -> Result<(), NodeError> {
let now_ms = Self::now_ms();
@@ -1140,10 +1178,16 @@ impl Node {
});
}
// 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, plaintext);
let inner_plaintext = fsp_prepend_inner_header(timestamp, msg_type, inner_flags, &port_payload);
// Determine whether coords fit within transport MTU.
// If not, send standalone CoordsWarmup before the data packet.
@@ -1151,7 +1195,7 @@ impl Node {
let src = self.tree_state.my_coords().clone();
let dst = self.get_dest_coords(dest_addr);
let coords_size = coords_wire_size(&src) + coords_wire_size(&dst);
let total_wire = FIPS_OVERHEAD as usize + coords_size + plaintext.len();
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 {
@@ -1226,7 +1270,7 @@ impl Node {
// Re-borrow after send (which borrowed &mut self)
if let Some(entry) = self.sessions.get_mut(dest_addr) {
entry.record_sent(plaintext.len());
entry.record_sent(payload.len());
if let Some(mmp) = entry.mmp_mut() {
mmp.sender.record_sent(counter, timestamp, ciphertext.len());
}
@@ -1236,6 +1280,24 @@ impl Node {
Ok(())
}
/// 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:
@@ -1525,7 +1587,7 @@ impl Node {
return;
}
}
if let Err(e) = self.send_session_data(&dest_addr, &ipv6_packet).await {
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;
@@ -1634,7 +1696,7 @@ impl Node {
None => return,
};
for packet in packets {
if let Err(e) = self.send_session_data(dest_addr, &packet).await {
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;
}
+19
View File
@@ -11,6 +11,17 @@
//! [ver+phase:1][flags:1][payload_len:2 LE]
//! ```
//!
//! ## DataPacket Port Multiplexing
//!
//! DataPacket (msg_type 0x10) payloads inside the AEAD envelope carry a 4-byte
//! port header for service dispatch:
//!
//! ```text
//! [src_port:2 LE][dst_port:2 LE][service payload...]
//! ```
//!
//! Port 256 (0x100) = IPv6 shim with header compression.
//!
//! ## Message Classes
//!
//! | Phase | U Flag | Type | Description |
@@ -58,6 +69,14 @@ const TAG_SIZE: usize = 16;
/// Minimum size for an encrypted FSP message: header + tag (no plaintext).
pub const FSP_ENCRYPTED_MIN_SIZE: usize = FSP_HEADER_SIZE + TAG_SIZE; // 28 bytes
// FSP DataPacket port header constants.
/// Size of the FSP DataPacket port header (src_port + dst_port).
pub const FSP_PORT_HEADER_SIZE: usize = 4;
/// FSP port: IPv6 shim service.
pub const FSP_PORT_IPV6_SHIM: u16 = 256;
// Cleartext flag bit constants (byte 1 of common prefix, phase 0x0 only).
/// Coords Present — source and destination coordinates follow the header.
+23 -9
View File
@@ -243,7 +243,7 @@ async fn test_session_direct_peer_data_transfer() {
let test_data = b"Hello, FIPS session!";
nodes[0]
.node
.send_session_data(&node1_addr, test_data)
.send_session_data(&node1_addr, 0, 0, test_data)
.await
.expect("send_session_data failed");
@@ -378,7 +378,7 @@ async fn test_session_3node_forwarded_data() {
let test_data = b"End-to-end through transit node B";
nodes[0]
.node
.send_session_data(&node2_addr, test_data)
.send_session_data(&node2_addr, 0, 0, test_data)
.await
.expect("send_session_data failed");
@@ -438,7 +438,7 @@ async fn test_session_send_data_no_session_fails() {
let mut node = make_node();
let fake_addr = make_node_addr(0xAA);
let result = node.send_session_data(&fake_addr, b"test").await;
let result = node.send_session_data(&fake_addr, 0, 0, b"test").await;
assert!(result.is_err(), "Should fail with no session");
}
@@ -618,11 +618,16 @@ async fn test_session_100_nodes() {
let dest_addr = all_info[dst].0;
let src_addr = all_info[src].0;
// Build IPv6 packets with pair index as payload
let src_fips = crate::FipsAddress::from_node_addr(&src_addr);
let dst_fips = crate::FipsAddress::from_node_addr(&dest_addr);
// Forward: initiator → responder
let fwd_payload = format!("fwd-{}", pair_idx).into_bytes();
let fwd_ipv6 = build_ipv6_packet(&src_fips, &dst_fips, &fwd_payload);
match nodes[src]
.node
.send_session_data(&dest_addr, &fwd_payload)
.send_ipv6_packet(&dest_addr, &fwd_ipv6)
.await
{
Ok(()) => send_forward_ok += 1,
@@ -634,9 +639,10 @@ async fn test_session_100_nodes() {
// Reverse: responder → initiator
// (Responder should already be Established after XK msg3)
let rev_payload = format!("rev-{}", pair_idx).into_bytes();
let rev_ipv6 = build_ipv6_packet(&dst_fips, &src_fips, &rev_payload);
match nodes[dst]
.node
.send_session_data(&src_addr, &rev_payload)
.send_ipv6_packet(&src_addr, &rev_ipv6)
.await
{
Ok(()) => send_reverse_ok += 1,
@@ -671,13 +677,21 @@ async fn test_session_100_nodes() {
let fwd_payload = format!("fwd-{}", pair_idx).into_bytes();
let rev_payload = format!("rev-{}", pair_idx).into_bytes();
if delivered_per_node[dst].contains(&fwd_payload) {
// After decompression, TUN receives full IPv6 packets.
// Check that delivered packet's upper-layer payload matches.
let fwd_found = delivered_per_node[dst]
.iter()
.any(|pkt| pkt.len() >= 40 && pkt[40..] == fwd_payload);
if fwd_found {
fwd_delivered += 1;
} else if fwd_missing.len() < 20 {
fwd_missing.push((src, dst));
}
if delivered_per_node[src].contains(&rev_payload) {
let rev_found = delivered_per_node[src]
.iter()
.any(|pkt| pkt.len() >= 40 && pkt[40..] == rev_payload);
if rev_found {
rev_delivered += 1;
} else if rev_missing.len() < 20 {
rev_missing.push((src, dst));
@@ -1836,7 +1850,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() {
// With coords (~66 extra), the wire could exceed B's recv buffer.
for _ in 0..5 {
let small = build_ipv6_packet(&src_fips, &dst_fips, &[0u8; 10]);
nodes[0].node.send_session_data(&node2_addr, &small).await.unwrap();
nodes[0].node.send_ipv6_packet(&node2_addr, &small).await.unwrap();
}
drain_to_quiescence(&mut nodes).await;
@@ -1855,7 +1869,7 @@ async fn test_multihop_pmtud_heterogeneous_mtu() {
// Send the oversized packet — B should fail to forward and send
// MtuExceeded signal back.
nodes[0].node.send_session_data(&node2_addr, &ipv6_packet).await.unwrap();
nodes[0].node.send_ipv6_packet(&node2_addr, &ipv6_packet).await.unwrap();
drain_to_quiescence(&mut nodes).await;
// Verify PathMtuState was updated on A