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(); 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 // Build and send msg3
let our_index = conn.our_index().unwrap_or(header.receiver_idx); let our_index = conn.our_index().unwrap_or(header.receiver_idx);
let wire_msg3 = build_msg3(our_index, header.sender_idx, &msg3_bytes); 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(); 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); let our_index = conn.our_index().unwrap_or(header.receiver_idx);
// Identity-based restart/rekey detection. // Identity-based restart/rekey detection.
+139 -183
View File
@@ -1,11 +1,11 @@
//! Node lifecycle management: start, stop, and peer connection initiation. //! Node lifecycle management: start, stop, and peer connection initiation.
use super::{Node, NodeError, NodeState}; use super::{Node, NodeError, NodeState};
use crate::node::wire::build_msg1;
use crate::peer::PeerConnection; use crate::peer::PeerConnection;
use crate::protocol::{Disconnect, DisconnectReason}; use crate::protocol::{Disconnect, DisconnectReason};
use crate::transport::{Link, LinkDirection, LinkId, TransportAddr, TransportId, packet_channel}; use crate::transport::{packet_channel, Link, LinkDirection, LinkId, TransportAddr, TransportId};
use crate::upper::tun::{TunDevice, TunState, run_tun_reader, shutdown_tun_interface}; use crate::upper::tun::{run_tun_reader, shutdown_tun_interface, TunDevice, TunState};
use crate::node::wire::build_msg1;
use crate::{NodeAddr, PeerIdentity}; use crate::{NodeAddr, PeerIdentity};
use std::thread; use std::thread;
use std::time::Duration; use std::time::Duration;
@@ -50,10 +50,7 @@ impl Node {
return; return;
} }
debug!( debug!(count = peer_configs.len(), "Initiating static peer connections");
count = peer_configs.len(),
"Initiating static peer connections"
);
for peer_config in peer_configs { for peer_config in peer_configs {
if let Err(e) = self.initiate_peer_connection(&peer_config).await { if let Err(e) = self.initiate_peer_connection(&peer_config).await {
@@ -70,16 +67,14 @@ impl Node {
/// Initiate a connection to a single peer. /// Initiate a connection to a single peer.
/// ///
/// Creates a link, starts the Noise handshake, and sends the first message. /// Creates a link, starts the Noise handshake, and sends the first message.
pub(super) async fn initiate_peer_connection( pub(super) async fn initiate_peer_connection(&mut self, peer_config: &crate::config::PeerConfig) -> Result<(), NodeError> {
&mut self,
peer_config: &crate::config::PeerConfig,
) -> Result<(), NodeError> {
// Parse the peer's npub to get their identity // Parse the peer's npub to get their identity
let peer_identity = let peer_identity = PeerIdentity::from_npub(&peer_config.npub).map_err(|e| {
PeerIdentity::from_npub(&peer_config.npub).map_err(|e| NodeError::InvalidPeerNpub { NodeError::InvalidPeerNpub {
npub: peer_config.npub.clone(), npub: peer_config.npub.clone(),
reason: e.to_string(), reason: e.to_string(),
})?; }
})?;
let peer_node_addr = *peer_identity.node_addr(); let peer_node_addr = *peer_identity.node_addr();
@@ -163,10 +158,7 @@ impl Node {
(tid, TransportAddr::from_string(&addr.addr)) (tid, TransportAddr::from_string(&addr.addr))
}; };
match self match self.initiate_connection(transport_id, remote_addr, Some(peer_identity)).await {
.initiate_connection(transport_id, remote_addr, peer_identity)
.await
{
Ok(()) => return Ok(()), Ok(()) => return Ok(()),
Err(e) => { Err(e) => {
debug!( debug!(
@@ -201,13 +193,9 @@ impl Node {
&mut self, &mut self,
transport_id: TransportId, transport_id: TransportId,
remote_addr: TransportAddr, remote_addr: TransportAddr,
peer_identity: PeerIdentity, peer_identity: Option<PeerIdentity>,
) -> Result<(), NodeError> { ) -> 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) .map(|t| t.transport_type().connection_oriented)
.unwrap_or(false); .unwrap_or(false);
@@ -243,13 +231,22 @@ impl Node {
if let Some(transport) = self.transports.get(&transport_id) { if let Some(transport) = self.transports.get(&transport_id) {
match transport.connect(&remote_addr).await { match transport.connect(&remote_addr).await {
Ok(()) => { Ok(()) => {
debug!( if let Some(ref id) = peer_identity {
peer = %self.peer_display_name(&peer_node_addr), debug!(
transport_id = %transport_id, peer = %self.peer_display_name(id.node_addr()),
remote_addr = %remote_addr, transport_id = %transport_id,
link_id = %link_id, remote_addr = %remote_addr,
"Transport connect initiated (non-blocking)" 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 { self.pending_connects.push(super::PendingConnect {
link_id, link_id,
transport_id, transport_id,
@@ -268,8 +265,7 @@ impl Node {
Ok(()) Ok(())
} else { } else {
// Connectionless: proceed with immediate handshake // Connectionless: proceed with immediate handshake
self.start_handshake(link_id, transport_id, remote_addr, peer_identity) self.start_handshake(link_id, transport_id, remote_addr, peer_identity).await
.await
} }
} }
@@ -282,16 +278,19 @@ impl Node {
link_id: LinkId, link_id: LinkId,
transport_id: TransportId, transport_id: TransportId,
remote_addr: TransportAddr, remote_addr: TransportAddr,
peer_identity: PeerIdentity, peer_identity: Option<PeerIdentity>,
) -> Result<(), NodeError> { ) -> Result<(), NodeError> {
let peer_node_addr = *peer_identity.node_addr(); // Create connection in handshake phase
// Create connection in handshake phase (outbound knows expected identity)
let current_time_ms = std::time::SystemTime::now() let current_time_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH) .duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64) .map(|d| d.as_millis() as u64)
.unwrap_or(0); .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 // Allocate a session index for this handshake
let our_index = match self.index_allocator.allocate() { let our_index = match self.index_allocator.allocate() {
@@ -306,17 +305,16 @@ impl Node {
// Start the Noise handshake and get message 1 // Start the Noise handshake and get message 1
let our_keypair = self.identity.keypair(); let our_keypair = self.identity.keypair();
let noise_msg1 = let noise_msg1 = match connection.start_handshake(our_keypair, self.startup_epoch, current_time_ms) {
match connection.start_handshake(our_keypair, self.startup_epoch, current_time_ms) { Ok(msg) => msg,
Ok(msg) => msg, Err(e) => {
Err(e) => { // Clean up the index and link
// Clean up the index and link let _ = self.index_allocator.free(our_index);
let _ = self.index_allocator.free(our_index); self.links.remove(&link_id);
self.links.remove(&link_id); self.addr_to_link.remove(&(transport_id, remote_addr));
self.addr_to_link.remove(&(transport_id, remote_addr)); return Err(NodeError::HandshakeFailed(e.to_string()));
return Err(NodeError::HandshakeFailed(e.to_string())); }
} };
};
// Set index and transport info on the connection // Set index and transport info on the connection
connection.set_our_index(our_index); connection.set_our_index(our_index);
@@ -326,22 +324,31 @@ impl Node {
// Build wire format msg1: [0x01][sender_idx:4 LE][noise_msg1:82] // Build wire format msg1: [0x01][sender_idx:4 LE][noise_msg1:82]
let wire_msg1 = build_msg1(our_index, &noise_msg1); let wire_msg1 = build_msg1(our_index, &noise_msg1);
debug!( if let Some(id) = connection.expected_identity() {
peer = %self.peer_display_name(&peer_node_addr), debug!(
transport_id = %transport_id, peer = %self.peer_display_name(id.node_addr()),
remote_addr = %remote_addr, transport_id = %transport_id,
link_id = %link_id, remote_addr = %remote_addr,
our_index = %our_index, link_id = %link_id,
"Connection initiated" 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 // Store msg1 for resend and schedule first resend
let resend_interval = self.config.node.rate_limit.handshake_resend_interval_ms; let resend_interval = self.config.node.rate_limit.handshake_resend_interval_ms;
connection.set_handshake_msg1(wire_msg1.clone(), current_time_ms + resend_interval); connection.set_handshake_msg1(wire_msg1.clone(), current_time_ms + resend_interval);
// Track in pending_outbound for msg2 dispatch // Track in pending_outbound for msg2 dispatch
self.pending_outbound self.pending_outbound.insert((transport_id, our_index.as_u32()), link_id);
.insert((transport_id, our_index.as_u32()), link_id);
self.connections.insert(link_id, connection); self.connections.insert(link_id, connection);
// Send the wire format handshake message // Send the wire format handshake message
@@ -380,7 +387,7 @@ impl Node {
/// newly discovered peers (if auto_connect is enabled). /// newly discovered peers (if auto_connect is enabled).
pub(super) async fn poll_transport_discovery(&mut self) { pub(super) async fn poll_transport_discovery(&mut self) {
// Collect discoveries first to avoid borrow conflict with 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 { for (transport_id, transport) in &self.transports {
if !transport.is_operational() { if !transport.is_operational() {
@@ -396,46 +403,59 @@ impl Node {
Err(_) => continue, Err(_) => continue,
}; };
for peer in discovered { for peer in discovered {
let pubkey = match peer.pubkey_hint { if let Some(pubkey) = peer.pubkey_hint {
Some(pk) => pk, // Identity known from discovery (e.g., config-based auto-connect)
None => continue, let identity = PeerIdentity::from_pubkey(pubkey);
}; let node_addr = *identity.node_addr();
let identity = PeerIdentity::from_pubkey(pubkey);
let node_addr = *identity.node_addr();
// Skip self // Skip self
if node_addr == *self.identity.node_addr() { if node_addr == *self.identity.node_addr() {
continue; continue;
} }
// Skip if already connected // Skip if already connected
if self.peers.contains_key(&node_addr) { if self.peers.contains_key(&node_addr) {
continue; continue;
} }
// Skip if connection already in progress // Skip if connection already in progress
let connecting = self.connections.values().any(|c| { let connecting = self.connections.values().any(|c| {
c.expected_identity() c.expected_identity()
.map(|id| id.node_addr() == &node_addr) .map(|id| id.node_addr() == &node_addr)
.unwrap_or(false) .unwrap_or(false)
}); });
if connecting { if connecting {
continue; 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 { for (transport_id, remote_addr, identity) in to_connect {
info!( if let Some(ref id) = identity {
peer = %self.peer_display_name(identity.node_addr()), info!(
transport_id = %transport_id, peer = %self.peer_display_name(id.node_addr()),
remote_addr = %remote_addr, transport_id = %transport_id,
"Auto-connecting to discovered peer" remote_addr = %remote_addr,
); "Auto-connecting to discovered peer"
if let Err(e) = self );
.initiate_connection(transport_id, remote_addr, identity) } else {
.await 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"); warn!(error = %e, "Failed to auto-connect to discovered peer");
} }
} }
@@ -489,7 +509,6 @@ impl Node {
} }
debug!( debug!(
peer = %self.peer_display_name(pending.peer_identity.node_addr()),
transport_id = %pending.transport_id, transport_id = %pending.transport_id,
remote_addr = %pending.remote_addr, remote_addr = %pending.remote_addr,
link_id = %pending.link_id, link_id = %pending.link_id,
@@ -497,15 +516,12 @@ impl Node {
); );
// Start the handshake now that the transport is connected // Start the handshake now that the transport is connected
if let Err(e) = self if let Err(e) = self.start_handshake(
.start_handshake( pending.link_id,
pending.link_id, pending.transport_id,
pending.transport_id, pending.remote_addr.clone(),
pending.remote_addr.clone(), pending.peer_identity,
pending.peer_identity, ).await {
)
.await
{
warn!( warn!(
link_id = %pending.link_id, link_id = %pending.link_id,
error = %e, error = %e,
@@ -517,7 +533,6 @@ impl Node {
} else { } else {
let reason = reason.unwrap_or_default(); let reason = reason.unwrap_or_default();
warn!( warn!(
peer = %self.peer_display_name(pending.peer_identity.node_addr()),
transport_id = %pending.transport_id, transport_id = %pending.transport_id,
remote_addr = %pending.remote_addr, remote_addr = %pending.remote_addr,
link_id = %pending.link_id, link_id = %pending.link_id,
@@ -528,11 +543,14 @@ impl Node {
// Clean up link and schedule retry // Clean up link and schedule retry
self.remove_link(&pending.link_id); self.remove_link(&pending.link_id);
self.links.remove(&pending.link_id); self.links.remove(&pending.link_id);
let now_ms = std::time::SystemTime::now() if let Some(ref id) = pending.peer_identity {
.duration_since(std::time::UNIX_EPOCH) let now_ms = std::time::SystemTime::now()
.map(|d| d.as_millis() as u64) .duration_since(std::time::UNIX_EPOCH)
.unwrap_or(0); .map(|d| d.as_millis() as u64)
self.schedule_retry(*pending.peer_identity.node_addr(), now_ms); .unwrap_or(0);
self.schedule_retry(*id.node_addr(), now_ms);
}
// Anonymous connections don't retry — they'll be rediscovered
} }
} }
} }
@@ -606,20 +624,6 @@ impl Node {
info!("effective MTU: {} bytes", effective_mtu); info!("effective MTU: {} bytes", effective_mtu);
debug!(" max TCP MSS: {} bytes", max_mss); 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) // Create writer (dups the fd for independent write access)
let (writer, tun_tx) = device.create_writer(max_mss)?; let (writer, tun_tx) = device.create_writer(max_mss)?;
@@ -637,28 +641,8 @@ impl Node {
// Spawn reader thread // Spawn reader thread
let transport_mtu = self.transport_mtu(); let transport_mtu = self.transport_mtu();
#[cfg(target_os = "macos")]
let reader_handle = thread::spawn(move || { let reader_handle = thread::spawn(move || {
run_tun_reader( run_tun_reader(device, mtu, our_addr, reader_tun_tx, outbound_tx, transport_mtu);
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,
);
}); });
self.tun_state = TunState::Active; self.tun_state = TunState::Active;
@@ -667,10 +651,6 @@ impl Node {
self.tun_outbound_rx = Some(outbound_rx); self.tun_outbound_rx = Some(outbound_rx);
self.tun_reader_handle = Some(reader_handle); self.tun_reader_handle = Some(reader_handle);
self.tun_writer_handle = Some(writer_handle); self.tun_writer_handle = Some(writer_handle);
#[cfg(target_os = "macos")]
{
self.tun_shutdown_fd = Some(shutdown_write_fd);
}
} }
Err(e) => { Err(e) => {
self.tun_state = TunState::Failed; self.tun_state = TunState::Failed;
@@ -687,19 +667,11 @@ impl Node {
let dns_channel_size = self.config.node.buffers.dns_channel; let dns_channel_size = self.config.node.buffers.dns_channel;
let (identity_tx, identity_rx) = tokio::sync::mpsc::channel(dns_channel_size); let (identity_tx, identity_rx) = tokio::sync::mpsc::channel(dns_channel_size);
let dns_ttl = self.config.dns.ttl(); let dns_ttl = self.config.dns.ttl();
let base_hosts = let base_hosts = crate::upper::hosts::HostMap::from_peer_configs(self.config.peers());
crate::upper::hosts::HostMap::from_peer_configs(self.config.peers()); let hosts_path = std::path::PathBuf::from(crate::upper::hosts::DEFAULT_HOSTS_PATH);
let hosts_path = let reloader = crate::upper::hosts::HostMapReloader::new(base_hosts, 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)"); 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( let handle = tokio::spawn(crate::upper::dns::run_dns_responder(socket, identity_tx, dns_ttl, reloader));
socket,
identity_tx,
dns_ttl,
reloader,
));
self.dns_identity_rx = Some(identity_rx); self.dns_identity_rx = Some(identity_rx);
self.dns_task = Some(handle); self.dns_task = Some(handle);
} }
@@ -735,8 +707,7 @@ impl Node {
} }
// Send disconnect notifications to all active peers before closing transports // Send disconnect notifications to all active peers before closing transports
self.send_disconnect_to_all_peers(DisconnectReason::Shutdown) self.send_disconnect_to_all_peers(DisconnectReason::Shutdown).await;
.await;
// Shutdown transports (they're packet producers) // Shutdown transports (they're packet producers)
let transport_ids: Vec<_> = self.transports.keys().cloned().collect(); 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 // Drop the tun_tx to signal the writer to stop
self.tun_tx.take(); 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 { if let Err(e) = shutdown_tun_interface(&name).await {
warn!(name = %name, error = %e, "Failed to shutdown TUN interface"); 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 // Wait for threads to finish
if let Some(handle) = self.tun_reader_handle.take() { if let Some(handle) = self.tun_reader_handle.take() {
let _ = handle.join(); let _ = handle.join();
@@ -810,9 +771,7 @@ impl Node {
let plaintext = disconnect.encode(); let plaintext = disconnect.encode();
// Collect node_addrs to avoid borrow conflict with send helper // Collect node_addrs to avoid borrow conflict with send helper
let peer_addrs: Vec<NodeAddr> = self let peer_addrs: Vec<NodeAddr> = self.peers.iter()
.peers
.iter()
.filter(|(_, peer)| peer.can_send() && peer.has_session()) .filter(|(_, peer)| peer.can_send() && peer.has_session())
.map(|(addr, _)| *addr) .map(|(addr, _)| *addr)
.collect(); .collect();
@@ -827,10 +786,7 @@ impl Node {
let mut sent = 0usize; let mut sent = 0usize;
for node_addr in &peer_addrs { for node_addr in &peer_addrs {
match self match self.send_encrypted_link_message(node_addr, &plaintext).await {
.send_encrypted_link_message(node_addr, &plaintext)
.await
{
Ok(()) => sent += 1, Ok(()) => sent += 1,
Err(e) => { Err(e) => {
debug!( debug!(
@@ -895,8 +851,8 @@ impl Node {
/// ///
/// Removes the peer and suppresses auto-reconnect. /// Removes the peer and suppresses auto-reconnect.
pub(crate) fn api_disconnect(&mut self, npub: &str) -> Result<serde_json::Value, String> { pub(crate) fn api_disconnect(&mut self, npub: &str) -> Result<serde_json::Value, String> {
let peer_identity = let peer_identity = PeerIdentity::from_npub(npub)
PeerIdentity::from_npub(npub).map_err(|e| format!("invalid npub '{npub}': {e}"))?; .map_err(|e| format!("invalid npub '{npub}': {e}"))?;
let node_addr = *peer_identity.node_addr(); let node_addr = *peer_identity.node_addr();
if !self.peers.contains_key(&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. /// The remote address being connected to.
remote_addr: TransportAddr, remote_addr: TransportAddr,
/// The peer identity (for handshake initiation). /// 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. /// A running FIPS node instance.
@@ -719,11 +721,9 @@ impl Node {
.map(|(name, config)| (name.map(|s| s.to_string()), config.clone())) .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
.collect(); .collect();
let xonly = self.identity.pubkey();
for (name, eth_config) in eth_instances { for (name, eth_config) in eth_instances {
let transport_id = self.allocate_transport_id(); let transport_id = self.allocate_transport_id();
let mut eth = EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone()); let eth = EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
eth.set_local_pubkey(xonly);
transports.push(TransportHandle::Ethernet(eth)); transports.push(TransportHandle::Ethernet(eth));
} }
} }
@@ -776,14 +776,13 @@ impl Node {
let mtu = ble_config.mtu(); let mtu = ble_config.mtu();
match crate::transport::ble::io::BluerIo::new(&adapter, mtu).await { match crate::transport::ble::io::BluerIo::new(&adapter, mtu).await {
Ok(io) => { Ok(io) => {
let mut ble = crate::transport::ble::BleTransport::new( let ble = crate::transport::ble::BleTransport::new(
transport_id, transport_id,
name, name,
ble_config, ble_config,
io, io,
packet_tx.clone(), packet_tx.clone(),
); );
ble.set_local_pubkey(self.identity.pubkey().serialize());
transports.push(TransportHandle::Ble(ble)); transports.push(TransportHandle::Ble(ble));
} }
Err(e) => { Err(e) => {
+7 -1
View File
@@ -278,6 +278,12 @@ async fn test_ble_discovery() {
}; };
let io = MockBleIo::new("hci0", addr.clone()); 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 (packet_tx, packet_rx) = packet_channel(256);
let mut transport = BleTransport::new(transport_id, None, config, io, packet_tx); let mut transport = BleTransport::new(transport_id, None, config, io, packet_tx);
transport.start_async().await.unwrap(); 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::time::advance(std::time::Duration::from_secs(6)).await;
tokio::task::yield_now().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(); let peers = transport.discover().unwrap();
assert_eq!(peers.len(), 2); 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). /// Create a new inbound connection (they are initiating).
/// ///
/// For inbound, we don't know who they are until we decrypt their /// 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. //! identity is exchanged during the Noise handshake.
use crate::transport::{DiscoveredPeer, TransportId}; use crate::transport::{DiscoveredPeer, TransportId};
use secp256k1::XOnlyPublicKey;
use std::sync::Mutex; use std::sync::Mutex;
use super::addr::BleAddr; use super::addr::BleAddr;
@@ -41,20 +40,6 @@ impl DiscoveryBuffer {
peers.push(peer); 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. /// Drain all discovered peers since the last call.
pub fn take(&self) -> Vec<DiscoveredPeer> { pub fn take(&self) -> Vec<DiscoveredPeer> {
let mut peers = self.peers.lock().unwrap(); let mut peers = self.peers.lock().unwrap();
+65 -265
View File
@@ -29,14 +29,12 @@ use super::{
TransportError, TransportId, TransportState, TransportType, TransportError, TransportId, TransportState, TransportType,
}; };
use crate::config::BleConfig; use crate::config::BleConfig;
use crate::identity::NodeAddr;
use addr::BleAddr; use addr::BleAddr;
use discovery::DiscoveryBuffer; use discovery::DiscoveryBuffer;
use io::{BleIo, BleScanner, BleStream}; use io::{BleIo, BleScanner, BleStream};
use pool::{BleConnection, ConnectionPool}; use pool::{BleConnection, ConnectionPool};
use stats::BleStats; use stats::BleStats;
use secp256k1::XOnlyPublicKey;
use std::collections::HashMap; use std::collections::HashMap;
use std::sync::Arc; use std::sync::Arc;
use tokio::sync::Mutex; use tokio::sync::Mutex;
@@ -58,6 +56,7 @@ pub type DefaultBleTransport = BleTransport<io::BluerIo>;
#[cfg(any(not(feature = "ble"), test))] #[cfg(any(not(feature = "ble"), test))]
pub type DefaultBleTransport = BleTransport<io::MockBleIo>; pub type DefaultBleTransport = BleTransport<io::MockBleIo>;
// ============================================================================ // ============================================================================
// BLE Transport // BLE Transport
// ============================================================================ // ============================================================================
@@ -92,13 +91,6 @@ pub struct BleTransport<I: BleIo> {
discovery_buffer: Arc<DiscoveryBuffer>, discovery_buffer: Arc<DiscoveryBuffer>,
/// Transport statistics. /// Transport statistics.
stats: Arc<BleStats>, 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. /// A pending background connection attempt.
@@ -129,7 +121,6 @@ impl<I: BleIo> BleTransport<I> {
scan_probe_task: None, scan_probe_task: None,
discovery_buffer: Arc::new(DiscoveryBuffer::new(transport_id)), discovery_buffer: Arc::new(DiscoveryBuffer::new(transport_id)),
stats: Arc::new(BleStats::new()), stats: Arc::new(BleStats::new()),
local_pubkey: None,
} }
} }
@@ -148,15 +139,6 @@ impl<I: BleIo> BleTransport<I> {
&self.io &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. /// Start the transport asynchronously.
pub async fn start_async(&mut self) -> Result<(), TransportError> { pub async fn start_async(&mut self) -> Result<(), TransportError> {
if !self.state.can_start() { if !self.state.can_start() {
@@ -167,13 +149,6 @@ impl<I: BleIo> BleTransport<I> {
let psm = self.config.psm(); let psm = self.config.psm();
let adapter = self.io.adapter_name().to_string(); 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 // Start L2CAP listener for inbound connections
if self.config.accept_connections() { if self.config.accept_connections() {
match self.io.listen(psm).await { match self.io.listen(psm).await {
@@ -191,9 +166,7 @@ impl<I: BleIo> BleTransport<I> {
transport_id, transport_id,
stats, stats,
max_conns, max_conns,
self.local_pubkey,
Arc::clone(&self.discovery_buffer), Arc::clone(&self.discovery_buffer),
local_node_addr,
))); )));
debug!(adapter = %adapter, psm = psm, "BLE accept loop started"); 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.pool),
Arc::clone(&self.discovery_buffer), Arc::clone(&self.discovery_buffer),
Arc::clone(&self.stats), Arc::clone(&self.stats),
self.local_pubkey,
self.config.psm(), self.config.psm(),
self.config.connect_timeout_ms(), self.config.connect_timeout_ms(),
self.config.probe_cooldown_secs(), self.config.probe_cooldown_secs(),
local_node_addr,
self.packet_tx.clone(), self.packet_tx.clone(),
self.transport_id, 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 self.promote_connection(addr, &ble_addr, stream).await
} }
@@ -469,8 +425,6 @@ impl<I: BleIo> BleTransport<I> {
let psm = self.config.psm(); let psm = self.config.psm();
let timeout_ms = self.config.connect_timeout_ms(); let timeout_ms = self.config.connect_timeout_ms();
let addr_clone = addr.clone(); 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 task = tokio::spawn(async move {
let result = tokio::time::timeout( let result = tokio::time::timeout(
@@ -484,23 +438,6 @@ impl<I: BleIo> BleTransport<I> {
match result { match result {
Ok(Ok(stream)) => { 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 send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu(); let recv_mtu = stream.recv_mtu();
let stream = Arc::new(stream); let stream = Arc::new(stream);
@@ -659,66 +596,7 @@ impl<I: BleIo> Transport for BleTransport<I> {
// Background Tasks // Background Tasks
// ============================================================================ // ============================================================================
/// Pre-handshake pubkey exchange prefix byte. /// Accept loop: accepts inbound L2CAP connections and adds to pool.
///
/// 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.
#[allow(clippy::too_many_arguments)] #[allow(clippy::too_many_arguments)]
async fn accept_loop<A>( async fn accept_loop<A>(
mut acceptor: A, mut acceptor: A,
@@ -727,9 +605,7 @@ async fn accept_loop<A>(
transport_id: TransportId, transport_id: TransportId,
stats: Arc<BleStats>, stats: Arc<BleStats>,
_max_conns: usize, _max_conns: usize,
local_pubkey: Option<[u8; 32]>,
discovery_buffer: Arc<DiscoveryBuffer>, discovery_buffer: Arc<DiscoveryBuffer>,
local_node_addr: Option<NodeAddr>,
) where ) where
A: io::BleAcceptor, A: io::BleAcceptor,
A::Stream: 'static, A::Stream: 'static,
@@ -752,33 +628,7 @@ async fn accept_loop<A>(
let send_mtu = stream.send_mtu(); let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu(); let recv_mtu = stream.recv_mtu();
// Pre-handshake pubkey exchange (temporary, pre-XX) discovery_buffer.add_peer(&addr);
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;
}
}
}
let stream = Arc::new(stream); let stream = Arc::new(stream);
@@ -885,11 +735,9 @@ async fn scan_probe_loop<I: io::BleIo>(
pool: Arc<Mutex<ConnectionPool<Arc<I::Stream>>>>, pool: Arc<Mutex<ConnectionPool<Arc<I::Stream>>>>,
buffer: Arc<DiscoveryBuffer>, buffer: Arc<DiscoveryBuffer>,
stats: Arc<BleStats>, stats: Arc<BleStats>,
local_pubkey: Option<[u8; 32]>,
psm: u16, psm: u16,
connect_timeout_ms: u64, connect_timeout_ms: u64,
cooldown_secs: u64, cooldown_secs: u64,
local_node_addr: Option<NodeAddr>,
packet_tx: PacketTx, packet_tx: PacketTx,
transport_id: TransportId, 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) // Record probe time (before attempt, so cooldown applies on failure too)
last_probed.insert(addr.clone(), tokio::time::Instant::now()); 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 // L2CAP connect
let stream = match tokio::time::timeout( let stream = match tokio::time::timeout(
std::time::Duration::from_millis(connect_timeout_ms), 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(); let ta = addr.to_transport_addr();
match pubkey_exchange(&stream, &our_pubkey).await { debug!(addr = %addr, "BLE probe complete");
Ok(peer_pubkey) => {
debug!(addr = %addr, "BLE probe complete");
// Cross-probe tie-breaker: smaller NodeAddr's outbound wins. let send_mtu = stream.send_mtu();
// If we lose, drop connection — accept_loop handles inbound. let recv_mtu = stream.recv_mtu();
if let Some(ref our_addr) = local_node_addr { let stream = Arc::new(stream);
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;
}
}
// Promote connection to pool — no second L2CAP connect needed let recv_task = tokio::spawn(receive_loop(
let send_mtu = stream.send_mtu(); Arc::clone(&stream),
let recv_mtu = stream.recv_mtu(); ta.clone(),
let stream = Arc::new(stream); Arc::clone(&pool),
packet_tx.clone(),
transport_id,
Arc::clone(&stats),
recv_mtu,
));
let recv_task = tokio::spawn(receive_loop( let conn = BleConnection {
Arc::clone(&stream), stream,
ta.clone(), recv_task: Some(recv_task),
Arc::clone(&pool), send_mtu,
packet_tx.clone(), recv_mtu,
transport_id, established_at: tokio::time::Instant::now(),
Arc::clone(&stats), is_static: false,
recv_mtu, addr: addr.clone(),
)); };
let conn = BleConnection { let mut pool_guard = pool.lock().await;
stream, match pool_guard.insert(ta.clone(), conn) {
recv_task: Some(recv_task), Ok(Some(evicted)) => {
send_mtu, stats.record_pool_eviction();
recv_mtu, debug!(addr = %ta, evicted = %evicted, "BLE probe promoted (evicted peer)");
established_at: tokio::time::Instant::now(), }
is_static: false, Ok(None) => {
addr: addr.clone(), debug!(addr = %ta, "BLE probe promoted to pool");
};
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);
} }
Err(e) => { 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( fn make_transport(
io: MockBleIo, 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 (tx, rx) = tokio::sync::mpsc::channel(64);
let config = BleConfig::default(); 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) (transport, rx)
} }
@@ -1122,6 +942,13 @@ mod tests {
#[tokio::test(start_paused = true)] #[tokio::test(start_paused = true)]
async fn test_scan_discovers_peers() { async fn test_scan_discovers_peers() {
let io = MockBleIo::new("hci0", test_addr(1)); 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); let (mut transport, _rx) = make_transport(io);
transport.start_async().await.unwrap(); transport.start_async().await.unwrap();
@@ -1136,7 +963,7 @@ mod tests {
// Let the expired entries get processed // Let the expired entries get processed
tokio::task::yield_now().await; 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(); let peers = transport.discovery_buffer.take();
assert_eq!(peers.len(), 2); assert_eq!(peers.len(), 2);
} }
@@ -1144,6 +971,12 @@ mod tests {
#[tokio::test(start_paused = true)] #[tokio::test(start_paused = true)]
async fn test_scan_deduplicates() { async fn test_scan_deduplicates() {
let io = MockBleIo::new("hci0", test_addr(1)); 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); let (mut transport, _rx) = make_transport(io);
transport.start_async().await.unwrap(); transport.start_async().await.unwrap();
@@ -1172,40 +1005,7 @@ mod tests {
let io = MockBleIo::new("hci0", test_addr(1)); let io = MockBleIo::new("hci0", test_addr(1));
let (transport, _rx) = make_transport(io); let (transport, _rx) = make_transport(io);
let addr = test_addr(2).to_transport_addr(); let addr = test_addr(2).to_transport_addr();
assert_eq!( assert_eq!(transport.connection_state_sync(&addr), ConnectionState::None);
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. //! Ethernet LAN discovery via broadcast beacons.
//! //!
//! Beacon format (34 bytes total): //! Beacon format (5 bytes total):
//! - `0x01` (1 byte): frame type = discovery announcement //! - Unified header (4 bytes): `[type:1][flags:1][length:2 LE]`
//! - `0x01` (1 byte): discovery protocol version //! - Version (1 byte): discovery protocol version
//! - x-only public key (32 bytes): node's Nostr identity
use crate::transport::{DiscoveredPeer, TransportAddr, TransportId}; use crate::transport::{DiscoveredPeer, TransportAddr, TransportId};
use secp256k1::XOnlyPublicKey;
use std::sync::Mutex; use std::sync::Mutex;
/// Discovery protocol version. /// Discovery protocol version.
@@ -18,32 +16,41 @@ pub const FRAME_TYPE_BEACON: u8 = 0x01;
/// Frame type prefix for FIPS data frames. /// Frame type prefix for FIPS data frames.
pub const FRAME_TYPE_DATA: u8 = 0x00; pub const FRAME_TYPE_DATA: u8 = 0x00;
/// Total beacon payload size: type(1) + version(1) + pubkey(32). /// Shared header size for all Ethernet frame types: type(1) + flags(1) + length(2).
pub const BEACON_SIZE: usize = 34; 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. /// 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]; let mut buf = [0u8; BEACON_SIZE];
buf[0] = FRAME_TYPE_BEACON; buf[0] = FRAME_TYPE_BEACON;
buf[1] = DISCOVERY_VERSION; buf[1] = 0x00; // flags (reserved)
buf[2..BEACON_SIZE].copy_from_slice(&pubkey.serialize()); buf[2..4].copy_from_slice(&(BEACON_PAYLOAD_SIZE as u16).to_le_bytes());
buf[4] = DISCOVERY_VERSION;
buf buf
} }
/// Parse a discovery announcement beacon payload. /// Parse a discovery announcement beacon payload.
/// ///
/// Returns the sender's public key, or None if the payload is invalid. /// Returns true if the payload is a valid beacon, false otherwise.
pub fn parse_beacon(data: &[u8]) -> Option<XOnlyPublicKey> { pub fn parse_beacon(data: &[u8]) -> bool {
if data.len() < BEACON_SIZE { if data.len() < BEACON_SIZE {
return None; return false;
} }
if data[0] != FRAME_TYPE_BEACON { if data[0] != FRAME_TYPE_BEACON {
return None; return false;
} }
if data[1] != DISCOVERY_VERSION { // flags byte data[1] accepted as any value for forward compatibility
return None; 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()`. /// Buffer for discovered peers, drained by `discover()`.
@@ -62,11 +69,10 @@ impl DiscoveryBuffer {
} }
/// Add a discovered peer from a received beacon. /// 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 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(); let mut peers = self.peers.lock().unwrap();
// Deduplicate by MAC address — keep the latest
peers.retain(|p| p.addr.as_bytes() != src_mac); peers.retain(|p| p.addr.as_bytes() != src_mac);
peers.push(peer); peers.push(peer);
} }
@@ -85,46 +91,38 @@ impl DiscoveryBuffer {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::*; 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] #[test]
fn test_build_parse_beacon() { fn test_build_parse_beacon() {
let pubkey = test_pubkey(); let beacon = build_beacon();
let beacon = build_beacon(&pubkey);
assert_eq!(beacon.len(), BEACON_SIZE); assert_eq!(beacon.len(), BEACON_SIZE);
assert_eq!(beacon[0], FRAME_TYPE_BEACON); 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!(parse_beacon(&beacon));
assert_eq!(parsed, pubkey);
} }
#[test] #[test]
fn test_parse_beacon_too_short() { fn test_parse_beacon_too_short() {
assert!(parse_beacon(&[0x01, 0x01]).is_none()); assert!(!parse_beacon(&[0x01, 0x00, 0x01, 0x00]));
assert!(parse_beacon(&[]).is_none()); assert!(!parse_beacon(&[]));
} }
#[test] #[test]
fn test_parse_beacon_wrong_type() { 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 beacon[0] = 0x00; // data frame, not beacon
assert!(parse_beacon(&beacon).is_none()); assert!(!parse_beacon(&beacon));
} }
#[test] #[test]
fn test_parse_beacon_wrong_version() { fn test_parse_beacon_wrong_version() {
let mut beacon = build_beacon(&test_pubkey()); let mut beacon = build_beacon();
beacon[1] = 0xFF; beacon[4] = 0xFF;
assert!(parse_beacon(&beacon).is_none()); assert!(!parse_beacon(&beacon));
} }
#[test] #[test]
@@ -133,18 +131,24 @@ mod tests {
assert_eq!(FRAME_TYPE_BEACON, 0x01); 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] #[test]
fn test_discovery_buffer() { fn test_discovery_buffer() {
let buffer = DiscoveryBuffer::new(TransportId::new(1)); let buffer = DiscoveryBuffer::new(TransportId::new(1));
let pubkey = test_pubkey();
let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff]; let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff];
buffer.add_peer(mac, pubkey); buffer.add_peer(mac);
let peers = buffer.take(); let peers = buffer.take();
assert_eq!(peers.len(), 1); assert_eq!(peers.len(), 1);
assert_eq!(peers[0].addr.as_bytes(), &mac); 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 // Second take should be empty
let peers = buffer.take(); let peers = buffer.take();
@@ -154,13 +158,17 @@ mod tests {
#[test] #[test]
fn test_discovery_buffer_dedup() { fn test_discovery_buffer_dedup() {
let buffer = DiscoveryBuffer::new(TransportId::new(1)); let buffer = DiscoveryBuffer::new(TransportId::new(1));
let pubkey = test_pubkey();
let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff]; let mac = [0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff];
buffer.add_peer(mac, pubkey); buffer.add_peer(mac);
buffer.add_peer(mac, pubkey); // same MAC again buffer.add_peer(mac); // same MAC again
let peers = buffer.take(); let peers = buffer.take();
assert_eq!(peers.len(), 1); 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 //! Ethernet Transport Implementation
//! //!
//! Provides raw Ethernet transport for FIPS peer communication. On Linux, //! Provides raw Ethernet transport for FIPS peer communication using
//! uses AF_PACKET/SOCK_DGRAM sockets; on macOS, uses BPF devices (`/dev/bpf*`). //! AF_PACKET sockets with SOCK_DGRAM. Works on wired Ethernet and WiFi
//! Works on wired Ethernet and WiFi interfaces (kernel mac80211 abstracts //! interfaces (kernel mac80211 abstracts 802.11 transparently).
//! 802.11 transparently on Linux).
pub mod discovery; pub mod discovery;
pub mod socket; pub mod socket;
@@ -14,11 +13,12 @@ use super::{
TransportId, TransportState, TransportType, TransportId, TransportState, TransportType,
}; };
use crate::config::EthernetConfig; use crate::config::EthernetConfig;
use discovery::{DiscoveryBuffer, FRAME_TYPE_BEACON, FRAME_TYPE_DATA, build_beacon, parse_beacon}; use discovery::{
use socket::{AsyncPacketSocket, ETHERNET_BROADCAST, PacketSocket}; build_beacon, parse_beacon, DiscoveryBuffer, FRAME_TYPE_BEACON, FRAME_TYPE_DATA,
};
use socket::{AsyncPacketSocket, PacketSocket, ETHERNET_BROADCAST};
use stats::EthernetStats; use stats::EthernetStats;
use secp256k1::XOnlyPublicKey;
use std::sync::Arc; use std::sync::Arc;
use tokio::task::JoinHandle; use tokio::task::JoinHandle;
use tracing::{debug, info, trace, warn}; use tracing::{debug, info, trace, warn};
@@ -49,14 +49,12 @@ pub struct EthernetTransport {
local_mac: Option<[u8; 6]>, local_mac: Option<[u8; 6]>,
/// Interface name (from config). /// Interface name (from config).
interface: String, interface: String,
/// Effective MTU (interface MTU - 1 for frame type prefix). /// Effective MTU (interface MTU - 4 for frame header).
effective_mtu: u16, effective_mtu: u16,
/// Discovery buffer for discovered peers. /// Discovery buffer for discovered peers.
discovery_buffer: Arc<DiscoveryBuffer>, discovery_buffer: Arc<DiscoveryBuffer>,
/// Transport-level statistics. /// Transport-level statistics.
stats: Arc<EthernetStats>, stats: Arc<EthernetStats>,
/// Node's public key for beacon construction.
local_pubkey: Option<XOnlyPublicKey>,
} }
impl EthernetTransport { impl EthernetTransport {
@@ -82,10 +80,9 @@ impl EthernetTransport {
beacon_task: None, beacon_task: None,
local_mac: None, local_mac: None,
interface, interface,
effective_mtu: 1499, // default, updated on start effective_mtu: 1496, // default, updated on start
discovery_buffer, discovery_buffer,
stats, stats,
local_pubkey: None,
} }
} }
@@ -104,13 +101,6 @@ impl EthernetTransport {
self.local_mac 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. /// Get a reference to the statistics.
pub fn stats(&self) -> &Arc<EthernetStats> { pub fn stats(&self) -> &Arc<EthernetStats> {
&self.stats &self.stats
@@ -134,13 +124,13 @@ impl EthernetTransport {
let local_mac = raw_socket.local_mac()?; let local_mac = raw_socket.local_mac()?;
let if_mtu = raw_socket.interface_mtu()?; let if_mtu = raw_socket.interface_mtu()?;
// Effective MTU: interface MTU minus 3 bytes for frame header // Effective MTU: interface MTU minus 4 bytes for frame header
// (1 byte frame type + 2 bytes LE payload length) // (1 byte frame type + 1 byte flags + 2 bytes LE payload length)
let effective_mtu = if let Some(configured_mtu) = self.config.mtu { let effective_mtu = if let Some(configured_mtu) = self.config.mtu {
// Config MTU cannot exceed interface MTU - 3 // Config MTU cannot exceed interface MTU - 4
configured_mtu.min(if_mtu.saturating_sub(3)) configured_mtu.min(if_mtu.saturating_sub(4))
} else { } else {
if_mtu.saturating_sub(3) if_mtu.saturating_sub(4)
}; };
self.effective_mtu = effective_mtu; self.effective_mtu = effective_mtu;
self.local_mac = Some(local_mac); self.local_mac = Some(local_mac);
@@ -179,34 +169,25 @@ impl EthernetTransport {
// Spawn beacon sender if announce is enabled // Spawn beacon sender if announce is enabled
if self.config.announce() { if self.config.announce() {
if let Some(pubkey) = self.local_pubkey { let beacon_socket = socket.clone();
let beacon_socket = socket.clone(); let interval_secs = self.config.beacon_interval_secs();
let interval_secs = self.config.beacon_interval_secs(); let beacon_stats = self.stats.clone();
let beacon_stats = self.stats.clone(); let beacon_transport_id = self.transport_id;
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_task = tokio::spawn(async move {
let beacon_ethertype = self.config.ethertype(); beacon_sender_loop(
beacon_socket,
let beacon_task = tokio::spawn(async move { interval_secs,
beacon_sender_loop( beacon_stats,
beacon_socket, beacon_transport_id,
pubkey, beacon_interface,
interval_secs, beacon_ethertype,
beacon_stats, )
beacon_transport_id, .await;
beacon_interface, });
beacon_ethertype, self.beacon_task = Some(beacon_task);
)
.await;
});
self.beacon_task = Some(beacon_task);
} else {
warn!(
transport_id = %self.transport_id,
"Announce enabled but no local pubkey set; beacons disabled"
);
}
} }
self.state = TransportState::Up; self.state = TransportState::Up;
@@ -239,30 +220,16 @@ impl EthernetTransport {
return Err(TransportError::NotStarted); return Err(TransportError::NotStarted);
} }
// Signal the socket to shut down. On macOS this writes to the // Abort beacon task
// 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.
if let Some(task) = self.beacon_task.take() { if let Some(task) = self.beacon_task.take() {
task.abort(); task.abort();
#[cfg(not(target_os = "macos"))] let _ = task.await;
{
let _ = task.await;
}
} }
// Abort receive task
if let Some(task) = self.recv_task.take() { if let Some(task) = self.recv_task.take() {
task.abort(); task.abort();
#[cfg(not(target_os = "macos"))] let _ = task.await;
{
let _ = task.await;
}
} }
// Drop socket // Drop socket
@@ -282,8 +249,7 @@ impl EthernetTransport {
/// Send a packet asynchronously. /// Send a packet asynchronously.
/// ///
/// The data is prepended with a FRAME_TYPE_DATA prefix byte before /// The data is prepended with a 4-byte frame header before transmission.
/// transmission.
pub async fn send_async( pub async fn send_async(
&self, &self,
addr: &TransportAddr, addr: &TransportAddr,
@@ -303,12 +269,13 @@ impl EthernetTransport {
let dest_mac = parse_mac_addr(addr)?; let dest_mac = parse_mac_addr(addr)?;
let socket = self.socket.as_ref().ok_or(TransportError::NotStarted)?; 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 // The length field lets the receiver trim Ethernet minimum-frame padding
// (NICs pad frames shorter than 46 bytes payload to 46 bytes with zeros, // (NICs pad frames shorter than 46 bytes payload to 46 bytes with zeros,
// which would otherwise corrupt AEAD ciphertext verification). // 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(FRAME_TYPE_DATA);
frame.push(0x00); // flags (reserved)
frame.extend_from_slice(&(data.len() as u16).to_le_bytes()); frame.extend_from_slice(&(data.len() as u16).to_le_bytes());
frame.extend_from_slice(data); frame.extend_from_slice(data);
@@ -322,8 +289,8 @@ impl EthernetTransport {
"Ethernet frame sent" "Ethernet frame sent"
); );
// Return the data bytes sent (excluding frame type prefix and length field) // Return the data bytes sent (excluding 4-byte frame header)
Ok(bytes_sent.saturating_sub(3)) Ok(bytes_sent.saturating_sub(4))
} }
} }
@@ -406,22 +373,23 @@ async fn ethernet_receive_loop(
let frame_type = buf[0]; let frame_type = buf[0];
match frame_type { match frame_type {
FRAME_TYPE_DATA => { FRAME_TYPE_DATA => {
// Data frame: [type:1][length:2 LE][payload:N] // Data frame: [type:1][flags:1][length:2 LE][payload:N]
// Use the length field to trim Ethernet minimum-frame padding. if len < 4 {
if len < 3 {
trace!("Data frame too short ({len} bytes), ignoring"); trace!("Data frame too short ({len} bytes), ignoring");
continue; continue;
} }
let payload_len = u16::from_le_bytes([buf[1], buf[2]]) as usize; // buf[1] is flags (reserved, ignored for now)
if payload_len > len - 3 { let payload_len =
u16::from_le_bytes([buf[2], buf[3]]) as usize;
if payload_len > len - 4 {
trace!( trace!(
"Data frame length field ({payload_len}) exceeds \ "Data frame length field ({payload_len}) exceeds \
available bytes ({}), ignoring", available bytes ({}), ignoring",
len - 3 len - 4
); );
continue; 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 addr = TransportAddr::from_bytes(&src_mac);
let packet = ReceivedPacket::new(transport_id, addr, data); let packet = ReceivedPacket::new(transport_id, addr, data);
@@ -443,8 +411,8 @@ async fn ethernet_receive_loop(
FRAME_TYPE_BEACON => { FRAME_TYPE_BEACON => {
stats.record_beacon_recv(); stats.record_beacon_recv();
if discovery_enabled && let Some(pubkey) = parse_beacon(&buf[..len]) { if discovery_enabled && parse_beacon(&buf[..len]) {
discovery_buffer.add_peer(src_mac, pubkey); discovery_buffer.add_peer(src_mac);
trace!( trace!(
transport_id = %transport_id, transport_id = %transport_id,
remote_mac = %format_mac(&src_mac), 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. /// failures, attempts to open a fresh socket on the same interface.
async fn beacon_sender_loop( async fn beacon_sender_loop(
mut socket: Arc<AsyncPacketSocket>, mut socket: Arc<AsyncPacketSocket>,
pubkey: XOnlyPublicKey,
interval_secs: u64, interval_secs: u64,
stats: Arc<EthernetStats>, stats: Arc<EthernetStats>,
transport_id: TransportId, transport_id: TransportId,
@@ -498,7 +465,7 @@ async fn beacon_sender_loop(
/// Number of consecutive ENXIO errors before attempting socket reopen. /// Number of consecutive ENXIO errors before attempting socket reopen.
const REOPEN_THRESHOLD: u32 = 3; const REOPEN_THRESHOLD: u32 = 3;
let beacon = build_beacon(&pubkey); let beacon = build_beacon();
let interval = tokio::time::Duration::from_secs(interval_secs); let interval = tokio::time::Duration::from_secs(interval_secs);
debug!( debug!(
@@ -712,28 +679,31 @@ mod tests {
#[test] #[test]
fn test_frame_type_data_prefix() { 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 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(FRAME_TYPE_DATA);
frame.push(0x00); // flags
frame.extend_from_slice(&(data.len() as u16).to_le_bytes()); frame.extend_from_slice(&(data.len() as u16).to_le_bytes());
frame.extend_from_slice(&data); frame.extend_from_slice(&data);
assert_eq!(frame[0], 0x00); // frame type assert_eq!(frame[0], 0x00); // frame type
assert_eq!(u16::from_le_bytes([frame[1], frame[2]]), 4); // length assert_eq!(frame[1], 0x00); // flags
assert_eq!(&frame[3..], &[1, 2, 3, 4]); // payload assert_eq!(u16::from_le_bytes([frame[2], frame[3]]), 4); // length
assert_eq!(&frame[4..], &[1, 2, 3, 4]); // payload
} }
#[test] #[test]
fn test_data_frame_padding_trimmed() { fn test_data_frame_padding_trimmed() {
// Simulate Ethernet minimum-frame padding: a 4-byte payload produces // 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 = vec![0xAA, 0xBB, 0xCC, 0xDD];
let payload_len = payload.len() as u16; let payload_len = payload.len() as u16;
// Build frame as sender would // 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(FRAME_TYPE_DATA);
frame.push(0x00); // flags
frame.extend_from_slice(&payload_len.to_le_bytes()); frame.extend_from_slice(&payload_len.to_le_bytes());
frame.extend_from_slice(&payload); frame.extend_from_slice(&payload);
@@ -741,13 +711,26 @@ mod tests {
frame.resize(46, 0x00); frame.resize(46, 0x00);
// Receiver extracts using length field // Receiver extracts using length field
let recv_len = u16::from_le_bytes([frame[1], frame[2]]) as usize; let recv_len = u16::from_le_bytes([frame[2], frame[3]]) as usize;
let extracted = &frame[3..3 + recv_len]; let extracted = &frame[4..4 + recv_len];
assert_eq!(extracted, &[0xAA, 0xBB, 0xCC, 0xDD]); assert_eq!(extracted, &[0xAA, 0xBB, 0xCC, 0xDD]);
} }
#[test] #[test]
fn test_beacon_size() { 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);
} }
} }