fix(transport/ble): bound the probe retry and explain the connect outcomes

Two defects in the scan/probe loop that share the same code and the same
field capture, so they are fixed together.

A discovered address that fails to connect was re-dialled every cooldown for
the life of the process. One unreachable peer therefore consumed a dial slot
forever, and because BLE hardware caps concurrent connections at roughly four
to ten, a handful of them starve discovery of everything behind them.
PendingProbes now backs a failing address off by powers of two up to
MAX_PROBE_BACKOFF_SHIFT, which at the 30 s default caps a failing address at
one attempt every sixteen minutes, and the book itself is capped at
MAX_PENDING_PROBES entries so that rotating private addresses cannot grow it
without bound.

The stats snapshot could not explain any of it. A probe that failed and one
that was never attempted were indistinguishable, so there was no way to tell a
peer out of range from a peer being dialled at the wrong PSM. Each connect
outcome now has its own counter and its own structured log line carrying the
role, the outcome, the PSM dialled and how long the peer took to go from
advertisement to conclusion.

The two touch the same arms because the outcome that needed counting most is
the one that also needed backing off: a dial failure now both records its
reason and forgets the peer's learned PSM, so a stale advertised PSM costs one
retry and is re-learned from the next advertisement rather than being retried
forever at a number that cannot work.

Behaviour on the deployed BlueZ path changes in one way worth naming: an
address that fails is dialled less often. Nothing about which peers are
reachable changes, and a peer that connects is unaffected.
This commit is contained in:
Arjen
2026-08-26 07:58:58 +01:00
committed by Johnathan Corgan
parent ae93c90908
commit 901947899f
2 changed files with 756 additions and 37 deletions
+694 -37
View File
@@ -391,12 +391,19 @@ impl<I: BleIo> BleTransport<I> {
{
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<I: BleIo> BleTransport<I> {
.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<I: BleIo> BleTransport<I> {
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<I: BleIo> BleTransport<I> {
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<A>(
{
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<A>(
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<A>(
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<S: BleStream + 'static>(
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<PendingProbe>,
cooldown: std::time::Duration,
}
impl PendingProbes {
fn new(cooldown: std::time::Duration) -> Self {
Self {
entries: Vec::new(),
cooldown,
}
}
fn position(&self, addr: &BleAddr) -> Option<usize> {
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<BleAddr> {
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<I: io::BleIo>(
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();
// 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<I: io::BleIo>(
// which is every peer that predates this and every backend that does not
// advertise service data.
let mut learned_psm: HashMap<BleAddr, u16> = 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<I: io::BleIo>(
_ = 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<I: io::BleIo>(
{
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<I: io::BleIo>(
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<I: io::BleIo>(
}
// 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<I: io::BleIo>(
// 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<I: io::BleIo>(
{
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<I: io::BleIo>(
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<I: io::BleIo>(
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<I: io::BleIo>(
{
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<I: io::BleIo>(
// 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<I: io::BleIo>(
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");
}
}
+62
View File
@@ -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,