Minimal shared-media beacons: 4-byte header, strip pubkey, remove BLE exchange

Redesign shared-media transport framing now that XX replaces IK:

- Ethernet: unified 4-byte header [type][flags][length:2 LE] for all
  frame types, effective MTU now if_mtu - 4
- Beacons: strip 32-byte pubkey (34 → 5 bytes), identity learned from
  XX handshake msg2/msg3 instead of beacon
- BLE: remove pre-handshake pubkey_exchange() and cross-probe
  tie-breaker, both unnecessary with XX
- Discovery: support anonymous connections (pubkey_hint: None) with
  address-based dedup for shared-media transports
- Post-handshake identity hooks in handle_msg2/msg3 for self-detection
  and future allow/deny list filtering (IDEA-0047)
This commit is contained in:
Johnathan Corgan
2026-04-11 08:16:01 +00:00
parent e06015dd43
commit 8162d3c237
9 changed files with 402 additions and 613 deletions
+21
View File
@@ -428,6 +428,16 @@ impl Node {
let peer_node_addr = *peer_identity.node_addr();
// Post-handshake identity filtering hook (IDEA-0047).
// With XX, shared-media transports discover peers without identity;
// this is the first point where the initiator knows the responder.
// Future: check allow/deny list here, abort if denied.
if peer_node_addr == *self.identity.node_addr() {
debug!(link_id = %link_id, "Discovered self via shared-media beacon, dropping");
self.connections.remove(&link_id);
return;
}
// Build and send msg3
let our_index = conn.our_index().unwrap_or(header.receiver_idx);
let wire_msg3 = build_msg3(our_index, header.sender_idx, &msg3_bytes);
@@ -756,6 +766,17 @@ impl Node {
};
let peer_node_addr = *peer_identity.node_addr();
// Post-handshake identity filtering hook (IDEA-0047).
// With XX, this is the first point where the responder knows
// the initiator's identity. Future: check allow/deny list here.
if peer_node_addr == *self.identity.node_addr() {
debug!(link_id = %link_id, "Received msg3 from self, dropping");
self.connections.remove(&link_id);
self.remove_link(&link_id);
return;
}
let our_index = conn.our_index().unwrap_or(header.receiver_idx);
// Identity-based restart/rekey detection.
+140 -184
View File
@@ -1,11 +1,11 @@
//! Node lifecycle management: start, stop, and peer connection initiation.
use super::{Node, NodeError, NodeState};
use crate::node::wire::build_msg1;
use crate::peer::PeerConnection;
use crate::protocol::{Disconnect, DisconnectReason};
use crate::transport::{Link, LinkDirection, LinkId, TransportAddr, TransportId, packet_channel};
use crate::upper::tun::{TunDevice, TunState, run_tun_reader, shutdown_tun_interface};
use crate::transport::{packet_channel, Link, LinkDirection, LinkId, TransportAddr, TransportId};
use crate::upper::tun::{run_tun_reader, shutdown_tun_interface, TunDevice, TunState};
use crate::node::wire::build_msg1;
use crate::{NodeAddr, PeerIdentity};
use std::thread;
use std::time::Duration;
@@ -50,10 +50,7 @@ impl Node {
return;
}
debug!(
count = peer_configs.len(),
"Initiating static peer connections"
);
debug!(count = peer_configs.len(), "Initiating static peer connections");
for peer_config in peer_configs {
if let Err(e) = self.initiate_peer_connection(&peer_config).await {
@@ -70,16 +67,14 @@ impl Node {
/// Initiate a connection to a single peer.
///
/// Creates a link, starts the Noise handshake, and sends the first message.
pub(super) async fn initiate_peer_connection(
&mut self,
peer_config: &crate::config::PeerConfig,
) -> Result<(), NodeError> {
pub(super) async fn initiate_peer_connection(&mut self, peer_config: &crate::config::PeerConfig) -> Result<(), NodeError> {
// Parse the peer's npub to get their identity
let peer_identity =
PeerIdentity::from_npub(&peer_config.npub).map_err(|e| NodeError::InvalidPeerNpub {
let peer_identity = PeerIdentity::from_npub(&peer_config.npub).map_err(|e| {
NodeError::InvalidPeerNpub {
npub: peer_config.npub.clone(),
reason: e.to_string(),
})?;
}
})?;
let peer_node_addr = *peer_identity.node_addr();
@@ -163,10 +158,7 @@ impl Node {
(tid, TransportAddr::from_string(&addr.addr))
};
match self
.initiate_connection(transport_id, remote_addr, peer_identity)
.await
{
match self.initiate_connection(transport_id, remote_addr, Some(peer_identity)).await {
Ok(()) => return Ok(()),
Err(e) => {
debug!(
@@ -201,13 +193,9 @@ impl Node {
&mut self,
transport_id: TransportId,
remote_addr: TransportAddr,
peer_identity: PeerIdentity,
peer_identity: Option<PeerIdentity>,
) -> Result<(), NodeError> {
let peer_node_addr = *peer_identity.node_addr();
let is_connection_oriented = self
.transports
.get(&transport_id)
let is_connection_oriented = self.transports.get(&transport_id)
.map(|t| t.transport_type().connection_oriented)
.unwrap_or(false);
@@ -243,13 +231,22 @@ impl Node {
if let Some(transport) = self.transports.get(&transport_id) {
match transport.connect(&remote_addr).await {
Ok(()) => {
debug!(
peer = %self.peer_display_name(&peer_node_addr),
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
"Transport connect initiated (non-blocking)"
);
if let Some(ref id) = peer_identity {
debug!(
peer = %self.peer_display_name(id.node_addr()),
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
"Transport connect initiated (non-blocking)"
);
} else {
debug!(
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
"Transport connect initiated (anonymous discovery)"
);
}
self.pending_connects.push(super::PendingConnect {
link_id,
transport_id,
@@ -268,8 +265,7 @@ impl Node {
Ok(())
} else {
// Connectionless: proceed with immediate handshake
self.start_handshake(link_id, transport_id, remote_addr, peer_identity)
.await
self.start_handshake(link_id, transport_id, remote_addr, peer_identity).await
}
}
@@ -282,16 +278,19 @@ impl Node {
link_id: LinkId,
transport_id: TransportId,
remote_addr: TransportAddr,
peer_identity: PeerIdentity,
peer_identity: Option<PeerIdentity>,
) -> Result<(), NodeError> {
let peer_node_addr = *peer_identity.node_addr();
// Create connection in handshake phase (outbound knows expected identity)
// Create connection in handshake phase
let current_time_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let mut connection = PeerConnection::outbound(link_id, peer_identity, current_time_ms);
let mut connection = if let Some(identity) = peer_identity {
PeerConnection::outbound(link_id, identity, current_time_ms)
} else {
// Anonymous discovery connection — identity learned from XX msg2
PeerConnection::outbound_anonymous(link_id, current_time_ms)
};
// Allocate a session index for this handshake
let our_index = match self.index_allocator.allocate() {
@@ -306,17 +305,16 @@ impl Node {
// Start the Noise handshake and get message 1
let our_keypair = self.identity.keypair();
let noise_msg1 =
match connection.start_handshake(our_keypair, self.startup_epoch, current_time_ms) {
Ok(msg) => msg,
Err(e) => {
// Clean up the index and link
let _ = self.index_allocator.free(our_index);
self.links.remove(&link_id);
self.addr_to_link.remove(&(transport_id, remote_addr));
return Err(NodeError::HandshakeFailed(e.to_string()));
}
};
let noise_msg1 = match connection.start_handshake(our_keypair, self.startup_epoch, current_time_ms) {
Ok(msg) => msg,
Err(e) => {
// Clean up the index and link
let _ = self.index_allocator.free(our_index);
self.links.remove(&link_id);
self.addr_to_link.remove(&(transport_id, remote_addr));
return Err(NodeError::HandshakeFailed(e.to_string()));
}
};
// Set index and transport info on the connection
connection.set_our_index(our_index);
@@ -326,22 +324,31 @@ impl Node {
// Build wire format msg1: [0x01][sender_idx:4 LE][noise_msg1:82]
let wire_msg1 = build_msg1(our_index, &noise_msg1);
debug!(
peer = %self.peer_display_name(&peer_node_addr),
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
our_index = %our_index,
"Connection initiated"
);
if let Some(id) = connection.expected_identity() {
debug!(
peer = %self.peer_display_name(id.node_addr()),
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
our_index = %our_index,
"Connection initiated"
);
} else {
debug!(
transport_id = %transport_id,
remote_addr = %remote_addr,
link_id = %link_id,
our_index = %our_index,
"Anonymous discovery connection initiated"
);
}
// Store msg1 for resend and schedule first resend
let resend_interval = self.config.node.rate_limit.handshake_resend_interval_ms;
connection.set_handshake_msg1(wire_msg1.clone(), current_time_ms + resend_interval);
// Track in pending_outbound for msg2 dispatch
self.pending_outbound
.insert((transport_id, our_index.as_u32()), link_id);
self.pending_outbound.insert((transport_id, our_index.as_u32()), link_id);
self.connections.insert(link_id, connection);
// Send the wire format handshake message
@@ -380,7 +387,7 @@ impl Node {
/// newly discovered peers (if auto_connect is enabled).
pub(super) async fn poll_transport_discovery(&mut self) {
// Collect discoveries first to avoid borrow conflict with self
let mut to_connect = Vec::new();
let mut to_connect: Vec<(TransportId, TransportAddr, Option<PeerIdentity>)> = Vec::new();
for (transport_id, transport) in &self.transports {
if !transport.is_operational() {
@@ -396,46 +403,59 @@ impl Node {
Err(_) => continue,
};
for peer in discovered {
let pubkey = match peer.pubkey_hint {
Some(pk) => pk,
None => continue,
};
let identity = PeerIdentity::from_pubkey(pubkey);
let node_addr = *identity.node_addr();
if let Some(pubkey) = peer.pubkey_hint {
// Identity known from discovery (e.g., config-based auto-connect)
let identity = PeerIdentity::from_pubkey(pubkey);
let node_addr = *identity.node_addr();
// Skip self
if node_addr == *self.identity.node_addr() {
continue;
}
// Skip if already connected
if self.peers.contains_key(&node_addr) {
continue;
}
// Skip if connection already in progress
let connecting = self.connections.values().any(|c| {
c.expected_identity()
.map(|id| id.node_addr() == &node_addr)
.unwrap_or(false)
});
if connecting {
continue;
}
// Skip self
if node_addr == *self.identity.node_addr() {
continue;
}
// Skip if already connected
if self.peers.contains_key(&node_addr) {
continue;
}
// Skip if connection already in progress
let connecting = self.connections.values().any(|c| {
c.expected_identity()
.map(|id| id.node_addr() == &node_addr)
.unwrap_or(false)
});
if connecting {
continue;
}
to_connect.push((*transport_id, peer.addr, identity));
to_connect.push((*transport_id, peer.addr, Some(identity)));
} else {
// Anonymous discovery (shared-media beacon without identity).
// Identity will be learned from XX handshake msg2.
// Dedup by transport address — skip if link already exists.
if self.addr_to_link.contains_key(&(*transport_id, peer.addr.clone())) {
continue;
}
to_connect.push((*transport_id, peer.addr, None));
}
}
}
for (transport_id, remote_addr, identity) in to_connect {
info!(
peer = %self.peer_display_name(identity.node_addr()),
transport_id = %transport_id,
remote_addr = %remote_addr,
"Auto-connecting to discovered peer"
);
if let Err(e) = self
.initiate_connection(transport_id, remote_addr, identity)
.await
{
if let Some(ref id) = identity {
info!(
peer = %self.peer_display_name(id.node_addr()),
transport_id = %transport_id,
remote_addr = %remote_addr,
"Auto-connecting to discovered peer"
);
} else {
info!(
transport_id = %transport_id,
remote_addr = %remote_addr,
"Auto-connecting to anonymous discovered peer"
);
}
if let Err(e) = self.initiate_connection(transport_id, remote_addr, identity).await {
warn!(error = %e, "Failed to auto-connect to discovered peer");
}
}
@@ -489,7 +509,6 @@ impl Node {
}
debug!(
peer = %self.peer_display_name(pending.peer_identity.node_addr()),
transport_id = %pending.transport_id,
remote_addr = %pending.remote_addr,
link_id = %pending.link_id,
@@ -497,15 +516,12 @@ impl Node {
);
// Start the handshake now that the transport is connected
if let Err(e) = self
.start_handshake(
pending.link_id,
pending.transport_id,
pending.remote_addr.clone(),
pending.peer_identity,
)
.await
{
if let Err(e) = self.start_handshake(
pending.link_id,
pending.transport_id,
pending.remote_addr.clone(),
pending.peer_identity,
).await {
warn!(
link_id = %pending.link_id,
error = %e,
@@ -517,7 +533,6 @@ impl Node {
} else {
let reason = reason.unwrap_or_default();
warn!(
peer = %self.peer_display_name(pending.peer_identity.node_addr()),
transport_id = %pending.transport_id,
remote_addr = %pending.remote_addr,
link_id = %pending.link_id,
@@ -528,11 +543,14 @@ impl Node {
// Clean up link and schedule retry
self.remove_link(&pending.link_id);
self.links.remove(&pending.link_id);
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.schedule_retry(*pending.peer_identity.node_addr(), now_ms);
if let Some(ref id) = pending.peer_identity {
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
self.schedule_retry(*id.node_addr(), now_ms);
}
// Anonymous connections don't retry — they'll be rediscovered
}
}
}
@@ -602,24 +620,10 @@ impl Node {
// Calculate max MSS for TCP clamping
let effective_mtu = self.effective_ipv6_mtu();
let max_mss = effective_mtu.saturating_sub(40).saturating_sub(20); // IPv6 + TCP headers
info!("effective MTU: {} bytes", effective_mtu);
debug!(" max TCP MSS: {} bytes", max_mss);
// On macOS, create a shutdown pipe. Writing to it unblocks the
// reader thread's select() loop without closing the TUN fd
// (which would cause a double-close when TunDevice drops).
#[cfg(target_os = "macos")]
let (shutdown_read_fd, shutdown_write_fd) = {
let mut fds = [0i32; 2];
if unsafe { libc::pipe(fds.as_mut_ptr()) } < 0 {
return Err(NodeError::Tun(crate::upper::tun::TunError::Configure(
"failed to create shutdown pipe".into(),
)));
}
(fds[0], fds[1])
};
// Create writer (dups the fd for independent write access)
let (writer, tun_tx) = device.create_writer(max_mss)?;
@@ -637,28 +641,8 @@ impl Node {
// Spawn reader thread
let transport_mtu = self.transport_mtu();
#[cfg(target_os = "macos")]
let reader_handle = thread::spawn(move || {
run_tun_reader(
device,
mtu,
our_addr,
reader_tun_tx,
outbound_tx,
transport_mtu,
shutdown_read_fd,
);
});
#[cfg(not(target_os = "macos"))]
let reader_handle = thread::spawn(move || {
run_tun_reader(
device,
mtu,
our_addr,
reader_tun_tx,
outbound_tx,
transport_mtu,
);
run_tun_reader(device, mtu, our_addr, reader_tun_tx, outbound_tx, transport_mtu);
});
self.tun_state = TunState::Active;
@@ -667,10 +651,6 @@ impl Node {
self.tun_outbound_rx = Some(outbound_rx);
self.tun_reader_handle = Some(reader_handle);
self.tun_writer_handle = Some(writer_handle);
#[cfg(target_os = "macos")]
{
self.tun_shutdown_fd = Some(shutdown_write_fd);
}
}
Err(e) => {
self.tun_state = TunState::Failed;
@@ -687,19 +667,11 @@ impl Node {
let dns_channel_size = self.config.node.buffers.dns_channel;
let (identity_tx, identity_rx) = tokio::sync::mpsc::channel(dns_channel_size);
let dns_ttl = self.config.dns.ttl();
let base_hosts =
crate::upper::hosts::HostMap::from_peer_configs(self.config.peers());
let hosts_path =
std::path::PathBuf::from(crate::upper::hosts::DEFAULT_HOSTS_PATH);
let reloader =
crate::upper::hosts::HostMapReloader::new(base_hosts, hosts_path);
let base_hosts = crate::upper::hosts::HostMap::from_peer_configs(self.config.peers());
let hosts_path = std::path::PathBuf::from(crate::upper::hosts::DEFAULT_HOSTS_PATH);
let reloader = crate::upper::hosts::HostMapReloader::new(base_hosts, hosts_path);
info!(bind = %bind, hosts = reloader.hosts().len(), "DNS responder started for .fips domain (auto-reload enabled)");
let handle = tokio::spawn(crate::upper::dns::run_dns_responder(
socket,
identity_tx,
dns_ttl,
reloader,
));
let handle = tokio::spawn(crate::upper::dns::run_dns_responder(socket, identity_tx, dns_ttl, reloader));
self.dns_identity_rx = Some(identity_rx);
self.dns_task = Some(handle);
}
@@ -735,8 +707,7 @@ impl Node {
}
// Send disconnect notifications to all active peers before closing transports
self.send_disconnect_to_all_peers(DisconnectReason::Shutdown)
.await;
self.send_disconnect_to_all_peers(DisconnectReason::Shutdown).await;
// Shutdown transports (they're packet producers)
let transport_ids: Vec<_> = self.transports.keys().cloned().collect();
@@ -770,21 +741,11 @@ impl Node {
// Drop the tun_tx to signal the writer to stop
self.tun_tx.take();
// Delete the interface (on Linux, causes reader to get EFAULT)
// Delete the interface (causes reader to get EFAULT)
if let Err(e) = shutdown_tun_interface(&name).await {
warn!(name = %name, error = %e, "Failed to shutdown TUN interface");
}
// On macOS, signal the reader thread to exit by writing to the
// shutdown pipe. The reader's select() will wake up and break.
#[cfg(target_os = "macos")]
if let Some(fd) = self.tun_shutdown_fd.take() {
unsafe {
libc::write(fd, b"x".as_ptr() as *const libc::c_void, 1);
libc::close(fd);
}
}
// Wait for threads to finish
if let Some(handle) = self.tun_reader_handle.take() {
let _ = handle.join();
@@ -810,9 +771,7 @@ impl Node {
let plaintext = disconnect.encode();
// Collect node_addrs to avoid borrow conflict with send helper
let peer_addrs: Vec<NodeAddr> = self
.peers
.iter()
let peer_addrs: Vec<NodeAddr> = self.peers.iter()
.filter(|(_, peer)| peer.can_send() && peer.has_session())
.map(|(addr, _)| *addr)
.collect();
@@ -827,10 +786,7 @@ impl Node {
let mut sent = 0usize;
for node_addr in &peer_addrs {
match self
.send_encrypted_link_message(node_addr, &plaintext)
.await
{
match self.send_encrypted_link_message(node_addr, &plaintext).await {
Ok(()) => sent += 1,
Err(e) => {
debug!(
@@ -895,8 +851,8 @@ impl Node {
///
/// Removes the peer and suppresses auto-reconnect.
pub(crate) fn api_disconnect(&mut self, npub: &str) -> Result<serde_json::Value, String> {
let peer_identity =
PeerIdentity::from_npub(npub).map_err(|e| format!("invalid npub '{npub}': {e}"))?;
let peer_identity = PeerIdentity::from_npub(npub)
.map_err(|e| format!("invalid npub '{npub}': {e}"))?;
let node_addr = *peer_identity.node_addr();
if !self.peers.contains_key(&node_addr) {
+5 -6
View File
@@ -234,7 +234,9 @@ struct PendingConnect {
/// The remote address being connected to.
remote_addr: TransportAddr,
/// The peer identity (for handshake initiation).
peer_identity: PeerIdentity,
/// None for anonymous discovery connections where identity isn't
/// known until the XX handshake completes.
peer_identity: Option<PeerIdentity>,
}
/// A running FIPS node instance.
@@ -719,11 +721,9 @@ impl Node {
.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);
let eth = EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
transports.push(TransportHandle::Ethernet(eth));
}
}
@@ -776,14 +776,13 @@ impl Node {
let mtu = ble_config.mtu();
match crate::transport::ble::io::BluerIo::new(&adapter, mtu).await {
Ok(io) => {
let mut ble = crate::transport::ble::BleTransport::new(
let ble = crate::transport::ble::BleTransport::new(
transport_id,
name,
ble_config,
io,
packet_tx.clone(),
);
ble.set_local_pubkey(self.identity.pubkey().serialize());
transports.push(TransportHandle::Ble(ble));
}
Err(e) => {
+7 -1
View File
@@ -278,6 +278,12 @@ async fn test_ble_discovery() {
};
let io = MockBleIo::new("hci0", addr.clone());
// Probe connect must succeed for peers to reach the discovery buffer
let local = addr.clone();
io.set_connect_handler(move |target, _psm| {
let (stream, _peer) = MockBleStream::pair(local.clone(), target.clone(), 2048);
Ok(stream)
});
let (packet_tx, packet_rx) = packet_channel(256);
let mut transport = BleTransport::new(transport_id, None, config, io, packet_tx);
transport.start_async().await.unwrap();
@@ -292,7 +298,7 @@ async fn test_ble_discovery() {
tokio::time::advance(std::time::Duration::from_secs(6)).await;
tokio::task::yield_now().await;
// Without pubkey set, peers appear as bare MACs in discovery buffer
// Peers appear as bare addresses in discovery buffer after probe
let peers = transport.discover().unwrap();
assert_eq!(peers.len(), 2);
+31
View File
@@ -177,6 +177,37 @@ impl PeerConnection {
}
}
/// Create a new outbound connection without pre-known identity.
///
/// Used for anonymous discovery on shared-media transports (Ethernet,
/// BLE) where the beacon doesn't carry identity. The peer's identity
/// is learned from XX msg2 during the handshake.
pub fn outbound_anonymous(link_id: LinkId, current_time_ms: u64) -> Self {
Self {
link_id,
direction: LinkDirection::Outbound,
handshake_state: HandshakeState::Initial,
expected_identity: None,
noise_handshake: None,
noise_session: None,
started_at: current_time_ms,
last_activity: current_time_ms,
link_stats: LinkStats::new(),
our_index: None,
their_index: None,
transport_id: None,
source_addr: None,
remote_epoch: None,
peer_profile: None,
agreed_bloom_size_class: None,
handshake_msg1: None,
handshake_msg2: None,
resend_count: 0,
next_resend_at_ms: 0,
}
}
/// Create a new inbound connection (they are initiating).
///
/// For inbound, we don't know who they are until we decrypt their
-15
View File
@@ -5,7 +5,6 @@
//! identity is exchanged during the Noise handshake.
use crate::transport::{DiscoveredPeer, TransportId};
use secp256k1::XOnlyPublicKey;
use std::sync::Mutex;
use super::addr::BleAddr;
@@ -41,20 +40,6 @@ impl DiscoveryBuffer {
peers.push(peer);
}
/// Add a discovered BLE peer with a known public key.
///
/// Used after the pre-handshake pubkey exchange confirms the peer's
/// identity. The pubkey_hint enables the node's auto-connect path
/// to initiate the XX handshake.
pub fn add_peer_with_pubkey(&self, addr: &BleAddr, pubkey: XOnlyPublicKey) {
let ta = addr.to_transport_addr();
let peer = DiscoveredPeer::with_hint(self.transport_id, ta.clone(), pubkey);
let mut peers = self.peers.lock().unwrap();
let addr_str = addr.to_string_repr();
peers.retain(|p| p.addr.as_str() != Some(addr_str.as_str()));
peers.push(peer);
}
/// Drain all discovered peers since the last call.
pub fn take(&self) -> Vec<DiscoveredPeer> {
let mut peers = self.peers.lock().unwrap();
+65 -265
View File
@@ -29,14 +29,12 @@ use super::{
TransportError, TransportId, TransportState, TransportType,
};
use crate::config::BleConfig;
use crate::identity::NodeAddr;
use addr::BleAddr;
use discovery::DiscoveryBuffer;
use io::{BleIo, BleScanner, BleStream};
use pool::{BleConnection, ConnectionPool};
use stats::BleStats;
use secp256k1::XOnlyPublicKey;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::Mutex;
@@ -58,6 +56,7 @@ pub type DefaultBleTransport = BleTransport<io::BluerIo>;
#[cfg(any(not(feature = "ble"), test))]
pub type DefaultBleTransport = BleTransport<io::MockBleIo>;
// ============================================================================
// BLE Transport
// ============================================================================
@@ -92,13 +91,6 @@ pub struct BleTransport<I: BleIo> {
discovery_buffer: Arc<DiscoveryBuffer>,
/// Transport statistics.
stats: Arc<BleStats>,
/// Our public key for pre-handshake identity exchange.
///
/// BLE advertisements carry only the FIPS UUID, not the pubkey.
/// After L2CAP connection, both sides exchange `[0x00][pubkey:32]`
/// so the node layer can initiate the XX handshake.
/// Temporary — removed when FMP switches to XX.
local_pubkey: Option<[u8; 32]>,
}
/// A pending background connection attempt.
@@ -129,7 +121,6 @@ impl<I: BleIo> BleTransport<I> {
scan_probe_task: None,
discovery_buffer: Arc::new(DiscoveryBuffer::new(transport_id)),
stats: Arc::new(BleStats::new()),
local_pubkey: None,
}
}
@@ -148,15 +139,6 @@ impl<I: BleIo> BleTransport<I> {
&self.io
}
/// Set the local public key for pre-handshake identity exchange.
///
/// Must be called before `start_async()`. Without this, BLE
/// connections skip the pubkey exchange and discovered peers
/// won't have identity information for auto-connect.
pub fn set_local_pubkey(&mut self, pubkey: [u8; 32]) {
self.local_pubkey = Some(pubkey);
}
/// Start the transport asynchronously.
pub async fn start_async(&mut self) -> Result<(), TransportError> {
if !self.state.can_start() {
@@ -167,13 +149,6 @@ impl<I: BleIo> BleTransport<I> {
let psm = self.config.psm();
let adapter = self.io.adapter_name().to_string();
// Pre-compute local NodeAddr for cross-probe tie-breaking
let local_node_addr = self.local_pubkey.and_then(|pk| {
XOnlyPublicKey::from_slice(&pk)
.ok()
.map(|xonly| NodeAddr::from_pubkey(&xonly))
});
// Start L2CAP listener for inbound connections
if self.config.accept_connections() {
match self.io.listen(psm).await {
@@ -191,9 +166,7 @@ impl<I: BleIo> BleTransport<I> {
transport_id,
stats,
max_conns,
self.local_pubkey,
Arc::clone(&self.discovery_buffer),
local_node_addr,
)));
debug!(adapter = %adapter, psm = psm, "BLE accept loop started");
}
@@ -225,11 +198,9 @@ impl<I: BleIo> BleTransport<I> {
Arc::clone(&self.pool),
Arc::clone(&self.discovery_buffer),
Arc::clone(&self.stats),
self.local_pubkey,
self.config.psm(),
self.config.connect_timeout_ms(),
self.config.probe_cooldown_secs(),
local_node_addr,
self.packet_tx.clone(),
self.transport_id,
)));
@@ -364,21 +335,6 @@ impl<I: BleIo> BleTransport<I> {
}
};
// Pre-handshake pubkey exchange (temporary, pre-XX)
if let Some(ref our_pubkey) = self.local_pubkey {
match pubkey_exchange(&stream, our_pubkey).await {
Ok(peer_pubkey) => {
debug!(addr = %addr, "BLE outbound pubkey exchange complete");
self.discovery_buffer
.add_peer_with_pubkey(&ble_addr, peer_pubkey);
}
Err(e) => {
warn!(addr = %addr, error = %e, "BLE outbound pubkey exchange failed");
return Err(e);
}
}
}
self.promote_connection(addr, &ble_addr, stream).await
}
@@ -469,8 +425,6 @@ impl<I: BleIo> BleTransport<I> {
let psm = self.config.psm();
let timeout_ms = self.config.connect_timeout_ms();
let addr_clone = addr.clone();
let local_pubkey = self.local_pubkey;
let discovery_buffer = Arc::clone(&self.discovery_buffer);
let task = tokio::spawn(async move {
let result = tokio::time::timeout(
@@ -484,23 +438,6 @@ impl<I: BleIo> BleTransport<I> {
match result {
Ok(Ok(stream)) => {
// Pre-handshake pubkey exchange (temporary, pre-XX)
if let Some(ref our_pubkey) = local_pubkey {
match pubkey_exchange(&stream, our_pubkey).await {
Ok(peer_pubkey) => {
debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete");
discovery_buffer.add_peer_with_pubkey(&ble_addr, peer_pubkey);
}
Err(e) => {
warn!(
addr = %addr_clone, error = %e,
"BLE outbound pubkey exchange failed"
);
return;
}
}
}
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
let stream = Arc::new(stream);
@@ -659,66 +596,7 @@ impl<I: BleIo> Transport for BleTransport<I> {
// Background Tasks
// ============================================================================
/// Pre-handshake pubkey exchange prefix byte.
///
/// Distinguishes the identity exchange from FMP packets (version ≥ 0x01).
/// Temporary — removed when FMP uses XX handshake for BLE.
const PUBKEY_EXCHANGE_PREFIX: u8 = 0x00;
/// Pre-handshake pubkey exchange message size: `[0x00][pubkey:32]`.
const PUBKEY_EXCHANGE_SIZE: usize = 33;
/// Timeout for pubkey exchange recv (seconds).
///
/// The peer should respond in milliseconds; 5s is generous. Without this,
/// a peer that connects but never sends its pubkey blocks the calling task
/// forever — killing scan_probe_loop, accept_loop, or the event loop.
const PUBKEY_EXCHANGE_TIMEOUT_SECS: u64 = 5;
/// Exchange public keys over a newly established L2CAP connection.
///
/// Both sides send `[0x00][our_pubkey:32]` and receive the peer's.
/// Returns the peer's XOnlyPublicKey on success.
async fn pubkey_exchange<S: BleStream>(
stream: &S,
local_pubkey: &[u8; 32],
) -> Result<XOnlyPublicKey, TransportError> {
// Send our pubkey
let mut msg = [0u8; PUBKEY_EXCHANGE_SIZE];
msg[0] = PUBKEY_EXCHANGE_PREFIX;
msg[1..].copy_from_slice(local_pubkey);
stream.send(&msg).await?;
// Receive peer's pubkey (with timeout to prevent indefinite blocking)
let mut buf = [0u8; PUBKEY_EXCHANGE_SIZE];
let timeout = std::time::Duration::from_secs(PUBKEY_EXCHANGE_TIMEOUT_SECS);
let n = match tokio::time::timeout(timeout, stream.recv(&mut buf)).await {
Ok(result) => result?,
Err(_) => return Err(TransportError::Timeout),
};
if n != PUBKEY_EXCHANGE_SIZE {
return Err(TransportError::RecvFailed(format!(
"pubkey exchange: expected {} bytes, got {}",
PUBKEY_EXCHANGE_SIZE, n
)));
}
if buf[0] != PUBKEY_EXCHANGE_PREFIX {
return Err(TransportError::RecvFailed(format!(
"pubkey exchange: bad prefix 0x{:02X}",
buf[0]
)));
}
XOnlyPublicKey::from_slice(&buf[1..])
.map_err(|e| TransportError::RecvFailed(format!("pubkey exchange: invalid key: {}", e)))
}
// Beacon loop removed — advertising is now continuous (started once
// in start_async, stopped in stop_async). BLE advertising overhead
// is negligible (~0.15% duty cycle on advertising channels).
/// Accept loop: accepts inbound L2CAP connections, exchanges pubkeys,
/// and adds to pool.
/// Accept loop: accepts inbound L2CAP connections and adds to pool.
#[allow(clippy::too_many_arguments)]
async fn accept_loop<A>(
mut acceptor: A,
@@ -727,9 +605,7 @@ async fn accept_loop<A>(
transport_id: TransportId,
stats: Arc<BleStats>,
_max_conns: usize,
local_pubkey: Option<[u8; 32]>,
discovery_buffer: Arc<DiscoveryBuffer>,
local_node_addr: Option<NodeAddr>,
) where
A: io::BleAcceptor,
A::Stream: 'static,
@@ -752,33 +628,7 @@ async fn accept_loop<A>(
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
// Pre-handshake pubkey exchange (temporary, pre-XX)
if let Some(ref our_pubkey) = local_pubkey {
match pubkey_exchange(&stream, our_pubkey).await {
Ok(peer_pubkey) => {
debug!(addr = %ta, "BLE inbound pubkey exchange complete");
discovery_buffer.add_peer_with_pubkey(&addr, peer_pubkey);
// Cross-probe tie-breaker: smaller NodeAddr's
// outbound wins. If we're smaller, our outbound
// should win — drop this inbound.
if let Some(ref our_addr) = local_node_addr {
let peer_addr = NodeAddr::from_pubkey(&peer_pubkey);
if our_addr < &peer_addr {
debug!(
addr = %ta,
"BLE inbound tie-breaker: dropping (our addr < peer, outbound wins)"
);
continue;
}
}
}
Err(e) => {
debug!(addr = %ta, error = %e, "BLE inbound pubkey exchange failed");
continue;
}
}
}
discovery_buffer.add_peer(&addr);
let stream = Arc::new(stream);
@@ -885,11 +735,9 @@ async fn scan_probe_loop<I: io::BleIo>(
pool: Arc<Mutex<ConnectionPool<Arc<I::Stream>>>>,
buffer: Arc<DiscoveryBuffer>,
stats: Arc<BleStats>,
local_pubkey: Option<[u8; 32]>,
psm: u16,
connect_timeout_ms: u64,
cooldown_secs: u64,
local_node_addr: Option<NodeAddr>,
packet_tx: PacketTx,
transport_id: TransportId,
) {
@@ -956,15 +804,6 @@ async fn scan_probe_loop<I: io::BleIo>(
// Record probe time (before attempt, so cooldown applies on failure too)
last_probed.insert(addr.clone(), tokio::time::Instant::now());
// Need pubkey for probe
let our_pubkey = match local_pubkey {
Some(pk) => pk,
None => {
buffer.add_peer(&addr);
continue;
}
};
// L2CAP connect
let stream = match tokio::time::timeout(
std::time::Duration::from_millis(connect_timeout_ms),
@@ -984,76 +823,54 @@ async fn scan_probe_loop<I: io::BleIo>(
}
};
// Pubkey exchange, then promote connection to pool
// Promote connection to pool
let ta = addr.to_transport_addr();
match pubkey_exchange(&stream, &our_pubkey).await {
Ok(peer_pubkey) => {
debug!(addr = %addr, "BLE probe complete");
debug!(addr = %addr, "BLE probe complete");
// Cross-probe tie-breaker: smaller NodeAddr's outbound wins.
// If we lose, drop connection — accept_loop handles inbound.
if let Some(ref our_addr) = local_node_addr {
let peer_addr = NodeAddr::from_pubkey(&peer_pubkey);
if our_addr >= &peer_addr {
debug!(
addr = %addr,
"BLE probe tie-breaker: yielding to peer's outbound"
);
buffer.add_peer_with_pubkey(&addr, peer_pubkey);
continue;
}
}
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
let stream = Arc::new(stream);
// Promote connection to pool — no second L2CAP connect needed
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
let stream = Arc::new(stream);
let recv_task = tokio::spawn(receive_loop(
Arc::clone(&stream),
ta.clone(),
Arc::clone(&pool),
packet_tx.clone(),
transport_id,
Arc::clone(&stats),
recv_mtu,
));
let recv_task = tokio::spawn(receive_loop(
Arc::clone(&stream),
ta.clone(),
Arc::clone(&pool),
packet_tx.clone(),
transport_id,
Arc::clone(&stats),
recv_mtu,
));
let conn = BleConnection {
stream,
recv_task: Some(recv_task),
send_mtu,
recv_mtu,
established_at: tokio::time::Instant::now(),
is_static: false,
addr: addr.clone(),
};
let conn = BleConnection {
stream,
recv_task: Some(recv_task),
send_mtu,
recv_mtu,
established_at: tokio::time::Instant::now(),
is_static: false,
addr: addr.clone(),
};
let mut pool_guard = pool.lock().await;
match pool_guard.insert(ta.clone(), conn) {
Ok(Some(evicted)) => {
stats.record_pool_eviction();
debug!(addr = %ta, evicted = %evicted, "BLE probe promoted (evicted peer)");
}
Ok(None) => {
debug!(addr = %ta, "BLE probe promoted to pool");
}
Err(e) => {
warn!(addr = %ta, error = %e, "BLE pool full, probe connection dropped");
stats.record_connection_rejected();
}
}
drop(pool_guard);
stats.record_connection_established();
pending_addrs.retain(|a| a != &addr);
// Report to node layer for auto-connect / handshake
buffer.add_peer_with_pubkey(&addr, peer_pubkey);
let mut pool_guard = pool.lock().await;
match pool_guard.insert(ta.clone(), conn) {
Ok(Some(evicted)) => {
stats.record_pool_eviction();
debug!(addr = %ta, evicted = %evicted, "BLE probe promoted (evicted peer)");
}
Ok(None) => {
debug!(addr = %ta, "BLE probe promoted to pool");
}
Err(e) => {
debug!(addr = %addr, error = %e, "BLE probe pubkey exchange failed");
warn!(addr = %ta, error = %e, "BLE pool full, probe connection dropped");
stats.record_connection_rejected();
}
}
drop(pool_guard);
stats.record_connection_established();
pending_addrs.retain(|a| a != &addr);
// Report to node layer for auto-connect / handshake
buffer.add_peer(&addr);
}
}
@@ -1075,13 +892,16 @@ mod tests {
fn make_transport(
io: MockBleIo,
) -> (
BleTransport<MockBleIo>,
tokio::sync::mpsc::Receiver<ReceivedPacket>,
) {
) -> (BleTransport<MockBleIo>, tokio::sync::mpsc::Receiver<ReceivedPacket>) {
let (tx, rx) = tokio::sync::mpsc::channel(64);
let config = BleConfig::default();
let transport = BleTransport::new(TransportId::new(1), None, config, io, tx);
let transport = BleTransport::new(
TransportId::new(1),
None,
config,
io,
tx,
);
(transport, rx)
}
@@ -1122,6 +942,13 @@ mod tests {
#[tokio::test(start_paused = true)]
async fn test_scan_discovers_peers() {
let io = MockBleIo::new("hci0", test_addr(1));
// Probe connect must succeed for peers to reach the discovery buffer
let local = test_addr(1);
io.set_connect_handler(move |addr, _psm| {
let (stream, _peer) =
io::MockBleStream::pair(local.clone(), addr.clone(), 2048);
Ok(stream)
});
let (mut transport, _rx) = make_transport(io);
transport.start_async().await.unwrap();
@@ -1136,7 +963,7 @@ mod tests {
// Let the expired entries get processed
tokio::task::yield_now().await;
// Without pubkey set, scan results go to discovery buffer as bare MACs
// Scan results go to discovery buffer as bare addresses after probe
let peers = transport.discovery_buffer.take();
assert_eq!(peers.len(), 2);
}
@@ -1144,6 +971,12 @@ mod tests {
#[tokio::test(start_paused = true)]
async fn test_scan_deduplicates() {
let io = MockBleIo::new("hci0", test_addr(1));
let local = test_addr(1);
io.set_connect_handler(move |addr, _psm| {
let (stream, _peer) =
io::MockBleStream::pair(local.clone(), addr.clone(), 2048);
Ok(stream)
});
let (mut transport, _rx) = make_transport(io);
transport.start_async().await.unwrap();
@@ -1172,40 +1005,7 @@ mod tests {
let io = MockBleIo::new("hci0", test_addr(1));
let (transport, _rx) = make_transport(io);
let addr = test_addr(2).to_transport_addr();
assert_eq!(
transport.connection_state_sync(&addr),
ConnectionState::None
);
assert_eq!(transport.connection_state_sync(&addr), ConnectionState::None);
}
/// Verify that the cross-probe tie-breaker follows the same convention
/// as `cross_connection_winner`: smaller NodeAddr's outbound wins.
#[test]
fn test_tiebreaker_convention() {
use secp256k1::{Secp256k1, SecretKey};
let secp = Secp256k1::new();
let sk_a = SecretKey::from_slice(&[1u8; 32]).unwrap();
let sk_b = SecretKey::from_slice(&[2u8; 32]).unwrap();
let (pk_a, _) = sk_a.public_key(&secp).x_only_public_key();
let (pk_b, _) = sk_b.public_key(&secp).x_only_public_key();
let addr_a = NodeAddr::from_pubkey(&pk_a);
let addr_b = NodeAddr::from_pubkey(&pk_b);
// Determine which is smaller
let (smaller, larger) = if addr_a < addr_b {
(addr_a, addr_b)
} else {
(addr_b, addr_a)
};
// scan_loop (outbound): promotes when our_addr < peer_addr
// Smaller node scanning larger → our_addr < peer_addr → promote (win)
assert!(smaller < larger, "test setup: smaller < larger");
// accept_loop (inbound): drops when our_addr < peer_addr
// Smaller node accepting from larger → drops inbound (outbound wins)
// This means: smaller always uses outbound, larger always uses inbound
}
}
+54 -46
View File
@@ -1,12 +1,10 @@
//! Ethernet LAN discovery via broadcast beacons.
//!
//! Beacon format (34 bytes total):
//! - `0x01` (1 byte): frame type = discovery announcement
//! - `0x01` (1 byte): discovery protocol version
//! - x-only public key (32 bytes): node's Nostr identity
//! Beacon format (5 bytes total):
//! - Unified header (4 bytes): `[type:1][flags:1][length:2 LE]`
//! - Version (1 byte): discovery protocol version
use crate::transport::{DiscoveredPeer, TransportAddr, TransportId};
use secp256k1::XOnlyPublicKey;
use std::sync::Mutex;
/// Discovery protocol version.
@@ -18,32 +16,41 @@ pub const FRAME_TYPE_BEACON: u8 = 0x01;
/// Frame type prefix for FIPS data frames.
pub const FRAME_TYPE_DATA: u8 = 0x00;
/// Total beacon payload size: type(1) + version(1) + pubkey(32).
pub const BEACON_SIZE: usize = 34;
/// Shared header size for all Ethernet frame types: type(1) + flags(1) + length(2).
pub const ETHERNET_HEADER_SIZE: usize = 4;
/// Beacon payload size: version(1).
pub const BEACON_PAYLOAD_SIZE: usize = 1;
/// Total beacon size: header(4) + payload(1).
pub const BEACON_SIZE: usize = ETHERNET_HEADER_SIZE + BEACON_PAYLOAD_SIZE;
/// Build a discovery announcement beacon payload.
pub fn build_beacon(pubkey: &XOnlyPublicKey) -> [u8; BEACON_SIZE] {
pub fn build_beacon() -> [u8; BEACON_SIZE] {
let mut buf = [0u8; BEACON_SIZE];
buf[0] = FRAME_TYPE_BEACON;
buf[1] = DISCOVERY_VERSION;
buf[2..BEACON_SIZE].copy_from_slice(&pubkey.serialize());
buf[1] = 0x00; // flags (reserved)
buf[2..4].copy_from_slice(&(BEACON_PAYLOAD_SIZE as u16).to_le_bytes());
buf[4] = DISCOVERY_VERSION;
buf
}
/// Parse a discovery announcement beacon payload.
///
/// Returns the sender's public key, or None if the payload is invalid.
pub fn parse_beacon(data: &[u8]) -> Option<XOnlyPublicKey> {
/// Returns true if the payload is a valid beacon, false otherwise.
pub fn parse_beacon(data: &[u8]) -> bool {
if data.len() < BEACON_SIZE {
return None;
return false;
}
if data[0] != FRAME_TYPE_BEACON {
return None;
return false;
}
if data[1] != DISCOVERY_VERSION {
return None;
// flags byte data[1] accepted as any value for forward compatibility
let length = u16::from_le_bytes([data[2], data[3]]);
if length < 1 {
return false;
}
XOnlyPublicKey::from_slice(&data[2..34]).ok()
data[4] == DISCOVERY_VERSION
}
/// Buffer for discovered peers, drained by `discover()`.
@@ -62,11 +69,10 @@ impl DiscoveryBuffer {
}
/// Add a discovered peer from a received beacon.
pub fn add_peer(&self, src_mac: [u8; 6], pubkey: XOnlyPublicKey) {
pub fn add_peer(&self, src_mac: [u8; 6]) {
let addr = TransportAddr::from_bytes(&src_mac);
let peer = DiscoveredPeer::with_hint(self.transport_id, addr, pubkey);
let peer = DiscoveredPeer::new(self.transport_id, addr);
let mut peers = self.peers.lock().unwrap();
// Deduplicate by MAC address — keep the latest
peers.retain(|p| p.addr.as_bytes() != src_mac);
peers.push(peer);
}
@@ -85,46 +91,38 @@ impl DiscoveryBuffer {
#[cfg(test)]
mod tests {
use super::*;
use secp256k1::{Secp256k1, SecretKey};
fn test_pubkey() -> XOnlyPublicKey {
let secp = Secp256k1::new();
let sk = SecretKey::from_slice(&[0x42; 32]).unwrap();
let (xonly, _) = sk.public_key(&secp).x_only_public_key();
xonly
}
#[test]
fn test_build_parse_beacon() {
let pubkey = test_pubkey();
let beacon = build_beacon(&pubkey);
let beacon = build_beacon();
assert_eq!(beacon.len(), BEACON_SIZE);
assert_eq!(beacon[0], FRAME_TYPE_BEACON);
assert_eq!(beacon[1], DISCOVERY_VERSION);
assert_eq!(beacon[1], 0x00); // flags
assert_eq!(u16::from_le_bytes([beacon[2], beacon[3]]), 1); // length
assert_eq!(beacon[4], DISCOVERY_VERSION);
let parsed = parse_beacon(&beacon).unwrap();
assert_eq!(parsed, pubkey);
assert!(parse_beacon(&beacon));
}
#[test]
fn test_parse_beacon_too_short() {
assert!(parse_beacon(&[0x01, 0x01]).is_none());
assert!(parse_beacon(&[]).is_none());
assert!(!parse_beacon(&[0x01, 0x00, 0x01, 0x00]));
assert!(!parse_beacon(&[]));
}
#[test]
fn test_parse_beacon_wrong_type() {
let mut beacon = build_beacon(&test_pubkey());
let mut beacon = build_beacon();
beacon[0] = 0x00; // data frame, not beacon
assert!(parse_beacon(&beacon).is_none());
assert!(!parse_beacon(&beacon));
}
#[test]
fn test_parse_beacon_wrong_version() {
let mut beacon = build_beacon(&test_pubkey());
beacon[1] = 0xFF;
assert!(parse_beacon(&beacon).is_none());
let mut beacon = build_beacon();
beacon[4] = 0xFF;
assert!(!parse_beacon(&beacon));
}
#[test]
@@ -133,18 +131,24 @@ mod tests {
assert_eq!(FRAME_TYPE_BEACON, 0x01);
}
#[test]
fn test_beacon_unified_header() {
let beacon = build_beacon();
assert_eq!(beacon[1], 0x00); // flags reserved, zero
assert_eq!(u16::from_le_bytes([beacon[2], beacon[3]]), 1); // length field = 1
}
#[test]
fn test_discovery_buffer() {
let buffer = DiscoveryBuffer::new(TransportId::new(1));
let pubkey = test_pubkey();
let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff];
buffer.add_peer(mac, pubkey);
buffer.add_peer(mac);
let peers = buffer.take();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].addr.as_bytes(), &mac);
assert_eq!(peers[0].pubkey_hint, Some(pubkey));
assert!(peers[0].pubkey_hint.is_none());
// Second take should be empty
let peers = buffer.take();
@@ -154,13 +158,17 @@ mod tests {
#[test]
fn test_discovery_buffer_dedup() {
let buffer = DiscoveryBuffer::new(TransportId::new(1));
let pubkey = test_pubkey();
let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff];
buffer.add_peer(mac, pubkey);
buffer.add_peer(mac, pubkey); // same MAC again
buffer.add_peer(mac);
buffer.add_peer(mac); // same MAC again
let peers = buffer.take();
assert_eq!(peers.len(), 1);
}
#[test]
fn test_beacon_size() {
assert_eq!(BEACON_SIZE, 5);
}
}
+79 -96
View File
@@ -1,9 +1,8 @@
//! Ethernet Transport Implementation
//!
//! Provides raw Ethernet transport for FIPS peer communication. On Linux,
//! uses AF_PACKET/SOCK_DGRAM sockets; on macOS, uses BPF devices (`/dev/bpf*`).
//! Works on wired Ethernet and WiFi interfaces (kernel mac80211 abstracts
//! 802.11 transparently on Linux).
//! Provides raw Ethernet transport for FIPS peer communication using
//! AF_PACKET sockets with SOCK_DGRAM. Works on wired Ethernet and WiFi
//! interfaces (kernel mac80211 abstracts 802.11 transparently).
pub mod discovery;
pub mod socket;
@@ -14,11 +13,12 @@ use super::{
TransportId, TransportState, TransportType,
};
use crate::config::EthernetConfig;
use discovery::{DiscoveryBuffer, FRAME_TYPE_BEACON, FRAME_TYPE_DATA, build_beacon, parse_beacon};
use socket::{AsyncPacketSocket, ETHERNET_BROADCAST, PacketSocket};
use discovery::{
build_beacon, parse_beacon, DiscoveryBuffer, FRAME_TYPE_BEACON, FRAME_TYPE_DATA,
};
use socket::{AsyncPacketSocket, PacketSocket, ETHERNET_BROADCAST};
use stats::EthernetStats;
use secp256k1::XOnlyPublicKey;
use std::sync::Arc;
use tokio::task::JoinHandle;
use tracing::{debug, info, trace, warn};
@@ -49,14 +49,12 @@ pub struct EthernetTransport {
local_mac: Option<[u8; 6]>,
/// Interface name (from config).
interface: String,
/// Effective MTU (interface MTU - 1 for frame type prefix).
/// Effective MTU (interface MTU - 4 for frame header).
effective_mtu: u16,
/// Discovery buffer for discovered peers.
discovery_buffer: Arc<DiscoveryBuffer>,
/// Transport-level statistics.
stats: Arc<EthernetStats>,
/// Node's public key for beacon construction.
local_pubkey: Option<XOnlyPublicKey>,
}
impl EthernetTransport {
@@ -82,10 +80,9 @@ impl EthernetTransport {
beacon_task: None,
local_mac: None,
interface,
effective_mtu: 1499, // default, updated on start
effective_mtu: 1496, // default, updated on start
discovery_buffer,
stats,
local_pubkey: None,
}
}
@@ -104,13 +101,6 @@ impl EthernetTransport {
self.local_mac
}
/// Set the node's public key for beacon construction.
///
/// Must be called before start if announce is enabled.
pub fn set_local_pubkey(&mut self, pubkey: XOnlyPublicKey) {
self.local_pubkey = Some(pubkey);
}
/// Get a reference to the statistics.
pub fn stats(&self) -> &Arc<EthernetStats> {
&self.stats
@@ -134,13 +124,13 @@ impl EthernetTransport {
let local_mac = raw_socket.local_mac()?;
let if_mtu = raw_socket.interface_mtu()?;
// Effective MTU: interface MTU minus 3 bytes for frame header
// (1 byte frame type + 2 bytes LE payload length)
// Effective MTU: interface MTU minus 4 bytes for frame header
// (1 byte frame type + 1 byte flags + 2 bytes LE payload length)
let effective_mtu = if let Some(configured_mtu) = self.config.mtu {
// Config MTU cannot exceed interface MTU - 3
configured_mtu.min(if_mtu.saturating_sub(3))
// Config MTU cannot exceed interface MTU - 4
configured_mtu.min(if_mtu.saturating_sub(4))
} else {
if_mtu.saturating_sub(3)
if_mtu.saturating_sub(4)
};
self.effective_mtu = effective_mtu;
self.local_mac = Some(local_mac);
@@ -179,34 +169,25 @@ impl EthernetTransport {
// Spawn beacon sender if announce is enabled
if self.config.announce() {
if let Some(pubkey) = self.local_pubkey {
let beacon_socket = socket.clone();
let interval_secs = self.config.beacon_interval_secs();
let beacon_stats = self.stats.clone();
let beacon_transport_id = self.transport_id;
let beacon_socket = socket.clone();
let interval_secs = self.config.beacon_interval_secs();
let beacon_stats = self.stats.clone();
let beacon_transport_id = self.transport_id;
let beacon_interface = self.config.interface.clone();
let beacon_ethertype = self.config.ethertype();
let beacon_interface = self.config.interface.clone();
let beacon_ethertype = self.config.ethertype();
let beacon_task = tokio::spawn(async move {
beacon_sender_loop(
beacon_socket,
pubkey,
interval_secs,
beacon_stats,
beacon_transport_id,
beacon_interface,
beacon_ethertype,
)
.await;
});
self.beacon_task = Some(beacon_task);
} else {
warn!(
transport_id = %self.transport_id,
"Announce enabled but no local pubkey set; beacons disabled"
);
}
let beacon_task = tokio::spawn(async move {
beacon_sender_loop(
beacon_socket,
interval_secs,
beacon_stats,
beacon_transport_id,
beacon_interface,
beacon_ethertype,
)
.await;
});
self.beacon_task = Some(beacon_task);
}
self.state = TransportState::Up;
@@ -239,30 +220,16 @@ impl EthernetTransport {
return Err(TransportError::NotStarted);
}
// Signal the socket to shut down. On macOS this writes to the
// shutdown pipe, waking the reader thread's select() immediately.
// On Linux this is a no-op (AsyncFd cancellation handles it).
if let Some(ref socket) = self.socket {
socket.shutdown();
}
// Abort tasks. On Linux, safe to await since all I/O is
// AsyncFd-based and cancellation-safe. On macOS, do NOT await —
// on a current_thread runtime the aborted task can't be polled
// while we're blocked on the JoinHandle, causing a deadlock.
// Abort beacon task
if let Some(task) = self.beacon_task.take() {
task.abort();
#[cfg(not(target_os = "macos"))]
{
let _ = task.await;
}
let _ = task.await;
}
// Abort receive task
if let Some(task) = self.recv_task.take() {
task.abort();
#[cfg(not(target_os = "macos"))]
{
let _ = task.await;
}
let _ = task.await;
}
// Drop socket
@@ -282,8 +249,7 @@ impl EthernetTransport {
/// Send a packet asynchronously.
///
/// The data is prepended with a FRAME_TYPE_DATA prefix byte before
/// transmission.
/// The data is prepended with a 4-byte frame header before transmission.
pub async fn send_async(
&self,
addr: &TransportAddr,
@@ -303,12 +269,13 @@ impl EthernetTransport {
let dest_mac = parse_mac_addr(addr)?;
let socket = self.socket.as_ref().ok_or(TransportError::NotStarted)?;
// Prepend frame type prefix and 2-byte LE payload length.
// Prepend 4-byte frame header: type(1) + flags(1) + length(2 LE).
// The length field lets the receiver trim Ethernet minimum-frame padding
// (NICs pad frames shorter than 46 bytes payload to 46 bytes with zeros,
// which would otherwise corrupt AEAD ciphertext verification).
let mut frame = Vec::with_capacity(3 + data.len());
let mut frame = Vec::with_capacity(4 + data.len());
frame.push(FRAME_TYPE_DATA);
frame.push(0x00); // flags (reserved)
frame.extend_from_slice(&(data.len() as u16).to_le_bytes());
frame.extend_from_slice(data);
@@ -322,8 +289,8 @@ impl EthernetTransport {
"Ethernet frame sent"
);
// Return the data bytes sent (excluding frame type prefix and length field)
Ok(bytes_sent.saturating_sub(3))
// Return the data bytes sent (excluding 4-byte frame header)
Ok(bytes_sent.saturating_sub(4))
}
}
@@ -406,22 +373,23 @@ async fn ethernet_receive_loop(
let frame_type = buf[0];
match frame_type {
FRAME_TYPE_DATA => {
// Data frame: [type:1][length:2 LE][payload:N]
// Use the length field to trim Ethernet minimum-frame padding.
if len < 3 {
// Data frame: [type:1][flags:1][length:2 LE][payload:N]
if len < 4 {
trace!("Data frame too short ({len} bytes), ignoring");
continue;
}
let payload_len = u16::from_le_bytes([buf[1], buf[2]]) as usize;
if payload_len > len - 3 {
// buf[1] is flags (reserved, ignored for now)
let payload_len =
u16::from_le_bytes([buf[2], buf[3]]) as usize;
if payload_len > len - 4 {
trace!(
"Data frame length field ({payload_len}) exceeds \
available bytes ({}), ignoring",
len - 3
len - 4
);
continue;
}
let data = buf[3..3 + payload_len].to_vec();
let data = buf[4..4 + payload_len].to_vec();
let addr = TransportAddr::from_bytes(&src_mac);
let packet = ReceivedPacket::new(transport_id, addr, data);
@@ -443,8 +411,8 @@ async fn ethernet_receive_loop(
FRAME_TYPE_BEACON => {
stats.record_beacon_recv();
if discovery_enabled && let Some(pubkey) = parse_beacon(&buf[..len]) {
discovery_buffer.add_peer(src_mac, pubkey);
if discovery_enabled && parse_beacon(&buf[..len]) {
discovery_buffer.add_peer(src_mac);
trace!(
transport_id = %transport_id,
remote_mac = %format_mac(&src_mac),
@@ -488,7 +456,6 @@ async fn ethernet_receive_loop(
/// failures, attempts to open a fresh socket on the same interface.
async fn beacon_sender_loop(
mut socket: Arc<AsyncPacketSocket>,
pubkey: XOnlyPublicKey,
interval_secs: u64,
stats: Arc<EthernetStats>,
transport_id: TransportId,
@@ -498,7 +465,7 @@ async fn beacon_sender_loop(
/// Number of consecutive ENXIO errors before attempting socket reopen.
const REOPEN_THRESHOLD: u32 = 3;
let beacon = build_beacon(&pubkey);
let beacon = build_beacon();
let interval = tokio::time::Duration::from_secs(interval_secs);
debug!(
@@ -712,28 +679,31 @@ mod tests {
#[test]
fn test_frame_type_data_prefix() {
// Verify data frames have type prefix + 2-byte LE length + payload
// Verify data frames have 4-byte header + payload
let data = vec![1, 2, 3, 4];
let mut frame = Vec::with_capacity(3 + data.len());
let mut frame = Vec::with_capacity(4 + data.len());
frame.push(FRAME_TYPE_DATA);
frame.push(0x00); // flags
frame.extend_from_slice(&(data.len() as u16).to_le_bytes());
frame.extend_from_slice(&data);
assert_eq!(frame[0], 0x00); // frame type
assert_eq!(u16::from_le_bytes([frame[1], frame[2]]), 4); // length
assert_eq!(&frame[3..], &[1, 2, 3, 4]); // payload
assert_eq!(frame[1], 0x00); // flags
assert_eq!(u16::from_le_bytes([frame[2], frame[3]]), 4); // length
assert_eq!(&frame[4..], &[1, 2, 3, 4]); // payload
}
#[test]
fn test_data_frame_padding_trimmed() {
// Simulate Ethernet minimum-frame padding: a 4-byte payload produces
// a 7-byte frame (type + len + payload), padded to 46 bytes by NIC.
// an 8-byte frame (header + payload), padded to 46 bytes by NIC.
let payload = vec![0xAA, 0xBB, 0xCC, 0xDD];
let payload_len = payload.len() as u16;
// Build frame as sender would
let mut frame = Vec::with_capacity(3 + payload.len());
let mut frame = Vec::with_capacity(4 + payload.len());
frame.push(FRAME_TYPE_DATA);
frame.push(0x00); // flags
frame.extend_from_slice(&payload_len.to_le_bytes());
frame.extend_from_slice(&payload);
@@ -741,13 +711,26 @@ mod tests {
frame.resize(46, 0x00);
// Receiver extracts using length field
let recv_len = u16::from_le_bytes([frame[1], frame[2]]) as usize;
let extracted = &frame[3..3 + recv_len];
let recv_len = u16::from_le_bytes([frame[2], frame[3]]) as usize;
let extracted = &frame[4..4 + recv_len];
assert_eq!(extracted, &[0xAA, 0xBB, 0xCC, 0xDD]);
}
#[test]
fn test_beacon_size() {
assert_eq!(discovery::BEACON_SIZE, 34);
assert_eq!(discovery::BEACON_SIZE, 5);
}
#[test]
fn test_unified_header_flags_byte() {
// Build a data frame and verify the flags byte at offset 1 is 0x00
let data = vec![0x42];
let mut frame = Vec::with_capacity(4 + data.len());
frame.push(FRAME_TYPE_DATA);
frame.push(0x00); // flags
frame.extend_from_slice(&(data.len() as u16).to_le_bytes());
frame.extend_from_slice(&data);
assert_eq!(frame[1], 0x00);
}
}