mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
Merge branch 'master' into next
Carries the BLE platform work up to the feature line: the reframing that
recovers packet boundaries from the FMP length prefix, the L2CAP PSM seam and
its BlueZ implementation, the io.rs / io_linux.rs / io_android.rs split, the
Android backend, the probe backoff and connect-outcome counters, the
stop_scanning seam, and the pool-refusal fix.
Nine conflicts, all in src/transport/ble/mod.rs, and all one collision rather
than nine: `next` removed the pre-handshake pubkey exchange and the cross-probe
tie-breaker in 8162d3c, because XX replaces IK and identity is learned from the
handshake rather than from the transport. Most of what master added since sits
on top of that exchange.
Resolved by taking master's module and re-applying next's removal over it.
`next`'s own BLE delta is purely subtractive, so the two are reconcilable
without inventing anything: what survives is every change that does not need
the exchange, and what goes is everything that does.
Dropped deliberately, not lost:
- The node-identity pool keying, which keyed the pool on the NodeAddr the
exchange learned. There is no identity source for it here, so `next` keeps
address-keyed dedup. The RPA-rotation defect it fixes on master therefore
stands on next, and closing it needs a decision about where BLE peer identity
comes from once the XX handshake owns it. That decision is not made here.
- The inbound handshake concurrency bound, which exists to run the pubkey
exchange off the accept loop. With no exchange there is nothing to run off
it, and ISSUE-2026-0171 already records next as NOT AFFECTED.
- The exchange itself, its tie-breaker, and the announced-address
canonicalisation built on them.
Five BLE counters went with them rather than being left to read zero forever:
pubkey_exchange_failures, tiebreaker_yields, tiebreaker_drops,
duplicate_node_declines and handshakes_aborted. Each counts a mechanism this
branch no longer has, and a metric that can only ever report zero is worse than
no metric.
The two scan tests take next's shape rather than master's: master's relied on
the no-local-pubkey shortcut to reach the neighbour buffer without dialling,
and that shortcut is gone, so they set a connect handler and let the probe
succeed.
Quartet green: rustfmt and clippy clean at -D warnings, 2478 tests passed,
0 failed. Six repo guards exit 0.
This commit is contained in:
@@ -1842,6 +1842,29 @@ with v0.4.x or earlier peers.
|
||||
doing so, which is a denial of new sessions rather than the unbounded memory
|
||||
growth it replaces.
|
||||
|
||||
- One inbound BLE connector can no longer stall every other inbound
|
||||
connection. The accept loop ran the pre-handshake pubkey exchange inline, so
|
||||
a peer that connected an L2CAP channel and then said nothing held the loop
|
||||
for the full 5-second exchange deadline and no other inbound connection was
|
||||
accepted in that window; the maximum-connections argument the loop was given
|
||||
was never used, so the effective concurrency was one. Each inbound connection
|
||||
now runs its handshake in its own task, up to eight in flight, and at that
|
||||
bound the oldest pending handshake is aborted to make room rather than the
|
||||
loop waiting for a slot: a healthy exchange is one round trip, so anything
|
||||
still pending under a flood is overwhelmingly the flooder's, and a genuinely
|
||||
slow peer that is aborted reconnects, which is better than never being
|
||||
accepted at all. Aborted handshakes are counted in the transport's stats as
|
||||
`handshakes_aborted`. The tasks live in a set owned by the accept loop, so
|
||||
stopping the transport stops them too and none can insert into a pool that
|
||||
stop has just drained. Separately, the send half of the pubkey exchange had
|
||||
no deadline at all while the receive half had one, so a peer that stopped
|
||||
draining its channel could park the write forever; it now shares the same
|
||||
5-second deadline, which also covers the outbound connect and scan-probe
|
||||
paths. **What this does not close**: eight simultaneous silent connectors
|
||||
still occupy the whole in-flight budget, and no BlueZ hardware was exercised,
|
||||
so the controller's own concurrent-link limit and accept backlog depth stay
|
||||
unmeasured.
|
||||
|
||||
- An accepted inbound TCP connection no longer holds a slot indefinitely
|
||||
without sending anything. The cap was tested at accept and the pool insert
|
||||
and counter bump followed with no read in between, while the frame reader's
|
||||
|
||||
@@ -43,7 +43,22 @@ fn main() {
|
||||
println!("cargo:rustc-check-cfg=cfg(bluer_available)");
|
||||
let target_os = std::env::var("CARGO_CFG_TARGET_OS").unwrap_or_default();
|
||||
let target_env = std::env::var("CARGO_CFG_TARGET_ENV").unwrap_or_default();
|
||||
if target_os == "linux" && target_env != "musl" {
|
||||
let bluer_available = target_os == "linux" && target_env != "musl";
|
||||
if bluer_available {
|
||||
println!("cargo:rustc-cfg=bluer_available");
|
||||
}
|
||||
|
||||
// Whether the BLE transport is compiled at all.
|
||||
//
|
||||
// This is the set of platforms that have a concrete `BleIo` backend, not
|
||||
// the set that could plausibly have one. A platform listed here with no
|
||||
// backend behind it does not get "BLE, degraded" — it gets an in-memory
|
||||
// transport that starts, reports itself up and never peers, with no error
|
||||
// anywhere. Add a platform here only in the same change that adds its
|
||||
// backend; `transport::ble` carries a compile-time tripwire that refuses
|
||||
// a build where the two disagree.
|
||||
println!("cargo:rustc-check-cfg=cfg(ble_available)");
|
||||
if bluer_available || target_os == "android" {
|
||||
println!("cargo:rustc-cfg=ble_available");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -755,8 +755,15 @@ entries become no-ops. Communicates with BlueZ via D-Bus through the
|
||||
|
||||
**Advertising and scanning.** When `advertise` is enabled, the transport
|
||||
advertises the FIPS service UUID continuously so that nearby nodes can
|
||||
discover and connect via L2CAP. When `scan` is enabled, the transport
|
||||
continuously scans for other FIPS nodes' advertisements. Discovered
|
||||
discover and connect via L2CAP, plus the L2CAP PSM its listener actually
|
||||
bound, as a service-data structure (see `src/transport/ble/psm.rs` for the
|
||||
wire layout and why platforms with OS-assigned PSMs need it). The
|
||||
advertisement carries no device name — alongside the PSM a name no longer
|
||||
fits the 31-byte legacy PDU, so the node shows up in generic Bluetooth
|
||||
scanners as an unnamed device with the FIPS UUID. When `scan` is enabled,
|
||||
the transport continuously scans for other FIPS nodes' advertisements and
|
||||
learns each peer's advertised PSM; a peer that advertises none is dialled
|
||||
at the configured `psm`. Discovered
|
||||
peers are probed immediately (L2CAP connect + pubkey exchange) with a
|
||||
cooldown (`probe_cooldown_secs`) to prevent rapid re-probing of the same
|
||||
address. If two nodes probe each other at the same time (cross-probe),
|
||||
|
||||
@@ -711,6 +711,12 @@ pub struct BleConfig {
|
||||
pub adapter: Option<String>,
|
||||
|
||||
/// L2CAP PSM for FIPS connections. Default: 0x0085 (133).
|
||||
///
|
||||
/// This is the PSM to request for this node's listener, and the PSM to
|
||||
/// dial a peer at when that peer advertises none. It is not always the
|
||||
/// PSM finally used: platforms that assign listener PSMs themselves
|
||||
/// report back what they bound and advertise that, and a peer that
|
||||
/// advertises its own PSM is dialled there instead.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub psm: Option<u16>,
|
||||
|
||||
|
||||
@@ -2710,7 +2710,7 @@ impl Node {
|
||||
}
|
||||
}
|
||||
} else if addr.transport == "ble" {
|
||||
#[cfg(bluer_available)]
|
||||
#[cfg(ble_available)]
|
||||
{
|
||||
match self.resolve_ble_addr(&addr.addr) {
|
||||
Ok(result) => result,
|
||||
@@ -2725,7 +2725,7 @@ impl Node {
|
||||
}
|
||||
}
|
||||
}
|
||||
#[cfg(not(bluer_available))]
|
||||
#[cfg(not(ble_available))]
|
||||
{
|
||||
debug!(transport = %addr.transport, "BLE transport not available on this build");
|
||||
continue;
|
||||
|
||||
+93
-2
@@ -560,6 +560,17 @@ pub struct Node {
|
||||
/// TUN interface name (for cleanup).
|
||||
tun_name: Option<String>,
|
||||
|
||||
/// Slot the embedder installs its BLE radio into, armed by
|
||||
/// [`Self::enable_app_owned_ble_radio`]. `None` unless armed.
|
||||
///
|
||||
/// Gated on the BLE transport existing *and* on its backend being the
|
||||
/// embedder-supplied one — the same condition
|
||||
/// `transport::ble::io_android` itself is compiled under, so the seam is
|
||||
/// absent on platforms whose radio is opened in process, and present in a
|
||||
/// test build so its contract is covered on an ordinary runner.
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
ble_radio: Option<Arc<crate::transport::ble::io_android::BleRadioSlot>>,
|
||||
|
||||
// === Index-Based Session Dispatch ===
|
||||
/// Allocator for session indices.
|
||||
index_allocator: IndexAllocator,
|
||||
@@ -859,6 +870,8 @@ impl Node {
|
||||
)),
|
||||
tun_state,
|
||||
tun_name: None,
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
ble_radio: None,
|
||||
index_allocator: IndexAllocator::new(),
|
||||
peers_by_index: HashMap::new(),
|
||||
pending_outbound: HashMap::new(),
|
||||
@@ -1029,6 +1042,8 @@ impl Node {
|
||||
)),
|
||||
tun_state,
|
||||
tun_name: None,
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
ble_radio: None,
|
||||
index_allocator: IndexAllocator::new(),
|
||||
peers_by_index: HashMap::new(),
|
||||
pending_outbound: HashMap::new(),
|
||||
@@ -1199,7 +1214,7 @@ impl Node {
|
||||
let transport_id = self.allocate_transport_id();
|
||||
let adapter = ble_config.adapter().to_string();
|
||||
let mtu = ble_config.mtu();
|
||||
match crate::transport::ble::io::BluerIo::new(&adapter, mtu).await {
|
||||
match crate::transport::ble::io_linux::BluerIo::new(&adapter, mtu).await {
|
||||
Ok(io) => {
|
||||
let ble = crate::transport::ble::BleTransport::new(
|
||||
transport_id,
|
||||
@@ -1223,6 +1238,32 @@ impl Node {
|
||||
}
|
||||
}
|
||||
|
||||
// Create BLE transport instances over an embedder-supplied radio.
|
||||
// Built whether or not a radio is installed yet: the backend resolves
|
||||
// the slot per operation, so one that arrives later is adopted in
|
||||
// place rather than needing the node rebuilt around it.
|
||||
#[cfg(all(target_os = "android", not(bluer_available), not(test)))]
|
||||
if let Some(slot) = self.ble_radio.clone() {
|
||||
let ble_instances: Vec<_> = self
|
||||
.config()
|
||||
.transports
|
||||
.ble
|
||||
.iter()
|
||||
.map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
|
||||
.collect();
|
||||
for (name, ble_config) in ble_instances {
|
||||
let transport_id = self.allocate_transport_id();
|
||||
let ble = crate::transport::ble::BleTransport::new(
|
||||
transport_id,
|
||||
name,
|
||||
ble_config,
|
||||
crate::transport::ble::io_android::AndroidIo::new(Arc::clone(&slot)),
|
||||
packet_tx.clone(),
|
||||
);
|
||||
transports.push(TransportHandle::Ble(ble));
|
||||
}
|
||||
}
|
||||
|
||||
transports
|
||||
}
|
||||
|
||||
@@ -1289,7 +1330,7 @@ impl Node {
|
||||
/// Resolve a BLE address string (`"adapter/AA:BB:CC:DD:EE:FF"`) to a
|
||||
/// (TransportId, TransportAddr) pair by finding the BLE transport
|
||||
/// instance matching the adapter name.
|
||||
#[cfg(bluer_available)]
|
||||
#[cfg(ble_available)]
|
||||
fn resolve_ble_addr(&self, addr_str: &str) -> Result<(TransportId, TransportAddr), NodeError> {
|
||||
let ta = TransportAddr::from_string(addr_str);
|
||||
let adapter = crate::transport::ble::addr::adapter_from_addr(&ta).ok_or_else(|| {
|
||||
@@ -3444,6 +3485,56 @@ impl Node {
|
||||
udp_fd_rx
|
||||
}
|
||||
|
||||
/// Set up an **app-owned BLE radio**: the embedder supplies the radio the
|
||||
/// BLE transport drives, because on this platform there is no
|
||||
/// Rust-reachable one to open. Call this after [`Node::new`] and
|
||||
/// **before** [`Self::start`] — the transport is built during `start`, and
|
||||
/// only a node armed by then has a slot to build it over.
|
||||
///
|
||||
/// Returns the slot. Installing, replacing and clearing a radio through it
|
||||
/// is safe at any time, from any thread, including long after the node is
|
||||
/// running:
|
||||
///
|
||||
/// ```no_run
|
||||
/// # async fn f(node: &mut fips::Node, radio: std::sync::Arc<dyn fips::transport::ble::io_android::AndroidRadio>)
|
||||
/// # -> Result<(), Box<dyn std::error::Error>> {
|
||||
/// let slot = node.enable_app_owned_ble_radio(); // after new(), before start()
|
||||
/// node.start().await?;
|
||||
/// // ...whenever the embedder's radio service comes up, and again each
|
||||
/// // time it restarts:
|
||||
/// slot.install(fips::transport::ble::io_android::AndroidBleBridge::new(radio));
|
||||
/// # Ok(())
|
||||
/// # }
|
||||
/// ```
|
||||
///
|
||||
/// The lateness is the point rather than a convenience. The radio belongs
|
||||
/// to a service whose lifetime is not the node's: the user can turn it on
|
||||
/// after the mesh is already running, and off and on again, and each start
|
||||
/// typically produces a fresh radio. A node that had to be built around an
|
||||
/// existing radio would make that mean "tear the node down and rebuild
|
||||
/// it", dropping every peer, session and route for as long as
|
||||
/// re-handshaking takes. So the transport is built and started whether or
|
||||
/// not a radio is installed, and resolves the slot per operation: it
|
||||
/// listens and scans against whichever radio is there, dials fail while
|
||||
/// there is none, and everything recovers on its own when one appears.
|
||||
/// Streams already open keep the radio they were opened on rather than
|
||||
/// migrating.
|
||||
///
|
||||
/// Deliberately narrow, and shaped like the [`Self::enable_app_owned_tun`]
|
||||
/// seam it sits beside: one call, no callbacks, and no lifecycle contract
|
||||
/// beyond the slot outliving the node. Arming twice returns the same slot,
|
||||
/// so a second call cannot orphan a radio installed through the first. The
|
||||
/// seam does not exist on platforms whose BLE backend is opened in
|
||||
/// process, since there is nothing there for an embedder to supply.
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
pub fn enable_app_owned_ble_radio(
|
||||
&mut self,
|
||||
) -> Arc<crate::transport::ble::io_android::BleRadioSlot> {
|
||||
Arc::clone(self.ble_radio.get_or_insert_with(|| {
|
||||
Arc::new(crate::transport::ble::io_android::BleRadioSlot::new())
|
||||
}))
|
||||
}
|
||||
|
||||
/// Address the built-in `.fips` DNS responder is listening on, or `None`
|
||||
/// when it is not running (`dns.enabled = false`, the bind failed, or the
|
||||
/// node is stopped).
|
||||
|
||||
@@ -6,7 +6,7 @@ use crate::utils::index::SessionIndex;
|
||||
use std::time::Duration;
|
||||
|
||||
mod acl;
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
mod ble;
|
||||
mod bloom;
|
||||
mod bloom_poison;
|
||||
|
||||
@@ -3152,6 +3152,115 @@ fn app_owned_udp_fd_seam_second_arm_replaces_the_first() {
|
||||
);
|
||||
}
|
||||
|
||||
/// The app-owned BLE radio seam. The slot is live from the moment it is armed
|
||||
/// — before `start()`, which is when the transport that reads it gets built —
|
||||
/// and installing a radio through it is a slot operation, not a node one.
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
#[test]
|
||||
fn app_owned_ble_radio_seam_hands_out_a_live_slot_before_start() {
|
||||
use crate::transport::ble::io_android::{AndroidBleBridge, BleRadioSlot};
|
||||
|
||||
let mut node = make_node();
|
||||
let slot: std::sync::Arc<BleRadioSlot> = node.enable_app_owned_ble_radio();
|
||||
|
||||
assert!(
|
||||
!slot.is_installed(),
|
||||
"arming supplies the slot, not a radio to put in it",
|
||||
);
|
||||
|
||||
slot.install(AndroidBleBridge::new(std::sync::Arc::new(
|
||||
test_radio::TestRadio,
|
||||
)));
|
||||
assert!(slot.is_installed(), "the embedder installs whenever it can");
|
||||
|
||||
slot.clear();
|
||||
assert!(!slot.is_installed(), "and can take it away again");
|
||||
}
|
||||
|
||||
/// Arming twice returns the same slot, so a second call cannot orphan a radio
|
||||
/// installed through the first. This is where the seam deliberately differs
|
||||
/// from `enable_app_owned_udp_fd`, whose last arming wins: a channel can be
|
||||
/// replaced because nothing was delivered on it yet, while a slot may already
|
||||
/// be holding the embedder's live radio.
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
#[test]
|
||||
fn app_owned_ble_radio_seam_second_arm_returns_the_same_slot() {
|
||||
use crate::transport::ble::io_android::AndroidBleBridge;
|
||||
|
||||
let mut node = make_node();
|
||||
let first = node.enable_app_owned_ble_radio();
|
||||
first.install(AndroidBleBridge::new(std::sync::Arc::new(
|
||||
test_radio::TestRadio,
|
||||
)));
|
||||
|
||||
let second = node.enable_app_owned_ble_radio();
|
||||
|
||||
assert!(
|
||||
std::sync::Arc::ptr_eq(&first, &second),
|
||||
"re-arming must not hand back a different slot",
|
||||
);
|
||||
assert!(
|
||||
second.is_installed(),
|
||||
"the radio installed through the first handle is still there",
|
||||
);
|
||||
}
|
||||
|
||||
/// The slot is per-node state, which is the whole reason it is not a process
|
||||
/// global: two nodes in one process each drive their own radio.
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
#[test]
|
||||
fn app_owned_ble_radio_slots_are_per_node() {
|
||||
use crate::transport::ble::io_android::AndroidBleBridge;
|
||||
|
||||
let mut node_a = make_node();
|
||||
let mut node_b = make_node();
|
||||
let slot_a = node_a.enable_app_owned_ble_radio();
|
||||
let slot_b = node_b.enable_app_owned_ble_radio();
|
||||
|
||||
slot_a.install(AndroidBleBridge::new(std::sync::Arc::new(
|
||||
test_radio::TestRadio,
|
||||
)));
|
||||
|
||||
assert!(slot_a.is_installed());
|
||||
assert!(
|
||||
!slot_b.is_installed(),
|
||||
"node B's radio is node B's — no shared or global slot",
|
||||
);
|
||||
}
|
||||
|
||||
/// A node that never armed the seam has no slot to hand the transport, which
|
||||
/// is how a build with the embedder-supplied backend distinguishes "no radio
|
||||
/// yet" from "this embedder does not supply one at all".
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
#[test]
|
||||
fn app_owned_ble_radio_seam_is_absent_until_armed() {
|
||||
let node = make_node();
|
||||
assert!(node.ble_radio.is_none());
|
||||
}
|
||||
|
||||
#[cfg(all(ble_available, any(target_os = "android", test)))]
|
||||
mod test_radio {
|
||||
use crate::transport::ble::addr::BleAddr;
|
||||
use crate::transport::ble::io_android::AndroidRadio;
|
||||
|
||||
/// A radio that does nothing. These tests are about the seam handing one
|
||||
/// over, not about what it then does — that is covered where the backend
|
||||
/// lives.
|
||||
pub(super) struct TestRadio;
|
||||
|
||||
impl AndroidRadio for TestRadio {
|
||||
fn listen(&self) -> u16 {
|
||||
0
|
||||
}
|
||||
fn connect(&self, _connect_id: i64, _addr: &BleAddr, _psm: u16) {}
|
||||
fn start_advertising(&self, _psm: u16) {}
|
||||
fn stop_advertising(&self) {}
|
||||
fn start_scanning(&self) {}
|
||||
fn stop_scanning(&self) {}
|
||||
fn close_channel(&self, _ch_id: i64) {}
|
||||
}
|
||||
}
|
||||
|
||||
/// The embedder-facing DNS contract, end to end.
|
||||
///
|
||||
/// An embedder that owns the TUN fd (Android `VpnService`) has no system DNS
|
||||
|
||||
+163
-406
@@ -1,8 +1,9 @@
|
||||
//! BLE I/O abstraction layer.
|
||||
//!
|
||||
//! Defines the `BleIo` trait that separates transport logic from the
|
||||
//! BlueZ/bluer stack. `BluerIo` (behind `cfg(bluer_available)`) provides
|
||||
//! the real implementation; `MockBleIo` provides an in-memory test double.
|
||||
//! Defines the `BleIo` seam that separates transport logic from any one
|
||||
//! radio stack, and the in-memory `MockBleIo` test double. Everything in
|
||||
//! this file is platform-neutral; each concrete backend lives beside it in
|
||||
//! its own `io_<platform>.rs` — [`super::io_linux`] for BlueZ.
|
||||
|
||||
use crate::transport::TransportError;
|
||||
|
||||
@@ -58,12 +59,53 @@ pub trait BleAcceptor: Send {
|
||||
) -> impl std::future::Future<Output = Result<Self::Stream, TransportError>> + Send;
|
||||
}
|
||||
|
||||
/// One advertisement observed by a scanner.
|
||||
///
|
||||
/// Carries what the backend could read from the advert, not what it wishes
|
||||
/// were there: a backend that cannot surface a field reports `None` for it,
|
||||
/// the way `TransportHandle::local_addr` already does for transports that
|
||||
/// have no address to give.
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub struct ScanAdvert {
|
||||
/// The advertiser's link address.
|
||||
pub addr: BleAddr,
|
||||
/// The L2CAP listener PSM the peer advertised, if it advertised one.
|
||||
///
|
||||
/// `None` for a legacy UUID-only advertiser, and for a backend that
|
||||
/// cannot read advertised service data. The dialer then falls back to the
|
||||
/// configured PSM. See [`super::psm`] for the wire layout.
|
||||
pub psm: Option<u16>,
|
||||
/// Received signal strength in dBm, if the backend reports it.
|
||||
pub rssi: Option<i16>,
|
||||
}
|
||||
|
||||
impl ScanAdvert {
|
||||
/// An advert carrying nothing but a link address — what a legacy
|
||||
/// UUID-only advertiser produces.
|
||||
pub fn new(addr: BleAddr) -> Self {
|
||||
Self {
|
||||
addr,
|
||||
psm: None,
|
||||
rssi: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// An advert carrying a listener PSM.
|
||||
pub fn with_psm(addr: BleAddr, psm: u16) -> Self {
|
||||
Self {
|
||||
addr,
|
||||
psm: Some(psm),
|
||||
rssi: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A scanner that yields discovered BLE devices advertising the FIPS UUID.
|
||||
pub trait BleScanner: Send {
|
||||
/// Wait for the next discovered device.
|
||||
/// Wait for the next observed advertisement.
|
||||
///
|
||||
/// Returns `None` when scanning is stopped.
|
||||
fn next(&mut self) -> impl std::future::Future<Output = Option<BleAddr>> + Send;
|
||||
fn next(&mut self) -> impl std::future::Future<Output = Option<ScanAdvert>> + Send;
|
||||
}
|
||||
|
||||
/// Core BLE I/O operations.
|
||||
@@ -79,11 +121,19 @@ pub trait BleIo: Send + Sync + 'static {
|
||||
/// The concrete scanner type.
|
||||
type Scanner: BleScanner + 'static;
|
||||
|
||||
/// Start listening for inbound L2CAP connections on the given PSM.
|
||||
/// Start listening for inbound L2CAP connections, and report the PSM
|
||||
/// actually bound.
|
||||
///
|
||||
/// `psm` is the PSM to request. Backends that let an application choose
|
||||
/// one (BlueZ) bind it and report it back unchanged. Backends whose
|
||||
/// platform assigns the PSM (Android, macOS) ignore the request and
|
||||
/// report what the OS gave them — which is why this returns a value
|
||||
/// rather than being assumed equal to the argument. The reported PSM is
|
||||
/// what gets advertised.
|
||||
fn listen(
|
||||
&self,
|
||||
psm: u16,
|
||||
) -> impl std::future::Future<Output = Result<Self::Acceptor, TransportError>> + Send;
|
||||
) -> impl std::future::Future<Output = Result<(Self::Acceptor, u16), TransportError>> + Send;
|
||||
|
||||
/// Connect to a remote BLE device on the given PSM.
|
||||
fn connect(
|
||||
@@ -92,9 +142,15 @@ pub trait BleIo: Send + Sync + 'static {
|
||||
psm: u16,
|
||||
) -> impl std::future::Future<Output = Result<Self::Stream, TransportError>> + Send;
|
||||
|
||||
/// Start advertising the FIPS service UUID.
|
||||
/// Start advertising the FIPS service UUID, and the listener PSM.
|
||||
///
|
||||
/// `psm` is the PSM this node's listener is bound to; see [`super::psm`]
|
||||
/// for the wire layout it should be advertised in. A backend that cannot
|
||||
/// put it in its advert ignores the argument, and peers dial it at their
|
||||
/// configured PSM as before.
|
||||
fn start_advertising(
|
||||
&self,
|
||||
psm: u16,
|
||||
) -> impl std::future::Future<Output = Result<(), TransportError>> + Send;
|
||||
|
||||
/// Stop advertising.
|
||||
@@ -107,6 +163,16 @@ pub trait BleIo: Send + Sync + 'static {
|
||||
&self,
|
||||
) -> impl std::future::Future<Output = Result<Self::Scanner, TransportError>> + Send;
|
||||
|
||||
/// Stop scanning.
|
||||
///
|
||||
/// The counterpart to [`Self::stop_advertising`], and needed for the same
|
||||
/// reason: dropping the transport's scan task stops *us* reading adverts,
|
||||
/// but on a backend whose radio is owned elsewhere it does not stop the
|
||||
/// radio. A backend whose scan ends when its `Scanner` is dropped
|
||||
/// implements this as a no-op and says so.
|
||||
fn stop_scanning(&self)
|
||||
-> impl std::future::Future<Output = Result<(), TransportError>> + Send;
|
||||
|
||||
/// Get the adapter's BLE address.
|
||||
fn local_addr(&self) -> Result<BleAddr, TransportError>;
|
||||
|
||||
@@ -114,390 +180,6 @@ pub trait BleIo: Send + Sync + 'static {
|
||||
fn adapter_name(&self) -> &str;
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// BluerIo — Production BLE I/O via BlueZ D-Bus
|
||||
// ============================================================================
|
||||
|
||||
#[cfg(bluer_available)]
|
||||
mod bluer_impl {
|
||||
use super::*;
|
||||
use crate::transport::TransportError;
|
||||
|
||||
use bluer::l2cap::{SeqPacket, SeqPacketListener, Socket, SocketAddr};
|
||||
use bluer::{
|
||||
AdapterEvent, AddressType, DiscoveryFilter, DiscoveryTransport, adv::Advertisement,
|
||||
};
|
||||
use futures::StreamExt;
|
||||
use std::collections::{BTreeSet, HashSet};
|
||||
use std::pin::Pin;
|
||||
use tokio::sync::Mutex;
|
||||
use tracing::{debug, trace};
|
||||
|
||||
/// FIPS BLE service UUID.
|
||||
///
|
||||
/// Derived from SHA-256("FIPS: welcome to cryptoanarchy") with UUID v4
|
||||
/// version/variant bits applied.
|
||||
pub const FIPS_SERVICE_UUID: bluer::Uuid =
|
||||
bluer::Uuid::from_u128(0x9c90_b790_2cc5_42c0_9f87_c9cc_4064_8f4c);
|
||||
|
||||
/// Map a bluer error to a TransportError.
|
||||
fn map_err(context: &str, e: bluer::Error) -> TransportError {
|
||||
TransportError::Io(std::io::Error::other(format!("{}: {}", context, e)))
|
||||
}
|
||||
|
||||
/// Map a std::io::Error to a TransportError.
|
||||
fn map_io_err(context: &str, e: std::io::Error) -> TransportError {
|
||||
TransportError::Io(std::io::Error::new(e.kind(), format!("{}: {}", context, e)))
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerStream
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// BLE stream wrapping a bluer L2CAP SeqPacket connection.
|
||||
pub struct BluerStream {
|
||||
conn: SeqPacket,
|
||||
remote: BleAddr,
|
||||
send_mtu: u16,
|
||||
recv_mtu: u16,
|
||||
}
|
||||
|
||||
impl BluerStream {
|
||||
/// Construct from a connected SeqPacket, querying MTU values.
|
||||
pub fn new(conn: SeqPacket, remote: BleAddr) -> Result<Self, TransportError> {
|
||||
let send_mtu = conn.send_mtu().map_err(|e| map_io_err("send_mtu", e))? as u16;
|
||||
let recv_mtu = conn.recv_mtu().map_err(|e| map_io_err("recv_mtu", e))? as u16;
|
||||
|
||||
// Log negotiated PHY for diagnostics (2M vs 1M)
|
||||
match conn.as_ref().phy() {
|
||||
Ok(phy) => {
|
||||
debug!(addr = %remote, phy, send_mtu, recv_mtu, "BLE connection established")
|
||||
}
|
||||
Err(_) => {
|
||||
debug!(addr = %remote, send_mtu, recv_mtu, "BLE connection established (PHY query unsupported)")
|
||||
}
|
||||
}
|
||||
|
||||
Ok(Self {
|
||||
conn,
|
||||
remote,
|
||||
send_mtu,
|
||||
recv_mtu,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl BleStream for BluerStream {
|
||||
async fn send(&self, data: &[u8]) -> Result<(), TransportError> {
|
||||
self.conn
|
||||
.send(data)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| TransportError::SendFailed(format!("{}", e)))
|
||||
}
|
||||
|
||||
async fn recv(&self, buf: &mut [u8]) -> Result<usize, TransportError> {
|
||||
self.conn
|
||||
.recv(buf)
|
||||
.await
|
||||
.map_err(|e| TransportError::RecvFailed(format!("{}", e)))
|
||||
}
|
||||
|
||||
fn send_mtu(&self) -> u16 {
|
||||
self.send_mtu
|
||||
}
|
||||
|
||||
fn recv_mtu(&self) -> u16 {
|
||||
self.recv_mtu
|
||||
}
|
||||
|
||||
fn remote_addr(&self) -> &BleAddr {
|
||||
&self.remote
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerAcceptor
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Acceptor wrapping a bluer L2CAP SeqPacketListener.
|
||||
pub struct BluerAcceptor {
|
||||
listener: SeqPacketListener,
|
||||
adapter_name: String,
|
||||
}
|
||||
|
||||
impl BleAcceptor for BluerAcceptor {
|
||||
type Stream = BluerStream;
|
||||
|
||||
async fn accept(&mut self) -> Result<BluerStream, TransportError> {
|
||||
let (conn, peer_sa) = self
|
||||
.listener
|
||||
.accept()
|
||||
.await
|
||||
.map_err(|e| map_io_err("accept", e))?;
|
||||
|
||||
let remote = BleAddr::from_bluer(peer_sa.addr, &self.adapter_name);
|
||||
BluerStream::new(conn, remote)
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerScanner
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Scanner wrapping a bluer discovery event stream.
|
||||
pub struct BluerScanner {
|
||||
events: Pin<Box<dyn futures::Stream<Item = AdapterEvent> + Send>>,
|
||||
adapter: bluer::Adapter,
|
||||
adapter_name: String,
|
||||
}
|
||||
|
||||
impl BleScanner for BluerScanner {
|
||||
async fn next(&mut self) -> Option<BleAddr> {
|
||||
loop {
|
||||
match self.events.next().await {
|
||||
Some(AdapterEvent::DeviceAdded(addr)) => {
|
||||
// Check if device advertises FIPS UUID
|
||||
if let Ok(device) = self.adapter.device(addr) {
|
||||
match device.uuids().await {
|
||||
Ok(Some(uuids)) if uuids.contains(&FIPS_SERVICE_UUID) => {
|
||||
let ble_addr = BleAddr::from_bluer(addr, &self.adapter_name);
|
||||
debug!(addr = %ble_addr, "BLE scanner: FIPS peer found");
|
||||
return Some(ble_addr);
|
||||
}
|
||||
Ok(_) => {
|
||||
trace!(addr = %addr, "BLE scanner: device without FIPS UUID");
|
||||
}
|
||||
Err(e) => {
|
||||
trace!(addr = %addr, error = %e, "BLE scanner: failed to read UUIDs");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Some(_) => continue,
|
||||
None => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerIo
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Production BLE I/O implementation via BlueZ D-Bus (bluer crate).
|
||||
pub struct BluerIo {
|
||||
#[allow(dead_code)] // Session must be kept alive for the adapter.
|
||||
session: bluer::Session,
|
||||
adapter: bluer::Adapter,
|
||||
adapter_name: String,
|
||||
adv_handle: Mutex<Option<bluer::adv::AdvertisementHandle>>,
|
||||
mtu: u16,
|
||||
}
|
||||
|
||||
impl BluerIo {
|
||||
/// Create a new BluerIo for the given adapter.
|
||||
///
|
||||
/// Connects to BlueZ via D-Bus and powers on the adapter.
|
||||
pub async fn new(adapter_name: &str, mtu: u16) -> Result<Self, TransportError> {
|
||||
let session = bluer::Session::new()
|
||||
.await
|
||||
.map_err(|e| map_err("Session::new", e))?;
|
||||
|
||||
let adapter = if adapter_name == "default" {
|
||||
session
|
||||
.default_adapter()
|
||||
.await
|
||||
.map_err(|e| map_err("default_adapter", e))?
|
||||
} else {
|
||||
session
|
||||
.adapter(adapter_name)
|
||||
.map_err(|e| map_err("adapter", e))?
|
||||
};
|
||||
|
||||
adapter
|
||||
.set_powered(true)
|
||||
.await
|
||||
.map_err(|e| map_err("set_powered", e))?;
|
||||
|
||||
let name = adapter.name().to_string();
|
||||
debug!(adapter = %name, "BluerIo initialized");
|
||||
|
||||
Ok(Self {
|
||||
session,
|
||||
adapter,
|
||||
adapter_name: name,
|
||||
adv_handle: Mutex::new(None),
|
||||
mtu,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl BleIo for BluerIo {
|
||||
type Stream = BluerStream;
|
||||
type Acceptor = BluerAcceptor;
|
||||
type Scanner = BluerScanner;
|
||||
|
||||
async fn listen(&self, psm: u16) -> Result<Self::Acceptor, TransportError> {
|
||||
let local_addr = self
|
||||
.adapter
|
||||
.address()
|
||||
.await
|
||||
.map_err(|e| map_err("address", e))?;
|
||||
|
||||
let sa = SocketAddr::new(local_addr, AddressType::LePublic, psm);
|
||||
let listener = SeqPacketListener::bind(sa)
|
||||
.await
|
||||
.map_err(|e| map_io_err("bind", e))?;
|
||||
|
||||
// Request high MTU for accepted connections
|
||||
listener
|
||||
.as_ref()
|
||||
.set_recv_mtu(self.mtu)
|
||||
.map_err(|e| map_io_err("set_recv_mtu", e))?;
|
||||
|
||||
// Prevent sniff mode to reduce latency during data transfer
|
||||
if let Err(e) = listener.as_ref().set_power_forced_active(true) {
|
||||
debug!(error = %e, "BLE listener: set_power_forced_active not supported");
|
||||
}
|
||||
|
||||
debug!(psm, mtu = self.mtu, "BLE listener bound");
|
||||
|
||||
Ok(BluerAcceptor {
|
||||
listener,
|
||||
adapter_name: self.adapter_name.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
async fn connect(&self, addr: &BleAddr, psm: u16) -> Result<Self::Stream, TransportError> {
|
||||
let target_sa = addr.to_socket_addr(psm);
|
||||
|
||||
let socket = Socket::<SeqPacket>::new_seq_packet()
|
||||
.map_err(|e| map_io_err("new_seq_packet", e))?;
|
||||
socket
|
||||
.bind(SocketAddr::any_le())
|
||||
.map_err(|e| map_io_err("bind", e))?;
|
||||
socket
|
||||
.set_recv_mtu(self.mtu)
|
||||
.map_err(|e| map_io_err("set_recv_mtu", e))?;
|
||||
|
||||
// Prevent sniff mode to reduce latency during data transfer
|
||||
if let Err(e) = socket.set_power_forced_active(true) {
|
||||
debug!(error = %e, "BLE connect: set_power_forced_active not supported");
|
||||
}
|
||||
|
||||
let conn = socket
|
||||
.connect(target_sa)
|
||||
.await
|
||||
.map_err(|e| map_io_err("connect", e))?;
|
||||
|
||||
let remote = addr.clone();
|
||||
BluerStream::new(conn, remote)
|
||||
}
|
||||
|
||||
async fn start_advertising(&self) -> Result<(), TransportError> {
|
||||
let adv = Advertisement {
|
||||
advertisement_type: bluer::adv::Type::Peripheral,
|
||||
service_uuids: {
|
||||
let mut s = BTreeSet::new();
|
||||
s.insert(FIPS_SERVICE_UUID);
|
||||
s
|
||||
},
|
||||
local_name: Some("fips".to_string()),
|
||||
min_interval: Some(std::time::Duration::from_millis(400)),
|
||||
max_interval: Some(std::time::Duration::from_millis(600)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let handle = self
|
||||
.adapter
|
||||
.advertise(adv)
|
||||
.await
|
||||
.map_err(|e| map_err("advertise", e))?;
|
||||
|
||||
*self.adv_handle.lock().await = Some(handle);
|
||||
debug!("BLE advertising started");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn stop_advertising(&self) -> Result<(), TransportError> {
|
||||
let _ = self.adv_handle.lock().await.take();
|
||||
debug!("BLE advertising stopped");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn start_scanning(&self) -> Result<Self::Scanner, TransportError> {
|
||||
// Clear cached devices so BlueZ fires DeviceAdded for every
|
||||
// advertisement. Without this, already-known devices only
|
||||
// produce PropertyChanged events (which bluer doesn't expose
|
||||
// at the device level), causing the scanner to miss peers
|
||||
// after a daemon restart.
|
||||
if let Ok(cached) = self.adapter.device_addresses().await {
|
||||
let count = cached.len();
|
||||
for addr in cached {
|
||||
let _ = self.adapter.remove_device(addr).await;
|
||||
}
|
||||
if count > 0 {
|
||||
debug!(count, "BLE scanner: cleared cached devices");
|
||||
}
|
||||
}
|
||||
|
||||
// Set discovery filter for LE transport with FIPS UUID
|
||||
let filter = DiscoveryFilter {
|
||||
transport: DiscoveryTransport::Le,
|
||||
uuids: {
|
||||
let mut s = HashSet::new();
|
||||
s.insert(FIPS_SERVICE_UUID);
|
||||
s
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.adapter
|
||||
.set_discovery_filter(filter)
|
||||
.await
|
||||
.map_err(|e| map_err("set_discovery_filter", e))?;
|
||||
|
||||
let events = self
|
||||
.adapter
|
||||
.discover_devices()
|
||||
.await
|
||||
.map_err(|e| map_err("discover_devices", e))?;
|
||||
|
||||
debug!("BLE scanning started");
|
||||
|
||||
Ok(BluerScanner {
|
||||
events: Box::pin(events),
|
||||
adapter: self.adapter.clone(),
|
||||
adapter_name: self.adapter_name.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
fn local_addr(&self) -> Result<BleAddr, TransportError> {
|
||||
// Use futures::executor::block_on since this is a sync method
|
||||
// but needs an async call. The adapter address is cached so
|
||||
// the D-Bus call is fast.
|
||||
let addr = futures::executor::block_on(self.adapter.address())
|
||||
.map_err(|e| map_err("address", e))?;
|
||||
Ok(BleAddr::from_bluer(addr, &self.adapter_name))
|
||||
}
|
||||
|
||||
fn adapter_name(&self) -> &str {
|
||||
&self.adapter_name
|
||||
}
|
||||
}
|
||||
|
||||
// Compile-time assertion that BluerIo satisfies Send + Sync.
|
||||
#[allow(dead_code)]
|
||||
fn _assert_bluer_io_send_sync() {
|
||||
fn require<T: Send + Sync>() {}
|
||||
require::<BluerIo>();
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(bluer_available)]
|
||||
pub use bluer_impl::{BluerAcceptor, BluerIo, BluerScanner, BluerStream, FIPS_SERVICE_UUID};
|
||||
|
||||
// ============================================================================
|
||||
// Mock BLE I/O (for testing without hardware)
|
||||
// ============================================================================
|
||||
@@ -583,13 +265,13 @@ impl BleAcceptor for MockBleAcceptor {
|
||||
}
|
||||
}
|
||||
|
||||
/// Mock BLE scanner backed by a channel of discovered addresses.
|
||||
/// Mock BLE scanner backed by a channel of observed adverts.
|
||||
pub struct MockBleScanner {
|
||||
rx: tokio::sync::mpsc::Receiver<BleAddr>,
|
||||
rx: tokio::sync::mpsc::Receiver<ScanAdvert>,
|
||||
}
|
||||
|
||||
impl BleScanner for MockBleScanner {
|
||||
async fn next(&mut self) -> Option<BleAddr> {
|
||||
async fn next(&mut self) -> Option<ScanAdvert> {
|
||||
self.rx.recv().await
|
||||
}
|
||||
}
|
||||
@@ -607,9 +289,17 @@ pub struct MockBleIo {
|
||||
local_addr: BleAddr,
|
||||
accept_tx: tokio::sync::mpsc::Sender<MockBleStream>,
|
||||
accept_rx: std::sync::Mutex<Option<tokio::sync::mpsc::Receiver<MockBleStream>>>,
|
||||
scan_tx: tokio::sync::mpsc::Sender<BleAddr>,
|
||||
scan_rx: std::sync::Mutex<Option<tokio::sync::mpsc::Receiver<BleAddr>>>,
|
||||
scan_tx: tokio::sync::mpsc::Sender<ScanAdvert>,
|
||||
scan_rx: std::sync::Mutex<Option<tokio::sync::mpsc::Receiver<ScanAdvert>>>,
|
||||
connect_handler: std::sync::Mutex<Option<ConnectHandler>>,
|
||||
/// PSM `listen` reports back, overriding the requested one.
|
||||
///
|
||||
/// Simulates a platform that assigns the PSM itself.
|
||||
bound_psm: std::sync::Mutex<Option<u16>>,
|
||||
/// PSM most recently passed to `start_advertising`.
|
||||
advertised_psm: std::sync::Mutex<Option<u16>>,
|
||||
/// Number of times `stop_scanning` has been called.
|
||||
stop_scans: std::sync::atomic::AtomicUsize,
|
||||
}
|
||||
|
||||
impl MockBleIo {
|
||||
@@ -625,17 +315,45 @@ impl MockBleIo {
|
||||
scan_tx,
|
||||
scan_rx: std::sync::Mutex::new(Some(scan_rx)),
|
||||
connect_handler: std::sync::Mutex::new(None),
|
||||
bound_psm: std::sync::Mutex::new(None),
|
||||
advertised_psm: std::sync::Mutex::new(None),
|
||||
stop_scans: std::sync::atomic::AtomicUsize::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
/// How many times the transport has asked this backend to stop scanning.
|
||||
pub fn stop_scan_calls(&self) -> usize {
|
||||
self.stop_scans.load(std::sync::atomic::Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Inject an inbound connection (simulates a remote device connecting).
|
||||
pub async fn inject_inbound(&self, stream: MockBleStream) {
|
||||
let _ = self.accept_tx.send(stream).await;
|
||||
}
|
||||
|
||||
/// Inject a scan result (simulates discovering a remote device).
|
||||
/// Inject a scan result (simulates discovering a legacy UUID-only
|
||||
/// advertiser, which carries no PSM).
|
||||
pub async fn inject_scan_result(&self, addr: BleAddr) {
|
||||
let _ = self.scan_tx.send(addr).await;
|
||||
self.inject_scan_advert(ScanAdvert::new(addr)).await;
|
||||
}
|
||||
|
||||
/// Inject an observed advertisement verbatim.
|
||||
pub async fn inject_scan_advert(&self, advert: ScanAdvert) {
|
||||
let _ = self.scan_tx.send(advert).await;
|
||||
}
|
||||
|
||||
/// Make `listen` report a PSM other than the one requested, the way a
|
||||
/// platform that assigns PSMs itself would.
|
||||
pub fn set_bound_psm(&self, psm: u16) {
|
||||
*self.bound_psm.lock().unwrap_or_else(|e| e.into_inner()) = Some(psm);
|
||||
}
|
||||
|
||||
/// The PSM most recently handed to `start_advertising`.
|
||||
pub fn advertised_psm(&self) -> Option<u16> {
|
||||
*self
|
||||
.advertised_psm
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
}
|
||||
|
||||
/// Set a handler for outbound connect calls.
|
||||
@@ -655,14 +373,19 @@ impl BleIo for MockBleIo {
|
||||
type Acceptor = MockBleAcceptor;
|
||||
type Scanner = MockBleScanner;
|
||||
|
||||
async fn listen(&self, _psm: u16) -> Result<Self::Acceptor, TransportError> {
|
||||
async fn listen(&self, psm: u16) -> Result<(Self::Acceptor, u16), TransportError> {
|
||||
let rx = self
|
||||
.accept_rx
|
||||
.lock()
|
||||
.unwrap()
|
||||
.take()
|
||||
.ok_or_else(|| TransportError::NotSupported("acceptor already taken".into()))?;
|
||||
Ok(MockBleAcceptor { rx })
|
||||
let bound = self
|
||||
.bound_psm
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.unwrap_or(psm);
|
||||
Ok((MockBleAcceptor { rx }, bound))
|
||||
}
|
||||
|
||||
async fn connect(&self, addr: &BleAddr, psm: u16) -> Result<Self::Stream, TransportError> {
|
||||
@@ -676,7 +399,11 @@ impl BleIo for MockBleIo {
|
||||
}
|
||||
}
|
||||
|
||||
async fn start_advertising(&self) -> Result<(), TransportError> {
|
||||
async fn start_advertising(&self, psm: u16) -> Result<(), TransportError> {
|
||||
*self
|
||||
.advertised_psm
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner()) = Some(psm);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -684,6 +411,12 @@ impl BleIo for MockBleIo {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn stop_scanning(&self) -> Result<(), TransportError> {
|
||||
self.stop_scans
|
||||
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn start_scanning(&self) -> Result<Self::Scanner, TransportError> {
|
||||
let rx = self
|
||||
.scan_rx
|
||||
@@ -751,7 +484,8 @@ mod tests {
|
||||
#[tokio::test]
|
||||
async fn test_mock_io_listen_accept() {
|
||||
let io = MockBleIo::new("hci0", test_addr(1));
|
||||
let mut acceptor = io.listen(0x0085).await.unwrap();
|
||||
let (mut acceptor, bound) = io.listen(0x0085).await.unwrap();
|
||||
assert_eq!(bound, 0x0085, "mock binds what it is asked for by default");
|
||||
|
||||
let (stream_a, _stream_b) = MockBleStream::pair(test_addr(1), test_addr(2), 2048);
|
||||
io.inject_inbound(stream_a).await;
|
||||
@@ -789,8 +523,8 @@ mod tests {
|
||||
io.inject_scan_result(test_addr(2)).await;
|
||||
io.inject_scan_result(test_addr(3)).await;
|
||||
|
||||
assert_eq!(scanner.next().await, Some(test_addr(2)));
|
||||
assert_eq!(scanner.next().await, Some(test_addr(3)));
|
||||
assert_eq!(scanner.next().await, Some(ScanAdvert::new(test_addr(2))));
|
||||
assert_eq!(scanner.next().await, Some(ScanAdvert::new(test_addr(3))));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -803,10 +537,33 @@ mod tests {
|
||||
#[tokio::test]
|
||||
async fn test_mock_io_advertising_noop() {
|
||||
let io = MockBleIo::new("hci0", test_addr(1));
|
||||
io.start_advertising().await.unwrap();
|
||||
io.start_advertising(0x0085).await.unwrap();
|
||||
assert_eq!(io.advertised_psm(), Some(0x0085));
|
||||
io.stop_advertising().await.unwrap();
|
||||
}
|
||||
|
||||
/// A backend whose platform assigns the PSM reports back something other
|
||||
/// than what was requested — the case the return value exists for.
|
||||
#[tokio::test]
|
||||
async fn test_mock_io_listen_reports_an_os_assigned_psm() {
|
||||
let io = MockBleIo::new("hci0", test_addr(1));
|
||||
io.set_bound_psm(0x00C1);
|
||||
let (_acceptor, bound) = io.listen(0x0085).await.unwrap();
|
||||
assert_eq!(bound, 0x00C1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_mock_io_scan_advert_carries_a_psm() {
|
||||
let io = MockBleIo::new("hci0", test_addr(1));
|
||||
let mut scanner = io.start_scanning().await.unwrap();
|
||||
io.inject_scan_advert(ScanAdvert::with_psm(test_addr(2), 0x00C1))
|
||||
.await;
|
||||
let advert = scanner.next().await.unwrap();
|
||||
assert_eq!(advert.addr, test_addr(2));
|
||||
assert_eq!(advert.psm, Some(0x00C1));
|
||||
assert_eq!(advert.rssi, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_mock_io_listen_twice_fails() {
|
||||
let io = MockBleIo::new("hci0", test_addr(1));
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,464 @@
|
||||
//! BlueZ backend for the BLE transport.
|
||||
//!
|
||||
//! The Linux implementation of the [`BleIo`](super::io::BleIo) seam, over the
|
||||
//! `bluer` crate's D-Bus binding to BlueZ. Everything platform-neutral lives
|
||||
//! in [`super::io`] and the modules beside it; this file holds only what
|
||||
//! speaks to BlueZ.
|
||||
|
||||
use super::io::*;
|
||||
use crate::transport::TransportError;
|
||||
|
||||
use bluer::l2cap::{SeqPacket, SeqPacketListener, Socket, SocketAddr};
|
||||
use bluer::{AdapterEvent, AddressType, DiscoveryFilter, DiscoveryTransport, adv::Advertisement};
|
||||
use futures::StreamExt;
|
||||
use std::collections::{BTreeMap, BTreeSet, HashSet};
|
||||
use std::pin::Pin;
|
||||
use tokio::sync::Mutex;
|
||||
use tracing::{debug, trace};
|
||||
|
||||
use super::addr::BleAddr;
|
||||
use super::psm;
|
||||
|
||||
/// FIPS BLE service UUID.
|
||||
///
|
||||
/// Derived from SHA-256("FIPS: welcome to cryptoanarchy") with UUID v4
|
||||
/// version/variant bits applied.
|
||||
pub const FIPS_SERVICE_UUID: bluer::Uuid =
|
||||
bluer::Uuid::from_u128(0x9c90_b790_2cc5_42c0_9f87_c9cc_4064_8f4c);
|
||||
|
||||
/// The PSM service-data key as a whole UUID.
|
||||
///
|
||||
/// BlueZ speaks in full UUIDs, so [`psm::PSM_SERVICE_DATA_UUID16`] is
|
||||
/// expanded through the Bluetooth base UUID
|
||||
/// (`00009C90-0000-1000-8000-00805F9B34FB`). The controller emits it
|
||||
/// back on the air as the 16-bit Service Data AD structure the wire
|
||||
/// layout in [`psm`] specifies.
|
||||
pub const PSM_SERVICE_DATA_UUID: bluer::Uuid = bluer::Uuid::from_u128(
|
||||
((psm::PSM_SERVICE_DATA_UUID16 as u128) << 96) | 0x0000_0000_0000_1000_8000_0080_5F9B_34FB,
|
||||
);
|
||||
|
||||
/// Map a bluer error to a TransportError.
|
||||
fn map_err(context: &str, e: bluer::Error) -> TransportError {
|
||||
TransportError::Io(std::io::Error::other(format!("{}: {}", context, e)))
|
||||
}
|
||||
|
||||
/// Map a std::io::Error to a TransportError.
|
||||
fn map_io_err(context: &str, e: std::io::Error) -> TransportError {
|
||||
TransportError::Io(std::io::Error::new(e.kind(), format!("{}: {}", context, e)))
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerStream
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// BLE stream wrapping a bluer L2CAP SeqPacket connection.
|
||||
pub struct BluerStream {
|
||||
conn: SeqPacket,
|
||||
remote: BleAddr,
|
||||
send_mtu: u16,
|
||||
recv_mtu: u16,
|
||||
}
|
||||
|
||||
impl BluerStream {
|
||||
/// Construct from a connected SeqPacket, querying MTU values.
|
||||
pub fn new(conn: SeqPacket, remote: BleAddr) -> Result<Self, TransportError> {
|
||||
let send_mtu = conn.send_mtu().map_err(|e| map_io_err("send_mtu", e))? as u16;
|
||||
let recv_mtu = conn.recv_mtu().map_err(|e| map_io_err("recv_mtu", e))? as u16;
|
||||
|
||||
// Log negotiated PHY for diagnostics (2M vs 1M)
|
||||
match conn.as_ref().phy() {
|
||||
Ok(phy) => {
|
||||
debug!(addr = %remote, phy, send_mtu, recv_mtu, "BLE connection established")
|
||||
}
|
||||
Err(_) => {
|
||||
debug!(addr = %remote, send_mtu, recv_mtu, "BLE connection established (PHY query unsupported)")
|
||||
}
|
||||
}
|
||||
|
||||
Ok(Self {
|
||||
conn,
|
||||
remote,
|
||||
send_mtu,
|
||||
recv_mtu,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl BleStream for BluerStream {
|
||||
async fn send(&self, data: &[u8]) -> Result<(), TransportError> {
|
||||
self.conn
|
||||
.send(data)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| TransportError::SendFailed(format!("{}", e)))
|
||||
}
|
||||
|
||||
async fn recv(&self, buf: &mut [u8]) -> Result<usize, TransportError> {
|
||||
self.conn
|
||||
.recv(buf)
|
||||
.await
|
||||
.map_err(|e| TransportError::RecvFailed(format!("{}", e)))
|
||||
}
|
||||
|
||||
fn send_mtu(&self) -> u16 {
|
||||
self.send_mtu
|
||||
}
|
||||
|
||||
fn recv_mtu(&self) -> u16 {
|
||||
self.recv_mtu
|
||||
}
|
||||
|
||||
fn remote_addr(&self) -> &BleAddr {
|
||||
&self.remote
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerAcceptor
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Acceptor wrapping a bluer L2CAP SeqPacketListener.
|
||||
pub struct BluerAcceptor {
|
||||
listener: SeqPacketListener,
|
||||
adapter_name: String,
|
||||
}
|
||||
|
||||
impl BleAcceptor for BluerAcceptor {
|
||||
type Stream = BluerStream;
|
||||
|
||||
async fn accept(&mut self) -> Result<BluerStream, TransportError> {
|
||||
let (conn, peer_sa) = self
|
||||
.listener
|
||||
.accept()
|
||||
.await
|
||||
.map_err(|e| map_io_err("accept", e))?;
|
||||
|
||||
let remote = BleAddr::from_bluer(peer_sa.addr, &self.adapter_name);
|
||||
BluerStream::new(conn, remote)
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerScanner
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Scanner wrapping a bluer discovery event stream.
|
||||
pub struct BluerScanner {
|
||||
events: Pin<Box<dyn futures::Stream<Item = AdapterEvent> + Send>>,
|
||||
adapter: bluer::Adapter,
|
||||
adapter_name: String,
|
||||
}
|
||||
|
||||
impl BleScanner for BluerScanner {
|
||||
/// Yields adverts with the PSM and RSSI when BlueZ can supply them.
|
||||
///
|
||||
/// The PSM comes out of the peer's Service Data AD structure (see
|
||||
/// `super::super::psm`); a peer that advertises none — a legacy
|
||||
/// UUID-only advertiser — yields `psm: None` and is dialled at the
|
||||
/// configured PSM, exactly as before.
|
||||
async fn next(&mut self) -> Option<ScanAdvert> {
|
||||
loop {
|
||||
match self.events.next().await {
|
||||
Some(AdapterEvent::DeviceAdded(addr)) => {
|
||||
// Check if device advertises FIPS UUID
|
||||
if let Ok(device) = self.adapter.device(addr) {
|
||||
match device.uuids().await {
|
||||
Ok(Some(uuids)) if uuids.contains(&FIPS_SERVICE_UUID) => {
|
||||
let ble_addr = BleAddr::from_bluer(addr, &self.adapter_name);
|
||||
let psm =
|
||||
device.service_data().await.ok().flatten().and_then(|sd| {
|
||||
sd.get(&PSM_SERVICE_DATA_UUID)
|
||||
.and_then(|data| psm::decode_psm(data))
|
||||
});
|
||||
let rssi = device.rssi().await.ok().flatten();
|
||||
debug!(addr = %ble_addr, ?psm, ?rssi, "BLE scanner: FIPS peer found");
|
||||
return Some(ScanAdvert {
|
||||
addr: ble_addr,
|
||||
psm,
|
||||
rssi,
|
||||
});
|
||||
}
|
||||
Ok(_) => {
|
||||
trace!(addr = %addr, "BLE scanner: device without FIPS UUID");
|
||||
}
|
||||
Err(e) => {
|
||||
trace!(addr = %addr, error = %e, "BLE scanner: failed to read UUIDs");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Some(_) => continue,
|
||||
None => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// BluerIo
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// Production BLE I/O implementation via BlueZ D-Bus (bluer crate).
|
||||
pub struct BluerIo {
|
||||
#[allow(dead_code)] // Session must be kept alive for the adapter.
|
||||
session: bluer::Session,
|
||||
adapter: bluer::Adapter,
|
||||
adapter_name: String,
|
||||
adv_handle: Mutex<Option<bluer::adv::AdvertisementHandle>>,
|
||||
mtu: u16,
|
||||
}
|
||||
|
||||
impl BluerIo {
|
||||
/// Create a new BluerIo for the given adapter.
|
||||
///
|
||||
/// Connects to BlueZ via D-Bus and powers on the adapter.
|
||||
pub async fn new(adapter_name: &str, mtu: u16) -> Result<Self, TransportError> {
|
||||
let session = bluer::Session::new()
|
||||
.await
|
||||
.map_err(|e| map_err("Session::new", e))?;
|
||||
|
||||
let adapter = if adapter_name == "default" {
|
||||
session
|
||||
.default_adapter()
|
||||
.await
|
||||
.map_err(|e| map_err("default_adapter", e))?
|
||||
} else {
|
||||
session
|
||||
.adapter(adapter_name)
|
||||
.map_err(|e| map_err("adapter", e))?
|
||||
};
|
||||
|
||||
adapter
|
||||
.set_powered(true)
|
||||
.await
|
||||
.map_err(|e| map_err("set_powered", e))?;
|
||||
|
||||
let name = adapter.name().to_string();
|
||||
debug!(adapter = %name, "BluerIo initialized");
|
||||
|
||||
Ok(Self {
|
||||
session,
|
||||
adapter,
|
||||
adapter_name: name,
|
||||
adv_handle: Mutex::new(None),
|
||||
mtu,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl BleIo for BluerIo {
|
||||
type Stream = BluerStream;
|
||||
type Acceptor = BluerAcceptor;
|
||||
type Scanner = BluerScanner;
|
||||
|
||||
/// Binds the requested PSM and reports it back unchanged.
|
||||
///
|
||||
/// BlueZ lets an application choose the PSM it binds, so the bound
|
||||
/// PSM is always the requested one. Backends whose platform assigns
|
||||
/// the PSM report something else; that is the reason for the return
|
||||
/// value, not anything BlueZ does.
|
||||
async fn listen(&self, psm: u16) -> Result<(Self::Acceptor, u16), TransportError> {
|
||||
let local_addr = self
|
||||
.adapter
|
||||
.address()
|
||||
.await
|
||||
.map_err(|e| map_err("address", e))?;
|
||||
|
||||
let sa = SocketAddr::new(local_addr, AddressType::LePublic, psm);
|
||||
let listener = SeqPacketListener::bind(sa)
|
||||
.await
|
||||
.map_err(|e| map_io_err("bind", e))?;
|
||||
|
||||
// Request high MTU for accepted connections
|
||||
listener
|
||||
.as_ref()
|
||||
.set_recv_mtu(self.mtu)
|
||||
.map_err(|e| map_io_err("set_recv_mtu", e))?;
|
||||
|
||||
// Prevent sniff mode to reduce latency during data transfer
|
||||
if let Err(e) = listener.as_ref().set_power_forced_active(true) {
|
||||
debug!(error = %e, "BLE listener: set_power_forced_active not supported");
|
||||
}
|
||||
|
||||
debug!(psm, mtu = self.mtu, "BLE listener bound");
|
||||
|
||||
Ok((
|
||||
BluerAcceptor {
|
||||
listener,
|
||||
adapter_name: self.adapter_name.clone(),
|
||||
},
|
||||
psm,
|
||||
))
|
||||
}
|
||||
|
||||
async fn connect(&self, addr: &BleAddr, psm: u16) -> Result<Self::Stream, TransportError> {
|
||||
let target_sa = addr.to_socket_addr(psm);
|
||||
|
||||
let socket =
|
||||
Socket::<SeqPacket>::new_seq_packet().map_err(|e| map_io_err("new_seq_packet", e))?;
|
||||
socket
|
||||
.bind(SocketAddr::any_le())
|
||||
.map_err(|e| map_io_err("bind", e))?;
|
||||
socket
|
||||
.set_recv_mtu(self.mtu)
|
||||
.map_err(|e| map_io_err("set_recv_mtu", e))?;
|
||||
|
||||
// Prevent sniff mode to reduce latency during data transfer
|
||||
if let Err(e) = socket.set_power_forced_active(true) {
|
||||
debug!(error = %e, "BLE connect: set_power_forced_active not supported");
|
||||
}
|
||||
|
||||
let conn = socket
|
||||
.connect(target_sa)
|
||||
.await
|
||||
.map_err(|e| map_io_err("connect", e))?;
|
||||
|
||||
let remote = addr.clone();
|
||||
BluerStream::new(conn, remote)
|
||||
}
|
||||
|
||||
/// Advertises the FIPS service UUID and the listener PSM.
|
||||
///
|
||||
/// The PSM rides the Service Data AD structure specified in
|
||||
/// `super::super::psm`. Emitting it costs the `local_name`: flags
|
||||
/// (3) + 128-bit UUID list (18) + service data (6) fill 27 of the
|
||||
/// 31-byte legacy PDU, and a name no longer fits. Peers that read
|
||||
/// the service data dial the advertised PSM; legacy peers keep
|
||||
/// dialling their configured one, which BlueZ listeners still bind.
|
||||
///
|
||||
/// `super::super::psm` requires the PSM to ride the primary
|
||||
/// advertisement, never the scan response, so a passive scanner
|
||||
/// still sees it. Nothing here enforces that: BlueZ takes a set of
|
||||
/// AD structures and chooses their placement itself. What keeps the
|
||||
/// requirement holding is the arithmetic above — 27 of 31 bytes
|
||||
/// used, so BlueZ has no reason to spill into the scan response —
|
||||
/// and dropping the name is what makes it hold.
|
||||
async fn start_advertising(&self, psm: u16) -> Result<(), TransportError> {
|
||||
let adv = Advertisement {
|
||||
advertisement_type: bluer::adv::Type::Peripheral,
|
||||
service_uuids: {
|
||||
let mut s = BTreeSet::new();
|
||||
s.insert(FIPS_SERVICE_UUID);
|
||||
s
|
||||
},
|
||||
service_data: {
|
||||
let mut m = BTreeMap::new();
|
||||
m.insert(PSM_SERVICE_DATA_UUID, psm::encode_psm(psm).to_vec());
|
||||
m
|
||||
},
|
||||
min_interval: Some(std::time::Duration::from_millis(400)),
|
||||
max_interval: Some(std::time::Duration::from_millis(600)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let handle = self
|
||||
.adapter
|
||||
.advertise(adv)
|
||||
.await
|
||||
.map_err(|e| map_err("advertise", e))?;
|
||||
|
||||
*self.adv_handle.lock().await = Some(handle);
|
||||
debug!(psm, "BLE advertising started");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn stop_advertising(&self) -> Result<(), TransportError> {
|
||||
let _ = self.adv_handle.lock().await.take();
|
||||
debug!("BLE advertising stopped");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// A no-op: BlueZ discovery ends when [`BluerScanner`]'s event stream is
|
||||
/// dropped, which happens when the transport drops the scanner. There is
|
||||
/// no separate adapter-level stop to issue.
|
||||
async fn stop_scanning(&self) -> Result<(), TransportError> {
|
||||
debug!("BLE scanning stops with the scanner");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn start_scanning(&self) -> Result<Self::Scanner, TransportError> {
|
||||
// Clear cached devices so BlueZ fires DeviceAdded for every
|
||||
// advertisement. Without this, already-known devices only
|
||||
// produce PropertyChanged events (which bluer doesn't expose
|
||||
// at the device level), causing the scanner to miss peers
|
||||
// after a daemon restart.
|
||||
if let Ok(cached) = self.adapter.device_addresses().await {
|
||||
let count = cached.len();
|
||||
for addr in cached {
|
||||
let _ = self.adapter.remove_device(addr).await;
|
||||
}
|
||||
if count > 0 {
|
||||
debug!(count, "BLE scanner: cleared cached devices");
|
||||
}
|
||||
}
|
||||
|
||||
// Set discovery filter for LE transport with FIPS UUID
|
||||
let filter = DiscoveryFilter {
|
||||
transport: DiscoveryTransport::Le,
|
||||
uuids: {
|
||||
let mut s = HashSet::new();
|
||||
s.insert(FIPS_SERVICE_UUID);
|
||||
s
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.adapter
|
||||
.set_discovery_filter(filter)
|
||||
.await
|
||||
.map_err(|e| map_err("set_discovery_filter", e))?;
|
||||
|
||||
let events = self
|
||||
.adapter
|
||||
.discover_devices()
|
||||
.await
|
||||
.map_err(|e| map_err("discover_devices", e))?;
|
||||
|
||||
debug!("BLE scanning started");
|
||||
|
||||
Ok(BluerScanner {
|
||||
events: Box::pin(events),
|
||||
adapter: self.adapter.clone(),
|
||||
adapter_name: self.adapter_name.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
fn local_addr(&self) -> Result<BleAddr, TransportError> {
|
||||
// Use futures::executor::block_on since this is a sync method
|
||||
// but needs an async call. The adapter address is cached so
|
||||
// the D-Bus call is fast.
|
||||
let addr = futures::executor::block_on(self.adapter.address())
|
||||
.map_err(|e| map_err("address", e))?;
|
||||
Ok(BleAddr::from_bluer(addr, &self.adapter_name))
|
||||
}
|
||||
|
||||
fn adapter_name(&self) -> &str {
|
||||
&self.adapter_name
|
||||
}
|
||||
}
|
||||
|
||||
// Compile-time assertion that BluerIo satisfies Send + Sync.
|
||||
#[allow(dead_code)]
|
||||
fn _assert_bluer_io_send_sync() {
|
||||
fn require<T: Send + Sync>() {}
|
||||
require::<BluerIo>();
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// A wrong shift in the base-UUID expansion would yield a plausible
|
||||
/// UUID that simply never matches any peer — silent discovery
|
||||
/// failure, not a build error. Companion to psm.rs's
|
||||
/// `test_key_is_the_leading_16_bits_of_the_fips_uuid`.
|
||||
#[test]
|
||||
fn psm_service_data_uuid_expands_the_key_over_the_base_uuid() {
|
||||
assert_eq!(
|
||||
PSM_SERVICE_DATA_UUID,
|
||||
"00009C90-0000-1000-8000-00805F9B34FB"
|
||||
.parse::<bluer::Uuid>()
|
||||
.unwrap()
|
||||
);
|
||||
}
|
||||
}
|
||||
+1090
-79
File diff suppressed because it is too large
Load Diff
@@ -1,8 +1,9 @@
|
||||
//! BLE neighbor detection via advertising and scanning.
|
||||
//!
|
||||
//! BLE advertisements carry a 128-bit FIPS service UUID for identification.
|
||||
//! Post-forklift, advertisements are UUID-only (no identity material);
|
||||
//! identity is exchanged during the Noise handshake.
|
||||
//! BLE advertisements carry a 128-bit FIPS service UUID for identification,
|
||||
//! and optionally the advertiser's L2CAP listener PSM (see `super::psm`).
|
||||
//! Post-forklift they carry no identity material; identity is exchanged
|
||||
//! during the Noise handshake.
|
||||
|
||||
use crate::transport::{DiscoveredPeer, TransportId};
|
||||
use std::sync::Mutex;
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
//! Advertising the L2CAP listener PSM.
|
||||
//!
|
||||
//! A dialer has to know which PSM a peer's L2CAP listener is bound to. On
|
||||
//! BlueZ an application can *choose* that number, so both ends can agree on a
|
||||
//! configured constant. BlueZ is the exception: Android's
|
||||
//! `listenUsingInsecureL2capChannel` and macOS's
|
||||
//! `CBPeripheralManager.publishL2CAPChannel` both return an **OS-assigned**
|
||||
//! PSM the application cannot request. A dialer cannot guess it, and before a
|
||||
//! connection exists there is no channel to be told it on other than the
|
||||
//! advertisement itself.
|
||||
//!
|
||||
//! This module is the wire specification for putting it there. It is
|
||||
//! deliberately state-free: learning and caching belong to the scan/probe
|
||||
//! loop, which already owns per-address state.
|
||||
//!
|
||||
//! # Wire layout
|
||||
//!
|
||||
//! The PSM rides a **Service Data — 16-bit UUID** AD structure (AD type
|
||||
//! `0x16`) keyed on [`PSM_SERVICE_DATA_UUID16`], carrying the PSM as two
|
||||
//! bytes little-endian.
|
||||
//!
|
||||
//! ## Why a 16-bit key, and not the FIPS service UUID
|
||||
//!
|
||||
//! A legacy advertising PDU carries 31 bytes of AD payload. Keying the
|
||||
//! service data on the full 128-bit FIPS service UUID does not fit:
|
||||
//!
|
||||
//! | AD structure | bytes |
|
||||
//! |-----------------------------------------------|-------|
|
||||
//! | Flags | 3 |
|
||||
//! | Complete list of 128-bit service UUIDs | 18 |
|
||||
//! | Service Data — **128-bit** UUID + 2-byte PSM | 20 |
|
||||
//! | **total** | **41** — over by 10 |
|
||||
//!
|
||||
//! Keying it on the 16-bit UUID [`PSM_SERVICE_DATA_UUID16`] does:
|
||||
//!
|
||||
//! | AD structure | bytes |
|
||||
//! |-----------------------------------------------|-------|
|
||||
//! | Flags | 3 |
|
||||
//! | Complete list of 128-bit service UUIDs | 18 |
|
||||
//! | Service Data — **16-bit** UUID + 2-byte PSM | 6 |
|
||||
//! | **total** | **27** — fits |
|
||||
//!
|
||||
//! `0x9C90` is the leading 16 bits of the FIPS service UUID, expanded through
|
||||
//! the Bluetooth base UUID (`00009C90-0000-1000-8000-00805F9B34FB`). The
|
||||
//! budget is asserted at compile time below, so a change that reverts to a
|
||||
//! 128-bit key fails the build rather than the radio. It also means an
|
||||
//! advertiser using this layout has no room left for a local name.
|
||||
//!
|
||||
//! ## Why the primary advertisement, not the scan response
|
||||
//!
|
||||
//! A scan response only arrives after a successful active-scan
|
||||
//! request/response round-trip, and that round-trip drops asymmetrically
|
||||
//! across chipsets. Peers that never answer a scan request would become
|
||||
//! undiscoverable rather than merely slower. The primary advertisement is
|
||||
//! received passively on every advertising interval, so the PSM must ride it.
|
||||
//!
|
||||
//! ## Compatibility
|
||||
//!
|
||||
//! A reader ignores trailing bytes, so the value can be extended without
|
||||
//! breaking older peers, and an advert with no service data at all decodes to
|
||||
//! `None` — which is what every legacy UUID-only advertiser produces, and
|
||||
//! what makes them keep working against the configured PSM.
|
||||
|
||||
/// Service-data key for the advertised L2CAP PSM.
|
||||
///
|
||||
/// The leading 16 bits of the FIPS service UUID, i.e. the Bluetooth
|
||||
/// base-range UUID `00009C90-0000-1000-8000-00805F9B34FB`. Backends that
|
||||
/// speak in whole UUIDs must expand it through the base UUID; backends that
|
||||
/// speak in AD structures emit it as AD type `0x16`.
|
||||
pub const PSM_SERVICE_DATA_UUID16: u16 = 0x9C90;
|
||||
|
||||
/// AD payload budget of a legacy advertising PDU, in bytes.
|
||||
const LEGACY_ADV_PAYLOAD_BYTES: usize = 31;
|
||||
|
||||
/// Flags AD structure: length + type + one byte of flags.
|
||||
const FLAGS_AD_BYTES: usize = 3;
|
||||
|
||||
/// Complete list of 128-bit service UUIDs: length + type + one UUID.
|
||||
const UUID128_LIST_AD_BYTES: usize = 2 + 16;
|
||||
|
||||
/// Service data keyed on a 16-bit UUID: length + type + key + PSM.
|
||||
const PSM_SERVICE_DATA_AD_BYTES: usize = 2 + 2 + PSM_ENCODED_LEN;
|
||||
|
||||
/// Encoded width of the PSM value itself.
|
||||
const PSM_ENCODED_LEN: usize = 2;
|
||||
|
||||
/// The layout above must fit a legacy advertising PDU. If this fails, the
|
||||
/// advert would be silently truncated or rejected by the controller.
|
||||
const _: () = assert!(
|
||||
FLAGS_AD_BYTES + UUID128_LIST_AD_BYTES + PSM_SERVICE_DATA_AD_BYTES <= LEGACY_ADV_PAYLOAD_BYTES,
|
||||
"PSM advert layout exceeds the 31-byte legacy advertising PDU"
|
||||
);
|
||||
|
||||
/// Encode a PSM as advertised service data: two bytes, little-endian.
|
||||
pub fn encode_psm(psm: u16) -> [u8; PSM_ENCODED_LEN] {
|
||||
psm.to_le_bytes()
|
||||
}
|
||||
|
||||
/// Decode a PSM from advertised service data.
|
||||
///
|
||||
/// Returns `None` for absent or truncated data — a legacy UUID-only
|
||||
/// advertiser, which the caller answers by dialling the configured PSM.
|
||||
/// Trailing bytes are ignored so the value can be extended later without
|
||||
/// breaking readers built against this version.
|
||||
pub fn decode_psm(data: &[u8]) -> Option<u16> {
|
||||
if data.len() < PSM_ENCODED_LEN {
|
||||
return None;
|
||||
}
|
||||
Some(u16::from_le_bytes([data[0], data[1]]))
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Tests
|
||||
// ============================================================================
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_encode_is_little_endian() {
|
||||
assert_eq!(encode_psm(0x0085), [0x85, 0x00]);
|
||||
assert_eq!(encode_psm(0x1234), [0x34, 0x12]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_round_trip() {
|
||||
for psm in [0u16, 1, 0x0085, 0x00FF, 0x1234, u16::MAX] {
|
||||
assert_eq!(decode_psm(&encode_psm(psm)), Some(psm), "psm {psm:#06x}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_absent_service_data_decodes_to_none() {
|
||||
// A legacy UUID-only advertiser carries no service data at all.
|
||||
assert_eq!(decode_psm(&[]), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_truncated_service_data_decodes_to_none() {
|
||||
assert_eq!(decode_psm(&[0x85]), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_trailing_bytes_are_ignored() {
|
||||
// Forward compatibility: a future advertiser may append fields.
|
||||
assert_eq!(decode_psm(&[0x85, 0x00, 0xFF, 0xFF]), Some(0x0085));
|
||||
}
|
||||
|
||||
/// The byte budget is a build-time assertion, not a comment. This test
|
||||
/// records the arithmetic it encodes so the numbers stay legible.
|
||||
#[test]
|
||||
fn test_advert_fits_the_legacy_pdu() {
|
||||
assert_eq!(FLAGS_AD_BYTES, 3);
|
||||
assert_eq!(UUID128_LIST_AD_BYTES, 18);
|
||||
assert_eq!(PSM_SERVICE_DATA_AD_BYTES, 6);
|
||||
assert_eq!(
|
||||
FLAGS_AD_BYTES + UUID128_LIST_AD_BYTES + PSM_SERVICE_DATA_AD_BYTES,
|
||||
27
|
||||
);
|
||||
assert_eq!(LEGACY_ADV_PAYLOAD_BYTES, 31);
|
||||
// A 128-bit service-data key would need 20 bytes, not 6 — the layout
|
||||
// this module exists to reject.
|
||||
assert_eq!(FLAGS_AD_BYTES + UUID128_LIST_AD_BYTES + 20, 41);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_key_is_the_leading_16_bits_of_the_fips_uuid() {
|
||||
// FIPS service UUID: 9c90b790-2cc5-42c0-9f87-c9cc40648f4c
|
||||
const FIPS_SERVICE_UUID_U128: u128 = 0x9c90_b790_2cc5_42c0_9f87_c9cc_4064_8f4c;
|
||||
assert_eq!(
|
||||
PSM_SERVICE_DATA_UUID16,
|
||||
(FIPS_SERVICE_UUID_U128 >> 112) as u16
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,12 @@
|
||||
//! BLE transport statistics.
|
||||
//!
|
||||
//! Counters reach an operator through `show_transports`, which serves them
|
||||
//! off the control socket. Each connect outcome also emits a `debug!` at the
|
||||
//! moment it happens, carrying a uniform field set — `addr`, `role`
|
||||
//! (`central` for a dial, `peripheral` for an accept), `outcome` (a stable
|
||||
//! kebab-case string matching the counter name), and `discovery_ms` where a
|
||||
//! probe stamp exists. The counters give the aggregate; the trace stream
|
||||
//! gives the same taxonomy per event and per peer.
|
||||
|
||||
use portable_atomic::{AtomicU64, Ordering};
|
||||
|
||||
@@ -20,6 +28,8 @@ pub struct BleStats {
|
||||
pub connections_accepted: AtomicU64,
|
||||
pub connections_rejected: AtomicU64,
|
||||
pub connect_timeouts: AtomicU64,
|
||||
/// Outbound connects that failed with an error rather than timing out.
|
||||
pub connect_errors: AtomicU64,
|
||||
pub pool_evictions: AtomicU64,
|
||||
pub advertisements_sent: AtomicU64,
|
||||
pub scan_results: AtomicU64,
|
||||
@@ -40,6 +50,7 @@ impl BleStats {
|
||||
connections_accepted: AtomicU64::new(0),
|
||||
connections_rejected: AtomicU64::new(0),
|
||||
connect_timeouts: AtomicU64::new(0),
|
||||
connect_errors: AtomicU64::new(0),
|
||||
pool_evictions: AtomicU64::new(0),
|
||||
advertisements_sent: AtomicU64::new(0),
|
||||
scan_results: AtomicU64::new(0),
|
||||
@@ -93,6 +104,15 @@ impl BleStats {
|
||||
self.connect_timeouts.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Record an outbound connect that failed with an error.
|
||||
///
|
||||
/// Kept separate from [`Self::record_connect_timeout`]: a refusal and a
|
||||
/// silence are different faults and blur into one useless number if
|
||||
/// merged.
|
||||
pub fn record_connect_error(&self) {
|
||||
self.connect_errors.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// Record a pool eviction (non-static peer displaced).
|
||||
pub fn record_pool_eviction(&self) {
|
||||
self.pool_evictions.fetch_add(1, Ordering::Relaxed);
|
||||
@@ -122,6 +142,7 @@ impl BleStats {
|
||||
connections_accepted: self.connections_accepted.load(Ordering::Relaxed),
|
||||
connections_rejected: self.connections_rejected.load(Ordering::Relaxed),
|
||||
connect_timeouts: self.connect_timeouts.load(Ordering::Relaxed),
|
||||
connect_errors: self.connect_errors.load(Ordering::Relaxed),
|
||||
pool_evictions: self.pool_evictions.load(Ordering::Relaxed),
|
||||
advertisements_sent: self.advertisements_sent.load(Ordering::Relaxed),
|
||||
scan_results: self.scan_results.load(Ordering::Relaxed),
|
||||
@@ -149,6 +170,7 @@ pub struct BleStatsSnapshot {
|
||||
pub connections_accepted: u64,
|
||||
pub connections_rejected: u64,
|
||||
pub connect_timeouts: u64,
|
||||
pub connect_errors: u64,
|
||||
pub pool_evictions: u64,
|
||||
pub advertisements_sent: u64,
|
||||
pub scan_results: u64,
|
||||
|
||||
@@ -0,0 +1,295 @@
|
||||
//! `AsyncRead` adapter over a `BleStream`.
|
||||
//!
|
||||
//! The BLE receive path used to treat one `recv()` as one whole FIPS
|
||||
//! packet. That is a property of a *SeqPacket* socket, not of L2CAP: a
|
||||
//! stream-oriented backend may return a fragment of a packet, or several
|
||||
//! packets coalesced, from a single read. This adapter turns the
|
||||
//! datagram-shaped [`BleStream`] into the [`AsyncRead`] that
|
||||
//! `crate::transport::framing::read_fmp_packet` expects, buffering bytes
|
||||
//! left over from one read into the next so packet boundaries are recovered
|
||||
//! from the FMP length prefix rather than trusted to the OS.
|
||||
//!
|
||||
//! Nothing here is backend-specific. On a boundary-preserving backend the
|
||||
//! adapter is a transparent pass-through: one `recv` fills the buffer, the
|
||||
//! framer consumes exactly it, and the next read hits the underlying stream
|
||||
//! again.
|
||||
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
use tokio::io::{AsyncRead, ReadBuf};
|
||||
|
||||
use crate::transport::TransportError;
|
||||
|
||||
use super::io::BleStream;
|
||||
|
||||
/// Smallest scratch buffer used for a single `recv`.
|
||||
///
|
||||
/// Guards against a backend reporting a degenerate receive MTU, which would
|
||||
/// otherwise make every `recv` return zero bytes and look like EOF.
|
||||
const MIN_RECV_CHUNK: usize = 64;
|
||||
|
||||
/// A pending `recv` that owns its scratch buffer and yields an owned `Vec`.
|
||||
///
|
||||
/// Owning the buffer is what makes the future `'static`, which is what lets
|
||||
/// it be held across `poll_read` calls when a read returns `Pending`.
|
||||
type RecvFuture = Pin<Box<dyn Future<Output = Result<Vec<u8>, TransportError>> + Send>>;
|
||||
|
||||
/// Buffered [`AsyncRead`] view of a [`BleStream`].
|
||||
pub struct BleStreamRead<S: BleStream + 'static> {
|
||||
stream: Arc<S>,
|
||||
/// Bytes received but not yet handed to the reader.
|
||||
chunk: Vec<u8>,
|
||||
/// Read cursor into `chunk`.
|
||||
pos: usize,
|
||||
/// Scratch size for one underlying `recv`.
|
||||
capacity: usize,
|
||||
/// In-flight `recv`, kept across polls.
|
||||
pending: Option<RecvFuture>,
|
||||
/// Set once the peer has closed the connection.
|
||||
eof: bool,
|
||||
}
|
||||
|
||||
impl<S: BleStream + 'static> BleStreamRead<S> {
|
||||
/// Wrap a stream, sizing the scratch buffer from its receive MTU.
|
||||
pub fn new(stream: Arc<S>, recv_mtu: u16) -> Self {
|
||||
Self {
|
||||
stream,
|
||||
chunk: Vec::new(),
|
||||
pos: 0,
|
||||
capacity: (recv_mtu as usize).max(MIN_RECV_CHUNK),
|
||||
pending: None,
|
||||
eof: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Number of bytes already received but not yet consumed.
|
||||
///
|
||||
/// Non-zero after a peer coalesces data behind an earlier message; the
|
||||
/// hand-off from the pubkey exchange to the framer must preserve them.
|
||||
#[cfg(test)]
|
||||
pub fn buffered(&self) -> usize {
|
||||
self.chunk.len() - self.pos
|
||||
}
|
||||
|
||||
fn start_recv(&self) -> RecvFuture {
|
||||
let stream = Arc::clone(&self.stream);
|
||||
let capacity = self.capacity;
|
||||
Box::pin(async move {
|
||||
let mut scratch = vec![0u8; capacity];
|
||||
let n = stream.recv(&mut scratch).await?;
|
||||
scratch.truncate(n);
|
||||
Ok(scratch)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// Map a transport error onto the `io::Error` `AsyncRead` must report.
|
||||
fn to_io(e: TransportError) -> std::io::Error {
|
||||
match e {
|
||||
TransportError::Io(e) => e,
|
||||
other => std::io::Error::other(other.to_string()),
|
||||
}
|
||||
}
|
||||
|
||||
impl<S: BleStream + 'static> AsyncRead for BleStreamRead<S> {
|
||||
fn poll_read(
|
||||
self: Pin<&mut Self>,
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut ReadBuf<'_>,
|
||||
) -> Poll<std::io::Result<()>> {
|
||||
let this = self.get_mut();
|
||||
loop {
|
||||
// Serve from the leftover buffer first.
|
||||
if this.pos < this.chunk.len() {
|
||||
let n = (this.chunk.len() - this.pos).min(buf.remaining());
|
||||
buf.put_slice(&this.chunk[this.pos..this.pos + n]);
|
||||
this.pos += n;
|
||||
if this.pos == this.chunk.len() {
|
||||
this.chunk.clear();
|
||||
this.pos = 0;
|
||||
}
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
// A closed connection stays closed: report EOF (a filled length
|
||||
// of zero) rather than re-polling a dead stream forever.
|
||||
if this.eof {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
|
||||
let mut fut = match this.pending.take() {
|
||||
Some(f) => f,
|
||||
None => this.start_recv(),
|
||||
};
|
||||
match fut.as_mut().poll(cx) {
|
||||
Poll::Pending => {
|
||||
this.pending = Some(fut);
|
||||
return Poll::Pending;
|
||||
}
|
||||
Poll::Ready(Ok(chunk)) => {
|
||||
// `recv` returning zero bytes is the peer-closed signal,
|
||||
// not an empty packet.
|
||||
if chunk.is_empty() {
|
||||
this.eof = true;
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
this.chunk = chunk;
|
||||
this.pos = 0;
|
||||
}
|
||||
Poll::Ready(Err(e)) => return Poll::Ready(Err(to_io(e))),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Tests
|
||||
// ============================================================================
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::transport::ble::addr::BleAddr;
|
||||
use crate::transport::ble::io::MockBleStream;
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::Mutex as TokioMutex;
|
||||
|
||||
fn test_addr(n: u8) -> BleAddr {
|
||||
BleAddr {
|
||||
adapter: "hci0".to_string(),
|
||||
device: [0xAA, 0xBB, 0xCC, 0xDD, 0xEE, n],
|
||||
}
|
||||
}
|
||||
|
||||
/// A stream that replays a fixed script of `recv` results and then
|
||||
/// returns `Ok(0)` forever — the "peer closed but socket still open"
|
||||
/// shape a channel-backed mock cannot produce.
|
||||
struct ScriptedStream {
|
||||
addr: BleAddr,
|
||||
chunks: TokioMutex<std::collections::VecDeque<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl ScriptedStream {
|
||||
fn new(chunks: Vec<Vec<u8>>) -> Self {
|
||||
Self {
|
||||
addr: test_addr(9),
|
||||
chunks: TokioMutex::new(chunks.into()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl BleStream for ScriptedStream {
|
||||
async fn send(&self, _data: &[u8]) -> Result<(), TransportError> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn recv(&self, buf: &mut [u8]) -> Result<usize, TransportError> {
|
||||
match self.chunks.lock().await.pop_front() {
|
||||
Some(chunk) => {
|
||||
let n = chunk.len().min(buf.len());
|
||||
buf[..n].copy_from_slice(&chunk[..n]);
|
||||
Ok(n)
|
||||
}
|
||||
None => Ok(0),
|
||||
}
|
||||
}
|
||||
|
||||
fn send_mtu(&self) -> u16 {
|
||||
2048
|
||||
}
|
||||
|
||||
fn recv_mtu(&self) -> u16 {
|
||||
2048
|
||||
}
|
||||
|
||||
fn remote_addr(&self) -> &BleAddr {
|
||||
&self.addr
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_fragmented_delivery_is_reassembled() {
|
||||
let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048);
|
||||
peer.send(b"abc").await.unwrap();
|
||||
peer.send(b"defg").await.unwrap();
|
||||
peer.send(b"hij").await.unwrap();
|
||||
|
||||
let mut reader = BleStreamRead::new(Arc::new(local), 2048);
|
||||
let mut out = [0u8; 10];
|
||||
reader.read_exact(&mut out).await.unwrap();
|
||||
assert_eq!(&out, b"abcdefghij");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_coalesced_delivery_keeps_the_tail() {
|
||||
let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048);
|
||||
peer.send(b"0123456789").await.unwrap();
|
||||
|
||||
let mut reader = BleStreamRead::new(Arc::new(local), 2048);
|
||||
let mut head = [0u8; 4];
|
||||
reader.read_exact(&mut head).await.unwrap();
|
||||
assert_eq!(&head, b"0123");
|
||||
assert_eq!(reader.buffered(), 6);
|
||||
|
||||
let mut tail = [0u8; 6];
|
||||
reader.read_exact(&mut tail).await.unwrap();
|
||||
assert_eq!(&tail, b"456789");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_peer_drop_surfaces_as_unexpected_eof() {
|
||||
let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048);
|
||||
peer.send(b"ab").await.unwrap();
|
||||
drop(peer);
|
||||
|
||||
let mut reader = BleStreamRead::new(Arc::new(local), 2048);
|
||||
let mut out = [0u8; 4];
|
||||
let err = reader.read_exact(&mut out).await.unwrap_err();
|
||||
assert_eq!(err.kind(), std::io::ErrorKind::UnexpectedEof);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_zero_length_recv_is_eof_not_readiness() {
|
||||
let stream = ScriptedStream::new(vec![b"xy".to_vec()]);
|
||||
let mut reader = BleStreamRead::new(Arc::new(stream), 2048);
|
||||
|
||||
let mut out = [0u8; 8];
|
||||
let n = reader.read(&mut out).await.unwrap();
|
||||
assert_eq!(&out[..n], b"xy");
|
||||
|
||||
// The scripted stream now returns Ok(0) forever. That must read as
|
||||
// EOF once and stay EOF, not as a spurious zero-length packet.
|
||||
assert_eq!(reader.read(&mut out).await.unwrap(), 0);
|
||||
assert_eq!(reader.read(&mut out).await.unwrap(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_read_smaller_than_chunk_leaves_remainder() {
|
||||
let stream = ScriptedStream::new(vec![b"abcdef".to_vec()]);
|
||||
let mut reader = BleStreamRead::new(Arc::new(stream), 2048);
|
||||
|
||||
let mut one = [0u8; 1];
|
||||
reader.read_exact(&mut one).await.unwrap();
|
||||
assert_eq!(&one, b"a");
|
||||
assert_eq!(reader.buffered(), 5);
|
||||
|
||||
let mut rest = [0u8; 5];
|
||||
reader.read_exact(&mut rest).await.unwrap();
|
||||
assert_eq!(&rest, b"bcdef");
|
||||
assert_eq!(reader.buffered(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_degenerate_recv_mtu_still_reads() {
|
||||
let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048);
|
||||
peer.send(b"hello").await.unwrap();
|
||||
|
||||
let mut reader = BleStreamRead::new(Arc::new(local), 0);
|
||||
let mut out = [0u8; 5];
|
||||
reader.read_exact(&mut out).await.unwrap();
|
||||
assert_eq!(&out, b"hello");
|
||||
}
|
||||
}
|
||||
+23
-23
@@ -15,10 +15,10 @@ pub mod udp;
|
||||
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||
pub mod ethernet;
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
pub mod ble;
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
use ble::DefaultBleTransport;
|
||||
#[cfg(any(target_os = "linux", target_os = "macos"))]
|
||||
use ethernet::EthernetTransport;
|
||||
@@ -673,7 +673,7 @@ pub enum TransportHandle {
|
||||
/// Nym mixnet transport (via SOCKS5).
|
||||
Nym(NymTransport),
|
||||
/// BLE L2CAP transport.
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
Ble(DefaultBleTransport),
|
||||
/// In-process loopback transport (test harness only).
|
||||
#[cfg(test)]
|
||||
@@ -690,7 +690,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.start_async().await,
|
||||
TransportHandle::Tor(t) => t.start_async().await,
|
||||
TransportHandle::Nym(t) => t.start_async().await,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.start_async().await,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.start_async().await,
|
||||
@@ -706,7 +706,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.stop_async().await,
|
||||
TransportHandle::Tor(t) => t.stop_async().await,
|
||||
TransportHandle::Nym(t) => t.stop_async().await,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.stop_async().await,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.stop_async().await,
|
||||
@@ -722,7 +722,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.send_async(addr, data).await,
|
||||
TransportHandle::Tor(t) => t.send_async(addr, data).await,
|
||||
TransportHandle::Nym(t) => t.send_async(addr, data).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.send_async(addr, data).await,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.send_async(addr, data).await,
|
||||
@@ -738,7 +738,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.transport_id(),
|
||||
TransportHandle::Tor(t) => t.transport_id(),
|
||||
TransportHandle::Nym(t) => t.transport_id(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.transport_id(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.transport_id(),
|
||||
@@ -754,7 +754,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.name(),
|
||||
TransportHandle::Tor(t) => t.name(),
|
||||
TransportHandle::Nym(t) => t.name(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.name(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => None,
|
||||
@@ -770,7 +770,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.transport_type(),
|
||||
TransportHandle::Tor(t) => t.transport_type(),
|
||||
TransportHandle::Nym(t) => t.transport_type(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.transport_type(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.transport_type(),
|
||||
@@ -786,7 +786,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.state(),
|
||||
TransportHandle::Tor(t) => t.state(),
|
||||
TransportHandle::Nym(t) => t.state(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.state(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.state(),
|
||||
@@ -802,7 +802,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.mtu(),
|
||||
TransportHandle::Tor(t) => t.mtu(),
|
||||
TransportHandle::Nym(t) => t.mtu(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.mtu(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.mtu(),
|
||||
@@ -821,7 +821,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.link_mtu(addr),
|
||||
TransportHandle::Tor(t) => t.link_mtu(addr),
|
||||
TransportHandle::Nym(t) => t.link_mtu(addr),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.link_mtu(addr),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.link_mtu(addr),
|
||||
@@ -837,7 +837,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.local_addr(),
|
||||
TransportHandle::Tor(_) => None,
|
||||
TransportHandle::Nym(_) => None,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(_) => None,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => None,
|
||||
@@ -857,7 +857,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(_) => None,
|
||||
TransportHandle::Tor(_) => None,
|
||||
TransportHandle::Nym(_) => None,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(_) => None,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => None,
|
||||
@@ -873,7 +873,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(_) => None,
|
||||
TransportHandle::Tor(_) => None,
|
||||
TransportHandle::Nym(_) => None,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(_) => None,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => None,
|
||||
@@ -913,7 +913,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.discover(),
|
||||
TransportHandle::Tor(t) => t.discover(),
|
||||
TransportHandle::Nym(t) => t.discover(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.discover(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.discover(),
|
||||
@@ -929,7 +929,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.auto_connect(),
|
||||
TransportHandle::Tor(t) => t.auto_connect(),
|
||||
TransportHandle::Nym(t) => t.auto_connect(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.auto_connect(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.auto_connect(),
|
||||
@@ -945,7 +945,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.accept_connections(),
|
||||
TransportHandle::Tor(t) => t.accept_connections(),
|
||||
TransportHandle::Nym(t) => t.accept_connections(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.accept_connections(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(t) => t.accept_connections(),
|
||||
@@ -967,7 +967,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.connect_async(addr).await,
|
||||
TransportHandle::Tor(t) => t.connect_async(addr).await,
|
||||
TransportHandle::Nym(t) => t.connect_async(addr).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.connect_async(addr).await,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => Ok(()), // connectionless
|
||||
@@ -987,7 +987,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.connection_state_sync(addr),
|
||||
TransportHandle::Tor(t) => t.connection_state_sync(addr),
|
||||
TransportHandle::Nym(t) => t.connection_state_sync(addr),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.connection_state_sync(addr),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => ConnectionState::Connected,
|
||||
@@ -1006,7 +1006,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(t) => t.close_connection_async(addr).await,
|
||||
TransportHandle::Tor(t) => t.close_connection_async(addr).await,
|
||||
TransportHandle::Nym(t) => t.close_connection_async(addr).await,
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => t.close_connection_async(addr).await,
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => {} // connectionless no-op
|
||||
@@ -1031,7 +1031,7 @@ impl TransportHandle {
|
||||
TransportHandle::Tcp(_) => TransportCongestion::default(),
|
||||
TransportHandle::Tor(_) => TransportCongestion::default(),
|
||||
TransportHandle::Nym(_) => TransportCongestion::default(),
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(_) => TransportCongestion::default(),
|
||||
#[cfg(test)]
|
||||
TransportHandle::Loopback(_) => TransportCongestion::default(),
|
||||
@@ -1059,7 +1059,7 @@ impl TransportHandle {
|
||||
TransportHandle::Nym(t) => {
|
||||
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
|
||||
}
|
||||
#[cfg(target_os = "linux")]
|
||||
#[cfg(ble_available)]
|
||||
TransportHandle::Ble(t) => {
|
||||
serde_json::to_value(t.stats().snapshot()).unwrap_or_default()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user