mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-06 11:38:24 +00:00
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:
@@ -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.
|
||||||
|
|||||||
+140
-184
@@ -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
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -602,24 +620,10 @@ impl Node {
|
|||||||
// Calculate max MSS for TCP clamping
|
// Calculate max MSS for TCP clamping
|
||||||
let effective_mtu = self.effective_ipv6_mtu();
|
let effective_mtu = self.effective_ipv6_mtu();
|
||||||
let max_mss = effective_mtu.saturating_sub(40).saturating_sub(20); // IPv6 + TCP headers
|
let max_mss = effective_mtu.saturating_sub(40).saturating_sub(20); // IPv6 + TCP headers
|
||||||
|
|
||||||
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
@@ -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) => {
|
||||||
|
|||||||
@@ -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);
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
@@ -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
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user