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