diff --git a/src/transport/ble/mod.rs b/src/transport/ble/mod.rs index 92d57062..16ab99db 100644 --- a/src/transport/ble/mod.rs +++ b/src/transport/ble/mod.rs @@ -385,12 +385,16 @@ impl BleTransport { let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); // Pre-handshake pubkey exchange (temporary, pre-XX) + let mut peer_node: Option = None; if let Some(ref our_pubkey) = self.local_pubkey { match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr, "BLE outbound pubkey exchange complete"); + let node = NodeAddr::from_pubkey(&peer_pubkey); + peer_node = Some(node); + let announced = announced_addr(&self.pool, &node, &ble_addr).await; self.neighbor_buffer - .add_peer_with_pubkey(&ble_addr, peer_pubkey); + .add_peer_with_pubkey(&announced, peer_pubkey); } Err(e) => { warn!(addr = %addr, error = %e, "BLE outbound pubkey exchange failed"); @@ -399,7 +403,7 @@ impl BleTransport { } } - self.promote_connection(addr, &ble_addr, stream, reader) + self.promote_connection(addr, &ble_addr, stream, reader, peer_node) .await } @@ -412,6 +416,7 @@ impl BleTransport { ble_addr: &BleAddr, stream: Arc, reader: BleStreamRead, + node_addr: Option, ) -> Result<(), TransportError> { let send_mtu = stream.send_mtu(); let recv_mtu = stream.recv_mtu(); @@ -434,6 +439,7 @@ impl BleTransport { established_at: tokio::time::Instant::now(), is_static: false, addr: ble_addr.clone(), + node_addr, }; let mut pool = self.pool.lock().await; @@ -512,11 +518,15 @@ impl BleTransport { let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); // Pre-handshake pubkey exchange (temporary, pre-XX) + let mut peer_node: Option = None; if let Some(ref our_pubkey) = local_pubkey { match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete"); - neighbor_buffer.add_peer_with_pubkey(&ble_addr, peer_pubkey); + let node = NodeAddr::from_pubkey(&peer_pubkey); + peer_node = Some(node); + let announced = announced_addr(&pool, &node, &ble_addr).await; + neighbor_buffer.add_peer_with_pubkey(&announced, peer_pubkey); } Err(e) => { warn!( @@ -546,6 +556,7 @@ impl BleTransport { established_at: tokio::time::Instant::now(), is_static: false, addr: ble_addr, + node_addr: peer_node, }; let mut pool = pool.lock().await; @@ -705,6 +716,36 @@ const PUBKEY_EXCHANGE_SIZE: usize = 33; /// forever — killing scan_probe_loop, accept_loop, or the event loop. const PUBKEY_EXCHANGE_TIMEOUT_SECS: u64 = 5; +/// The link address a completed pubkey exchange should be announced under. +/// +/// A peer using resolvable private addresses presents a different link +/// address on every rotation, so the address an exchange happened on is a +/// transient alias for the peer, not a durable way to name it. Announcing the +/// alias makes a consumer that compares addresses treat it as a *new path* to +/// a peer it is already connected to and dial it; the duplicate is declined +/// here, so it never reaches the pool, so nothing upstream remembers the +/// conclusion and the next discovery round pays the same connect and exchange +/// again. `scan_probe_loop` breaks that cycle for its own probes, but callers +/// that reach `connect_async` directly never consult it. +/// +/// So when the peer is already connected, report the address its link is +/// actually on: same peer, named by the address that works. When it is not, +/// there is no incumbent and the observed address stands. +/// +/// Canonicalising rather than withholding matters: suppressing the +/// announcement would also stop the peer being offered at all, and consumers +/// legitimately re-probe a peer whose link has gone idle to recover it. +async fn announced_addr( + pool: &Mutex>, + node: &NodeAddr, + observed: &BleAddr, +) -> BleAddr { + pool.lock() + .await + .live_addr_of_node(node) + .unwrap_or_else(|| observed.clone()) +} + /// Exchange public keys over a newly established L2CAP connection. /// /// Both sides send `[0x00][our_pubkey:32]` and receive the peer's. @@ -794,24 +835,51 @@ async fn accept_loop( let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); // Pre-handshake pubkey exchange (temporary, pre-XX) + let mut peer_node_addr: Option = None; if let Some(ref our_pubkey) = local_pubkey { match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %ta, "BLE inbound pubkey exchange complete"); - neighbor_buffer.add_peer_with_pubkey(&addr, peer_pubkey); + let peer_node = NodeAddr::from_pubkey(&peer_pubkey); + peer_node_addr = Some(peer_node); + let announced = announced_addr(&pool, &peer_node, &addr).await; + neighbor_buffer.add_peer_with_pubkey(&announced, peer_pubkey); + + // Already linked to this peer on another address? + // A peer using resolvable private addresses rotates + // continually, and every rotation dials in looking + // like a new device. Admitting those would put one + // peer in several pool slots and evict real ones. + // The incumbent link is kept: it is known-good, and + // a genuinely dead one is already reaped by the + // send-error and receive-loop paths. + let dup = { + let pool_guard = pool.lock().await; + pool_guard.find_by_node(&peer_node) + }; + if let Some(existing) = dup + && existing != ta + { + debug!( + addr = %ta, + existing = %existing, + "BLE inbound: peer already connected on another address, dropping duplicate" + ); + stats.record_duplicate_node_decline(); + continue; + } // Cross-probe tie-breaker: smaller NodeAddr's // outbound wins. If we're smaller, our outbound // should win — drop this inbound. - if let Some(ref our_addr) = local_node_addr { - let peer_addr = NodeAddr::from_pubkey(&peer_pubkey); - if our_addr < &peer_addr { - debug!( - addr = %ta, - "BLE inbound tie-breaker: dropping (our addr < peer, outbound wins)" - ); - continue; - } + if let Some(ref our_addr) = local_node_addr + && our_addr < &peer_node + { + debug!( + addr = %ta, + "BLE inbound tie-breaker: dropping (our addr < peer, outbound wins)" + ); + continue; } } Err(e) => { @@ -840,6 +908,7 @@ async fn accept_loop( established_at: tokio::time::Instant::now(), is_static: false, addr, + node_addr: peer_node_addr, }; let mut pool_guard = pool.lock().await; @@ -942,6 +1011,13 @@ async fn scan_probe_loop( // 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(); + // 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 + // exchange to rediscover that every cooldown. Rotation means this grows by + // one per rotation, so entries are dropped once their node is no longer in + // the pool — a peer that genuinely goes away is probed again normally. + let mut known_node_of: 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); @@ -997,6 +1073,23 @@ async fn scan_probe_loop( continue; } + // Skip an address already known to belong to a peer we are connected + // to. Without this the loop re-dials every rotated address of a live + // peer once per cooldown, forever: the duplicate is declined so it + // never enters the pool, so the pool-keyed guard above never sees it. + if let Some(node) = known_node_of.get(&addr) { + let still_connected = { + let pool_guard = pool.lock().await; + pool_guard.find_by_node(node).is_some() + }; + if still_connected { + pending_addrs.retain(|a| a != &addr); + continue; + } + // That peer is gone — forget the mapping and probe normally. + known_node_of.remove(&addr); + } + // Record probe time (before attempt, so cooldown applies on failure too) last_probed.insert(addr.clone(), tokio::time::Instant::now()); @@ -1038,19 +1131,50 @@ async fn scan_probe_loop( match pubkey_exchange(stream.as_ref(), &mut reader, &our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr, "BLE probe complete"); + let peer_node = NodeAddr::from_pubkey(&peer_pubkey); // 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 { - debug!( - addr = %addr, - "BLE probe tie-breaker: yielding to peer's outbound" - ); - buffer.add_peer_with_pubkey(&addr, peer_pubkey); - continue; - } + if let Some(ref our_addr) = local_node_addr + && our_addr >= &peer_node + { + debug!( + addr = %addr, + "BLE probe tie-breaker: yielding to peer's outbound" + ); + let announced = announced_addr(&pool, &peer_node, &addr).await; + buffer.add_peer_with_pubkey(&announced, peer_pubkey); + continue; + } + + // Same duplicate guard as the inbound path: a rotated address + // for a peer we already hold a link to must not become a + // second pool entry. Checked after the tie-breaker so the two + // decisions stay independent. + let dup = { + let pool_guard = pool.lock().await; + pool_guard.find_by_node(&peer_node) + }; + if let Some(existing) = dup + && existing != ta + { + debug!( + addr = %ta, + existing = %existing, + "BLE probe: peer already connected on another address, dropping duplicate" + ); + stats.record_duplicate_node_decline(); + // Remember what this address resolved to, so the next + // cooldown skips it outright rather than paying another + // connect and exchange to reach the same conclusion. + known_node_of.insert(addr.clone(), peer_node); + // Report the peer under the address its live link is on, + // so the node layer is not handed an alias with no + // 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); + continue; } // Promote connection to pool — no second L2CAP connect needed @@ -1072,6 +1196,7 @@ async fn scan_probe_loop( established_at: tokio::time::Instant::now(), is_static: false, addr: addr.clone(), + node_addr: Some(peer_node), }; let mut pool_guard = pool.lock().await; @@ -1371,6 +1496,7 @@ mod tests { established_at: tokio::time::Instant::now(), is_static: false, addr: test_addr(2), + node_addr: None, }, ) .unwrap(); @@ -1434,6 +1560,181 @@ mod tests { assert_eq!(got.serialize(), peer_pk); } + // ------------------------------------------------------------------ + // Node identity vs. rotating link address + // ------------------------------------------------------------------ + + /// Two pubkeys, returned as `(smaller_node_addr, larger_node_addr)`. + /// + /// The cross-probe tie-breaker is decided by `NodeAddr` ordering, so a + /// test that wants a connection admitted has to know which side it is. + fn pubkeys_ordered_by_node_addr() -> ([u8; 32], [u8; 32]) { + let a = test_pubkey(1); + let b = test_pubkey(2); + let na = NodeAddr::from_pubkey(&XOnlyPublicKey::from_slice(&a).unwrap()); + let nb = NodeAddr::from_pubkey(&XOnlyPublicKey::from_slice(&b).unwrap()); + if na < nb { (a, b) } else { (b, a) } + } + + /// Run the peer half of the pubkey exchange over a mock stream end. + async fn peer_side_exchange(peer: &MockBleStream, peer_pubkey: &[u8; 32]) { + let mut msg = [0u8; PUBKEY_EXCHANGE_SIZE]; + msg[0] = PUBKEY_EXCHANGE_PREFIX; + msg[1..].copy_from_slice(peer_pubkey); + peer.send(&msg).await.unwrap(); + let mut buf = [0u8; PUBKEY_EXCHANGE_SIZE]; + let n = peer.recv(&mut buf).await.unwrap(); + assert_eq!(n, PUBKEY_EXCHANGE_SIZE); + } + + fn identity_test_config() -> BleConfig { + BleConfig { + adapter: Some("hci0".to_string()), + scan: Some(false), + advertise: Some(false), + accept_connections: Some(true), + probe_cooldown_secs: Some(1), + ..Default::default() + } + } + + /// Let spawned loops make progress. + async fn settle() { + for _ in 0..64 { + tokio::task::yield_now().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. + #[tokio::test] + async fn test_inbound_rotation_is_declined_and_keeps_the_incumbent() { + let (smaller, larger) = pubkeys_ordered_by_node_addr(); + 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); + // We take the larger node address, so the inbound tie-breaker admits + // rather than drops — the duplicate guard is what is under test. + transport.set_local_pubkey(larger); + transport.start_async().await.unwrap(); + + // First inbound, on link address 2. + let (ours, peer_a) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + transport.io.inject_inbound(ours).await; + peer_side_exchange(&peer_a, &smaller).await; + settle().await; + + assert_eq!(transport.pool.lock().await.len(), 1); + assert!( + transport + .pool + .lock() + .await + .contains(&test_addr(2).to_transport_addr()) + ); + + // The same node dials in again after rotating to link address 3. + let (ours2, peer_b) = MockBleStream::pair(test_addr(1), test_addr(3), 2048); + transport.io.inject_inbound(ours2).await; + peer_side_exchange(&peer_b, &smaller).await; + settle().await; + + let pool = transport.pool.lock().await; + assert_eq!(pool.len(), 1, "the rotation must not become a second link"); + assert!( + pool.contains(&test_addr(2).to_transport_addr()), + "the incumbent link is kept" + ); + assert!(!pool.contains(&test_addr(3).to_transport_addr())); + drop(pool); + + assert_eq!(transport.stats.snapshot().duplicate_node_declines, 1); + + // Discovery names the peer by the address its link is actually on, + // not by the alias the rotation arrived from — otherwise the node + // layer is handed an address with no connection behind it. + let peers = transport.neighbor_buffer.take(); + assert_eq!(peers.len(), 1); + assert_eq!(peers[0].addr, test_addr(2).to_transport_addr()); + + transport.stop_async().await.unwrap(); + drop((peer_a, peer_b)); + } + + /// Once a rotated alias has been resolved to a peer that holds a live + /// link, the scan loop stops paying a connect and exchange to reach that + /// same conclusion every cooldown. + #[tokio::test(start_paused = true)] + async fn test_scan_loop_stops_reprobing_a_resolved_alias() { + use std::sync::Mutex as StdMutex; + + let (smaller, larger) = pubkeys_ordered_by_node_addr(); + let io = MockBleIo::new("hci0", test_addr(1)); + + let connects: Arc>> = Arc::new(StdMutex::new(Vec::new())); + let (peer_tx, mut peer_rx) = tokio::sync::mpsc::unbounded_channel(); + { + let connects = Arc::clone(&connects); + io.set_connect_handler(move |addr, _psm| { + let (ours, theirs) = MockBleStream::pair(test_addr(1), addr.clone(), 2048); + connects.lock().unwrap().push(addr.clone()); + peer_tx + .send(theirs) + .map_err(|_| TransportError::ConnectionRefused)?; + Ok(ours) + }); + } + + // The remote answers every probe with one identity, whichever link + // address the probe went to. + tokio::spawn(async move { + let mut alive = Vec::new(); + while let Some(theirs) = peer_rx.recv().await { + peer_side_exchange(&theirs, &larger).await; + alive.push(theirs); + } + }); + + let config = BleConfig { + scan: Some(true), + accept_connections: Some(false), + ..identity_test_config() + }; + let (tx, _rx) = tokio::sync::mpsc::channel(64); + let mut transport = BleTransport::new(TransportId::new(1), None, config, io, tx); + // We take the smaller node address, so our outbound wins the + // tie-breaker and the probe is promoted. + transport.set_local_pubkey(smaller); + transport.start_async().await.unwrap(); + + transport.io.inject_scan_result(test_addr(2)).await; + settle().await; + assert_eq!(transport.pool.lock().await.len(), 1); + assert_eq!(connects.lock().unwrap().len(), 1); + + // The peer rotates to address 3. That probe is paid once and declined. + transport.io.inject_scan_result(test_addr(3)).await; + settle().await; + assert_eq!(connects.lock().unwrap().len(), 2); + assert_eq!(transport.stats.snapshot().duplicate_node_declines, 1); + assert_eq!(transport.pool.lock().await.len(), 1); + + // The alias is advertised again after the cooldown expires. It must + // not be dialled a third time: the loop already knows whose it is. + tokio::time::advance(std::time::Duration::from_secs(5)).await; + transport.io.inject_scan_result(test_addr(3)).await; + settle().await; + assert_eq!( + connects.lock().unwrap().len(), + 2, + "a resolved alias of a live peer must not be re-dialled" + ); + + transport.stop_async().await.unwrap(); + } + /// A peer that opens with something other than the exchange prefix is /// rejected before the framer ever sees the bytes. #[tokio::test] diff --git a/src/transport/ble/pool.rs b/src/transport/ble/pool.rs index ffdf45db..837b1c0c 100644 --- a/src/transport/ble/pool.rs +++ b/src/transport/ble/pool.rs @@ -8,6 +8,7 @@ use std::collections::HashMap; use tokio::task::JoinHandle; +use crate::identity::NodeAddr; use crate::transport::{TransportAddr, TransportError}; use super::addr::BleAddr; @@ -28,6 +29,16 @@ pub struct BleConnection { pub is_static: bool, /// Parsed remote address. pub addr: BleAddr, + /// The peer's node address, once the pubkey exchange has learned it. + /// + /// The pool is keyed by *link* address, but a BLE link address is not a + /// stable identity: peers using resolvable private addresses rotate theirs + /// continually, and each rotation looks like a brand-new device. This + /// field carries the identity that does not rotate, so + /// [`ConnectionPool::find_by_node`] can recognise a peer already connected + /// under an address never seen before. `None` for a connection whose peer + /// is not yet identified. + pub node_addr: Option, } impl BleConnection { @@ -95,6 +106,38 @@ impl ConnectionPool { self.connections.contains_key(addr) } + /// Find an existing connection to `node`, whatever link address it + /// arrived on. + /// + /// This is the identity check [`Self::contains`] cannot make. A peer using + /// resolvable private addresses presents a different link address every + /// rotation, so an address-keyed lookup reports "not connected" for a peer + /// that is very much connected — and the caller then opens a second link + /// to it, and a third. Callers that know the peer's node address should + /// ask this before admitting a connection. + /// + /// Only connections whose pubkey exchange has completed carry a node + /// address, so an unidentified connection is never matched. + pub fn find_by_node(&self, node: &NodeAddr) -> Option { + self.connections + .iter() + .find(|(_, c)| c.node_addr.as_ref() == Some(node)) + .map(|(addr, _)| addr.clone()) + } + + /// The live link address for `node`, if it is connected. + /// + /// [`Self::find_by_node`] answers with the pool key; this answers with the + /// `BleAddr` the link is actually on, which is what a caller needs when it + /// has to *name* the peer's current address rather than merely test for + /// one. + pub fn live_addr_of_node(&self, node: &NodeAddr) -> Option { + self.connections + .values() + .find(|c| c.node_addr.as_ref() == Some(node)) + .map(|c| c.addr.clone()) + } + /// Try to insert a connection, evicting if necessary. /// /// Returns `Ok(evicted_addr)` on success (with optional evicted peer), @@ -186,6 +229,13 @@ mod tests { } } + /// A distinct node identity per `n` — the identity that does NOT rotate. + fn test_node(n: u8) -> NodeAddr { + let mut bytes = [0u8; 16]; + bytes[0] = n; + NodeAddr::from_bytes(bytes) + } + fn test_conn(n: u8, is_static: bool) -> BleConnection<()> { BleConnection { stream: (), @@ -195,6 +245,7 @@ mod tests { established_at: tokio::time::Instant::now(), is_static, addr: test_ble_addr(n), + node_addr: None, } } @@ -287,4 +338,111 @@ mod tests { addrs.sort_by(|a, b| a.as_str().cmp(&b.as_str())); assert_eq!(addrs.len(), 2); } + + /// A node address is found regardless of which link address it arrived on + /// — the whole point of the lookup, since the link address rotates. + #[test] + fn test_find_by_node_matches_across_a_rotated_link_address() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + let node = test_node(1); + let mut conn = test_conn(1, false); + conn.node_addr = Some(node); + pool.insert(test_addr(1), conn).unwrap(); + + // Found under the address it was inserted with... + assert_eq!(pool.find_by_node(&node), Some(test_addr(1))); + assert!(pool.contains(&test_addr(1))); + // ...and a rotated address for the same peer is NOT found by + // `contains`, which is exactly the gap `find_by_node` closes. + assert!(!pool.contains(&test_addr(99))); + assert_eq!(pool.find_by_node(&node), Some(test_addr(1))); + } + + #[test] + fn test_find_by_node_ignores_unidentified_connections() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + // No pubkey exchange yet, so no node address. + pool.insert(test_addr(1), test_conn(1, false)).unwrap(); + assert_eq!(pool.find_by_node(&test_node(1)), None); + } + + #[test] + fn test_find_by_node_returns_none_for_an_unconnected_node() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + let mut conn = test_conn(1, false); + conn.node_addr = Some(test_node(1)); + pool.insert(test_addr(1), conn).unwrap(); + assert_eq!(pool.find_by_node(&test_node(2)), None); + } + + /// Distinct nodes do not alias: each resolves to its own link address. + #[test] + fn test_find_by_node_distinguishes_two_nodes() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + let mut a = test_conn(1, false); + a.node_addr = Some(test_node(1)); + let mut b = test_conn(2, false); + b.node_addr = Some(test_node(2)); + pool.insert(test_addr(1), a).unwrap(); + pool.insert(test_addr(2), b).unwrap(); + + assert_eq!(pool.find_by_node(&test_node(1)), Some(test_addr(1))); + assert_eq!(pool.find_by_node(&test_node(2)), Some(test_addr(2))); + } + + /// The live link address is reported for a peer found under any of its + /// rotated aliases — what a caller needs when it has to name the peer's + /// current address rather than merely test for one. + #[test] + fn test_live_addr_of_node_reports_the_incumbent_not_the_alias() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + let node = test_node(1); + let mut conn = test_conn(1, false); + conn.node_addr = Some(node); + pool.insert(test_addr(1), conn).unwrap(); + + assert_eq!(pool.live_addr_of_node(&node), Some(test_ble_addr(1))); + assert!(!pool.contains(&test_addr(99))); + // An unconnected node has no incumbent, so the caller keeps whatever + // address it observed. + assert_eq!(pool.live_addr_of_node(&test_node(2)), None); + } + + #[test] + fn test_live_addr_of_node_ignores_unidentified_connections() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + pool.insert(test_addr(1), test_conn(1, false)).unwrap(); + assert_eq!(pool.live_addr_of_node(&test_node(1)), None); + } + + /// The regression this guards: without a node-identity check, N rotated + /// addresses for ONE peer become N pool entries and evict real peers. + /// With it, the caller sees the peer is already present and declines. + #[test] + fn test_rotated_addresses_would_otherwise_fill_the_pool() { + let mut pool: ConnectionPool<()> = ConnectionPool::new(7); + let node = test_node(1); + + let mut first = test_conn(1, false); + first.node_addr = Some(node); + pool.insert(test_addr(1), first).unwrap(); + + // Ten rotations arrive. Each is a distinct link address, so `contains` + // says "new" every time — but `find_by_node` recognises all of them. + for n in 2..12u8 { + assert!( + !pool.contains(&test_addr(n)), + "rotation {n} looks new by address" + ); + assert_eq!( + pool.find_by_node(&node), + Some(test_addr(1)), + "rotation {n} is recognised as the peer already connected", + ); + } + // Nothing was admitted, so the pool still holds exactly one link, and + // it is the incumbent — the first one, not the newest. + assert_eq!(pool.len(), 1); + assert!(pool.contains(&test_addr(1))); + } } diff --git a/src/transport/ble/stats.rs b/src/transport/ble/stats.rs index f4fbccf5..9e43a82d 100644 --- a/src/transport/ble/stats.rs +++ b/src/transport/ble/stats.rs @@ -23,6 +23,9 @@ pub struct BleStats { pub pool_evictions: AtomicU64, pub advertisements_sent: AtomicU64, pub scan_results: AtomicU64, + /// Connections declined because the peer was already linked on another + /// link address (see `ConnectionPool::find_by_node`). + pub duplicate_node_declines: AtomicU64, } impl BleStats { @@ -43,6 +46,7 @@ impl BleStats { pool_evictions: AtomicU64::new(0), advertisements_sent: AtomicU64::new(0), scan_results: AtomicU64::new(0), + duplicate_node_declines: AtomicU64::new(0), } } @@ -108,6 +112,16 @@ impl BleStats { self.scan_results.fetch_add(1, Ordering::Relaxed); } + /// Record a connection declined as a duplicate of a peer already linked + /// under a different link address. + /// + /// A peer using resolvable private addresses rotates continually, so a + /// climbing count against a busy mesh is that rotation being absorbed — + /// not an error. + pub fn record_duplicate_node_decline(&self) { + self.duplicate_node_declines.fetch_add(1, Ordering::Relaxed); + } + /// Take a snapshot of all counters. pub fn snapshot(&self) -> BleStatsSnapshot { BleStatsSnapshot { @@ -125,6 +139,7 @@ impl BleStats { pool_evictions: self.pool_evictions.load(Ordering::Relaxed), advertisements_sent: self.advertisements_sent.load(Ordering::Relaxed), scan_results: self.scan_results.load(Ordering::Relaxed), + duplicate_node_declines: self.duplicate_node_declines.load(Ordering::Relaxed), } } } @@ -152,4 +167,5 @@ pub struct BleStatsSnapshot { pub pool_evictions: u64, pub advertisements_sent: u64, pub scan_results: u64, + pub duplicate_node_declines: u64, }