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"); + } +}