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 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; if batch_bytes + packet.len() + PACKET_OVERHEAD > MAX_BATCH_SIZE_BYTES {
94 if batch.is_empty() {
95 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(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 if let Some(data) = this.recv_buffer.pop_front() {
210 return Poll::Ready(Some(Ok(data)));
211 }
212
213 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; }
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 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 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 for _ in 0..32 {
325 let size = channel.write(&data).expect("sync write to succeed");
326 assert_eq!(size, data.len());
327 }
328
329 let result = channel.write(&data);
331 assert_eq!(result, Err(zx::Status::SHOULD_WAIT));
332
333 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 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 let data = vec![1; 11];
370 let result = channel.write(&data);
371 assert_eq!(result, Err(zx::Status::OUT_OF_RANGE));
372
373 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 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 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 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 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 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 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 assert!(
490 FidlClientConnection::collect_packets_for_batch(&buffer, &flush_task_queued).is_none()
491 );
492 assert!(!flush_task_queued.load(Ordering::Relaxed));
493
494 flush_task_queued.store(true, Ordering::Relaxed);
496
497 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 let p1 = vec![0; 60 * 1024 - 24]; 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 let p_large = vec![0; 60 * 1024 + 10];
527 buffer.lock().push_back(p_large.clone()); 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}