Add FMP node profile negotiation with mixed-profile integration tests

NodeProfile enum (Full/NonRouting/Leaf) with FMP feature bitfield:
bits 0-2 profile, bits 3-6 MMP wants/provides, bit 7 bloom filter
size negotiable. Bloom size TLV always sent with min=max=1KB.

Config mapping: leaf_only -> Leaf, disable_routing -> NonRouting.
Profile and agreed bloom size stored on PeerConnection and
ActivePeer. Handshake msg2/msg3 carry FMP negotiation payload
with profile validation and bloom size agreement.

MMP report sending gated by profile wants/provides. Parent
selection and routing constraints skip non-full peers. One-way
bloom filters for non-routing peers (F inserts N as dependent).
Leaf mode: single-peer enforcement, suppress tree announces and
discovery forwarding.

Mixed-profile integration test: A(Full) + B(Full) + C(NonRouting)
+ D(Leaf) with 9 connectivity assertions.
This commit is contained in:
Johnathan Corgan
2026-04-11 08:16:01 +00:00
parent 8357200b0e
commit f4278f1dd7
20 changed files with 1323 additions and 299 deletions
+128 -122
View File
@@ -5,43 +5,42 @@
//! Bloom filters, coordinate caches, transports, links, and peers.
mod bloom;
mod discovery_rate_limit;
mod handlers;
mod lifecycle;
mod rate_limit;
mod retry;
mod discovery_rate_limit;
mod rate_limit;
mod routing_error_rate_limit;
pub(crate) mod session;
pub(crate) mod session_wire;
pub(crate) mod wire;
pub(crate) mod stats;
mod tree;
#[cfg(test)]
mod tests;
mod tree;
pub(crate) mod wire;
use crate::bloom::BloomState;
use crate::protocol::NodeProfile;
use crate::cache::CoordCache;
use crate::utils::index::IndexAllocator;
use crate::node::session::SessionEntry;
use crate::peer::{ActivePeer, PeerConnection};
use self::discovery_rate_limit::{DiscoveryBackoff, DiscoveryForwardRateLimiter};
use self::rate_limit::HandshakeRateLimiter;
use self::routing_error_rate_limit::RoutingErrorRateLimiter;
use self::wire::{
FLAG_CE, FLAG_KEY_EPOCH, FLAG_SP, build_encrypted, build_established_header,
prepend_inner_header,
};
use crate::bloom::BloomState;
use crate::cache::CoordCache;
use crate::node::session::SessionEntry;
use crate::peer::{ActivePeer, PeerConnection};
use crate::transport::ethernet::EthernetTransport;
use crate::transport::tcp::TcpTransport;
use crate::transport::tor::TorTransport;
use crate::transport::udp::UdpTransport;
use crate::transport::{
Link, LinkId, PacketRx, PacketTx, TransportAddr, TransportError, TransportHandle, TransportId,
};
use crate::transport::udp::UdpTransport;
use crate::transport::tcp::TcpTransport;
use crate::transport::tor::TorTransport;
#[cfg(target_os = "linux")]
use crate::transport::ethernet::EthernetTransport;
use crate::tree::TreeState;
use crate::upper::hosts::HostMap;
use crate::upper::icmp_rate_limit::IcmpRateLimiter;
use crate::upper::tun::{TunError, TunOutboundRx, TunState, TunTx};
use crate::utils::index::IndexAllocator;
use self::wire::{build_encrypted, build_established_header, prepend_inner_header, FLAG_CE, FLAG_KEY_EPOCH, FLAG_SP};
use crate::{Config, ConfigError, Identity, IdentityError, NodeAddr, PeerIdentity};
use rand::Rng;
use std::collections::{HashMap, VecDeque};
@@ -108,11 +107,7 @@ pub enum NodeError {
SendFailed { node_addr: NodeAddr, reason: String },
#[error("mtu exceeded forwarding to {node_addr}: packet {packet_size} > mtu {mtu}")]
MtuExceeded {
node_addr: NodeAddr,
packet_size: usize,
mtu: u16,
},
MtuExceeded { node_addr: NodeAddr, packet_size: usize, mtu: u16 },
#[error("config error: {0}")]
Config(#[from] ConfigError),
@@ -278,6 +273,9 @@ pub struct Node {
/// Whether this is a leaf-only node.
is_leaf_only: bool,
/// Node profile derived from config (Full, NonRouting, Leaf).
node_profile: NodeProfile,
// === Spanning Tree ===
/// Local spanning tree state.
tree_state: TreeState,
@@ -370,10 +368,6 @@ pub struct Node {
tun_reader_handle: Option<JoinHandle<()>>,
/// TUN writer thread handle.
tun_writer_handle: Option<JoinHandle<()>>,
/// Shutdown pipe: writing to this fd unblocks the TUN reader thread on macOS.
/// On Linux, deleting the interface via netlink serves the same purpose.
#[cfg(target_os = "macos")]
tun_shutdown_fd: Option<std::os::unix::io::RawFd>,
// === DNS Responder ===
/// Receiver for resolved identities from the DNS responder.
@@ -454,6 +448,7 @@ impl Node {
let identity = config.create_identity()?;
let node_addr = *identity.node_addr();
let is_leaf_only = config.is_leaf_only();
let node_profile = config.node_profile();
let mut startup_epoch = [0u8; 8];
rand::rng().fill_bytes(&mut startup_epoch);
@@ -516,6 +511,7 @@ impl Node {
config,
state: NodeState::Created,
is_leaf_only,
node_profile,
tree_state,
bloom_state,
coord_cache,
@@ -544,8 +540,6 @@ impl Node {
tun_outbound_rx: None,
tun_reader_handle: None,
tun_writer_handle: None,
#[cfg(target_os = "macos")]
tun_shutdown_fd: None,
dns_identity_rx: None,
dns_task: None,
index_allocator: IndexAllocator::new(),
@@ -558,7 +552,10 @@ impl Node {
coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
std::time::Duration::from_millis(coords_response_interval_ms),
),
discovery_backoff: DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs),
discovery_backoff: DiscoveryBackoff::with_params(
backoff_base_secs,
backoff_max_secs,
),
discovery_forward_limiter: DiscoveryForwardRateLimiter::with_interval(
std::time::Duration::from_secs(forward_min_interval_secs),
),
@@ -626,6 +623,7 @@ impl Node {
config,
state: NodeState::Created,
is_leaf_only: false,
node_profile: NodeProfile::Full,
tree_state,
bloom_state,
coord_cache,
@@ -654,8 +652,6 @@ impl Node {
tun_outbound_rx: None,
tun_reader_handle: None,
tun_writer_handle: None,
#[cfg(target_os = "macos")]
tun_shutdown_fd: None,
dns_identity_rx: None,
dns_task: None,
index_allocator: IndexAllocator::new(),
@@ -685,6 +681,7 @@ impl Node {
pub fn leaf_only(config: Config) -> Result<Self, NodeError> {
let mut node = Self::new(config)?;
node.is_leaf_only = true;
node.node_profile = NodeProfile::Leaf;
node.bloom_state = BloomState::leaf_only(*node.identity.node_addr());
Ok(node)
}
@@ -712,19 +709,23 @@ impl Node {
}
// Create Ethernet transport instances
let eth_instances: Vec<_> = self
.config
.transports
.ethernet
.iter()
.map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
.collect();
let xonly = self.identity.pubkey();
for (name, eth_config) in eth_instances {
let transport_id = self.allocate_transport_id();
let mut eth = EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
eth.set_local_pubkey(xonly);
transports.push(TransportHandle::Ethernet(eth));
#[cfg(target_os = "linux")]
{
let eth_instances: Vec<_> = self
.config
.transports
.ethernet
.iter()
.map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
.collect();
let xonly = self.identity.pubkey();
for (name, eth_config) in eth_instances {
let transport_id = self.allocate_transport_id();
let mut eth = EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
eth.set_local_pubkey(xonly);
transports.push(TransportHandle::Ethernet(eth));
}
}
// Create TCP transport instances
@@ -794,9 +795,7 @@ impl Node {
#[cfg(any(not(feature = "ble"), test))]
if !ble_instances.is_empty() {
#[cfg(not(test))]
tracing::warn!(
"BLE transport configured but 'ble' feature not enabled at compile time"
);
tracing::warn!("BLE transport configured but 'ble' feature not enabled at compile time");
}
}
@@ -846,9 +845,18 @@ impl Node {
))
})?;
// Parse the MAC address
#[cfg(target_os = "linux")]
let mac = crate::transport::ethernet::parse_mac_string(mac_str).map_err(|e| {
NodeError::NoTransportForType(format!("invalid MAC in '{}': {}", addr_str, e))
})?;
#[cfg(not(target_os = "linux"))]
let mac: [u8; 6] = {
let _ = mac_str;
return Err(NodeError::NoTransportForType(
"Ethernet transport not available on this platform".into(),
));
};
Ok((transport_id, TransportAddr::from_bytes(&mac)))
}
@@ -857,9 +865,13 @@ impl Node {
/// (TransportId, TransportAddr) pair by finding the BLE transport
/// instance matching the adapter name.
#[cfg(target_os = "linux")]
fn resolve_ble_addr(&self, addr_str: &str) -> Result<(TransportId, TransportAddr), NodeError> {
fn resolve_ble_addr(
&self,
addr_str: &str,
) -> Result<(TransportId, TransportAddr), NodeError> {
let ta = TransportAddr::from_string(addr_str);
let adapter = crate::transport::ble::addr::adapter_from_addr(&ta).ok_or_else(|| {
let adapter = crate::transport::ble::addr::adapter_from_addr(&ta)
.ok_or_else(|| {
NodeError::NoTransportForType(format!(
"invalid BLE address format '{}': expected 'adapter/mac'",
addr_str
@@ -870,7 +882,9 @@ impl Node {
let transport_id = self
.transports
.iter()
.find(|(_, handle)| handle.transport_type().name == "ble" && handle.is_operational())
.find(|(_, handle)| {
handle.transport_type().name == "ble" && handle.is_operational()
})
.map(|(id, _)| *id)
.ok_or_else(|| {
NodeError::NoTransportForType(format!(
@@ -987,6 +1001,22 @@ impl Node {
self.is_leaf_only
}
/// Get the node's profile (Full, NonRouting, Leaf).
pub fn node_profile(&self) -> NodeProfile {
self.node_profile
}
/// Collect the set of peers that are not full nodes (non-routing/leaf).
///
/// Used by tree and routing functions to skip non-transit peers.
fn non_full_peers(&self) -> std::collections::HashSet<NodeAddr> {
self.peers
.iter()
.filter(|(_, p)| p.peer_profile() != NodeProfile::Full)
.map(|(addr, _)| *addr)
.collect()
}
// === Tree State ===
/// Get the tree state.
@@ -1066,10 +1096,9 @@ impl Node {
let now = std::time::Instant::now();
let should_log = match self.last_mesh_size_log {
None => true,
Some(last) => {
now.duration_since(last)
>= std::time::Duration::from_secs(self.config.node.mmp.log_interval_secs)
}
Some(last) => now.duration_since(last) >= std::time::Duration::from_secs(
self.config.node.mmp.log_interval_secs,
),
};
if should_log {
tracing::debug!(
@@ -1118,6 +1147,7 @@ impl Node {
self.tun_name.as_deref()
}
// === Resource Limits ===
/// Set the maximum number of connections (handshake phase).
@@ -1198,17 +1228,14 @@ impl Node {
/// Add a link.
pub fn add_link(&mut self, link: Link) -> Result<(), NodeError> {
if self.max_links > 0 && self.links.len() >= self.max_links {
return Err(NodeError::MaxLinksExceeded {
max: self.max_links,
});
return Err(NodeError::MaxLinksExceeded { max: self.max_links });
}
let link_id = link.link_id();
let transport_id = link.transport_id();
let remote_addr = link.remote_addr().clone();
self.links.insert(link_id, link);
self.addr_to_link
.insert((transport_id, remote_addr), link_id);
self.addr_to_link.insert((transport_id, remote_addr), link_id);
Ok(())
}
@@ -1223,14 +1250,8 @@ impl Node {
}
/// Find link ID by transport address.
pub fn find_link_by_addr(
&self,
transport_id: TransportId,
addr: &TransportAddr,
) -> Option<LinkId> {
self.addr_to_link
.get(&(transport_id, addr.clone()))
.copied()
pub fn find_link_by_addr(&self, transport_id: TransportId, addr: &TransportAddr) -> Option<LinkId> {
self.addr_to_link.get(&(transport_id, addr.clone())).copied()
}
/// Remove a link.
@@ -1376,14 +1397,11 @@ impl Node {
pub(crate) fn register_identity(&mut self, node_addr: NodeAddr, pubkey: secp256k1::PublicKey) {
let mut prefix = [0u8; 15];
prefix.copy_from_slice(&node_addr.as_bytes()[0..15]);
self.identity_cache
.insert(prefix, (node_addr, pubkey, Self::now_ms()));
self.identity_cache.insert(prefix, (node_addr, pubkey, Self::now_ms()));
// LRU eviction
let max = self.config.node.cache.identity_size;
if self.identity_cache.len() > max
&& let Some(oldest_key) = self
.identity_cache
.iter()
&& let Some(oldest_key) = self.identity_cache.iter()
.min_by_key(|(_, (_, _, ts))| *ts)
.map(|(k, _)| *k)
{
@@ -1392,10 +1410,7 @@ impl Node {
}
/// Look up a destination by FipsAddress prefix (bytes 1-15 of the IPv6 address).
pub(crate) fn lookup_by_fips_prefix(
&mut self,
prefix: &[u8; 15],
) -> Option<(NodeAddr, secp256k1::PublicKey)> {
pub(crate) fn lookup_by_fips_prefix(&mut self, prefix: &[u8; 15]) -> Option<(NodeAddr, secp256k1::PublicKey)> {
if let Some(entry) = self.identity_cache.get_mut(prefix) {
entry.2 = Self::now_ms(); // LRU touch
Some((entry.0, entry.1))
@@ -1434,7 +1449,9 @@ impl Node {
/// has declared us as their parent (making them our child).
pub(crate) fn is_tree_peer(&self, peer_addr: &NodeAddr) -> bool {
// Peer is our parent
if !self.tree_state.is_root() && self.tree_state.my_declaration().parent_id() == peer_addr {
if !self.tree_state.is_root()
&& self.tree_state.my_declaration().parent_id() == peer_addr
{
return true;
}
// Peer is our child (their declaration names us as parent)
@@ -1483,22 +1500,17 @@ impl Node {
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let dest_coords = self
.coord_cache
.get_and_touch(dest_node_addr, now_ms)?
.clone();
let dest_coords = self.coord_cache.get_and_touch(dest_node_addr, now_ms)?.clone();
// 3. Bloom filter candidates — requires dest_coords for loop-free selection.
// If no candidate is strictly closer, fall through to tree routing.
// 3. Bloom filter candidates — requires dest_coords for loop-free selection
let candidates: Vec<&ActivePeer> = self.destination_in_filters(dest_node_addr);
if !candidates.is_empty()
&& let Some(peer) = self.select_best_candidate(&candidates, &dest_coords)
{
return Some(peer);
if !candidates.is_empty() {
return self.select_best_candidate(&candidates, &dest_coords);
}
// 4. Greedy tree routing fallback
let next_hop_id = self.tree_state.find_next_hop(&dest_coords)?;
// 4. Greedy tree routing fallback (skip non-routing/leaf peers)
let skip = self.non_full_peers();
let next_hop_id = self.tree_state.find_next_hop(&dest_coords, &skip)?;
self.peers.get(&next_hop_id).filter(|p| p.can_send())
}
@@ -1558,9 +1570,15 @@ impl Node {
best.map(|(peer, _, _)| peer)
}
/// Check if a destination is in any peer's bloom filter.
/// Check if a destination is in any full peer's bloom filter.
///
/// Skips non-routing and leaf peers (defensive — they shouldn't have
/// filters claiming transit reachability, but guard against propagation bugs).
pub fn destination_in_filters(&self, dest: &NodeAddr) -> Vec<&ActivePeer> {
self.peers.values().filter(|p| p.may_reach(dest)).collect()
self.peers
.values()
.filter(|p| p.peer_profile() == NodeProfile::Full && p.may_reach(dest))
.collect()
}
/// Get the TUN packet sender channel.
@@ -1588,8 +1606,7 @@ impl Node {
node_addr: &NodeAddr,
plaintext: &[u8],
) -> Result<(), NodeError> {
self.send_encrypted_link_message_with_ce(node_addr, plaintext, false)
.await
self.send_encrypted_link_message_with_ce(node_addr, plaintext, false).await
}
/// Like `send_encrypted_link_message` but allows setting the FMP CE flag.
@@ -1601,9 +1618,7 @@ impl Node {
plaintext: &[u8],
ce_flag: bool,
) -> Result<(), NodeError> {
let peer = self
.peers
.get_mut(node_addr)
let peer = self.peers.get_mut(node_addr)
.ok_or(NodeError::PeerNotFound(*node_addr))?;
let their_index = peer.their_index().ok_or_else(|| NodeError::SendFailed {
@@ -1614,19 +1629,18 @@ impl Node {
node_addr: *node_addr,
reason: "no transport_id".into(),
})?;
let remote_addr = peer
.current_addr()
.cloned()
.ok_or_else(|| NodeError::SendFailed {
node_addr: *node_addr,
reason: "no current_addr".into(),
})?;
let remote_addr = peer.current_addr().cloned().ok_or_else(|| NodeError::SendFailed {
node_addr: *node_addr,
reason: "no current_addr".into(),
})?;
// Prepend 4-byte session-relative timestamp (inner header)
let timestamp_ms = peer.session_elapsed_ms();
// MMP: read spin bit value before entering session borrow
let sp_flag = peer.mmp().map(|mmp| mmp.spin_bit.tx_bit()).unwrap_or(false);
let sp_flag = peer.mmp()
.map(|mmp| mmp.spin_bit.tx_bit())
.unwrap_or(false);
let mut flags = if sp_flag { FLAG_SP } else { 0 };
if ce_flag {
flags |= FLAG_CE;
@@ -1635,12 +1649,10 @@ impl Node {
flags |= FLAG_KEY_EPOCH;
}
let session = peer
.noise_session_mut()
.ok_or_else(|| NodeError::SendFailed {
node_addr: *node_addr,
reason: "no noise session".into(),
})?;
let session = peer.noise_session_mut().ok_or_else(|| NodeError::SendFailed {
node_addr: *node_addr,
reason: "no noise session".into(),
})?;
// Inner plaintext: [timestamp:4 LE][msg_type][payload...]
let inner_plaintext = prepend_inner_header(timestamp_ms, plaintext);
@@ -1651,24 +1663,18 @@ impl Node {
let header = build_established_header(their_index, counter, flags, payload_len);
// Encrypt with AAD binding to the outer header
let ciphertext = session
.encrypt_with_aad(&inner_plaintext, &header)
.map_err(|e| NodeError::SendFailed {
node_addr: *node_addr,
reason: format!("encryption failed: {}", e),
})?;
let ciphertext = session.encrypt_with_aad(&inner_plaintext, &header).map_err(|e| NodeError::SendFailed {
node_addr: *node_addr,
reason: format!("encryption failed: {}", e),
})?;
let wire_packet = build_encrypted(&header, &ciphertext);
// Re-borrow peer for stats update after sending
let transport = self
.transports
.get(&transport_id)
let transport = self.transports.get(&transport_id)
.ok_or(NodeError::TransportNotFound(transport_id))?;
let bytes_sent = transport
.send(&remote_addr, &wire_packet)
.await
let bytes_sent = transport.send(&remote_addr, &wire_packet).await
.map_err(|e| match e {
TransportError::MtuExceeded { packet_size, mtu } => NodeError::MtuExceeded {
node_addr: *node_addr,