diff --git a/src/transport/ble/mod.rs b/src/transport/ble/mod.rs index 61d6c916..077953b6 100644 --- a/src/transport/ble/mod.rs +++ b/src/transport/ble/mod.rs @@ -391,12 +391,19 @@ impl BleTransport { { Ok(Ok(stream)) => stream, Ok(Err(e)) => { - debug!(addr = %addr, error = %e, "BLE connect-on-send failed"); + self.stats.record_connect_error(); + debug!( + addr = %addr, role = "central", outcome = "connect-error", error = %e, + "BLE connect-on-send failed" + ); return Err(TransportError::ConnectionRefused); } Err(_) => { self.stats.record_connect_timeout(); - debug!(addr = %addr, "BLE connect-on-send timeout"); + debug!( + addr = %addr, role = "central", outcome = "connect-timeout", + "BLE connect-on-send timeout" + ); return Err(TransportError::Timeout); } }; @@ -421,7 +428,11 @@ impl BleTransport { .add_peer_with_pubkey(&announced, peer_pubkey); } Err(e) => { - warn!(addr = %addr, error = %e, "BLE outbound pubkey exchange failed"); + self.stats.record_pubkey_exchange_failure(); + warn!( + addr = %addr, role = "central", outcome = "pubkey-exchange-failed", + error = %e, "BLE outbound pubkey exchange failed" + ); return Err(e); } } @@ -553,8 +564,12 @@ impl BleTransport { neighbor_buffer.add_peer_with_pubkey(&announced, peer_pubkey); } Err(e) => { + stats.record_pubkey_exchange_failure(); warn!( - addr = %addr_clone, error = %e, + addr = %addr_clone, + role = "central", + outcome = "pubkey-exchange-failed", + error = %e, "BLE outbound pubkey exchange failed" ); return; @@ -601,11 +616,18 @@ impl BleTransport { stats.record_connection_established(); } Ok(Err(e)) => { - debug!(addr = %addr_clone, error = %e, "BLE connect failed"); + stats.record_connect_error(); + debug!( + addr = %addr_clone, role = "central", outcome = "connect-error", + error = %e, "BLE connect failed" + ); } Err(_) => { stats.record_connect_timeout(); - debug!(addr = %addr_clone, "BLE connect timeout"); + debug!( + addr = %addr_clone, role = "central", outcome = "connect-timeout", + "BLE connect timeout" + ); } } }); @@ -886,6 +908,8 @@ async fn accept_loop( { debug!( addr = %ta, + role = "peripheral", + outcome = "duplicate-node-decline", existing = %existing, "BLE inbound: peer already connected on another address, dropping duplicate" ); @@ -899,15 +923,23 @@ async fn accept_loop( if let Some(ref our_addr) = local_node_addr && our_addr < &peer_node { + stats.record_tiebreaker_drop(); debug!( addr = %ta, + role = "peripheral", + outcome = "tiebreaker-drop", "BLE inbound tie-breaker: dropping (our addr < peer, outbound wins)" ); continue; } } Err(e) => { - debug!(addr = %ta, error = %e, "BLE inbound pubkey exchange failed"); + stats.record_pubkey_exchange_failure(); + debug!( + addr = %ta, role = "peripheral", + outcome = "pubkey-exchange-failed", error = %e, + "BLE inbound pubkey exchange failed" + ); continue; } } @@ -945,8 +977,11 @@ async fn accept_loop( info!(addr = %ta, send_mtu, recv_mtu, "BLE inbound connection accepted"); } Err(e) => { - warn!(addr = %ta, error = %e, "BLE pool full, inbound connection rejected"); stats.record_connection_rejected(); + warn!( + addr = %ta, role = "peripheral", outcome = "pool-rejected", + error = %e, "BLE pool full, inbound connection rejected" + ); continue; } } @@ -1003,6 +1038,145 @@ async fn receive_loop( pool.remove(&addr); } +/// Consecutive-failure backoff ceiling for a pending address, as a power of +/// two multiple of the base cooldown. At the 30 s default this caps a failing +/// address at one dial attempt every 16 minutes. +const MAX_PROBE_BACKOFF_SHIFT: u32 = 5; + +/// Ceiling on how many discovered-but-unconnected addresses are kept for +/// retry. Resolvable private addresses rotate, so without a bound the book +/// grows for the life of the process; with one, the total retry dial rate is +/// bounded too (at most one dial per retry tick, spread over the book). +const MAX_PENDING_PROBES: usize = 32; + +/// One discovered address awaiting a successful probe. +#[derive(Debug, Clone)] +struct PendingProbe { + addr: BleAddr, + /// Consecutive failed probes. Reset only by removal from the book, which + /// every conclusive outcome (connected, duplicate declined, already + /// pooled) performs. + failures: u32, + /// Earliest instant at which this address may be dialled again. + next_attempt: tokio::time::Instant, +} + +/// The retry book for addresses the scanner has offered but which are not yet +/// connected. +/// +/// Exists because a scanner is not a reliable repeat source: BlueZ emits +/// `DeviceAdded` once per address per discovery session, so an address the +/// probe loop forgets is never offered again. Everything here therefore +/// throttles rather than discards — an entry leaves the book on a *conclusive* +/// outcome, or when [`MAX_PENDING_PROBES`] other addresses compete for its +/// slot, never because it failed. +/// +/// Two properties matter: +/// +/// - **Consecutive failures back an address off exponentially.** A dead +/// address is retried on a doubling interval up to +/// [`MAX_PROBE_BACKOFF_SHIFT`], instead of being re-dialled every cooldown +/// forever. Addresses rotate and links are lossy, so a handful of failures +/// is normal and must not retire a peer that is still there. +/// - **The retry tick rotates.** Probing only the head of the book let one +/// slow or dead address starve every other pending address behind it, which +/// on a busy radio is most of them. +#[derive(Debug)] +struct PendingProbes { + entries: Vec, + cooldown: std::time::Duration, +} + +impl PendingProbes { + fn new(cooldown: std::time::Duration) -> Self { + Self { + entries: Vec::new(), + cooldown, + } + } + + fn position(&self, addr: &BleAddr) -> Option { + self.entries.iter().position(|e| &e.addr == addr) + } + + /// Record a sighting. A previously unseen address becomes immediately + /// eligible; a known one keeps whatever backoff it has earned, so a + /// scanner that re-reports the same address many times a second cannot + /// wash out the backoff. + /// + /// When the book is full the *most-failed* entry is evicted to make room, + /// which is the entry least likely to still have a peer behind it. + fn observe(&mut self, addr: &BleAddr, now: tokio::time::Instant) { + if self.position(addr).is_some() { + return; + } + if self.entries.len() >= MAX_PENDING_PROBES + && let Some(worst) = self + .entries + .iter() + .enumerate() + .max_by_key(|(_, e)| (e.failures, e.next_attempt)) + .map(|(i, _)| i) + { + self.entries.remove(worst); + } + self.entries.push(PendingProbe { + addr: addr.clone(), + failures: 0, + next_attempt: now, + }); + } + + /// Whether `addr` may be dialled now. An address that is not in the book + /// has no history to hold it back. + fn is_due(&self, addr: &BleAddr, now: tokio::time::Instant) -> bool { + match self.position(addr) { + Some(i) => self.entries[i].next_attempt <= now, + None => true, + } + } + + /// Note that a probe is starting: hold the address for one base cooldown + /// so the attempt in flight is not duplicated. + fn mark_attempt(&mut self, addr: &BleAddr, now: tokio::time::Instant) { + if let Some(i) = self.position(addr) { + self.entries[i].next_attempt = now + self.cooldown; + } + } + + /// Note that a probe failed. Returns the new consecutive-failure count. + fn record_failure(&mut self, addr: &BleAddr, now: tokio::time::Instant) -> u32 { + let Some(i) = self.position(addr) else { + return 0; + }; + let e = &mut self.entries[i]; + e.failures = e.failures.saturating_add(1); + let shift = (e.failures - 1).min(MAX_PROBE_BACKOFF_SHIFT); + e.next_attempt = now + self.cooldown * 2u32.pow(shift); + e.failures + } + + /// Drop an address that reached a conclusive outcome. + fn resolve(&mut self, addr: &BleAddr) { + self.entries.retain(|e| &e.addr != addr); + } + + /// Drop every address for which `connected` reports a live pool entry. + fn drop_connected(&mut self, connected: impl Fn(&BleAddr) -> bool) { + self.entries.retain(|e| !connected(&e.addr)); + } + + /// The next address due for a retry, rotated to the back of the book so + /// the following tick starts after it rather than on it. + fn next_due(&mut self, now: tokio::time::Instant) -> Option { + let i = self.entries.iter().position(|e| e.next_attempt <= now)?; + let entry = self.entries.remove(i); + let addr = entry.addr.clone(); + self.entries.push(entry); + Some(addr) + } +} + /// Combined scan + probe loop. /// /// Scanner events arrive continuously (both sides advertise continuously). @@ -1030,11 +1204,12 @@ async fn scan_probe_loop( 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(); + // Addresses discovered but not yet connected — retried after cooldown even + // if the scanner doesn't fire again (BlueZ deduplicates), on a per-address + // backoff that widens with consecutive failures. Also the cooldown record: + // an address leaves the book the moment it reaches a conclusive outcome, + // after which the pool and `known_node_of` guards below cover it. + let mut pending = PendingProbes::new(std::time::Duration::from_secs(cooldown_secs)); // Link addresses already resolved to a node identity by a completed pubkey // exchange. Lets the loop skip an address it has *already* learned belongs // to a peer it is connected to, instead of paying a full connect and @@ -1049,7 +1224,6 @@ async fn scan_probe_loop( // which is every peer that predates this and every backend that does not // advertise service data. let mut learned_psm: HashMap = HashMap::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 @@ -1075,12 +1249,13 @@ async fn scan_probe_loop( _ = 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())); + pending.drop_connected(|a| pool_guard.contains(&a.to_transport_addr())); drop(pool_guard); - if let Some(a) = pending_addrs.first().cloned() { - a - } else { - continue; + // Rotating rather than always taking the head is what stops one + // slow or dead address from starving every other pending one. + match pending.next_due(tokio::time::Instant::now()) { + Some(a) => a, + None => continue, } } }; @@ -1092,21 +1267,17 @@ 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); + pending.resolve(&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()); - } + let now = tokio::time::Instant::now(); + pending.observe(&addr, now); - // Skip if in cooldown - if last_probed - .get(&addr) - .is_some_and(|last| last.elapsed() < cooldown) - { + // Skip if in cooldown, or backed off after consecutive failures + if !pending.is_due(&addr, now) { continue; } @@ -1120,7 +1291,7 @@ async fn scan_probe_loop( pool_guard.find_by_node(node).is_some() }; if still_connected { - pending_addrs.retain(|a| a != &addr); + pending.resolve(&addr); continue; } // That peer is gone — forget the mapping and probe normally. @@ -1128,7 +1299,7 @@ async fn scan_probe_loop( } // Record probe time (before attempt, so cooldown applies on failure too) - last_probed.insert(addr.clone(), tokio::time::Instant::now()); + pending.mark_attempt(&addr, now); // Need pubkey for probe let our_pubkey = match local_pubkey { @@ -1141,6 +1312,9 @@ async fn scan_probe_loop( // L2CAP connect, at whatever PSM this peer advertised. let dial_psm = learned_psm.get(&addr).copied().unwrap_or(configured_psm); + // Stamped here so every outcome below can report how long the peer + // took to go from advertisement to conclusion. + let probe_started = tokio::time::Instant::now(); let stream = match tokio::time::timeout( std::time::Duration::from_millis(connect_timeout_ms), io.connect(&addr, dial_psm), @@ -1149,7 +1323,13 @@ async fn scan_probe_loop( { Ok(Ok(s)) => s, Ok(Err(e)) => { - debug!(addr = %addr, psm = dial_psm, error = %e, "BLE probe connect failed"); + stats.record_connect_error(); + let failures = pending.record_failure(&addr, tokio::time::Instant::now()); + debug!( + addr = %addr, role = "central", outcome = "connect-error", + psm = dial_psm, discovery_ms = probe_started.elapsed().as_millis() as u64, + failures, error = %e, "BLE probe connect failed" + ); // A learned PSM that does not answer is stale — forget it, so // the next advert re-learns it and the fallback applies in the // meantime. Costs one retry. @@ -1157,8 +1337,13 @@ async fn scan_probe_loop( continue; } Err(_) => { - debug!(addr = %addr, psm = dial_psm, "BLE probe connect timeout"); stats.record_connect_timeout(); + let failures = pending.record_failure(&addr, tokio::time::Instant::now()); + debug!( + addr = %addr, role = "central", outcome = "connect-timeout", + psm = dial_psm, discovery_ms = probe_started.elapsed().as_millis() as u64, + failures, "BLE probe connect timeout" + ); learned_psm.remove(&addr); continue; } @@ -1181,10 +1366,21 @@ async fn scan_probe_loop( if let Some(ref our_addr) = local_node_addr && our_addr >= &peer_node { + stats.record_tiebreaker_yield(); debug!( addr = %addr, + role = "central", + outcome = "tiebreaker-yield", + discovery_ms = probe_started.elapsed().as_millis() as u64, "BLE probe tie-breaker: yielding to peer's outbound" ); + // Same reasoning as the duplicate-decline path below: the + // exchange has resolved this address to a node, so once + // that node holds a link the next cooldown can skip the + // address outright instead of paying another connect and + // exchange to yield again. The tie-breaker decision itself + // is unchanged — only the cost of re-reaching it. + known_node_of.insert(addr.clone(), peer_node); let announced = announced_addr(&pool, &peer_node, &addr).await; buffer.add_peer_with_pubkey(&announced, peer_pubkey); continue; @@ -1203,7 +1399,10 @@ async fn scan_probe_loop( { debug!( addr = %ta, + role = "central", + outcome = "duplicate-node-decline", existing = %existing, + discovery_ms = probe_started.elapsed().as_millis() as u64, "BLE probe: peer already connected on another address, dropping duplicate" ); stats.record_duplicate_node_decline(); @@ -1216,7 +1415,7 @@ async fn scan_probe_loop( // connection behind it. let announced = announced_addr(&pool, &peer_node, &addr).await; buffer.add_peer_with_pubkey(&announced, peer_pubkey); - pending_addrs.retain(|a| a != &addr); + pending.resolve(&addr); continue; } @@ -1249,22 +1448,35 @@ async fn scan_probe_loop( debug!(addr = %ta, evicted = %evicted, "BLE probe promoted (evicted peer)"); } Ok(None) => { - debug!(addr = %ta, "BLE probe promoted to pool"); + debug!( + addr = %ta, role = "central", outcome = "connected", + discovery_ms = probe_started.elapsed().as_millis() as u64, + "BLE probe promoted to pool" + ); } Err(e) => { - warn!(addr = %ta, error = %e, "BLE pool full, probe connection dropped"); stats.record_connection_rejected(); + warn!( + addr = %ta, role = "central", outcome = "pool-rejected", + error = %e, "BLE pool full, probe connection dropped" + ); } } drop(pool_guard); stats.record_connection_established(); - pending_addrs.retain(|a| a != &addr); + pending.resolve(&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"); + stats.record_pubkey_exchange_failure(); + let failures = pending.record_failure(&addr, tokio::time::Instant::now()); + debug!( + addr = %addr, role = "central", outcome = "pubkey-exchange-failed", + discovery_ms = probe_started.elapsed().as_millis() as u64, + failures, error = %e, "BLE probe pubkey exchange failed" + ); } } } @@ -1281,6 +1493,211 @@ mod tests { use io::{MockBleIo, MockBleStream}; use secp256k1::{Secp256k1, SecretKey}; + // ------------------------------------------------------------------ + // PendingProbes — the retry/backoff policy for discovered addresses + // ------------------------------------------------------------------ + + const TEST_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(30); + + fn probes() -> PendingProbes { + PendingProbes::new(TEST_COOLDOWN) + } + + fn a(n: u8) -> BleAddr { + BleAddr::parse(&format!("ble0/AA:BB:CC:DD:EE:{:02X}", n)).unwrap() + } + + /// A fresh sighting is dialled straight away — discovery must not wait a + /// cooldown to try a peer it has never met. + #[test] + fn a_newly_seen_address_is_due_immediately() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + assert!(p.is_due(&a(1), t0)); + } + + /// The regression this policy exists for: an address that keeps failing + /// must be dialled exponentially less often, not once per cooldown for as + /// long as the process lives. + #[test] + fn consecutive_failures_back_an_address_off_exponentially() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + + for expected_shift in 0..MAX_PROBE_BACKOFF_SHIFT { + let n = p.record_failure(&a(1), t0); + assert_eq!(n, expected_shift + 1); + let wait = TEST_COOLDOWN * 2u32.pow(expected_shift); + assert!( + !p.is_due(&a(1), t0 + wait - std::time::Duration::from_millis(1)), + "due too early after {n} failures" + ); + assert!(p.is_due(&a(1), t0 + wait), "not due after {n} failures"); + } + + // And the interval stops growing at the ceiling rather than running + // away to hours. + let capped = TEST_COOLDOWN * 2u32.pow(MAX_PROBE_BACKOFF_SHIFT); + for _ in 0..8 { + p.record_failure(&a(1), t0); + assert!(p.is_due(&a(1), t0 + capped)); + } + } + + /// Under the old policy an address failing every 30 s for 37 minutes was + /// dialled 49 times. Pin the improvement rather than just the formula. + #[test] + fn a_dead_address_is_dialled_a_handful_of_times_an_hour() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + + let mut dials = 0; + let mut now = t0; + let deadline = t0 + std::time::Duration::from_secs(37 * 60); + // Tick at the retry interval, exactly as the loop does. + while now <= deadline { + if p.is_due(&a(1), now) { + p.mark_attempt(&a(1), now); + p.record_failure(&a(1), now); + dials += 1; + } + now += TEST_COOLDOWN; + } + assert!( + (1..=10).contains(&dials), + "expected a handful of dials in 37 minutes, got {dials}" + ); + } + + /// A scanner that re-reports the same address many times a second (which + /// Android does, at roughly 52/min) must not wash the backoff out. + #[test] + fn repeated_sightings_do_not_reset_the_backoff() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + for _ in 0..4 { + p.record_failure(&a(1), t0); + } + let still_blocked = t0 + TEST_COOLDOWN; + for _ in 0..100 { + p.observe(&a(1), still_blocked); + } + assert!(!p.is_due(&a(1), still_blocked)); + assert_eq!(p.entries.len(), 1); + } + + /// Failing never removes an address. This is what keeps a BlueZ node + /// recoverable: BlueZ emits `DeviceAdded` once per address per discovery + /// session, so an address dropped from the book would never be offered + /// again and the peer behind it would be unreachable for the life of the + /// process. + #[test] + fn failures_never_evict_the_address_itself() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + for _ in 0..500 { + p.record_failure(&a(1), t0); + } + assert_eq!(p.entries.len(), 1); + // Still reachable: once the (capped) backoff elapses it is dialled + // again, so a peer that comes back is picked up without a new sighting. + let capped = TEST_COOLDOWN * 2u32.pow(MAX_PROBE_BACKOFF_SHIFT); + assert_eq!(p.next_due(t0 + capped), Some(a(1))); + } + + /// Conclusive outcomes clear the address *and* its failure history, so a + /// peer that reconnects later starts from a clean slate. + #[test] + fn resolving_clears_the_failure_history() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(1), t0); + for _ in 0..5 { + p.record_failure(&a(1), t0); + } + p.resolve(&a(1)); + assert!(p.entries.is_empty()); + p.observe(&a(1), t0); + assert!(p.is_due(&a(1), t0)); + } + + /// The head-of-line half of the bug: probing only the first entry let one + /// address monopolise the retry tick. Rotation gives every due address a + /// turn. + #[test] + fn the_retry_tick_rotates_across_due_addresses() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + for n in 0..3 { + p.observe(&a(n), t0); + } + let order: Vec<_> = (0..6).filter_map(|_| p.next_due(t0)).collect(); + assert_eq!(order, vec![a(0), a(1), a(2), a(0), a(1), a(2)]); + } + + /// A backed-off address is skipped by the tick rather than blocking the + /// addresses behind it. + #[test] + fn a_backed_off_address_does_not_block_the_others() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(0), t0); + p.observe(&a(1), t0); + p.record_failure(&a(0), t0); + assert_eq!(p.next_due(t0), Some(a(1))); + assert_eq!(p.next_due(t0), Some(a(1))); + } + + /// Nothing is due when everything is backed off — the tick idles rather + /// than dialling something it just said it would not. + #[test] + fn next_due_yields_nothing_when_all_are_backed_off() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + p.observe(&a(0), t0); + p.record_failure(&a(0), t0); + assert_eq!(p.next_due(t0), None); + } + + /// Addresses rotate, so the book is capacity-bounded. Eviction is by + /// failure count, so the entry least likely to have a peer behind it goes + /// first and a healthy address is never displaced by a dead one. + #[test] + fn a_full_book_evicts_the_most_failed_address() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + for n in 0..MAX_PENDING_PROBES as u8 { + p.observe(&a(n), t0); + } + // One entry is much worse than the rest. + for _ in 0..3 { + p.record_failure(&a(7), t0); + } + p.observe(&a(200), t0); + assert_eq!(p.entries.len(), MAX_PENDING_PROBES); + assert!(p.position(&a(7)).is_none(), "the worst entry should go"); + assert!(p.position(&a(200)).is_some(), "the new entry should land"); + assert!(p.position(&a(0)).is_some(), "healthy entries should stay"); + } + + /// Pool membership clears entries in bulk on the retry tick. + #[test] + fn connected_addresses_leave_the_book() { + let mut p = probes(); + let t0 = tokio::time::Instant::now(); + for n in 0..3 { + p.observe(&a(n), t0); + } + p.drop_connected(|addr| addr == &a(1)); + assert_eq!(p.entries.len(), 2); + assert!(p.position(&a(1)).is_none()); + } + /// Deterministic x-only pubkey for exchange tests. fn test_pubkey(seed: u8) -> [u8; 32] { let secp = Secp256k1::new(); @@ -1642,12 +2059,48 @@ mod tests { } /// Let spawned loops make progress. + /// + /// Cooperative only: this hands the scheduler control, it does not move + /// the clock. Anything gated on a `tokio::time` timer needs + /// [`wait_for`] instead. async fn settle() { for _ in 0..64 { tokio::task::yield_now().await; } } + /// Wait until `cond` holds, or fail the test. + /// + /// A fixed number of `yield_now()` calls is not a wait, it is a race + /// against the clock, and it loses whenever a loop under test is parked + /// on a timer rather than on a channel. `scan_probe_loop` is: before it + /// reaches its `select!` it consumes the retry interval's first tick, + /// and tokio rounds a timer deadline up to the next whole millisecond of + /// its wheel — so unless the runtime clock happens to sit exactly on a + /// millisecond boundary, that tick cannot fire until real time crosses + /// the next one. No number of yields makes real time pass, so whether a + /// yield budget covers the gap depends on how long a yield takes on the + /// host: comfortably on a slow one, not at all on a fast one. + /// + /// Polling the condition with a sleep between attempts removes the + /// dependency entirely — the sleep is what lets the timer fire, and the + /// condition is what ends the wait. The already-satisfied case still + /// costs only a `settle`, so nothing that passes today gets slower. + async fn wait_for(what: &str, mut cond: impl FnMut() -> bool) { + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5); + loop { + settle().await; + if cond() { + return; + } + assert!( + tokio::time::Instant::now() < deadline, + "timed out waiting for {what}" + ); + tokio::time::sleep(std::time::Duration::from_millis(1)).await; + } + } + /// A second inbound connection from a rotated address for a peer already /// in the pool is declined, the incumbent link is kept, and the peer is /// still announced — under the address its live link is on. @@ -1938,4 +2391,208 @@ mod tests { .unwrap_err(); assert!(matches!(err, TransportError::RecvFailed(_))); } + + // ------------------------------------------------------------------ + // Connect outcome counters + // ------------------------------------------------------------------ + + /// A dial that errors is counted as an error, not as a timeout. The two + /// are different faults and blur into one useless number if merged. + #[tokio::test(start_paused = true)] + async fn test_a_refused_dial_counts_as_an_error_not_a_timeout() { + let dials: DialLog = Arc::new(std::sync::Mutex::new(Vec::new())); + let (mut transport, _rx) = psm_probe_transport(Arc::clone(&dials)); + transport.start_async().await.unwrap(); + + transport.io.inject_scan_result(test_addr(2)).await; + settle().await; + + let snap = transport.stats.snapshot(); + assert_eq!(snap.connect_errors, 1); + assert_eq!(snap.connect_timeouts, 0); + assert_eq!(snap.connections_established, 0); + transport.stop_async().await.unwrap(); + } + + /// A peer that connects and then sends a bad exchange is counted as a + /// pubkey-exchange failure — the link came up and produced nothing + /// usable, which is a different fault from never connecting. + #[tokio::test] + async fn test_a_bad_exchange_counts_as_a_pubkey_exchange_failure() { + let io = MockBleIo::new("hci0", test_addr(1)); + let (tx, _rx) = tokio::sync::mpsc::channel(64); + let mut transport = + BleTransport::new(TransportId::new(1), None, identity_test_config(), io, tx); + transport.set_local_pubkey(test_pubkey(1)); + transport.start_async().await.unwrap(); + + let (ours, peer) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + transport.io.inject_inbound(ours).await; + let mut wire = vec![0xFFu8]; + wire.extend_from_slice(&test_pubkey(2)); + peer.send(&wire).await.unwrap(); + settle().await; + + let snap = transport.stats.snapshot(); + assert_eq!(snap.pubkey_exchange_failures, 1); + assert_eq!(snap.connections_accepted, 0); + assert_eq!(transport.pool.lock().await.len(), 0); + transport.stop_async().await.unwrap(); + } + + /// The tie-breaker pair. Its convention is deterministic in source, but + /// nothing recorded whether two nodes actually agreed at runtime, and a + /// disagreement leaves every existing counter at zero. Across a pair, one + /// yield and one drop is agreement. + #[tokio::test] + async fn test_tiebreaker_records_one_yield_and_one_drop_across_a_pair() { + let (smaller, larger) = pubkeys_ordered_by_node_addr(); + + // The node with the LARGER address accepts an inbound from the + // smaller: its inbound wins, so nothing is stood down here. Invert it + // — the SMALLER node accepting from the larger stands its inbound + // down, because its own outbound is meant to win. + let io = MockBleIo::new("hci0", test_addr(1)); + let (tx, _rx) = tokio::sync::mpsc::channel(64); + let mut inbound_side = + BleTransport::new(TransportId::new(1), None, identity_test_config(), io, tx); + inbound_side.set_local_pubkey(smaller); + inbound_side.start_async().await.unwrap(); + + let (ours, peer) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + inbound_side.io.inject_inbound(ours).await; + peer_side_exchange(&peer, &larger).await; + { + let stats = Arc::clone(&inbound_side.stats); + wait_for("the inbound tie-breaker to conclude", || { + stats.snapshot().tiebreaker_drops == 1 + }) + .await; + } + + let snap = inbound_side.stats.snapshot(); + assert_eq!(snap.tiebreaker_drops, 1, "our inbound stood down"); + assert_eq!(snap.tiebreaker_yields, 0); + assert_eq!(inbound_side.pool.lock().await.len(), 0); + inbound_side.stop_async().await.unwrap(); + + // The other side of the same pair: the node with the LARGER address + // probing outbound stands its dial down, because the smaller node's + // outbound is meant to win. + let dials: DialLog = Arc::new(std::sync::Mutex::new(Vec::new())); + let io2 = MockBleIo::new("hci0", test_addr(2)); + let (peer_tx, mut peer_rx) = tokio::sync::mpsc::unbounded_channel(); + io2.set_connect_handler(move |addr, psm| { + dials.lock().unwrap().push((addr.clone(), psm)); + let (ours, theirs) = MockBleStream::pair(test_addr(2), addr.clone(), 2048); + peer_tx + .send(theirs) + .map_err(|_| TransportError::ConnectionRefused)?; + Ok(ours) + }); + tokio::spawn(async move { + let mut alive = Vec::new(); + while let Some(theirs) = peer_rx.recv().await { + peer_side_exchange(&theirs, &smaller).await; + alive.push(theirs); + } + }); + + let config = BleConfig { + scan: Some(true), + accept_connections: Some(false), + ..identity_test_config() + }; + let (tx2, _rx2) = tokio::sync::mpsc::channel(64); + let mut outbound_side = BleTransport::new(TransportId::new(2), None, config, io2, tx2); + outbound_side.set_local_pubkey(larger); + outbound_side.start_async().await.unwrap(); + + outbound_side.io.inject_scan_result(test_addr(1)).await; + { + let stats = Arc::clone(&outbound_side.stats); + wait_for("the outbound tie-breaker to conclude", || { + stats.snapshot().tiebreaker_yields == 1 + }) + .await; + } + + let snap = outbound_side.stats.snapshot(); + assert_eq!(snap.tiebreaker_yields, 1, "our outbound stood down"); + assert_eq!(snap.tiebreaker_drops, 0); + assert_eq!(outbound_side.pool.lock().await.len(), 0); + outbound_side.stop_async().await.unwrap(); + } + + /// An oversized packet is a caller bug, not a property of the peer's + /// link. Folding it into `send_errors` would make that number useless as + /// evidence. + #[tokio::test] + async fn test_mtu_rejection_does_not_count_as_a_send_error() { + let io = MockBleIo::new("hci0", test_addr(1)); + let (tx, _rx) = tokio::sync::mpsc::channel(64); + let transport = + BleTransport::new(TransportId::new(1), None, identity_test_config(), io, tx); + + let ta = test_addr(2).to_transport_addr(); + let (parked, _peer) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + transport + .pool + .lock() + .await + .insert( + ta.clone(), + BleConnection { + stream: Arc::new(parked), + recv_task: None, + send_mtu: 64, + recv_mtu: 64, + established_at: tokio::time::Instant::now(), + is_static: false, + addr: test_addr(2), + node_addr: None, + }, + ) + .unwrap(); + + let err = transport.send_async(&ta, &[0u8; 128]).await.unwrap_err(); + assert!(matches!(err, TransportError::MtuExceeded { .. })); + + let snap = transport.stats.snapshot(); + assert_eq!(snap.mtu_exceeded, 1); + assert_eq!(snap.send_errors, 0); + } + + /// The snapshot is the control-socket contract. Pin every key so a field + /// cannot be dropped or renamed without a test saying so. + #[test] + fn test_snapshot_carries_every_counter() { + let value = serde_json::to_value(BleStats::new().snapshot()).unwrap(); + let object = value.as_object().unwrap(); + let expected = [ + "packets_sent", + "bytes_sent", + "packets_recv", + "bytes_recv", + "send_errors", + "recv_errors", + "mtu_exceeded", + "connections_established", + "connections_accepted", + "connections_rejected", + "connect_timeouts", + "connect_errors", + "pubkey_exchange_failures", + "tiebreaker_yields", + "tiebreaker_drops", + "pool_evictions", + "advertisements_sent", + "scan_results", + "duplicate_node_declines", + ]; + for key in expected { + assert!(object.contains_key(key), "snapshot lost `{key}`"); + } + assert_eq!(object.len(), expected.len(), "snapshot gained a key"); + } } diff --git a/src/transport/ble/stats.rs b/src/transport/ble/stats.rs index 9e43a82d..d84cc3ef 100644 --- a/src/transport/ble/stats.rs +++ b/src/transport/ble/stats.rs @@ -1,4 +1,12 @@ //! BLE transport statistics. +//! +//! Counters reach an operator through `show_transports`, which serves them +//! off the control socket. Each connect outcome also emits a `debug!` at the +//! moment it happens, carrying a uniform field set — `addr`, `role` +//! (`central` for a dial, `peripheral` for an accept), `outcome` (a stable +//! kebab-case string matching the counter name), and `discovery_ms` where a +//! probe stamp exists. The counters give the aggregate; the trace stream +//! gives the same taxonomy per event and per peer. use portable_atomic::{AtomicU64, Ordering}; @@ -20,6 +28,14 @@ pub struct BleStats { pub connections_accepted: AtomicU64, pub connections_rejected: AtomicU64, pub connect_timeouts: AtomicU64, + /// Outbound connects that failed with an error rather than timing out. + pub connect_errors: AtomicU64, + /// Connections dropped because the pre-handshake pubkey exchange failed. + pub pubkey_exchange_failures: AtomicU64, + /// Outbound connections stood down by the cross-probe tie-breaker. + pub tiebreaker_yields: AtomicU64, + /// Inbound connections stood down by the cross-probe tie-breaker. + pub tiebreaker_drops: AtomicU64, pub pool_evictions: AtomicU64, pub advertisements_sent: AtomicU64, pub scan_results: AtomicU64, @@ -43,6 +59,10 @@ impl BleStats { connections_accepted: AtomicU64::new(0), connections_rejected: AtomicU64::new(0), connect_timeouts: AtomicU64::new(0), + connect_errors: AtomicU64::new(0), + pubkey_exchange_failures: AtomicU64::new(0), + tiebreaker_yields: AtomicU64::new(0), + tiebreaker_drops: AtomicU64::new(0), pool_evictions: AtomicU64::new(0), advertisements_sent: AtomicU64::new(0), scan_results: AtomicU64::new(0), @@ -97,6 +117,40 @@ impl BleStats { self.connect_timeouts.fetch_add(1, Ordering::Relaxed); } + /// Record an outbound connect that failed with an error. + /// + /// Kept separate from [`Self::record_connect_timeout`]: a refusal and a + /// silence are different faults and blur into one useless number if + /// merged. + pub fn record_connect_error(&self) { + self.connect_errors.fetch_add(1, Ordering::Relaxed); + } + + /// Record a failed pre-handshake pubkey exchange. + /// + /// The link came up and then produced nothing usable — a different fault + /// from never connecting at all. + pub fn record_pubkey_exchange_failure(&self) { + self.pubkey_exchange_failures + .fetch_add(1, Ordering::Relaxed); + } + + /// Record an outbound connection stood down by the cross-probe + /// tie-breaker. + /// + /// Read together with [`Self::record_tiebreaker_drop`] across a pair of + /// nodes: one yield and one drop is the two sides agreeing; two yields or + /// two drops is the disagreement that leaves no other evidence. + pub fn record_tiebreaker_yield(&self) { + self.tiebreaker_yields.fetch_add(1, Ordering::Relaxed); + } + + /// Record an inbound connection stood down by the cross-probe + /// tie-breaker. See [`Self::record_tiebreaker_yield`]. + pub fn record_tiebreaker_drop(&self) { + self.tiebreaker_drops.fetch_add(1, Ordering::Relaxed); + } + /// Record a pool eviction (non-static peer displaced). pub fn record_pool_eviction(&self) { self.pool_evictions.fetch_add(1, Ordering::Relaxed); @@ -136,6 +190,10 @@ impl BleStats { connections_accepted: self.connections_accepted.load(Ordering::Relaxed), connections_rejected: self.connections_rejected.load(Ordering::Relaxed), connect_timeouts: self.connect_timeouts.load(Ordering::Relaxed), + connect_errors: self.connect_errors.load(Ordering::Relaxed), + pubkey_exchange_failures: self.pubkey_exchange_failures.load(Ordering::Relaxed), + tiebreaker_yields: self.tiebreaker_yields.load(Ordering::Relaxed), + tiebreaker_drops: self.tiebreaker_drops.load(Ordering::Relaxed), pool_evictions: self.pool_evictions.load(Ordering::Relaxed), advertisements_sent: self.advertisements_sent.load(Ordering::Relaxed), scan_results: self.scan_results.load(Ordering::Relaxed), @@ -164,6 +222,10 @@ pub struct BleStatsSnapshot { pub connections_accepted: u64, pub connections_rejected: u64, pub connect_timeouts: u64, + pub connect_errors: u64, + pub pubkey_exchange_failures: u64, + pub tiebreaker_yields: u64, + pub tiebreaker_drops: u64, pub pool_evictions: u64, pub advertisements_sent: u64, pub scan_results: u64,