BLE transport reliability: probe promotion, send fail-fast, pubkey timeout

- 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.
This commit is contained in:
Johnathan Corgan
2026-03-26 14:24:22 +00:00
parent 8f1494853a
commit db9549885a
4 changed files with 174 additions and 32 deletions
+17 -1
View File
@@ -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
+15
View File
@@ -421,6 +421,21 @@ mod bluer_impl {
}
async fn start_scanning(&self) -> Result<Self::Scanner, TransportError> {
// 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,
+141 -31
View File
@@ -231,6 +231,8 @@ impl<I: BleIo> BleTransport<I> {
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<I: BleIo> BleTransport<I> {
/// 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<usize, TransportError> {
// 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<I: BleIo> BleTransport<I> {
/// 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<I: BleIo> BleTransport<I> {
let psm = self.config.psm();
let timeout_ms = self.config.connect_timeout_ms();
let addr_clone = addr.clone();
let local_pubkey = self.local_pubkey;
let discovery_buffer = Arc::clone(&self.discovery_buffer);
let task = tokio::spawn(async move {
let result = tokio::time::timeout(
@@ -481,6 +485,24 @@ impl<I: BleIo> BleTransport<I> {
match result {
Ok(Ok(stream)) => {
// Pre-handshake pubkey exchange (temporary, pre-XX)
if let Some(ref our_pubkey) = local_pubkey {
match pubkey_exchange(&stream, our_pubkey).await {
Ok(peer_pubkey) => {
debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete");
discovery_buffer
.add_peer_with_pubkey(&ble_addr, peer_pubkey);
}
Err(e) => {
warn!(
addr = %addr_clone, error = %e,
"BLE outbound pubkey exchange failed"
);
return;
}
}
}
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
let stream = Arc::new(stream);
@@ -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<S: BleStream>(
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<S: BleStream>(
///
/// 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<I: io::BleIo>(
connect_timeout_ms: u64,
cooldown_secs: u64,
local_node_addr: Option<NodeAddr>,
packet_tx: PacketTx,
transport_id: TransportId,
) {
// Track last probe time per address for cooldown
let mut last_probed: HashMap<BleAddr, tokio::time::Instant> = 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<BleAddr> = 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<I: io::BleIo>(
{
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<I: io::BleIo>(
}
};
// 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<I: io::BleIo>(
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)
}
}
+1
View File
@@ -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 {