fuchsia_bluetooth/types/channel/
fidl_client.rs1use 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
26pub struct FidlClientConnection {
28 proxy: fidl_bt::ChannelProxy,
29 receive_stream: HangingGetStream<fidl_bt::ChannelProxy, Vec<fidl_bt::Packet>>,
31 send_buffer: Arc<Mutex<VecDeque<Vec<u8>>>>,
33 waker: Arc<Mutex<Option<std::task::Waker>>>,
35 recv_buffer: VecDeque<Vec<u8>>,
37 flush_scope: fasync::Scope,
39 flush_task_queued: Arc<AtomicBool>,
41 flush_finished: Mutex<Option<BoxFuture<'static, ()>>>,
43 max_tx_size: usize,
45 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 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(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 if let Some(data) = this.recv_buffer.pop_front() {
208 return Poll::Ready(Some(Ok(data)));
209 }
210
211 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; }
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 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 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 for _ in 0..32 {
323 let size = channel.write(&data).expect("sync write to succeed");
324 assert_eq!(size, data.len());
325 }
326
327 let result = channel.write(&data);
329 assert_eq!(result, Err(zx::Status::SHOULD_WAIT));
330
331 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 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 let data = vec![1; 11];
368 let result = channel.write(&data);
369 assert_eq!(result, Err(zx::Status::OUT_OF_RANGE));
370
371 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 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 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 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 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 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 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 assert!(
488 FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued).is_none()
489 );
490 assert!(!flush_task_queued.load(Ordering::Relaxed));
491
492 flush_task_queued.store(true, Ordering::Relaxed);
494
495 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 let p1 = vec![0; 60 * 1024 - 24]; 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 let p_large = vec![0; 60 * 1024 + 10];
525 buffer.lock().push_back(p_large.clone()); 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}