fix(transport/ble): recognise a peer by node identity, not by its rotating link address

A peer using resolvable private addresses rotates its BLE address
continually, and every rotation presents as a brand-new device. This is
not an exotic case: RPA rotation is a BLE privacy feature that modern
phones use by default, so a BlueZ node scanning one hits this today.
Field capture from a two-node mesh recorded one peer dialling in 28 times
from 28 distinct addresses in twenty minutes.

Every identity check in this transport keys on the link address, so none
of them could tell. `ConnectionPool` is a `HashMap<TransportAddr, _>`,
and all three already-connected guards — one in `accept_loop`, two in
`scan_probe_loop` — ask `pool.contains(&addr.to_transport_addr())`. For a
rotated address the answer is always "not connected", so the caller opens
another link, and another.

Those 28 rotations were harmless only by luck. The cross-probe
tie-breaker happened to reject every one of them, and which side it
protects is decided by a byte comparison of two node addresses. Had they
sorted the other way, the same 28 inbound dials would have been admitted
into a pool that holds seven, evicting genuine peers roughly four times
over. The tie-breaker is not the problem and is unchanged here — the
problem is that link identity was being used as node identity.

`BleConnection` now carries the peer's `NodeAddr` once the pubkey
exchange has learned it, and `ConnectionPool::find_by_node` looks a peer
up by the identity that does not rotate. Both admission points consult it
after the exchange and decline a duplicate, keeping the incumbent link:
it is known-good, and a genuinely dead one is already reaped by the
send-error and receive-loop paths.

Two consequences follow and are handled here rather than left as sequels.

The scan loop remembers what a declined address resolved to and skips it
while that node is still connected. Without that it re-dials every
rotated alias of a live peer once per cooldown, forever — the duplicate
is declined, so the alias never enters the pool, so the pool-keyed guard
above never sees it. The mapping is dropped as soon as the node leaves
the pool, so a peer that genuinely goes away is probed normally again.

And a completed exchange is announced under the address the peer's link
is actually on, not the rotated alias it happened on. Announcing the
alias makes a consumer that compares addresses treat it as a new path to
a peer it already holds a link to, and dial it; the duplicate is
declined, so nothing upstream remembers the conclusion and the next
round repeats it. `ConnectionPool::live_addr_of_node` answers with the
`BleAddr` the link is on, where `find_by_node` answers with the pool key
— the caller here has to name the address, not merely test for one.
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.

Declines are counted rather than silent, so absorption of a rotating peer
shows up in `show_transports` as a climbing `duplicate_node_declines`
instead of as an absence of log lines.

Coverage: all shared code driven by `MockBleIo`, so `cargo test` and
`cargo clippy --all-targets` cover every branch on any platform where
`transport::ble` compiles. Seven new pool tests, including one pinning
the regression — ten rotations of one peer leave the pool holding exactly
one link — plus two integration tests driving the inbound decline and the
scan loop's alias suppression end to end. Nothing here is inside
`cfg(bluer_available)`.
This commit is contained in:
Arjen
2026-08-26 07:26:23 +01:00
committed by Johnathan Corgan
parent 814ceef680
commit 8ba8076dbb
3 changed files with 498 additions and 23 deletions
+324 -23
View File
@@ -385,12 +385,16 @@ impl<I: BleIo> BleTransport<I> {
let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu);
// Pre-handshake pubkey exchange (temporary, pre-XX)
let mut peer_node: Option<NodeAddr> = 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<I: BleIo> BleTransport<I> {
}
}
self.promote_connection(addr, &ble_addr, stream, reader)
self.promote_connection(addr, &ble_addr, stream, reader, peer_node)
.await
}
@@ -412,6 +416,7 @@ impl<I: BleIo> BleTransport<I> {
ble_addr: &BleAddr,
stream: Arc<I::Stream>,
reader: BleStreamRead<I::Stream>,
node_addr: Option<NodeAddr>,
) -> Result<(), TransportError> {
let send_mtu = stream.send_mtu();
let recv_mtu = stream.recv_mtu();
@@ -434,6 +439,7 @@ impl<I: BleIo> BleTransport<I> {
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<I: BleIo> BleTransport<I> {
let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu);
// Pre-handshake pubkey exchange (temporary, pre-XX)
let mut peer_node: Option<NodeAddr> = 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<I: BleIo> BleTransport<I> {
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<S>(
pool: &Mutex<ConnectionPool<S>>,
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<A>(
let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu);
// Pre-handshake pubkey exchange (temporary, pre-XX)
let mut peer_node_addr: Option<NodeAddr> = 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<A>(
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<I: io::BleIo>(
// 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();
// 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<BleAddr, NodeAddr> = 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<I: io::BleIo>(
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<I: io::BleIo>(
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<I: io::BleIo>(
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<StdMutex<Vec<BleAddr>>> = 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]
+158
View File
@@ -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<S> {
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<NodeAddr>,
}
impl<S> BleConnection<S> {
@@ -95,6 +106,38 @@ impl<S> ConnectionPool<S> {
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<TransportAddr> {
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<BleAddr> {
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)));
}
}
+16
View File
@@ -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,
}