From db9549885afdeb6bcd1bbd3456c33059fa92441b Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 26 Mar 2026 14:24:22 +0000 Subject: [PATCH] BLE transport reliability: probe promotion, send fail-fast, pubkey timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Promote probe connections directly into pool instead of dropping and reconnecting. Eliminates fragile two-phase connect pattern (probe → disconnect → reconnect) that caused race conditions on restart. - send_async fails fast when no connection exists, triggering a background connect_async instead of blocking the event loop for up to 10s on inline L2CAP connect. Prevents control socket query timeouts and MMP processing stalls. - Add 5-second timeout to pubkey_exchange recv. Without this, a peer that connects but never sends its pubkey blocks the calling task forever, killing the scan_probe_loop or accept_loop entirely. - Add retry timer in scan_probe_loop for addresses that failed probe but won't get another DeviceAdded from BlueZ (deduplication). - Clear BlueZ cached devices before starting scan so fresh advertisements trigger DeviceAdded after daemon restart. - connect_async now performs pubkey exchange before promoting to pool, matching the accept_loop's expectation on inbound connections. - BLE tests updated to pre-establish connections via connect_async since send_async no longer does inline connect. --- src/node/tests/ble.rs | 18 +++- src/transport/ble/io.rs | 15 ++++ src/transport/ble/mod.rs | 172 ++++++++++++++++++++++++++++++------- src/transport/ble/stats.rs | 1 + 4 files changed, 174 insertions(+), 32 deletions(-) diff --git a/src/node/tests/ble.rs b/src/node/tests/ble.rs index c3c974a..656e486 100644 --- a/src/node/tests/ble.rs +++ b/src/node/tests/ble.rs @@ -119,6 +119,18 @@ fn install_connect_handler(nodes: &[TestNode], i: usize, bank: &StreamBank) { } } +/// Establish a BLE connection from node `i` to node `j` via connect_async. +/// +/// Must be called after `wire_ble_connection` and `install_connect_handler`. +/// BLE send_async fails fast if no connection exists, so connections must +/// be pre-established before initiating handshakes. +async fn establish_ble_connection(nodes: &[TestNode], i: usize, j: usize) { + let transport = nodes[i].node.transports.get(&nodes[i].transport_id).unwrap(); + transport.connect(&nodes[j].addr).await.unwrap(); + // Let the background connect task complete + tokio::task::yield_now().await; +} + /// Two BLE nodes complete a Noise handshake and establish bidirectional peering. #[tokio::test] async fn test_ble_two_node_handshake() { @@ -128,8 +140,9 @@ async fn test_ble_two_node_handshake() { let bank: StreamBank = Arc::new(StdMutex::new(HashMap::new())); wire_ble_connection(&nodes, 0, 1, &bank).await; install_connect_handler(&nodes, 0, &bank); + establish_ble_connection(&nodes, 0, 1).await; - // Initiate handshake (connect-on-send creates the BLE connection) + // Initiate handshake initiate_handshake(&mut nodes, 0, 1).await; // Drain all packets (handshake + TreeAnnounce exchange) @@ -167,6 +180,8 @@ async fn test_ble_three_node_chain() { wire_ble_connection(&nodes, 1, 2, &bank).await; install_connect_handler(&nodes, 0, &bank); install_connect_handler(&nodes, 1, &bank); + establish_ble_connection(&nodes, 0, 1).await; + establish_ble_connection(&nodes, 1, 2).await; initiate_handshake(&mut nodes, 0, 1).await; initiate_handshake(&mut nodes, 1, 2).await; @@ -212,6 +227,7 @@ async fn test_ble_mixed_transport() { let bank: StreamBank = Arc::new(StdMutex::new(HashMap::new())); wire_ble_connection(&nodes, 2, 3, &bank).await; install_connect_handler(&nodes, 2, &bank); + establish_ble_connection(&nodes, 2, 3).await; // Handshake within each component initiate_handshake(&mut nodes, 0, 1).await; // UDP pair diff --git a/src/transport/ble/io.rs b/src/transport/ble/io.rs index 1968b74..b05b910 100644 --- a/src/transport/ble/io.rs +++ b/src/transport/ble/io.rs @@ -421,6 +421,21 @@ mod bluer_impl { } async fn start_scanning(&self) -> Result { + // Clear cached devices so BlueZ fires DeviceAdded for every + // advertisement. Without this, already-known devices only + // produce PropertyChanged events (which bluer doesn't expose + // at the device level), causing the scanner to miss peers + // after a daemon restart. + if let Ok(cached) = self.adapter.device_addresses().await { + let count = cached.len(); + for addr in cached { + let _ = self.adapter.remove_device(addr).await; + } + if count > 0 { + debug!(count, "BLE scanner: cleared cached devices"); + } + } + // Set discovery filter for LE transport with FIPS UUID let filter = DiscoveryFilter { transport: DiscoveryTransport::Le, diff --git a/src/transport/ble/mod.rs b/src/transport/ble/mod.rs index f87b739..c5369f8 100644 --- a/src/transport/ble/mod.rs +++ b/src/transport/ble/mod.rs @@ -231,6 +231,8 @@ impl BleTransport { self.config.connect_timeout_ms(), self.config.probe_cooldown_secs(), local_node_addr, + self.packet_tx.clone(), + self.transport_id, ))); debug!(adapter = %adapter, "BLE scan+probe loop started"); } @@ -283,27 +285,26 @@ impl BleTransport { /// Send data to a remote BLE address. /// - /// If no connection exists to the target, attempts connect-on-send - /// (inline connection with timeout), matching TCP transport behavior. + /// If no connection exists, triggers a background connect and fails + /// fast. The next send retry (typically 1s later for handshake msg1) + /// will find the connection established. This avoids blocking the + /// event loop on L2CAP connect (up to 10s). pub async fn send_async( &self, addr: &TransportAddr, data: &[u8], ) -> Result { - // Get existing connection or connect inline - let has_conn = { - let pool = self.pool.lock().await; - pool.contains(addr) - }; - - if !has_conn { - self.connect_inline(addr).await?; - } - let pool = self.pool.lock().await; - let conn = pool - .get(addr) - .ok_or_else(|| TransportError::SendFailed("not connected".into()))?; + let conn = match pool.get(addr) { + Some(c) => c, + None => { + // Drop pool lock before triggering background connect + drop(pool); + // Fire-and-forget: connect_async spawns a background task + let _ = self.connect_async(addr).await; + return Err(TransportError::SendFailed("not connected".into())); + } + }; // MTU check let mtu = conn.effective_mtu() as usize; @@ -334,8 +335,9 @@ impl BleTransport { /// Connect to a remote BLE device inline (blocking the caller). /// - /// Used by connect-on-send. Connects with timeout, promotes to pool, - /// and spawns the receive loop. + /// Not used in normal operation (send_async fails fast instead). + /// Retained for manual debugging / testing scenarios. + #[allow(dead_code)] async fn connect_inline(&self, addr: &TransportAddr) -> Result<(), TransportError> { let ble_addr = BleAddr::parse( addr.as_str() @@ -468,6 +470,8 @@ impl BleTransport { let psm = self.config.psm(); let timeout_ms = self.config.connect_timeout_ms(); let addr_clone = addr.clone(); + let local_pubkey = self.local_pubkey; + let discovery_buffer = Arc::clone(&self.discovery_buffer); let task = tokio::spawn(async move { let result = tokio::time::timeout( @@ -481,6 +485,24 @@ impl BleTransport { match result { Ok(Ok(stream)) => { + // Pre-handshake pubkey exchange (temporary, pre-XX) + if let Some(ref our_pubkey) = local_pubkey { + match pubkey_exchange(&stream, our_pubkey).await { + Ok(peer_pubkey) => { + debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete"); + discovery_buffer + .add_peer_with_pubkey(&ble_addr, peer_pubkey); + } + Err(e) => { + warn!( + addr = %addr_clone, error = %e, + "BLE outbound pubkey exchange failed" + ); + return; + } + } + } + let send_mtu = stream.send_mtu(); let recv_mtu = stream.recv_mtu(); let stream = Arc::new(stream); @@ -648,6 +670,13 @@ 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. @@ -662,9 +691,13 @@ async fn pubkey_exchange( msg[1..].copy_from_slice(local_pubkey); stream.send(&msg).await?; - // Receive peer's pubkey + // Receive peer's pubkey (with timeout to prevent indefinite blocking) let mut buf = [0u8; PUBKEY_EXCHANGE_SIZE]; - let n = stream.recv(&mut buf).await?; + 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 {}", @@ -839,8 +872,10 @@ async fn receive_loop( /// /// Scanner events arrive continuously (both sides advertise continuously). /// Each scan result is probed immediately unless the address is in cooldown -/// (recently probed) or already connected. On successful probe, the peer -/// is reported to the discovery buffer for the node layer to auto-connect. +/// (recently probed) or already connected. On successful probe, the +/// connection is promoted directly into the pool (no second L2CAP connect +/// needed) and the peer is reported to the discovery buffer for the node +/// layer to auto-connect. /// /// Cooldown prevents rapid re-probing of the same address: after any probe /// attempt (success or failure), the address is suppressed for @@ -857,17 +892,41 @@ async fn scan_probe_loop( connect_timeout_ms: u64, cooldown_secs: u64, local_node_addr: Option, + packet_tx: PacketTx, + transport_id: TransportId, ) { // Track last probe time per address for cooldown let mut last_probed: HashMap = HashMap::new(); + // Addresses discovered but not yet connected — retried after cooldown + // even if the scanner doesn't fire again (BlueZ deduplicates). + let mut pending_addrs: Vec = Vec::new(); let cooldown = std::time::Duration::from_secs(cooldown_secs); + let retry_interval = tokio::time::interval(std::time::Duration::from_secs(cooldown_secs)); + tokio::pin!(retry_interval); + retry_interval.tick().await; // consume initial tick loop { - let addr = match scanner.next().await { - Some(a) => a, - None => { - debug!("BLE scanner ended"); - break; + // Either a scanner event or the retry timer fires + let addr = tokio::select! { + result = scanner.next() => { + match result { + Some(a) => a, + None => { + debug!("BLE scanner ended"); + break; + } + } + } + _ = retry_interval.tick() => { + // Re-probe pending addresses that aren't connected + let pool_guard = pool.lock().await; + pending_addrs.retain(|a| !pool_guard.contains(&a.to_transport_addr())); + drop(pool_guard); + if let Some(a) = pending_addrs.first().cloned() { + a + } else { + continue; + } } }; @@ -878,10 +937,16 @@ async fn scan_probe_loop( { let pool_guard = pool.lock().await; if pool_guard.contains(&addr.to_transport_addr()) { + pending_addrs.retain(|a| a != &addr); continue; } } + // Track for retry in case probe fails and scanner doesn't re-fire + if !pending_addrs.contains(&addr) { + pending_addrs.push(addr.clone()); + } + // Skip if in cooldown if last_probed .get(&addr) @@ -921,12 +986,14 @@ async fn scan_probe_loop( } }; - // Pubkey exchange, then close the L2CAP connection + // Pubkey exchange, then promote connection to pool + let ta = addr.to_transport_addr(); match pubkey_exchange(&stream, &our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr, "BLE probe complete"); - // Cross-probe tie-breaker: smaller NodeAddr's outbound wins + // Cross-probe tie-breaker: smaller NodeAddr's outbound wins. + // If we lose, drop connection — accept_loop handles inbound. if let Some(ref our_addr) = local_node_addr { let peer_addr = NodeAddr::from_pubkey(&peer_pubkey); if our_addr >= &peer_addr { @@ -934,18 +1001,61 @@ async fn scan_probe_loop( addr = %addr, "BLE probe tie-breaker: yielding to peer's outbound" ); + buffer.add_peer_with_pubkey(&addr, peer_pubkey); + continue; } } - // Report to node layer — auto-connect will establish - // a persistent connection via send_async/connect_inline + // Promote connection to pool — no second L2CAP connect needed + let send_mtu = stream.send_mtu(); + let recv_mtu = stream.recv_mtu(); + let stream = Arc::new(stream); + + let recv_task = tokio::spawn(receive_loop( + Arc::clone(&stream), + ta.clone(), + Arc::clone(&pool), + packet_tx.clone(), + transport_id, + Arc::clone(&stats), + recv_mtu, + )); + + let conn = BleConnection { + stream, + recv_task: Some(recv_task), + send_mtu, + recv_mtu, + established_at: tokio::time::Instant::now(), + is_static: false, + addr: addr.clone(), + }; + + let mut pool_guard = pool.lock().await; + match pool_guard.insert(ta.clone(), conn) { + Ok(Some(evicted)) => { + stats.record_pool_eviction(); + debug!(addr = %ta, evicted = %evicted, "BLE probe promoted (evicted peer)"); + } + Ok(None) => { + debug!(addr = %ta, "BLE probe promoted to pool"); + } + Err(e) => { + warn!(addr = %ta, error = %e, "BLE pool full, probe connection dropped"); + stats.record_connection_rejected(); + } + } + drop(pool_guard); + stats.record_connection_established(); + pending_addrs.retain(|a| a != &addr); + + // Report to node layer for auto-connect / handshake buffer.add_peer_with_pubkey(&addr, peer_pubkey); } Err(e) => { debug!(addr = %addr, error = %e, "BLE probe pubkey exchange failed"); } } - // L2CAP connection dropped here (stream goes out of scope) } } diff --git a/src/transport/ble/stats.rs b/src/transport/ble/stats.rs index bae7461..cc0e89f 100644 --- a/src/transport/ble/stats.rs +++ b/src/transport/ble/stats.rs @@ -108,6 +108,7 @@ impl BleStats { self.scan_results.fetch_add(1, Ordering::Relaxed); } + /// Take a snapshot of all counters. pub fn snapshot(&self) -> BleStatsSnapshot { BleStatsSnapshot {