From 7ecd91ef4200a7182e3e294c8d7089f3fffe47e4 Mon Sep 17 00:00:00 2001 From: webbeef Date: Mon, 21 Sep 2026 17:43:37 -0700 Subject: [PATCH] p2p: improve BLE reliability Signed-off-by: webbeef --- crates/beaver_p2p/Cargo.toml | 14 +- .../{src/main.rs => examples/pair.rs} | 5 + crates/beaver_p2p/src/discovery/ble.rs | 301 ++++++++++-------- crates/beaver_p2p/src/discovery/mod.rs | 52 +-- crates/beaver_p2p/src/lib.rs | 178 ++++------- crates/beaver_p2p/src/packet.rs | 14 +- crates/beaver_p2p/src/pairing_protocol.rs | 140 ++++---- crates/beaver_p2p/src/state.rs | 72 ++--- crates/beaver_p2p/tests/pairing.rs | 78 ++--- ui/system/mediacenter/peers.js | 19 +- 10 files changed, 398 insertions(+), 475 deletions(-) rename crates/beaver_p2p/{src/main.rs => examples/pair.rs} (95%) diff --git a/crates/beaver_p2p/Cargo.toml b/crates/beaver_p2p/Cargo.toml index 15927f2..3a16ab1 100644 --- a/crates/beaver_p2p/Cargo.toml +++ b/crates/beaver_p2p/Cargo.toml @@ -4,14 +4,9 @@ version = "0.1.0" edition = "2024" [features] -# Bluetooth LE as a second transport alongside IP + mDNS. Off by default: it pulls -# in a BlueZ/CoreBluetooth stack, and `iroh`'s custom-transport API is explicitly -# exempt from semver, so a bad iroh bump should degrade to mDNS-only rather than -# break the build. ble = ["dep:iroh-ble-transport", "iroh/unstable-custom-transports"] [dependencies] -env_logger = "0.11" iroh = { workspace = true } iroh-ble-transport = { workspace = true, optional = true } iroh-mdns-address-lookup = { workspace = true } @@ -19,11 +14,14 @@ iroh-persist = { workspace = true } log = "0.4" n0-error = "1.0" n0-future = "0.3" -parking_lot = "0.12" -petname = { workspace = true } postcard = "1.1" rand = "0.9" serde = "1.0" thiserror = "2.0" -tokio = { version = "1.50", features = ["signal"] } +tokio = { version = "1.50", features = ["macros", "rt", "rt-multi-thread", "sync", "time"] } tokio-stream = { workspace = true } + +[dev-dependencies] +env_logger = "0.11" +parking_lot = "0.12" +petname = { workspace = true } diff --git a/crates/beaver_p2p/src/main.rs b/crates/beaver_p2p/examples/pair.rs similarity index 95% rename from crates/beaver_p2p/src/main.rs rename to crates/beaver_p2p/examples/pair.rs index 6f17b12..b3d5f19 100644 --- a/crates/beaver_p2p/src/main.rs +++ b/crates/beaver_p2p/examples/pair.rs @@ -1,3 +1,8 @@ +/* SPDX Id: AGPL-3.0-or-later */ + +//! Two-device pairing harness: run on both machines, pass any argument on the one +//! that should dial. + use std::sync::mpsc::channel; use beaver_p2p::{PairingManager, PairingResult, PeerEvent}; diff --git a/crates/beaver_p2p/src/discovery/ble.rs b/crates/beaver_p2p/src/discovery/ble.rs index 04efcc4..914bede 100644 --- a/crates/beaver_p2p/src/discovery/ble.rs +++ b/crates/beaver_p2p/src/discovery/ble.rs @@ -5,11 +5,10 @@ //! `iroh-ble-transport` has no discovery stream of its own, and its //! `BlePeerInfo` view carries no display name and no signal strength. The radio //! layer underneath has both, and the transport builder lets us hand it a -//! `Central` we keep a reference to -- so discovery reads the same scan the -//! transport is driving, via an independent subscription (blew's event stream is a -//! broadcast, so subscribing twice is fine). +//! `Central` we keep a reference to. Discovery can read the same scan the +//! transport is driving, via an independent subscription. //! -//! Two things make this source more than a translation layer: +//! This source does: //! //! - **Prefix resolution.** An advertisement identifies a peer by the first 12 //! bytes of its public key, not a full `EndpointId`, and every dial path needs @@ -17,10 +16,8 @@ //! asked directly, over a GATT read of its IDENTITY characteristic. That read //! requires a connection, which makes the peer stop advertising, so a //! resolution is announced immediately rather than on the next sighting. -//! - **Synthesised expiry.** BLE has no "advertisement stopped" event -- a peer -//! simply goes quiet. Without a `Lost` of our own, `seen_on` would keep `Ble` -//! forever and an mDNS expiry could never evict a peer that BLE had once seen, -//! which would regress LAN behaviour. So we age sightings out on a tick. +//! - **Synthesised expiry.** BLE has no "advertisement stopped" event, instead a peer +//! simply goes quiet. use std::collections::HashMap; use std::sync::Arc; @@ -42,49 +39,21 @@ use crate::discovery::{DiscoveredPeer, DiscoveryEvent, DiscoveryTransport}; use crate::state::SharedState; /// How long a peer may go unseen before we call it lost. -/// -/// Advertising intervals are well under a second, so this looks generous -- but the -/// silence we have to tolerate is largely self-inflicted: a peripheral stops -/// advertising while connected, and reading a peer's identity connects to it. A -/// measured round trip was 37s from announce to the next sighting, so 30s produced a -/// spurious `Lost` followed immediately by a re-announce. A spurious `Lost` also -/// costs a paired peer its `PairedConnected` status, so erring long is the cheaper -/// mistake; the price is that a peer which really has left lingers this long. const SIGHTING_TTL: Duration = Duration::from_secs(60); /// How often sightings are aged out. Also the worst-case lateness of a `Lost`. const SWEEP_INTERVAL: Duration = Duration::from_secs(5); /// How many sweeps between health lines while no peer is visible. -/// -/// Everything here logs at `info!` because `beaver_shell` enables -/// `log/release_max_level_info` by default, so `debug!` is compiled out of the builds -/// that actually run on a device. That makes volume a design constraint rather than -/// something `RUST_LOG` can fix, hence the throttle. const HEALTH_EVERY_SWEEPS: u32 = 12; /// Backstop on the GATT connect-and-read that turns a prefix into an identity. -/// -/// Deliberately longer than blew's own connect timeout (see `bring_up`), so that -/// bound fires first: blew cancels the pending connect and reports an error, where -/// this one only abandons the future and would leave the connect queued in the OS. -/// This is here for the case where the read itself wedges after connecting. const IDENTITY_READ_TIMEOUT: Duration = Duration::from_secs(20); /// How long to leave a peer alone between identity reads. -/// -/// Doubles as the in-flight guard: a read cannot outlive `IDENTITY_READ_TIMEOUT`, -/// which is shorter than this, so a peer still being read is always inside its own -/// retry window and cannot be dialled twice. Without a gate at all, a peer that -/// cannot answer -- an older build with no IDENTITY characteristic, or one simply -/// refusing connections -- would be re-dialled on every sighting. const IDENTITY_RETRY_AFTER: Duration = Duration::from_secs(60); /// The radio handles and the transport built on top of them. -/// -/// We construct `Central`/`Peripheral` rather than letting the transport do it, so -/// discovery can read the same scan (see the module docs), and so there is something -/// left to switch the radio off with at teardown. #[derive(Clone)] pub(crate) struct Ble { pub central: Arc, @@ -97,11 +66,7 @@ impl Ble { /// /// Has to be explicit: `BleTransport` has no shutdown of its own and holds its /// own `Arc` on both handles, so dropping ours would leave the scan and the - /// advertisement running for the life of the process -- the single largest power - /// draw this transport adds on a phone. - /// - /// Failures are logged rather than propagated. This is the teardown path, and a - /// radio that has already gone away is not worth failing over. + /// advertisement running. pub(crate) async fn shut_down(&self) { if let Err(err) = self.central.stop_scan().await { warn!("[BLE] failed to stop scanning: {err}"); @@ -114,24 +79,12 @@ impl Ble { /// Bring up the radio and the BLE transport, advertising as `name`. pub(crate) async fn bring_up(name: &str, endpoint_id: EndpointId) -> BleResult { - // `with_config` rather than `new()`: `Central::new()` leaves `connect_timeout` - // unset, and an unset timeout means the backend awaits a connect forever without - // ever cancelling it. Dialling an address a peer has stopped advertising -- which - // happens routinely, since peers rotate BLE addresses -- would then wedge that - // connect for the life of the process. The configured path cancels it and reports - // an error instead. - // - // Independent: each opens its own session and resolves the default adapter. let (central, peripheral) = tokio::try_join!( Central::with_config(CentralConfig::default()), Peripheral::new() )?; let (central, peripheral) = (Arc::new(central), Arc::new(peripheral)); - // Fail fast when the adapter is off, which on a desktop is the common case. - // `build()` below would otherwise spend `wait_ready`'s two 5s timeouts back to - // back discovering the same thing, and the caller holds the pairing-service lock - // throughout, so that time is paid by the pairing UI. if !central.is_powered().await? { return Err(BleError::AdapterOff); } @@ -141,9 +94,7 @@ pub(crate) async fn bring_up(name: &str, endpoint_id: EndpointId) -> BleResult BleResult Duration { + match failures { + 0 | 1 => Duration::from_secs(4), + 2 => Duration::from_secs(16), + _ => IDENTITY_RETRY_AFTER, + } +} + +/// Where a device's identity read has got to. +enum Attempt { + /// A read is running: the prefix it is asking about, and a handle to stop it. + Reading { + prefix: KeyPrefix, + abort: tokio::task::AbortHandle, + failures: u32, + }, + /// Nothing running; the last read ended at `since`. + Idle { since: Instant, failures: u32 }, +} + +impl Attempt { + fn failures(&self) -> u32 { + match self { + Self::Reading { failures, .. } | Self::Idle { failures, .. } => *failures, + } + } + + /// Stop reading if this device was being asked about `prefix`, because another + /// address already answered for it. + fn give_up_on(&mut self, prefix: &KeyPrefix) -> bool { + let failures = match self { + Self::Reading { + prefix: reading, + abort, + failures, + } if reading == prefix => { + abort.abort(); + *failures + }, + _ => return false, + }; + // Not counted as a failure: nothing is wrong with this address, the question + // just stopped being worth asking. + *self = Self::Idle { + since: Instant::now(), + failures, + }; + true + } + + /// A read ended without an identity. Success needs no bookkeeping: the caller + /// memoizes the prefix and stops asking. + fn failed(&mut self) { + let failures = self.failures().saturating_add(1); + *self = Self::Idle { + since: Instant::now(), + failures, + }; + } + + /// Leave this device alone: a read is in flight, or the last one failed too + /// recently to try again. + fn still_waiting(&self) -> bool { + match self { + Self::Reading { .. } => true, + Self::Idle { since, failures } => since.elapsed() < identity_retry_delay(*failures), + } + } + + /// Worth keeping in the map. An in-flight read always is. + fn still_tracked(&self) -> bool { + match self { + Self::Reading { .. } => true, + Self::Idle { since, .. } => since.elapsed() < IDENTITY_RETRY_AFTER, + } + } +} + +/// The outcome of one identity read, reported back to the discovery loop. +enum ReadDone { + Resolved(Resolution), + Failed(DeviceId), +} + /// A completed identity read, carrying the advertisement that triggered it. -/// -/// The rssi travels with the read rather than being looked up when it finishes, -/// because a success has to be announced straight away: the connect the read needs -/// makes the peer stop advertising, so the sighting we would otherwise announce from -/// may never come. The name comes back from the peer itself. struct Resolution { prefix: KeyPrefix, id: EndpointId, @@ -171,9 +202,6 @@ struct Resolution { /// Record the sighting and hand the peer upward. Returns false once the consumer is /// gone, which is the loop's signal to stop. -/// -/// The only place a `DiscoveredPeer` is built and the only place `seen` is written: -/// both announce paths go through here so the two cannot drift. fn announce( tx: &mpsc::UnboundedSender, seen: &mut HashMap, @@ -181,8 +209,6 @@ fn announce( name: Option, rssi: Option, ) -> bool { - // `Found` goes out on every sighting -- `State::on_discovery` absorbs the repeats - // -- but only the first is logged, so the line means "newly visible". if seen.insert(id, Instant::now()).is_none() { info!("[BLE] announcing {id} as {name:?}"); } @@ -202,8 +228,6 @@ fn announce( /// Render a key prefix as lowercase hex, so a log line can be matched against the /// `EndpointId` it resolves to: `EndpointId`'s `Display` is hex of all 32 bytes, so /// the prefix appears verbatim at the head of it. -/// -/// Formats into the log writer rather than building a `String`. fn hex_prefix(prefix: &KeyPrefix) -> impl std::fmt::Display + '_ { struct Hex<'a>(&'a KeyPrefix); impl std::fmt::Display for Hex<'_> { @@ -215,58 +239,57 @@ fn hex_prefix(prefix: &KeyPrefix) -> impl std::fmt::Display + '_ { } /// Read a device's `EndpointId` over GATT, off the discovery loop. -/// -/// Spawned rather than awaited inline: this connects to the device, and a peer that -/// has just gone out of range would otherwise stall discovery for everything else. fn spawn_identity_read( transport: Arc, - resolved_tx: mpsc::UnboundedSender, + done_tx: mpsc::UnboundedSender, device_id: DeviceId, prefix: KeyPrefix, rssi: Option, -) { +) -> tokio::task::AbortHandle { tokio::spawn(async move { let read = time::timeout(IDENTITY_READ_TIMEOUT, transport.read_identity(&device_id)); - match read.await { + let done = match read.await { Ok(Ok(Some(peer))) => { info!( "[BLE] {device_id:?} identified as {} named {:?}", peer.endpoint_id, peer.name ); - let _ = resolved_tx.send(Resolution { + ReadDone::Resolved(Resolution { prefix, id: peer.endpoint_id, name: peer.name, rssi, - }); + }) }, // IDENTITY was not a 32-byte key, so there is nothing to dial. - Ok(Ok(None)) => info!("[BLE] {device_id:?} published no usable identity"), - Ok(Err(err)) => warn!("[BLE] identity read failed for {device_id:?}: {err}"), - Err(_) => warn!("[BLE] identity read timed out for {device_id:?}"), - } - }); + Ok(Ok(None)) => { + info!("[BLE] {device_id:?} published no usable identity"); + ReadDone::Failed(device_id) + }, + Ok(Err(err)) => { + warn!("[BLE] identity read failed for {device_id:?}: {err}"); + ReadDone::Failed(device_id) + }, + Err(_) => { + warn!("[BLE] identity read timed out for {device_id:?}"); + ReadDone::Failed(device_id) + }, + }; + let _ = done_tx.send(done); + }) + .abort_handle() } /// Subscribe to Bluetooth discovery. -/// -/// Runs the mapping in a task rather than as a stream adapter because resolving a -/// key prefix needs `State`, which is behind an async lock. pub(crate) async fn subscribe( central: Arc, transport: Arc, state: SharedState, ) -> BoxStream { let mut events = central.events(); - // Unbounded on purpose: blocking this task drops advertisements, because blew's - // event stream discards `Lagged` and the broadcast behind it holds only 256 -- at - // the hundreds of events a second a real scan delivers, a stalled consumer loses - // sightings outright. let (tx, rx) = mpsc::unbounded_channel(); - // Identity reads report back here rather than mutating shared state, so the loop - // below stays the single owner of every map it consults. - let (resolved_tx, mut resolved_rx) = mpsc::unbounded_channel::(); + let (resolved_tx, mut resolved_rx) = mpsc::unbounded_channel::(); tokio::spawn(async move { // When each announced peer was last sighted, keyed by the endpoint we @@ -278,13 +301,7 @@ pub(crate) async fn subscribe( // alone. Facts about a key rather than about a sighting, so only ever added to. let mut resolved: HashMap)> = HashMap::new(); // When we last asked each device who it is, whether or not it answered. - // - // Keyed by device rather than by key prefix, because what we dial is an - // address. A peer rotates its BLE address, and a connect to an address it has - // stopped advertising does not fail -- it queues until the timeout. Gating by - // prefix would let the first address we tried lock out the live one for the - // whole backoff, so a peer that moved mid-read could never be resolved. - let mut attempted: HashMap = HashMap::new(); + let mut attempted: HashMap = HashMap::new(); // Non-Beaver sightings, counted rather than logged: in a busy RF environment // they are the overwhelming majority, and the count is all that is needed to // show the scan is alive. @@ -294,11 +311,37 @@ pub(crate) async fn subscribe( loop { tokio::select! { - resolution = resolved_rx.recv() => { - let Some(r) = resolution else { return }; - resolved.insert(r.prefix, (r.id, Some(r.name.clone()))); - if !announce(&tx, &mut seen, r.id, Some(r.name), r.rssi) { - return; + done = resolved_rx.recv() => { + match done { + None => return, + Some(ReadDone::Resolved(r)) => { + resolved.insert(r.prefix, (r.id, Some(r.name.clone()))); + // Every other address advertising this prefix is now + // answered for; see `Attempt::give_up_on`. + let mut dropped = 0usize; + for attempt in attempted.values_mut() { + if attempt.give_up_on(&r.prefix) { + dropped += 1; + } + } + if dropped > 0 { + info!( + "[BLE] stopped {dropped} redundant identity read(s) for {}", + hex_prefix(&r.prefix) + ); + } + if !announce(&tx, &mut seen, r.id, Some(r.name), r.rssi) { + return; + } + }, + // Clears the in-flight mark and starts the retry clock, so the + // next sighting can try again rather than waiting out a gate + // sized for a read that is no longer running. + Some(ReadDone::Failed(device_id)) => { + if let Some(attempt) = attempted.get_mut(&device_id) { + attempt.failed(); + } + }, } } event = events.next() => { @@ -318,12 +361,10 @@ pub(crate) async fn subscribe( continue; }; // While backing off from this device, skip the lookups below as - // well as the read: they walk the transport's peer table and take - // the `State` lock, and neither is likely to have changed since a - // sighting moments ago. Blocking here costs advertisements. + // well as the read. let backing_off = attempted .get(&device.id) - .is_some_and(|at| at.elapsed() < IDENTITY_RETRY_AFTER); + .is_some_and(Attempt::still_waiting); // Deliberately not `device.name`: over BLE that is whatever the // central chose to report, which on Linux is often the peer's // operating-system name rather than the one it published. The @@ -353,33 +394,35 @@ pub(crate) async fn subscribe( device.rssi, hex_prefix(&prefix) ); - attempted.insert(device.id.clone(), Instant::now()); - spawn_identity_read( + let failures = + attempted.get(&device.id).map_or(0, Attempt::failures); + let abort = spawn_identity_read( Arc::clone(&transport), resolved_tx.clone(), - device.id, + device.id.clone(), prefix, device.rssi, ); + attempted.insert( + device.id, + Attempt::Reading { + prefix, + abort, + failures, + }, + ); } continue; }; - // Memoize however we got here, so later sightings of this peer are - // a map hit rather than a `snapshot_peers` walk and a `State` lock. resolved.insert(prefix, (id, name.clone())); if !announce(&tx, &mut seen, id, name, device.rssi) { return; } } - // Skipped while there is nothing to age, so an idle radio does not - // wake this task every few seconds for the life of the process. + // Skipped while there is nothing to age. _ = sweep.tick(), if !seen.is_empty() || !attempted.is_empty() || others_seen > 0 => { - // Only while there is nothing to show, which is when "why do I - // see no peers?" is the live question, and throttled because this - // cannot be a `debug!`. Distinguishes a dead scan (zero non-iroh - // devices) from a working one that has found nothing of ours. sweeps = sweeps.wrapping_add(1); - if seen.is_empty() && sweeps % HEALTH_EVERY_SWEEPS == 0 { + if seen.is_empty() && sweeps.is_multiple_of(HEALTH_EVERY_SWEEPS) { info!( "[BLE] scan alive: 0 announced, {} in identity backoff, {} non-iroh devices in the last {}s", attempted.len(), @@ -391,11 +434,21 @@ pub(crate) async fn subscribe( // An elapsed attempt and an absent one mean the same thing, so // dropping these is behaviour-neutral and keeps the map bounded. - attempted.retain(|_, at| at.elapsed() < IDENTITY_RETRY_AFTER); + // Dropping an entry forgets its failure count, which is right: + // a peer unseen for this long is effectively new again. + attempted.retain(|_, attempt| attempt.still_tracked()); + // One clock for the whole sweep, so every peer is aged against the + // same instant. + let now = Instant::now(); let mut lost = Vec::new(); seen.retain(|id, last| { - let alive = last.elapsed() <= SIGHTING_TTL; + // Silence from a peer we hold a pipe to is not absence (see + // `SIGHTING_TTL`), and it is asymmetric. + if transport.is_endpoint_linked(id) { + *last = now; + } + let alive = now.saturating_duration_since(*last) <= SIGHTING_TTL; if !alive { lost.push(*id); } @@ -423,12 +476,6 @@ pub(crate) async fn subscribe( /// Turn an advertised key prefix into a full `EndpointId`, without touching the /// radio. -/// -/// The two lookups that are not a plain map hit: the transport's own verified set -/// (populated once a connection has completed), then a peer we have on file whose key -/// starts with this prefix -- a restored pairing, say. The caller checks its own -/// prefix cache first, so reaching here means that already missed; a miss here means -/// the caller has to go and ask the device. async fn resolve( transport: &BleTransport, state: &SharedState, diff --git a/crates/beaver_p2p/src/discovery/mod.rs b/crates/beaver_p2p/src/discovery/mod.rs index c1dec40..eab1d76 100644 --- a/crates/beaver_p2p/src/discovery/mod.rs +++ b/crates/beaver_p2p/src/discovery/mod.rs @@ -1,13 +1,6 @@ /* SPDX Id: AGPL-3.0-or-later */ //! Peer discovery, abstracted away from any one transport. -//! -//! Discovery used to be mDNS and nothing else, and -//! `iroh_mdns_address_lookup::DiscoveryEvent` leaked all the way up into the -//! constellation. The types here are Beaver's own and `pub(crate)`, so a second -//! source (Bluetooth) can be added without either crate's type reaching past -//! `State`, which resolves raw events into the peer transitions the rest of the -//! system cares about. use iroh::{EndpointAddr, EndpointId}; use n0_future::IterExt; @@ -18,18 +11,12 @@ pub(crate) mod ble; pub(crate) mod mdns; /// How much a name can be trusted to be the one Beaver asked for. -/// -/// Ordered least to most trustworthy; `Ord` derives from declaration order, so a -/// higher variant must never be replaced by a lower one. #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] pub(crate) enum NameSource { - /// Read from a BLE advertisement. On some platforms this is the operating - /// system's own device name rather than anything Beaver set: BlueZ prefers a - /// dual-mode peer's Classic EIR name, and the name may not fit in the - /// advertisement alongside the service UUID in the first place. - Advertised, - /// A field Beaver populated itself -- the mDNS TXT record. Unauthenticated, but - /// it is our name and no OS can substitute for it. + /// Nobody has supplied a name yet, so anything else outranks this. Must stay + /// first: the ordering is the trust ranking. + Unknown, + /// A field Beaver populated itself like the mDNS TXT record. Reported, /// Sent in the pairing handshake, over a connection whose TLS handshake has /// already bound the peer's key. Attributable to the peer we are pairing with. @@ -37,9 +24,6 @@ pub(crate) enum NameSource { } /// Which transport surfaced a peer. -/// -/// A peer can be visible on more than one at a time, so this is tracked as a set -/// per peer rather than a single value. #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] pub(crate) enum DiscoveryTransport { /// Found on the local network. @@ -49,33 +33,16 @@ pub(crate) enum DiscoveryTransport { Ble, } -impl DiscoveryTransport { - /// How far a name from this transport can be trusted. - pub(crate) fn name_source(self) -> NameSource { - match self { - Self::Mdns => NameSource::Reported, - #[cfg(feature = "ble")] - // Read from the peer's NAME characteristic over GATT, never from the - // advertisement -- so it is Beaver's own name, same as the mDNS field. - Self::Ble => NameSource::Reported, - } - } -} - /// A peer as reported by one discovery source. #[derive(Clone, Debug)] pub(crate) struct DiscoveredPeer { pub id: EndpointId, /// The peer's self-reported display name. /// - /// `None` when the source cannot carry one, which must not be confused with - /// "the peer is called nothing": a source that reports `None` should never - /// overwrite a name another source already supplied. + /// `None` when the source cannot carry one. pub name: Option, pub addr: EndpointAddr, - /// Received signal strength in dBm, when the source measures it. `None` for - /// mDNS, which has no equivalent -- so an absent value means "not measurable - /// here", not "far away". + /// Received signal strength in dBm, when the source measures it. pub rssi: Option, pub transport: DiscoveryTransport, } @@ -85,8 +52,7 @@ pub(crate) struct DiscoveredPeer { pub(crate) enum DiscoveryEvent { /// A peer appeared, or its advertised data changed. Found(DiscoveredPeer), - /// A peer is no longer visible *on this transport*. It may still be visible on - /// another, which is why the transport is part of the event. + /// A peer is no longer visible on this transport. Lost { id: EndpointId, transport: DiscoveryTransport, @@ -95,10 +61,6 @@ pub(crate) enum DiscoveryEvent { /// Fan several sources into a single stream, so there is one consumer task /// regardless of how many transports are enabled. -/// -/// Each source is a free function returning its own stream rather than an impl of -/// a shared trait: they have genuinely different setup (mDNS cannot fail, a radio -/// can), and nothing holds a source after it has been subscribed. pub(crate) fn merge(streams: Vec>) -> BoxStream { Box::pin(streams.merge()) } diff --git a/crates/beaver_p2p/src/lib.rs b/crates/beaver_p2p/src/lib.rs index 927eb65..c9af6ba 100644 --- a/crates/beaver_p2p/src/lib.rs +++ b/crates/beaver_p2p/src/lib.rs @@ -36,13 +36,6 @@ use crate::state::{EndpointDescription, PairingCommand, SharedState, State}; pub use crate::state::{EndpointStatus, PeerEvent}; /// How long to wait for [`Endpoint::connect`] before giving up. -/// -/// iroh only resolves a pending connect once a path is found or the address-lookup -/// stream finishes. mDNS's stream does finish, so today an unreachable peer fails on -/// its own and this bound is belt-and-braces. A BLE address lookup's stream never -/// finishes by design, so once a second transport is wired in an unreachable peer -/// would park the caller forever - and `broadcast_message` walks peers serially, so a -/// single dead peer would stall every other one behind it. const CONNECT_TIMEOUT: Duration = Duration::from_secs(15); /// Why a bounded connect did not produce a connection. Flattened into whichever @@ -83,10 +76,6 @@ async fn connect_bounded( } /// Evict every guest and announce that guest mode ended. -/// -/// Shared by [`PairingManager::disable_guest_mode`] and the expiry task so the two -/// paths cannot drift: an expiry that forgot to close connections, or a disable that -/// forgot to notify, would both be silent bugs. fn end_guest_mode(state: &mut State) { let removed = state.disable_guest_mode(); for peer_id in &removed { @@ -125,7 +114,7 @@ pub enum PairingError { } /// The outcome of a pairing request. -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum PairingResult { /// The remote peer accepted the pairing request. Accepted, @@ -164,9 +153,7 @@ pub struct PairingManagerInner { /// mDNS-only, and a BLE peer would otherwise have no way to learn it. local_name: String, - /// The Bluetooth radio handles, kept only so `stop()` can switch them off. The - /// transport holds its own `Arc`s on these and offers no shutdown, so if we drop - /// ours the scan and the advertisement outlive the service. + /// The Bluetooth radio handles, kept only so `stop()` can switch them off. #[cfg(feature = "ble")] ble: Option, @@ -186,9 +173,7 @@ impl PairingManager { let user_data: UserData = name.to_owned().try_into().unwrap(); // Bluetooth is brought up before the endpoint: its `build()` is async and - // yields the hook and the custom transport the builder needs. A radio that - // will not start is not fatal -- log it and carry on IP-only, since the - // common case on a desktop is simply that Bluetooth is off. + // yields the hook and the custom transport the builder needs. #[cfg(feature = "ble")] let ble = match ble::bring_up(name, secret_key.public()).await { Ok(ble) => Some(ble), @@ -202,9 +187,6 @@ impl PairingManager { .secret_key(secret_key) .user_data_for_address_lookup(user_data); - // Deliberately no `clear_ip_transports()`: IP + mDNS coexisting with BLE is - // the point. Deliberately no `PathSelector` either -- the default biased-RTT - // selector already prefers IP when both paths exist. #[cfg(feature = "ble")] let builder = match &ble { Some(ble) => builder @@ -229,7 +211,6 @@ impl PairingManager { // One task drains a merged stream of every discovery source, so adding a // second transport does not add a second consumer of `State`. - // `mut` is only needed when the BLE branch below is compiled in. #[cfg_attr(not(feature = "ble"), allow(unused_mut))] let mut streams = vec![mdns::subscribe(mdns).await]; #[cfg(feature = "ble")] @@ -277,10 +258,7 @@ impl PairingManager { // Send a pairing request to an endpoint. pub async fn request_pairing(&self, id: &EndpointId) -> Result { info!("Pairing with {id}"); - let Some(ref inner) = self.inner else { - error!("Not initialized"); - return Err(PairingError::InvalidState); - }; + let inner = self.inner()?; let addr = { let state = inner.state.lock().await; @@ -302,24 +280,15 @@ impl PairingManager { } } - let result = self.do_pairing_handshake(inner, addr).await; + let request = PairingCommand::Request { + name: inner.local_name.clone(), + }; + let result = self.do_pairing_handshake(inner, addr, request).await; // Always clean up the requested state, regardless of success or failure. let mut state = inner.state.lock().await; state.remove_pairing_requested(id); - - match &result { - Ok(PairingResult::Accepted) => { - state.notify(PeerEvent::PairingAccepted(*id)); - state.set_status(id, EndpointStatus::PairedConnected); - }, - Ok(PairingResult::Rejected) => { - state.notify(PeerEvent::PairingRejected(*id)); - }, - Err(_) => { - state.notify(PeerEvent::PairingFailed(*id)); - }, - } + state.settle_pairing(id, result.as_ref().ok().copied()); result } @@ -332,9 +301,7 @@ impl PairingManager { pin: &str, ) -> Result { info!("Guest pairing with {id}"); - let Some(ref inner) = self.inner else { - return Err(PairingError::NotInitialized); - }; + let inner = self.inner()?; let addr = { let state = inner.state.lock().await; @@ -344,40 +311,12 @@ impl PairingManager { remote.addr() }; - let connection = connect_bounded(inner.router.endpoint(), addr, PAIRING_ALPN).await?; - let (mut sender, mut receiver) = connection.open_bi().await?; - let request = PairingCommand::GuestRequest { pin: pin.to_string(), name: inner.local_name.clone(), }; - PostcardPacket::send(request, &mut sender).await?; - - let command: PairingCommand = PostcardPacket::recv(&mut receiver).await?; - let result = match command { - PairingCommand::Accept => PairingResult::Accepted, - PairingCommand::Reject => PairingResult::Rejected, - _ => { - error!("Unexpected response to guest pairing: {command:?}"); - return Err(PairingError::InvalidState); - }, - }; - - PostcardPacket::send(PairingCommand::Ack, &mut sender).await?; - sender.finish()?; - let _ = sender.stopped().await; - - let mut state = inner.state.lock().await; - match &result { - PairingResult::Accepted => { - state.notify(PeerEvent::PairingAccepted(*id)); - state.set_status(id, EndpointStatus::PairedConnected); - }, - PairingResult::Rejected => { - state.notify(PeerEvent::PairingRejected(*id)); - }, - } - + let result = self.do_pairing_handshake(inner, addr, request).await?; + inner.state.lock().await.settle_pairing(id, Some(result)); Ok(result) } @@ -387,15 +326,14 @@ impl PairingManager { &self, inner: &PairingManagerInner, addr: EndpointAddr, + request: PairingCommand, ) -> Result { - let connection = connect_bounded(inner.router.endpoint(), addr, PAIRING_ALPN).await?; + let connection = + connect_bounded(inner.router.endpoint(), addr.clone(), PAIRING_ALPN).await?; // Send the request. let (mut sender, mut receiver) = connection.open_bi().await?; - let request = PairingCommand::Request { - name: inner.local_name.clone(), - }; PostcardPacket::send(request, &mut sender).await?; // Wait for accept or reject @@ -418,6 +356,12 @@ impl PairingManager { sender.finish()?; let _ = sender.stopped().await; + // Here rather than in the caller: the caller sets `PairedConnected` after this + // returns, by which point `connection` has already been dropped. + if result == PairingResult::Accepted { + self.prime_message_connection(inner, &addr).await; + } + Ok(result) } @@ -428,10 +372,7 @@ impl PairingManager { command: PairingCommand, ) -> Result<(), PairingError> { info!("Sending pairing response: {command:?}"); - let Some(ref inner) = self.inner else { - error!("Not initialized"); - return Err(PairingError::InvalidState); - }; + let inner = self.inner()?; if inner.state.lock().await.by_id(from).is_none() { error!("Can't send pairing response to unknown remote {from}"); @@ -464,19 +405,8 @@ impl PairingManager { // Accept a pairing request pub async fn accept_pairing(&self, from: &EndpointId) -> Result<(), PairingError> { - let Some(ref inner) = self.inner else { - error!("Not initialized"); - return Err(PairingError::InvalidState); - }; - - let res = self - .send_pairing_response(from, PairingCommand::Accept) - .await; - let mut state = inner.state.lock().await; - res.map(|_| { - // Add the endpoint to the set of pending acks if successfully sending. - state.set_pending_ack(from); - }) + self.send_pairing_response(from, PairingCommand::Accept) + .await } // Reject a pairing request @@ -485,10 +415,18 @@ impl PairingManager { .await } + /// The running service, or `NotInitialized` if `start` has not run yet, or `stop` + /// has already run. + fn inner(&self) -> Result<&PairingManagerInner, PairingError> { + self.inner.as_ref().ok_or_else(|| { + error!("[P2P] Pairing manager not initialized"); + PairingError::NotInitialized + }) + } + // Get the list of current known peers. pub async fn peers(&self) -> Vec { - let Some(ref inner) = self.inner else { - error!("Not initialized"); + let Ok(inner) = self.inner() else { return vec![]; }; @@ -506,10 +444,7 @@ impl PairingManager { // Caches the QUIC connection per peer and opens a fresh uni stream per message. // If sending fails (stale connection), retries once with a fresh connection. pub async fn send_message(&self, to: &EndpointId, message: &[u8]) -> Result<(), MessageError> { - let Some(ref inner) = self.inner else { - error!("Not initialized"); - return Err(MessageError::InvalidState); - }; + let inner = self.inner().map_err(|_| MessageError::InvalidState)?; let addr = { let state = inner.state.lock().await; @@ -540,6 +475,20 @@ impl PairingManager { Self::try_send(&connection, message).await } + /// Open the message connection while the pairing connection is still up. + /// + /// Pairing and messaging are separate connections to the same peer, so closing the + /// first before opening the second leaves the peer with none in between. Over BLE + /// that gap is expensive. + async fn prime_message_connection(&self, inner: &PairingManagerInner, addr: &EndpointAddr) { + if let Err(err) = self.get_or_connect(inner, &addr.id, addr).await { + info!( + "[P2P] could not pre-open the message connection to {}: {err}", + addr.id + ); + } + } + async fn get_or_connect( &self, inner: &PairingManagerInner, @@ -570,8 +519,7 @@ impl PairingManager { } pub async fn stop(&mut self) { - let Some(ref inner) = self.inner else { - error!("Not initialized"); + let Ok(inner) = self.inner() else { return; }; @@ -593,8 +541,7 @@ impl PairingManager { } pub async fn set_status(&self, endpoint: &EndpointId, status: EndpointStatus) { - let Some(ref inner) = self.inner else { - error!("Not initialized"); + let Ok(inner) = self.inner() else { return; }; @@ -605,8 +552,7 @@ impl PairingManager { /// Creates the endpoint entry with PairedDisconnected status and an empty address. /// When the peer is discovered via mDNS, its address will be updated. pub async fn register_paired_peer(&self, id: &EndpointId, name: &str) { - let Some(ref inner) = self.inner else { - error!("Not initialized"); + let Ok(inner) = self.inner() else { return; }; @@ -621,8 +567,7 @@ impl PairingManager { /// Remove a paired peer from the runtime state. pub async fn remove_paired_peer(&self, id: &EndpointId) { - let Some(ref inner) = self.inner else { - error!("Not initialized"); + let Ok(inner) = self.inner() else { return; }; @@ -634,14 +579,9 @@ impl PairingManager { /// Enable guest mode with a random 6-digit PIN. /// Returns the PIN to display to the user. /// - /// The session is torn down automatically after `duration_minutes`. Before that - /// timer existed, `expires_at` was only consulted when validating an incoming - /// PIN, so an expired session kept its guests connected indefinitely and hosts - /// kept displaying a PIN that no longer worked. + /// The session is torn down automatically after `duration_minutes`. pub async fn enable_guest_mode(&self, duration_minutes: u32) -> Result { - let Some(ref inner) = self.inner else { - return Err(PairingError::NotInitialized); - }; + let inner = self.inner()?; let pin: String = { let mut rng = rand::rng(); @@ -678,19 +618,13 @@ impl PairingManager { /// Disable guest mode and remove all guest peers. pub async fn disable_guest_mode(&self) -> Result<(), PairingError> { - let Some(ref inner) = self.inner else { - return Err(PairingError::NotInitialized); - }; + let inner = self.inner()?; end_guest_mode(&mut *inner.state.lock().await); Ok(()) } /// The PIN of the running guest session, if there is one. - /// - /// Lets a surface that did not enable guest mode itself -- a media center page - /// that reloaded, say -- recover the current state rather than waiting for the - /// next change event. pub async fn guest_pin(&self) -> Option { let inner = self.inner.as_ref()?; inner.state.lock().await.guest_pin().map(str::to_owned) diff --git a/crates/beaver_p2p/src/packet.rs b/crates/beaver_p2p/src/packet.rs index 36f34a6..7d16cfb 100644 --- a/crates/beaver_p2p/src/packet.rs +++ b/crates/beaver_p2p/src/packet.rs @@ -15,8 +15,13 @@ pub enum PacketError { Postcard(#[from] postcard::Error), #[error("I/O error")] Io(#[from] std::io::Error), + #[error("Declared payload of {0} bytes is over the limit")] + TooLarge(u32), } +/// Largest payload one framed packet may declare. +const MAX_PAYLOAD: u32 = 2 * 1024 * 1024 * 1024; + pub(crate) struct BasePacket {} impl BasePacket { @@ -31,9 +36,12 @@ impl BasePacket { let mut len_buf: [u8; 4] = [0; 4]; stream.read_exact(&mut len_buf).await?; let len = u32::from_be_bytes(len_buf); - let mut payload = Vec::with_capacity(len as _); - let _remaining = payload.spare_capacity_mut(); - unsafe { payload.set_len(len as _) }; + if len > MAX_PAYLOAD { + return Err(PacketError::TooLarge(len)); + } + // Zeroed, not `set_len` over uninitialised capacity: that handed `read_exact` + // uninit bytes to write into, and left them readable if it failed partway. + let mut payload = vec![0u8; len as usize]; stream.read_exact(&mut payload).await?; Ok(payload) } diff --git a/crates/beaver_p2p/src/pairing_protocol.rs b/crates/beaver_p2p/src/pairing_protocol.rs index 06efe58..6fbd3d9 100644 --- a/crates/beaver_p2p/src/pairing_protocol.rs +++ b/crates/beaver_p2p/src/pairing_protocol.rs @@ -1,13 +1,12 @@ /* SPDX Id: AGPL-3.0-or-later */ -use std::collections::BTreeSet; - -use iroh::EndpointAddr; use iroh::endpoint::Connection; use iroh::protocol::{AcceptError, ProtocolHandler}; +use iroh::{EndpointAddr, EndpointId}; use log::{error, info}; use tokio::sync::mpsc::channel as tokio_channel; +use crate::PairingResult; use crate::packet::PostcardPacket; pub use crate::state::PeerEvent; use crate::state::{EndpointStatus, PairingCommand, SharedState}; @@ -21,6 +20,20 @@ impl PairingProtocol { pub(crate) fn new(state: SharedState) -> Self { Self { state } } + + /// Tell the rest of the system this handshake did not complete. + async fn fail_pairing(&self, remote_id: EndpointId) { + self.state + .lock() + .await + .notify(PeerEvent::PairingFailed(remote_id)); + } +} + +/// The handler's only error shape: the peer did something the protocol does not allow, +/// so the connection is refused. +fn protocol_error(message: impl std::fmt::Display) -> AcceptError { + AcceptError::from(n0_error::AnyError::from(message.to_string())) } /// Pairing handshake: @@ -41,17 +54,11 @@ impl ProtocolHandler for PairingProtocol { let paths = connection.paths(); let Some(path_info) = paths.into_iter().find(|path| path.is_selected()) else { error!("Failed to get selected path for {connection:?}"); - return Err(AcceptError::from(n0_error::AnyError::from(format!( + return Err(protocol_error(format!( "Failed to get selected path for {connection:?}" - )))); - }; - let remote_addr = path_info.remote_addr(); - let mut addrs = BTreeSet::new(); - addrs.insert(remote_addr.clone()); - let addr = EndpointAddr { - id: remote_id, - addrs, + ))); }; + let addr = EndpointAddr::from_parts(remote_id, [path_info.remote_addr().clone()]); let (mut sender, mut receiver) = connection.accept_bi().await?; @@ -66,42 +73,18 @@ impl ProtocolHandler for PairingProtocol { state.notify(PeerEvent::PairingRequest(remote_id, name)); }, PairingCommand::GuestRequest { pin, name } => { - let mut state = self.state.lock().await; - state.introduce(remote_id, &addr, &name); - if state.validate_guest_pin(&pin) { - state.set_status(&remote_id, EndpointStatus::GuestConnected); - state.add_guest_peer(&remote_id); - drop(state); - - PostcardPacket::send(PairingCommand::Accept, &mut sender) - .await - .expect("Failed to send guest accept"); - - match PostcardPacket::recv(&mut receiver).await { - Ok(PairingCommand::Ack) => { - self.state - .lock() - .await - .notify(PeerEvent::GuestAccepted(remote_id)); - info!("[P2P] Guest paired: {remote_id}"); - }, - Ok(cmd) => { - error!("Unexpected command after guest accept: {cmd:?}"); - self.state - .lock() - .await - .notify(PeerEvent::PairingFailed(remote_id)); - }, - Err(e) => { - error!("Failed to read guest Ack: {e}"); - self.state - .lock() - .await - .notify(PeerEvent::PairingFailed(remote_id)); - }, + let accepted = { + let mut state = self.state.lock().await; + state.introduce(remote_id, &addr, &name); + let accepted = state.validate_guest_pin(&pin); + if accepted { + state.set_status(&remote_id, EndpointStatus::GuestConnected); + state.add_guest_peer(&remote_id); } - return Ok(()); - } else { + accepted + }; + + if !accepted { info!("[P2P] Guest pairing rejected: invalid PIN from {remote_id}"); PostcardPacket::send(PairingCommand::Reject, &mut sender) .await @@ -112,12 +95,33 @@ impl ProtocolHandler for PairingProtocol { .notify(PeerEvent::PairingRejected(remote_id)); return Ok(()); } + + PostcardPacket::send(PairingCommand::Accept, &mut sender) + .await + .expect("Failed to send guest accept"); + + match PostcardPacket::recv(&mut receiver).await { + Ok(PairingCommand::Ack) => { + self.state + .lock() + .await + .notify(PeerEvent::GuestAccepted(remote_id)); + info!("[P2P] Guest paired: {remote_id}"); + }, + Ok(cmd) => { + error!("Unexpected command after guest accept: {cmd:?}"); + self.fail_pairing(remote_id).await; + }, + Err(e) => { + error!("Failed to read guest Ack: {e}"); + self.fail_pairing(remote_id).await; + }, + } + return Ok(()); }, _ => { error!("Unexpected command: {command:?}"); - return Err(AcceptError::from(n0_error::AnyError::from(format!( - "Unexpected command: {command:?}" - )))); + return Err(protocol_error(format!("Unexpected command: {command:?}"))); }, } @@ -145,40 +149,30 @@ impl ProtocolHandler for PairingProtocol { Ok(cmd) => cmd, Err(e) => { error!("Failed to read Ack: {e}"); - self.state - .lock() - .await - .notify(PeerEvent::PairingFailed(remote_id)); + self.fail_pairing(remote_id).await; let _ = ack_sender.send(false).await; - return Err(AcceptError::from(n0_error::AnyError::from(format!( - "Failed to read Ack: {e}" - )))); + return Err(protocol_error(format!("Failed to read Ack: {e}"))); }, }; match command { PairingCommand::Ack => { - { - let mut state = self.state.lock().await; - if accepted { - state.notify(PeerEvent::PairingAccepted(remote_id)); - state.set_status(&remote_id, EndpointStatus::PairedConnected); - } else { - state.notify(PeerEvent::PairingRejected(remote_id)); - } - } - let _ = ack_sender.send(true).await; - }, - _ => { + let outcome = if accepted { + PairingResult::Accepted + } else { + PairingResult::Rejected + }; self.state .lock() .await - .notify(PeerEvent::PairingFailed(remote_id)); + .settle_pairing(&remote_id, Some(outcome)); + let _ = ack_sender.send(true).await; + }, + _ => { + self.fail_pairing(remote_id).await; let _ = ack_sender.send(false).await; error!("Unexpected command: {command:?}"); - return Err(AcceptError::from(n0_error::AnyError::from(format!( - "Unexpected command: {command:?}" - )))); + return Err(protocol_error(format!("Unexpected command: {command:?}"))); }, } diff --git a/crates/beaver_p2p/src/state.rs b/crates/beaver_p2p/src/state.rs index 2e2947c..70890f8 100644 --- a/crates/beaver_p2p/src/state.rs +++ b/crates/beaver_p2p/src/state.rs @@ -12,6 +12,7 @@ use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; use tokio::sync::mpsc::{Receiver as TokioReceiver, Sender as TokioSender}; +use crate::PairingResult; use crate::discovery::{DiscoveredPeer, DiscoveryEvent, DiscoveryTransport, NameSource}; pub(crate) type SharedState = Arc>; @@ -32,8 +33,7 @@ pub enum EndpointStatus { /// What one discovery transport currently reports about a peer. #[derive(Debug)] pub(crate) struct Sighting { - /// Signal strength in dBm, for transports that measure it. `None` means this - /// transport cannot measure it (mDNS), not that the peer is distant. + /// Signal strength in dBm, for transports that measure it. rssi: Option, } @@ -41,22 +41,13 @@ pub(crate) struct Sighting { pub(crate) struct EndpointProxy { name: String, /// Where `name` came from, so a less trustworthy source cannot overwrite a more - /// trustworthy one. Without this a BLE advertisement carrying the peer's - /// *operating system* name silently replaces the name Beaver was told over mDNS, - /// or worse, the one the peer introduced itself with during pairing. + /// trustworthy one. name_source: NameSource, id: EndpointId, addr: EndpointAddr, status: EndpointStatus, /// What each discovery transport currently sees of this peer. A peer is only - /// really gone once this empties, so one transport expiring cannot evict a peer - /// that another still sees. Empty for peers we learned about some other way -- a - /// restored pairing, or an inbound connection -- until a source reports them. - /// - /// Keyed by transport rather than held flat because a measurement belongs to the - /// sighting that produced it: when a transport stops seeing the peer its readings - /// go with it, instead of lingering as a stale value on a peer another transport - /// still sees. + /// really gone once this empties. sightings: HashMap, /// Cached QUIC connection for sending messages. connection: Option, @@ -71,9 +62,7 @@ impl EndpointProxy { ) -> Self { Self { name: name.to_owned(), - // The floor, so nothing is ever locked out. Callers that know better raise - // it: `merge_discovered` and `introduce` both do. - name_source: NameSource::Advertised, + name_source: NameSource::Unknown, id: id.to_owned(), addr, status, @@ -101,19 +90,12 @@ impl EndpointProxy { } /// Fold in what a discovery source just reported. - /// - /// Addresses are unioned rather than replaced: two transports can each know a - /// different way to reach the same peer, and dropping one because the other - /// reported first would lose a working path. - /// - /// The name is taken only when the source supplied one *and* is at least as - /// trustworthy as what we already have. An absent name must not clobber a good - /// one, and neither must a worse-sourced one: a BLE advertisement often carries - /// the peer's operating system name rather than its Beaver name. Equal rank does - /// replace, so a rename still propagates. pub(crate) fn merge_discovered(&mut self, peer: &DiscoveredPeer) { self.addr.addrs.extend(peer.addr.addrs.iter().cloned()); - let source = peer.transport.name_source(); + // Both sources hand us a name Beaver itself published: mDNS from its TXT + // record, BLE from the peer's NAME characteristic (never the advertisement, + // which may carry the operating system's name instead). + let source = NameSource::Reported; // Re-announcements carry the same name every time; skip the reallocation. if let Some(name) = &peer.name && source >= self.name_source && @@ -186,10 +168,6 @@ pub enum PeerEvent { #[derive(Serialize, Deserialize, Debug, PartialEq)] pub(crate) enum PairingCommand { /// Pairing request, carrying the dialer's display name. - /// - /// The name rides along with the request rather than arriving in a separate - /// introduction: it is only ever meaningful as part of one, and folding it in - /// keeps the responder to a single read. Request { name: String, }, @@ -224,19 +202,13 @@ impl From<&EndpointProxy> for EndpointDescription { } } -/// Guest mode state — active when the host is accepting temporary guest pairings. +/// Guest mode state, active when the host is accepting temporary guest pairings. #[derive(Debug)] pub(crate) struct GuestModeState { pub pin: String, pub expires_at: Instant, pub guest_peers: HashSet, /// Distinguishes this session from any earlier one. - /// - /// The expiry task captures the generation it was spawned for and checks it on - /// wake, so a timer left over from a session that was disabled — or replaced by - /// a re-enable — cannot end the session that is currently running. A superseded - /// task sleeps out its deadline and then does nothing, which is cheaper than - /// tracking abort handles for something that fires at most once per session. pub generation: u64, } @@ -252,7 +224,6 @@ pub(crate) struct State { pairing_responders: HashMap, TokioReceiver)>, /// The set of endpoints that we are waiting to receive ack from. - pending_ack: HashSet, /// The sender side of the channel used to receive high level events. sender: Sender, @@ -275,7 +246,6 @@ impl State { endpoints: HashMap::new(), pairing_requested: HashSet::new(), pairing_responders: HashMap::new(), - pending_ack: HashSet::new(), sender, guest_mode: None, guest_generation: 0, @@ -318,7 +288,6 @@ impl State { pub(crate) fn remove_endpoint(&mut self, id: &EndpointId) { self.endpoints.remove(id); self.pairing_requested.remove(id); - self.pending_ack.remove(id); self.pairing_responders.remove(id); self.close_incoming_connection(id); } @@ -331,12 +300,21 @@ impl State { self.pairing_requested.insert(*id); } - pub(crate) fn remove_pairing_requested(&mut self, id: &EndpointId) { - self.pairing_requested.remove(id); + /// Record how a pairing handshake ended: notify, and promote the peer on success. + /// `None` means it failed before either side answered. + pub(crate) fn settle_pairing(&mut self, id: &EndpointId, outcome: Option) { + match outcome { + Some(PairingResult::Accepted) => { + self.notify(PeerEvent::PairingAccepted(*id)); + self.set_status(id, EndpointStatus::PairedConnected); + }, + Some(PairingResult::Rejected) => self.notify(PeerEvent::PairingRejected(*id)), + None => self.notify(PeerEvent::PairingFailed(*id)), + } } - pub(crate) fn set_pending_ack(&mut self, id: &EndpointId) { - self.pending_ack.insert(*id); + pub(crate) fn remove_pairing_requested(&mut self, id: &EndpointId) { + self.pairing_requested.remove(id); } pub(crate) fn set_pairing_responder( @@ -389,9 +367,7 @@ impl State { } /// Record a peer under the name it introduced itself with over an authenticated - /// connection, registering it if this is the first we have heard of it -- over - /// Bluetooth that is the normal case, since an advertisement alone never yields a - /// full key. + /// connection. pub(crate) fn introduce(&mut self, id: EndpointId, addr: &EndpointAddr, name: &str) { match self.endpoints.get_mut(&id) { Some(existing) => { diff --git a/crates/beaver_p2p/tests/pairing.rs b/crates/beaver_p2p/tests/pairing.rs index 37bcc81..9f6fad3 100644 --- a/crates/beaver_p2p/tests/pairing.rs +++ b/crates/beaver_p2p/tests/pairing.rs @@ -8,7 +8,11 @@ use parking_lot::Mutex; type PairingReceiver = Receiver; // Helper: create a named pairing manager. +// +// `try_init` because these tests share one process: `env_logger::init` panics on the +// second call, which failed every test after whichever happened to run first. async fn create_manager(name: &str) -> (PairingManager, PairingReceiver) { + let _ = env_logger::try_init(); let (sender, receiver) = channel(); let secret = iroh_persist::KeyRetriever::new("test") .lenient() @@ -18,48 +22,44 @@ async fn create_manager(name: &str) -> (PairingManager, PairingReceiver) { (manager, receiver) } +/// Spawn a thread that consumes events until `pick` returns a value. +/// +/// Every wait in this file is the same loop over one receiver, matching one variant; +/// this is that loop, once. +fn expect_event( + receiver: Arc>, + mut pick: impl FnMut(PeerEvent) -> Option + Send + 'static, +) -> std::thread::JoinHandle { + std::thread::spawn(move || { + loop { + match receiver.lock().recv() { + Ok(event) => { + if let Some(value) = pick(event) { + return value; + } + }, + Err(err) => panic!("event channel closed before the expected event: {err:?}"), + } + } + }) +} + // Helper: wait for 2 managers to have discovered each other. // Returns the endpoints discovered by each receiver. fn wait_for_discovery( receiver1: Arc>, receiver2: Arc>, ) -> (EndpointId, EndpointId) { - let handle1 = std::thread::spawn(move || { - loop { - match receiver1.lock().recv() { - Ok(event) => match event { - PeerEvent::PeerDiscovered(id, _) => { - return id; - }, - _ => {}, - }, - Err(err) => { - panic!("Should not error! {err:?}"); - }, - } + fn discovered(event: PeerEvent) -> Option { + match event { + PeerEvent::PeerDiscovered(id, _) => Some(id), + _ => None, } - }); - - let handle2 = std::thread::spawn(move || { - loop { - match receiver2.lock().recv() { - Ok(event) => match event { - PeerEvent::PeerDiscovered(id, _) => { - return id; - }, - _ => {}, - }, - Err(err) => { - panic!("Should not error! {err:?}"); - }, - } - } - }); - - let endpoint1 = handle1.join().unwrap(); - let endpoint2 = handle2.join().unwrap(); + } - (endpoint1, endpoint2) + let handle1 = expect_event(receiver1, discovered); + let handle2 = expect_event(receiver2, discovered); + (handle1.join().unwrap(), handle2.join().unwrap()) } // Discover and stop. @@ -112,8 +112,6 @@ async fn discover_and_shutdown() { // Discover, and wait for second peer to disappear #[tokio::test(flavor = "multi_thread")] async fn discover_and_expire() { - env_logger::init(); - let (mut manager1, receiver1) = create_manager("test-1").await; let (mut manager2, receiver2) = create_manager("test-2").await; @@ -158,8 +156,6 @@ async fn discover_and_expire() { // Peer1 requests pairing, peer2 rejects it. #[tokio::test(flavor = "multi_thread")] async fn reject_pairing() { - env_logger::init(); - let (mut manager1, receiver1) = create_manager("test-1").await; let (mut manager2, receiver2) = create_manager("test-2").await; @@ -240,8 +236,6 @@ async fn reject_pairing() { // Peer1 requests pairing, peer2 accepts it. #[tokio::test(flavor = "multi_thread")] async fn accept_pairing() { - env_logger::init(); - let (mut manager1, receiver1) = create_manager("test-1").await; let (mut manager2, receiver2) = create_manager("test-2").await; @@ -323,8 +317,6 @@ async fn accept_pairing() { // Try to send a message to an unpaired peer #[tokio::test(flavor = "multi_thread")] async fn message_to_unpaired() { - env_logger::init(); - let (mut manager1, receiver1) = create_manager("test-1").await; let (mut manager2, receiver2) = create_manager("test-2").await; @@ -350,8 +342,6 @@ async fn message_to_unpaired() { // Exchange messages #[tokio::test(flavor = "multi_thread")] async fn message_to_paired() { - env_logger::init(); - let (mut manager1, receiver1) = create_manager("test-1").await; let (mut manager2, receiver2) = create_manager("test-2").await; diff --git a/ui/system/mediacenter/peers.js b/ui/system/mediacenter/peers.js index c652183..ce7445e 100644 --- a/ui/system/mediacenter/peers.js +++ b/ui/system/mediacenter/peers.js @@ -82,11 +82,20 @@ export function initPeers({ controller }) { return; } - await pairing - .setName("Beaver Media Center") - .catch((e) => console.error("[P2P] setName failed:", e)); - - await pairing.stop().catch((e) => console.error("[P2P] stop failed:", e)); + // The service auto-starts from persisted config, so on every boot after the first + // our name is already set and there is nothing to restart for. Restarting anyway + // is not free: stopping only halts scanning and advertising, leaving the previous + // BLE stack's GATT server and L2CAP listener registered. A peer then reads a PSM + // belonging to a listener whose endpoint is gone, dials it, and waits for an + // answer that cannot come. + const NAME = "Beaver Media Center"; + const local = await pairing.local().catch(() => null); + if (local?.displayName !== NAME) { + await pairing + .setName(NAME) + .catch((e) => console.error("[P2P] setName failed:", e)); + await pairing.stop().catch((e) => console.error("[P2P] stop failed:", e)); + } await pairing.start().catch((e) => console.error("[P2P] start failed:", e)); -- 2.51.2