Skip to main content

fuchsia_bluetooth/types/channel/
fidl_client.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use async_utils::hanging_get::client::HangingGetStream;
6use fidl::endpoints::Proxy;
7use fuchsia_async as fasync;
8use fuchsia_sync::Mutex;
9use futures::future::BoxFuture;
10use futures::sink::Sink;
11use futures::stream::Stream;
12use futures::{Future, StreamExt, ready};
13use log::{trace, warn};
14use std::collections::VecDeque;
15use std::pin::Pin;
16use std::sync::Arc;
17use std::sync::atomic::{AtomicBool, Ordering};
18
19use std::task::{Context, Poll};
20
21use super::{Connection, ConnectionBackendType};
22use fidl_fuchsia_bluetooth as fidl_bt;
23use fidl_fuchsia_bluetooth_bredr as bredr;
24use zx;
25
26/// A client-side implementation of the Bluetooth channel transport using the FIDL protocol.
27pub struct FidlClientConnection {
28    proxy: fidl_bt::ChannelProxy,
29    /// Hanging-get stream for receiving packets from the remote server end.
30    receive_stream: HangingGetStream<fidl_bt::ChannelProxy, Vec<fidl_bt::Packet>>,
31    /// The buffer for outgoing packets that are waiting to be sent.
32    send_buffer: Arc<Mutex<VecDeque<Vec<u8>>>>,
33    /// Waker to notify if the queue has space.
34    waker: Arc<Mutex<Option<std::task::Waker>>>,
35    /// The buffer for incoming packets that are waiting to be polled by the stream consumer.
36    recv_buffer: VecDeque<Vec<u8>>,
37    /// Scope for managing background flush tasks.
38    flush_scope: fasync::Scope,
39    /// Flag indicating if a background flush task is currently queued or running.
40    flush_task_queued: Arc<AtomicBool>,
41    /// Future that waits for all tasks in `flush_scope` to complete.
42    flush_finished: Mutex<Option<BoxFuture<'static, ()>>>,
43    /// Max TX packet size.
44    max_tx_size: usize,
45    /// Terminal error if the background flush task fails.
46    terminal_error: Arc<Mutex<Option<zx::Status>>>,
47}
48
49impl std::fmt::Debug for FidlClientConnection {
50    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
51        f.debug_struct("FidlClientConnection")
52            .field("flush_task_queued", &self.flush_task_queued.load(Ordering::Relaxed))
53            .field("send_buffer_len", &self.send_buffer.lock().len())
54            .finish()
55    }
56}
57
58impl FidlClientConnection {
59    const SEND_BUFFER_SIZE: usize = 32;
60
61    pub fn new(proxy: fidl_bt::ChannelProxy, max_tx_size: usize) -> Self {
62        let receive_stream =
63            HangingGetStream::new_with_fn_ptr(proxy.clone(), fidl_bt::ChannelProxy::receive);
64
65        Self {
66            proxy,
67            receive_stream,
68            send_buffer: Arc::new(Mutex::new(VecDeque::with_capacity(Self::SEND_BUFFER_SIZE))),
69            waker: Arc::new(Mutex::new(None)),
70            recv_buffer: VecDeque::new(),
71            flush_scope: fasync::Scope::new(),
72            flush_task_queued: Arc::new(AtomicBool::new(false)),
73            flush_finished: Mutex::new(None),
74            max_tx_size,
75            terminal_error: Arc::new(Mutex::new(None)),
76        }
77    }
78
79    fn collect_packets_for_batch(
80        send_buffer: &Mutex<VecDeque<Vec<u8>>>,
81        flush_task_queued: &AtomicBool,
82    ) -> Option<Vec<fidl_bt::Packet>> {
83        let mut buffer = send_buffer.lock();
84        if buffer.is_empty() {
85            flush_task_queued.store(false, Ordering::SeqCst);
86            return None;
87        }
88        let mut batch = Vec::new();
89        let mut batch_bytes = 0;
90        while let Some(packet) = buffer.front() {
91            if batch_bytes + packet.len() + super::PACKET_OVERHEAD > super::MAX_BATCH_SIZE_BYTES {
92                if batch.is_empty() {
93                    // If a single packet is larger than the limit, we still try to send it.
94                    let item = buffer.pop_front().unwrap();
95                    batch.push(fidl_bt::Packet { packet: item });
96                }
97                break;
98            }
99            let item = buffer.pop_front().unwrap();
100            batch_bytes += item.len() + super::PACKET_OVERHEAD;
101            batch.push(fidl_bt::Packet { packet: item });
102        }
103        Some(batch)
104    }
105
106    fn ensure_flush_task_running(&self) {
107        if self.send_buffer.lock().is_empty() {
108            return;
109        }
110        if self.flush_task_queued.swap(true, Ordering::SeqCst) {
111            return;
112        }
113        let send_buffer = self.send_buffer.clone();
114        let proxy = self.proxy.clone();
115        let flush_task_queued = self.flush_task_queued.clone();
116        let terminal_error = self.terminal_error.clone();
117        let waker = self.waker.clone();
118        let _ = self.flush_scope.spawn(async move {
119            loop {
120                let Some(packets) =
121                    Self::collect_packets_for_batch(&send_buffer, &flush_task_queued)
122                else {
123                    break;
124                };
125                if let Some(w) = waker.lock().take() {
126                    w.wake();
127                }
128                trace!(
129                    "FidlClientConnection: Sending batch of {} packets to FIDL server",
130                    packets.len()
131                );
132                if let Err(e) = proxy.send_(&packets).await {
133                    warn!(e:?; "FIDL Send error");
134                    let status =
135                        if e.is_closed() { zx::Status::PEER_CLOSED } else { zx::Status::INTERNAL };
136                    *terminal_error.lock() = Some(status);
137                    flush_task_queued.store(false, Ordering::SeqCst);
138                    break;
139                }
140            }
141        });
142    }
143
144    fn enqueue_packet(&self, packet: Vec<u8>) -> Result<(), zx::Status> {
145        if packet.len() > self.max_tx_size {
146            warn!("Packet size {} exceeds max_tx_size {}", packet.len(), self.max_tx_size);
147            return Err(zx::Status::OUT_OF_RANGE);
148        }
149        let mut buffer = self.send_buffer.lock();
150        if buffer.len() >= Self::SEND_BUFFER_SIZE {
151            warn!(
152                "FidlClientConnection: send buffer is full ({}/{})",
153                buffer.len(),
154                Self::SEND_BUFFER_SIZE
155            );
156            return Err(zx::Status::SHOULD_WAIT);
157        }
158        let len = packet.len();
159        buffer.push_back(packet);
160        trace!("FidlClientConnection: Enqueued packet of size {}", len);
161        Ok(())
162    }
163}
164
165impl Connection for FidlClientConnection {
166    fn closed<'a>(&'a self) -> Pin<Box<dyn Future<Output = Result<(), zx::Status>> + 'a>> {
167        let proxy_cloned = self.proxy.clone();
168        Box::pin(async move {
169            let _ = proxy_cloned.on_closed().await;
170            Ok(())
171        })
172    }
173
174    fn connection_type(&self) -> ConnectionBackendType {
175        ConnectionBackendType::FidlClient
176    }
177
178    fn write(&self, bytes: &[u8]) -> Result<usize, zx::Status> {
179        self.enqueue_packet(bytes.to_vec())?;
180        self.ensure_flush_task_running();
181        Ok(bytes.len())
182    }
183
184    fn is_closed(&self) -> bool {
185        self.proxy.is_closed()
186    }
187
188    fn into_fidl_channel(self: Box<Self>) -> Result<bredr::Channel, zx::Status> {
189        let this = *self;
190        // Drop any ongoing active hanging-get Receive request or Send request.
191        drop(this.receive_stream);
192        drop(this.flush_scope);
193
194        let client_end = this.proxy.into_client_end().map_err(|_| zx::Status::UNAVAILABLE)?;
195        Ok(bredr::Channel { connection: Some(client_end), ..Default::default() })
196    }
197}
198
199impl Stream for FidlClientConnection {
200    type Item = Result<Vec<u8>, zx::Status>;
201
202    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
203        let this = self.get_mut();
204
205        loop {
206            // Return buffered items first
207            if let Some(data) = this.recv_buffer.pop_front() {
208                return Poll::Ready(Some(Ok(data)));
209            }
210
211            // No items in the buffer.
212            let res = ready!(this.receive_stream.poll_next_unpin(cx));
213            match res {
214                Some(Ok(packets)) => {
215                    trace!(
216                        "FidlClientConnection: Received {} packets from FIDL server",
217                        packets.len()
218                    );
219                    for packet in packets {
220                        this.recv_buffer.push_back(packet.packet);
221                    }
222                    continue; // Loop again to return the first buffered item
223                }
224                Some(Err(e)) if e.is_closed() => {
225                    trace!("FIDL channel is closed");
226                    return Poll::Ready(None);
227                }
228                Some(Err(e)) => {
229                    warn!("FIDL Receive error: {:?}", e);
230                    return Poll::Ready(Some(Err(zx::Status::INTERNAL)));
231                }
232                None => {
233                    trace!("FIDL channel terminated");
234                    return Poll::Ready(None);
235                }
236            }
237        }
238    }
239}
240
241impl Sink<Vec<u8>> for FidlClientConnection {
242    type Error = zx::Status;
243
244    fn poll_ready(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
245        let this = self.get_mut();
246        if let Some(err) = *this.terminal_error.lock() {
247            return Poll::Ready(Err(err));
248        }
249        // Limit queue size
250        let buffer = this.send_buffer.lock();
251        if buffer.len() >= Self::SEND_BUFFER_SIZE {
252            *this.waker.lock() = Some(cx.waker().clone());
253            return Poll::Pending;
254        }
255        Poll::Ready(Ok(()))
256    }
257
258    fn start_send(self: Pin<&mut Self>, item: Vec<u8>) -> Result<(), Self::Error> {
259        let this = self.get_mut();
260        if let Some(err) = *this.terminal_error.lock() {
261            return Err(err);
262        }
263        this.enqueue_packet(item)?;
264        Ok(())
265    }
266
267    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
268        let this = self.get_mut();
269        loop {
270            if let Some(err) = *this.terminal_error.lock() {
271                return Poll::Ready(Err(err));
272            }
273
274            this.ensure_flush_task_running();
275
276            let mut flush_finished_guard = this.flush_finished.lock();
277            let flush_finished = flush_finished_guard.get_or_insert_with(|| {
278                let handle = this.flush_scope.to_handle();
279                Box::pin(async move {
280                    handle.on_no_tasks().await;
281                })
282            });
283
284            ready!(flush_finished.as_mut().poll(cx));
285
286            *flush_finished_guard = None;
287
288            if let Some(err) = *this.terminal_error.lock() {
289                return Poll::Ready(Err(err));
290            }
291
292            if this.send_buffer.lock().is_empty() {
293                return Poll::Ready(Ok(()));
294            }
295        }
296    }
297
298    // TODO(https://fxbug.dev/414410187) Consider actually closing the underlying channel.
299    fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
300        Sink::poll_flush(self, cx)
301    }
302}
303
304#[cfg(test)]
305mod tests {
306    use super::*;
307    use crate::types::Channel;
308    use fidl::endpoints::create_proxy_and_stream;
309    use fuchsia_async as fasync;
310    use futures::stream::FusedStream;
311    use futures::{SinkExt, StreamExt};
312
313    #[test]
314    fn channel_sync_write() {
315        let mut exec = fasync::TestExecutor::new();
316        let (proxy, mut stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
317        let channel = Channel::from_fidl_client(proxy, Channel::DEFAULT_MAX_TX);
318
319        let data = vec![1, 2, 3];
320
321        // Fill up the send buffer (32 writes).
322        for _ in 0..32 {
323            let size = channel.write(&data).expect("sync write to succeed");
324            assert_eq!(size, data.len());
325        }
326
327        // The 33rd write should fail with SHOULD_WAIT because the buffer is full.
328        let result = channel.write(&data);
329        assert_eq!(result, Err(zx::Status::SHOULD_WAIT));
330
331        // Check the stream (server end) for the request.
332        let mut stream_fut = stream.next();
333        match exec.run_until_stalled(&mut stream_fut) {
334            Poll::Ready(Some(Ok(fidl_bt::ChannelRequest::Send_ { packets, responder }))) => {
335                assert_eq!(packets.len(), 32);
336                assert_eq!(packets[0].packet, data);
337                responder.send().expect("responder to send successfully");
338            }
339            x => panic!("Expected Send_ request, got {:?}", x),
340        }
341
342        // And now we should be able to write again because one packet was processed and acknowledged.
343        let size = channel.write(&data).expect("sync write to succeed again");
344        assert_eq!(size, data.len());
345    }
346
347    #[test]
348    fn channel_into_fidl() {
349        let _exec = fasync::TestExecutor::new();
350        let (proxy, _stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
351        let conn = FidlClientConnection::new(proxy, Channel::DEFAULT_MAX_TX);
352
353        let fidl_channel =
354            Box::new(conn).into_fidl_channel().expect("into_fidl_channel to succeed");
355        println!("FIDL Channel: {:?}", fidl_channel);
356        assert!(fidl_channel.connection.is_some());
357        assert!(fidl_channel.socket.is_none());
358    }
359
360    #[test]
361    fn channel_write_too_large() {
362        let mut exec = fasync::TestExecutor::new();
363        let (proxy, _stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
364        let mut channel = Channel::from_fidl_client(proxy, 10);
365
366        // Test Connection::write
367        let data = vec![1; 11];
368        let result = channel.write(&data);
369        assert_eq!(result, Err(zx::Status::OUT_OF_RANGE));
370
371        // Test SinkExt::send
372        let mut send_fut = channel.send(data);
373        let result = exec.run_until_stalled(&mut send_fut);
374        assert_eq!(result, Poll::Ready(Err(zx::Status::OUT_OF_RANGE)));
375    }
376
377    #[test]
378    fn channel_closed() {
379        let mut exec = fasync::TestExecutor::new();
380        let (proxy, stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
381        let channel = Channel::from_fidl_client(proxy, Channel::DEFAULT_MAX_TX);
382
383        let mut closed_fut = channel.closed();
384        assert!(exec.run_until_stalled(&mut closed_fut).is_pending());
385        assert!(!channel.is_closed());
386
387        drop(stream);
388
389        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
390        assert!(exec.run_until_stalled(&mut closed_fut).is_ready());
391        assert!(channel.is_closed());
392    }
393
394    #[fuchsia::test]
395    fn channel_sink() {
396        let mut exec = fasync::TestExecutor::new();
397        let (proxy, mut stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
398        let mut channel = Channel::from_fidl_client(proxy, Channel::DEFAULT_MAX_TX);
399
400        let data = vec![1, 2, 3];
401        let mut send_fut = channel.send(data.clone());
402        assert!(exec.run_until_stalled(&mut send_fut).is_pending());
403
404        // Now expect Send_ request
405        let mut stream_fut = stream.next();
406        match exec.run_until_stalled(&mut stream_fut) {
407            Poll::Ready(Some(Ok(fidl_bt::ChannelRequest::Send_ { packets, responder }))) => {
408                assert_eq!(packets.len(), 1);
409                assert_eq!(packets[0].packet, data);
410                responder.send().expect("responder to send successfully");
411            }
412            x => panic!("Expected Send_ request, got {:?}", x),
413        }
414
415        assert!(exec.run_until_stalled(&mut send_fut).is_ready());
416
417        // Verify start_send fails with SHOULD_WAIT when the buffer is full (32 writes).
418        for _ in 0..32 {
419            Pin::new(&mut channel).start_send(data.clone()).expect("start_send to succeed");
420        }
421        let result = Pin::new(&mut channel).start_send(data.clone());
422        assert_eq!(result, Err(zx::Status::SHOULD_WAIT));
423
424        // Verify Sink batching: flushing the full buffer sends all 32 packets together.
425        let mut flush_fut = channel.flush();
426        assert!(exec.run_until_stalled(&mut flush_fut).is_pending());
427
428        let mut stream_fut = stream.next();
429        match exec.run_until_stalled(&mut stream_fut) {
430            Poll::Ready(Some(Ok(fidl_bt::ChannelRequest::Send_ { packets, responder }))) => {
431                assert_eq!(packets.len(), 32);
432                assert_eq!(packets[0].packet, data);
433                responder.send().expect("responder to send successfully");
434            }
435            x => panic!("Expected Send_ request, got {:?}", x),
436        }
437        assert!(exec.run_until_stalled(&mut flush_fut).is_ready());
438    }
439
440    #[test]
441    fn channel_stream() {
442        let mut exec = fasync::TestExecutor::new();
443        let (proxy, mut stream) = create_proxy_and_stream::<fidl_bt::ChannelMarker>();
444        let mut channel = Channel::from_fidl_client(proxy, Channel::DEFAULT_MAX_TX);
445
446        let data = vec![4, 5, 6];
447
448        // Trigger initial Receive request
449        let mut next_fut = channel.next();
450        assert!(exec.run_until_stalled(&mut next_fut).is_pending());
451
452        let mut stream_fut = stream.next();
453
454        match exec.run_until_stalled(&mut stream_fut) {
455            Poll::Ready(Some(Ok(fidl_bt::ChannelRequest::Receive { responder }))) => {
456                let packets = vec![fidl_bt::Packet { packet: data.clone() }];
457                responder.send(&packets).expect("responder to send packets successfully");
458            }
459            x => panic!("Expected Receive request, got {:?}", x),
460        }
461
462        match exec.run_until_stalled(&mut next_fut) {
463            Poll::Ready(Some(Ok(received))) => {
464                assert_eq!(received, data);
465            }
466            x => panic!("Expected data from stream, got {:?}", x),
467        }
468
469        // After the sender is dropped, the stream should terminate.
470        drop(stream);
471
472        let mut next_fut = channel.next();
473        let Poll::Ready(None) = exec.run_until_stalled(&mut next_fut) else {
474            panic!("Expected None from the stream")
475        };
476
477        // It should continue to report terminated.
478        assert!(channel.is_terminated());
479    }
480
481    #[test]
482    fn test_collect_packets_for_batch() {
483        let buffer = Mutex::new(VecDeque::new());
484        let flush_task_queued = AtomicBool::new(true);
485
486        // When buffer is empty, collect_packets_for_batch should return None
487        assert!(
488            FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued).is_none()
489        );
490        assert!(!flush_task_queued.load(Ordering::Relaxed));
491
492        // Reset flag
493        flush_task_queued.store(true, Ordering::Relaxed);
494
495        // When buffer has packets, it should return Some(packets)
496        buffer.lock().push_back(vec![1, 2, 3]);
497        buffer.lock().push_back(vec![4, 5, 6]);
498
499        let packets = FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued)
500            .expect("should return packets");
501        assert_eq!(packets.len(), 2);
502        assert_eq!(packets[0].packet, vec![1, 2, 3]);
503        assert_eq!(packets[1].packet, vec![4, 5, 6]);
504        assert!(buffer.lock().is_empty());
505        assert!(flush_task_queued.load(Ordering::Relaxed));
506
507        // Check the max batch size logic.
508        // MAX_BATCH_SIZE_BYTES = 60 * 1024 = 61440.
509        // PACKET_OVERHEAD = 24.
510        // Let's create two packets: one that is close to the limit, and one that exceeds the limit when combined.
511        let p1 = vec![0; 60 * 1024 - 24]; // Exactly 61416 bytes packet. Total size with overhead = 61440.
512        let p2 = vec![0; 10];
513        buffer.lock().push_back(p1.clone());
514        buffer.lock().push_back(p2.clone());
515
516        let packets = FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued)
517            .expect("should return packets");
518        assert_eq!(packets.len(), 1);
519        assert_eq!(packets[0].packet, p1);
520        assert_eq!(buffer.lock().len(), 1);
521        assert_eq!(buffer.lock().front().unwrap(), &p2);
522
523        // When the first packet exceeds the limit by itself, it should still be popped and returned.
524        let p_large = vec![0; 60 * 1024 + 10];
525        buffer.lock().push_back(p_large.clone()); // buffer is [p2, p_large]
526
527        let packets = FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued)
528            .expect("should return packets");
529        assert_eq!(packets.len(), 1);
530        assert_eq!(packets[0].packet, p2);
531
532        let packets = FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued)
533            .expect("should return packets");
534        assert_eq!(packets.len(), 1);
535        assert_eq!(packets[0].packet, p_large);
536        assert!(buffer.lock().is_empty());
537    }
538}