From 814ceef680c65ff9f2dd220964d5a75c3c49a021 Mon Sep 17 00:00:00 2001 From: Arjen <18398758+Origami74@users.noreply.github.com> Date: Sun, 23 Aug 2026 17:51:58 +0100 Subject: [PATCH] fix(transport/ble): recover packet boundaries from the byte stream MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The BLE receive path assumes one `recv()` returns exactly one whole FIPS packet. That holds only for BlueZ's `SOCK_SEQPACKET`, which preserves SDU boundaries. It is not a property of L2CAP — it is a property of one backend's socket type. A stream-oriented L2CAP backend (Android's `BluetoothSocket` input stream, CoreBluetooth's `CBL2CAPChannel`) may return a fragment of a packet or several packets coalesced in a single read. Under the current loop a fragment ships up as a runt that FMP and Noise reject, and a coalesced tail is silently truncated and dropped — packets are lost and the transport thrashes with no error to show for it. Recover the boundaries from the bytes instead of trusting the OS to preserve them. FIPS packets are self-delimiting via the 4-byte FMP common prefix, and `transport::framing::read_fmp_packet` already parses exactly that, shared by every stream-oriented transport. All that was missing is an adapter: a new `BleStreamRead` turns the datagram-shaped `BleStream` into the `AsyncRead` the framer expects, buffering bytes left over from one read into the next. Its pending-read future owns its scratch buffer and yields an owned `Vec`, so it is `'static` and can be held across `poll_read` calls. One reader is threaded through both phases of a connection — the pre-handshake pubkey exchange and then the receive loop — so bytes a peer coalesces after the 33-byte pubkey stay buffered rather than being dropped at the hand-off. The exchange now reads via `read_exact`, which also reassembles a fragmented pubkey, instead of a single `recv()` with an exact-length check that a fragmenting backend can never satisfy. That ordering is load-bearing and worth stating plainly: the pubkey message's `0x00` prefix decodes as FMP version 0, phase 0 (established), with a payload length read out of the pubkey's own bytes. If the framer ever saw the exchange it would mis-frame badly. The exchange must be fully consumed before the framer starts, which is exactly what threading one reader guarantees. The constant now says so. On BlueZ this is a transparent pass-through — one `recv` already is one packet, so the adapter serves it whole and the framer takes it whole. On a stream backend it reassembles. Either way the layer above sees one complete packet per read, identically on every platform. Coverage: all of this is shared code exercised by `MockBleIo` on the host, so `cargo test` and `cargo clippy --all-targets` cover it in full on any platform where `transport::ble` compiles. Nothing here is inside `cfg(bluer_available)`, so no part of the change depends on a Bluetooth adapter to be checked. --- src/transport/ble/mod.rs | 335 ++++++++++++++++++++++++++----- src/transport/ble/stream_read.rs | 295 +++++++++++++++++++++++++++ 2 files changed, 584 insertions(+), 46 deletions(-) create mode 100644 src/transport/ble/stream_read.rs diff --git a/src/transport/ble/mod.rs b/src/transport/ble/mod.rs index 098868fd..92d57062 100644 --- a/src/transport/ble/mod.rs +++ b/src/transport/ble/mod.rs @@ -1,9 +1,18 @@ //! BLE L2CAP Transport Implementation //! -//! Provides BLE-based transport for FIPS peer communication using L2CAP -//! Connection-Oriented Channels (CoC) in SeqPacket mode. L2CAP CoC -//! preserves message boundaries (unlike TCP byte streams), so no FMP -//! framing is needed — each send/recv is one FIPS packet. +//! Provides BLE-based transport for FIPS peer communication over L2CAP +//! Connection-Oriented Channels. +//! +//! ## Packet boundaries +//! +//! Message-boundary preservation is a property of the *socket type* a +//! backend uses, not of L2CAP. BlueZ's `SOCK_SEQPACKET` preserves SDU +//! boundaries; other backends expose an L2CAP channel as a byte stream and +//! may return a fragment of a packet or several packets coalesced from one +//! read. The receive path therefore recovers boundaries from the FMP length +//! prefix via [`stream_read::BleStreamRead`] and +//! [`crate::transport::framing::read_fmp_packet`], which is a transparent +//! pass-through on a boundary-preserving backend. //! //! ## Architecture //! @@ -23,7 +32,9 @@ pub mod io; pub mod neighbor; pub mod pool; pub mod stats; +pub mod stream_read; +use super::framing::{StreamError, read_fmp_packet}; use super::{ ConnectionState, DiscoveredPeer, PacketTx, ReceivedPacket, Transport, TransportAddr, TransportError, TransportId, TransportState, TransportType, @@ -35,10 +46,12 @@ use io::{BleIo, BleScanner, BleStream}; use neighbor::NeighborBuffer; use pool::{BleConnection, ConnectionPool}; use stats::BleStats; +use stream_read::BleStreamRead; use secp256k1::XOnlyPublicKey; use std::collections::HashMap; use std::sync::Arc; +use tokio::io::AsyncReadExt; use tokio::sync::Mutex; use tokio::task::JoinHandle; use tracing::{debug, info, trace, warn}; @@ -364,9 +377,16 @@ impl BleTransport { } }; + // One reader for the life of the connection: the pubkey exchange and + // the receive loop must share it, or bytes the peer coalesced behind + // the exchange are dropped at the hand-off. + let stream = Arc::new(stream); + let recv_mtu = stream.recv_mtu(); + let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); + // Pre-handshake pubkey exchange (temporary, pre-XX) if let Some(ref our_pubkey) = self.local_pubkey { - match pubkey_exchange(&stream, our_pubkey).await { + match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr, "BLE outbound pubkey exchange complete"); self.neighbor_buffer @@ -379,7 +399,8 @@ impl BleTransport { } } - self.promote_connection(addr, &ble_addr, stream).await + self.promote_connection(addr, &ble_addr, stream, reader) + .await } /// Promote a newly established stream into the connection pool. @@ -389,14 +410,14 @@ impl BleTransport { &self, addr: &TransportAddr, ble_addr: &BleAddr, - stream: I::Stream, + stream: Arc, + reader: BleStreamRead, ) -> Result<(), TransportError> { let send_mtu = stream.send_mtu(); let recv_mtu = stream.recv_mtu(); - let stream = Arc::new(stream); let recv_task = tokio::spawn(receive_loop( - Arc::clone(&stream), + reader, addr.clone(), Arc::clone(&self.pool), self.packet_tx.clone(), @@ -484,9 +505,15 @@ impl BleTransport { match result { Ok(Ok(stream)) => { + let send_mtu = stream.send_mtu(); + let recv_mtu = stream.recv_mtu(); + let stream = Arc::new(stream); + // One reader across both phases — see `pubkey_exchange`. + let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); + // Pre-handshake pubkey exchange (temporary, pre-XX) if let Some(ref our_pubkey) = local_pubkey { - match pubkey_exchange(&stream, our_pubkey).await { + match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete"); neighbor_buffer.add_peer_with_pubkey(&ble_addr, peer_pubkey); @@ -501,12 +528,8 @@ impl BleTransport { } } - let send_mtu = stream.send_mtu(); - let recv_mtu = stream.recv_mtu(); - let stream = Arc::new(stream); - let recv_task = tokio::spawn(receive_loop( - Arc::clone(&stream), + reader, addr_clone.clone(), Arc::clone(&pool), packet_tx, @@ -663,6 +686,13 @@ impl Transport for BleTransport { /// /// Distinguishes the identity exchange from FMP packets (version ≥ 0x01). /// Temporary — removed when FMP switches from IK to XX handshake. +/// +/// Caution: this prefix is *not* distinguishable from an FMP packet by the +/// framer. `0x00` decodes as FMP version 0, phase 0 (established), with a +/// payload length read out of the pubkey's own bytes — i.e. arbitrary. Any +/// code that runs the framer over a connection before the exchange has been +/// fully consumed will mis-frame badly. Threading one reader through both +/// phases is what guarantees the ordering. const PUBKEY_EXCHANGE_PREFIX: u8 = 0x00; /// Pre-handshake pubkey exchange message size: `[0x00][pubkey:32]`. @@ -679,8 +709,16 @@ const PUBKEY_EXCHANGE_TIMEOUT_SECS: u64 = 5; /// /// Both sides send `[0x00][our_pubkey:32]` and receive the peer's. /// Returns the peer's XOnlyPublicKey on success. -async fn pubkey_exchange( +/// +/// Reads through the connection's `BleStreamRead` rather than calling +/// `recv` directly, for two reasons. It reassembles an exchange a +/// stream-oriented backend fragmented, which a single `recv` with an +/// exact-length check can never do. And anything the peer coalesced behind +/// the exchange stays buffered in the reader that the receive loop then +/// takes over, instead of being discarded at the hand-off. +async fn pubkey_exchange( stream: &S, + reader: &mut BleStreamRead, local_pubkey: &[u8; 32], ) -> Result { // Send our pubkey @@ -692,15 +730,15 @@ async fn pubkey_exchange( // Receive peer's pubkey (with timeout to prevent indefinite blocking) let mut buf = [0u8; PUBKEY_EXCHANGE_SIZE]; let timeout = std::time::Duration::from_secs(PUBKEY_EXCHANGE_TIMEOUT_SECS); - let n = match tokio::time::timeout(timeout, stream.recv(&mut buf)).await { - Ok(result) => result?, + match tokio::time::timeout(timeout, reader.read_exact(&mut buf)).await { + Ok(Ok(_)) => {} + Ok(Err(e)) => { + return Err(TransportError::RecvFailed(format!( + "pubkey exchange: {}", + e + ))); + } Err(_) => return Err(TransportError::Timeout), - }; - if n != PUBKEY_EXCHANGE_SIZE { - return Err(TransportError::RecvFailed(format!( - "pubkey exchange: expected {} bytes, got {}", - PUBKEY_EXCHANGE_SIZE, n - ))); } if buf[0] != PUBKEY_EXCHANGE_PREFIX { return Err(TransportError::RecvFailed(format!( @@ -751,10 +789,13 @@ async fn accept_loop( let send_mtu = stream.send_mtu(); let recv_mtu = stream.recv_mtu(); + let stream = Arc::new(stream); + // One reader across both phases — see `pubkey_exchange`. + let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); // Pre-handshake pubkey exchange (temporary, pre-XX) if let Some(ref our_pubkey) = local_pubkey { - match pubkey_exchange(&stream, our_pubkey).await { + match pubkey_exchange(stream.as_ref(), &mut reader, our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %ta, "BLE inbound pubkey exchange complete"); neighbor_buffer.add_peer_with_pubkey(&addr, peer_pubkey); @@ -780,11 +821,9 @@ async fn accept_loop( } } - let stream = Arc::new(stream); - // Spawn receive loop let recv_task = tokio::spawn(receive_loop( - Arc::clone(&stream), + reader, ta.clone(), Arc::clone(&pool), packet_tx.clone(), @@ -829,8 +868,14 @@ async fn accept_loop( } /// Receive loop: reads packets from a BLE stream and delivers to node. -async fn receive_loop( - stream: Arc, +/// +/// Takes the connection's `BleStreamRead` — already positioned past the +/// pubkey exchange, and still holding anything the peer coalesced behind it +/// — and pulls whole FIPS packets out of it using the FMP length prefix. +/// Boundaries come from the bytes, not from the backend's socket type, so a +/// fragment is reassembled and a coalesced tail is not lost. +async fn receive_loop( + mut reader: BleStreamRead, addr: TransportAddr, pool: Arc>>>, packet_tx: PacketTx, @@ -838,21 +883,20 @@ async fn receive_loop( stats: Arc, recv_mtu: u16, ) { - let mut buf = vec![0u8; recv_mtu as usize]; loop { - match stream.recv(&mut buf).await { - Ok(0) => { - debug!(addr = %addr, "BLE connection closed by peer"); - break; - } - Ok(n) => { - stats.record_recv(n); - let packet = ReceivedPacket::new(transport_id, addr.clone(), buf[..n].to_vec()); + match read_fmp_packet(&mut reader, recv_mtu).await { + Ok(data) => { + stats.record_recv(data.len()); + let packet = ReceivedPacket::new(transport_id, addr.clone(), data); if packet_tx.send(packet).await.is_err() { trace!("BLE packet_tx closed, stopping receive loop"); break; } } + Err(StreamError::Io(e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => { + debug!(addr = %addr, "BLE connection closed by peer"); + break; + } Err(e) => { debug!(addr = %addr, error = %e, "BLE receive error"); stats.record_recv_error(); @@ -986,7 +1030,12 @@ async fn scan_probe_loop( // Pubkey exchange, then promote connection to pool let ta = addr.to_transport_addr(); - match pubkey_exchange(&stream, &our_pubkey).await { + let send_mtu = stream.send_mtu(); + let recv_mtu = stream.recv_mtu(); + let stream = Arc::new(stream); + // One reader across both phases — see `pubkey_exchange`. + let mut reader = BleStreamRead::new(Arc::clone(&stream), recv_mtu); + match pubkey_exchange(stream.as_ref(), &mut reader, &our_pubkey).await { Ok(peer_pubkey) => { debug!(addr = %addr, "BLE probe complete"); @@ -1005,12 +1054,8 @@ async fn scan_probe_loop( } // Promote connection to pool — no second L2CAP connect needed - let send_mtu = stream.send_mtu(); - let recv_mtu = stream.recv_mtu(); - let stream = Arc::new(stream); - let recv_task = tokio::spawn(receive_loop( - Arc::clone(&stream), + reader, ta.clone(), Arc::clone(&pool), packet_tx.clone(), @@ -1064,7 +1109,42 @@ async fn scan_probe_loop( #[cfg(test)] mod tests { use super::*; - use io::MockBleIo; + use crate::transport::framing::build_established_frame; + use io::{MockBleIo, MockBleStream}; + use secp256k1::{Secp256k1, SecretKey}; + + /// Deterministic x-only pubkey for exchange tests. + fn test_pubkey(seed: u8) -> [u8; 32] { + let secp = Secp256k1::new(); + let sk = SecretKey::from_slice(&[seed; 32]).unwrap(); + sk.public_key(&secp).x_only_public_key().0.serialize() + } + + /// Handles a receive-loop test needs to observe: the task, the packets + /// it delivers, and the pool it reaps its entry from. + type ReceiveLoopHarness = ( + JoinHandle<()>, + tokio::sync::mpsc::Receiver, + Arc>>>, + ); + + /// Wire up a receive loop over one end of a mock stream pair. + fn spawn_receive_loop(local: MockBleStream) -> ReceiveLoopHarness { + let addr = test_addr(2).to_transport_addr(); + let pool = Arc::new(Mutex::new(ConnectionPool::new(7))); + let (tx, rx) = tokio::sync::mpsc::channel(16); + let reader = BleStreamRead::new(Arc::new(local), 2048); + let task = tokio::spawn(receive_loop( + reader, + addr, + Arc::clone(&pool), + tx, + TransportId::new(1), + Arc::new(BleStats::new()), + 2048, + )); + (task, rx, pool) + } fn test_addr(n: u8) -> BleAddr { BleAddr { @@ -1208,4 +1288,167 @@ mod tests { // Smaller node accepting from larger → drops inbound (outbound wins) // This means: smaller always uses outbound, larger always uses inbound } + + // ------------------------------------------------------------------ + // Packet boundary recovery + // ------------------------------------------------------------------ + + /// Two whole FMP packets delivered in one `recv` must both arrive. + /// Before reframing the tail was silently truncated and lost. + #[tokio::test] + async fn test_receive_loop_splits_coalesced_packets() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let (task, mut rx, _pool) = spawn_receive_loop(local); + + let first = build_established_frame(16); + let second = build_established_frame(48); + let mut both = first.clone(); + both.extend_from_slice(&second); + peer.send(&both).await.unwrap(); + + assert_eq!(rx.recv().await.unwrap().data, first); + assert_eq!(rx.recv().await.unwrap().data, second); + task.abort(); + } + + /// One FMP packet split across three `recv`s arrives once, whole — + /// not as three runts that FMP and Noise would reject. + #[tokio::test] + async fn test_receive_loop_reassembles_fragmented_packet() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let (task, mut rx, _pool) = spawn_receive_loop(local); + + let frame = build_established_frame(64); + let third = frame.len() / 3; + peer.send(&frame[..third]).await.unwrap(); + peer.send(&frame[third..2 * third]).await.unwrap(); + peer.send(&frame[2 * third..]).await.unwrap(); + + assert_eq!(rx.recv().await.unwrap().data, frame); + assert!(rx.try_recv().is_err(), "no runt packets"); + task.abort(); + } + + /// One `send` per packet still yields one packet per `send`, byte for + /// byte — the boundary-preserving backend regression. + #[tokio::test] + async fn test_receive_loop_passes_through_whole_packets() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let (task, mut rx, _pool) = spawn_receive_loop(local); + + let frames: Vec> = [8u16, 0, 512] + .iter() + .map(|n| build_established_frame(*n)) + .collect(); + for f in &frames { + peer.send(f).await.unwrap(); + } + for f in &frames { + assert_eq!(&rx.recv().await.unwrap().data, f); + } + task.abort(); + } + + /// A malformed frame closes the connection and drops it from the pool + /// rather than spinning the loop. + #[tokio::test] + async fn test_receive_loop_drops_connection_on_bad_frame() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let ta = test_addr(2).to_transport_addr(); + let (task, _rx, pool) = spawn_receive_loop(local); + + // Put a pool entry in place so its removal is observable. + let (parked, _other) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + pool.lock() + .await + .insert( + ta.clone(), + BleConnection { + stream: Arc::new(parked), + recv_task: None, + send_mtu: 2048, + recv_mtu: 2048, + established_at: tokio::time::Instant::now(), + is_static: false, + addr: test_addr(2), + }, + ) + .unwrap(); + assert!(pool.lock().await.contains(&ta)); + + // 0x16 is a TLS ClientHello record type; it parses as FMP version 1. + peer.send(&[0x16, 0x03, 0x01, 0x00]).await.unwrap(); + + // The loop exits and clears the pool entry. + for _ in 0..50 { + if !pool.lock().await.contains(&ta) { + break; + } + tokio::task::yield_now().await; + } + assert!(!pool.lock().await.contains(&ta)); + assert!(task.await.is_ok(), "loop exited cleanly"); + } + + /// A peer that coalesces its first data packet behind the 33-byte + /// pubkey exchange must not lose it at the hand-off to the framer. + #[tokio::test] + async fn test_pubkey_exchange_preserves_coalesced_data() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let local = Arc::new(local); + let mut reader = BleStreamRead::new(Arc::clone(&local), 2048); + + let peer_pk = test_pubkey(2); + let frame = build_established_frame(24); + let mut wire = vec![PUBKEY_EXCHANGE_PREFIX]; + wire.extend_from_slice(&peer_pk); + wire.extend_from_slice(&frame); + peer.send(&wire).await.unwrap(); + + let got = pubkey_exchange(local.as_ref(), &mut reader, &test_pubkey(1)) + .await + .unwrap(); + assert_eq!(got.serialize(), peer_pk); + + let packet = read_fmp_packet(&mut reader, 2048).await.unwrap(); + assert_eq!(packet, frame); + } + + /// A fragmented pubkey exchange completes. The old exact-length `recv` + /// check could never satisfy this. + #[tokio::test] + async fn test_pubkey_exchange_reassembles_fragments() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let local = Arc::new(local); + let mut reader = BleStreamRead::new(Arc::clone(&local), 2048); + + let peer_pk = test_pubkey(3); + let mut wire = vec![PUBKEY_EXCHANGE_PREFIX]; + wire.extend_from_slice(&peer_pk); + peer.send(&wire[..17]).await.unwrap(); + peer.send(&wire[17..]).await.unwrap(); + + let got = pubkey_exchange(local.as_ref(), &mut reader, &test_pubkey(1)) + .await + .unwrap(); + assert_eq!(got.serialize(), peer_pk); + } + + /// A peer that opens with something other than the exchange prefix is + /// rejected before the framer ever sees the bytes. + #[tokio::test] + async fn test_pubkey_exchange_rejects_bad_prefix() { + let (peer, local) = MockBleStream::pair(test_addr(1), test_addr(2), 2048); + let local = Arc::new(local); + let mut reader = BleStreamRead::new(Arc::clone(&local), 2048); + + let mut wire = vec![0xFFu8]; + wire.extend_from_slice(&test_pubkey(4)); + peer.send(&wire).await.unwrap(); + + let err = pubkey_exchange(local.as_ref(), &mut reader, &test_pubkey(1)) + .await + .unwrap_err(); + assert!(matches!(err, TransportError::RecvFailed(_))); + } } diff --git a/src/transport/ble/stream_read.rs b/src/transport/ble/stream_read.rs new file mode 100644 index 00000000..59b12d66 --- /dev/null +++ b/src/transport/ble/stream_read.rs @@ -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, TransportError>> + Send>>; + +/// Buffered [`AsyncRead`] view of a [`BleStream`]. +pub struct BleStreamRead { + stream: Arc, + /// Bytes received but not yet handed to the reader. + chunk: Vec, + /// Read cursor into `chunk`. + pos: usize, + /// Scratch size for one underlying `recv`. + capacity: usize, + /// In-flight `recv`, kept across polls. + pending: Option, + /// Set once the peer has closed the connection. + eof: bool, +} + +impl BleStreamRead { + /// Wrap a stream, sizing the scratch buffer from its receive MTU. + pub fn new(stream: Arc, 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 AsyncRead for BleStreamRead { + fn poll_read( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &mut ReadBuf<'_>, + ) -> Poll> { + 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>>, + } + + impl ScriptedStream { + fn new(chunks: Vec>) -> 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 { + 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"); + } +}