1use anyhow::Context as _;
6use bt_avdtp::{
7 self as avdtp, MediaCodecType, ServiceCapability, ServiceCategory, StreamEndpoint,
8 StreamEndpointId,
9};
10use fidl_fuchsia_bluetooth::ChannelParameters;
11use fidl_fuchsia_bluetooth_avrcp as avrcp;
12use fidl_fuchsia_bluetooth_bredr::{
13 ConnectParameters, L2capParameters, PSM_AVDTP, ProfileDescriptor, ProfileProxy,
14};
15use fuchsia_async as fasync;
16use fuchsia_bluetooth::inspect::DebugExt;
17use fuchsia_bluetooth::types::{Channel, PeerId};
18use fuchsia_inspect as inspect;
19use fuchsia_inspect_derive::{AttachError, Inspect};
20use fuchsia_sync::Mutex;
21use futures::channel::mpsc;
22use futures::future::{BoxFuture, Either};
23use futures::stream::FuturesUnordered;
24use futures::task::{Context, Poll, Waker};
25use futures::{Future, FutureExt, StreamExt, select};
26use log::{debug, info, trace, warn};
27use std::collections::{BTreeMap, HashMap, HashSet};
28use std::pin::Pin;
29use std::sync::{Arc, Weak};
30
31mod controller;
33pub use controller::ControllerPool;
34mod volume_relay;
35pub(crate) use volume_relay::run_avrcp_volume_relay;
36
37use crate::codec::MediaCodecConfig;
38use crate::media_task::MediaTaskStatus;
39use crate::permits::{Permit, Permits};
40use crate::stream::{Stream, Streams};
41
42#[derive(Inspect)]
45pub struct Peer {
46 id: PeerId,
48 #[inspect(forward)]
50 inner: Arc<Mutex<PeerInner>>,
51 profile: ProfileProxy,
53 descriptor: Mutex<Option<ProfileDescriptor>>,
55 closed_wakers: Arc<Mutex<Option<Vec<Waker>>>>,
59 metrics: bt_metrics::MetricsLogger,
61}
62
63const STREAM_DWELL: zx::MonotonicDuration = zx::MonotonicDuration::from_millis(500);
66
67const START_RETRY_DELAY: zx::MonotonicDuration = zx::MonotonicDuration::from_seconds(1);
69
70const START_ATTEMPTS: usize = 3;
73
74#[derive(Clone)]
78struct StreamPermits {
79 permits: Permits,
80 open_streams: Arc<Mutex<HashMap<StreamEndpointId, Permit>>>,
81 reserved_streams: Arc<Mutex<HashSet<StreamEndpointId>>>,
82 inner: Weak<Mutex<PeerInner>>,
83 peer_id: PeerId,
84 sender: mpsc::UnboundedSender<BoxFuture<'static, StreamPermit>>,
85}
86
87#[derive(Debug)]
88struct StreamPermit {
89 local_id: StreamEndpointId,
90 open_streams: Arc<Mutex<HashMap<StreamEndpointId, Permit>>>,
91}
92
93impl StreamPermit {
94 fn local_id(&self) -> &StreamEndpointId {
95 &self.local_id
96 }
97
98 fn is_held(&self) -> bool {
100 self.open_streams.lock().contains_key(&self.local_id)
101 }
102}
103
104impl Drop for StreamPermit {
105 fn drop(&mut self) {
106 let _ = self.open_streams.lock().remove(&self.local_id);
107 }
108}
109
110impl StreamPermits {
111 fn new(
112 inner: Weak<Mutex<PeerInner>>,
113 peer_id: PeerId,
114 permits: Permits,
115 ) -> (Self, mpsc::UnboundedReceiver<BoxFuture<'static, StreamPermit>>) {
116 let (sender, reservations_receiver) = futures::channel::mpsc::unbounded();
117 (
118 Self {
119 inner,
120 permits,
121 peer_id,
122 sender,
123 open_streams: Default::default(),
124 reserved_streams: Default::default(),
125 },
126 reservations_receiver,
127 )
128 }
129
130 fn label_for(&self, local_id: &StreamEndpointId) -> String {
131 format!("{} {}", self.peer_id, local_id)
132 }
133
134 fn get(&self, local_id: StreamEndpointId) -> Option<StreamPermit> {
137 let revoke_fn = self.make_revocation_fn(&local_id);
138 let Some(permit) = self.permits.get_revokable(revoke_fn) else {
139 info!("No permits available: {:?}", self.permits);
140 return None;
141 };
142 permit.relabel(self.label_for(&local_id));
143 if let Some(_) = self.open_streams.lock().insert(local_id.clone(), permit) {
144 warn!(id:% = self.peer_id; "Started stream {local_id:?} twice, dropping previous permit");
145 }
146 Some(StreamPermit { local_id, open_streams: self.open_streams.clone() })
147 }
148
149 fn setup_reservation_for(&self, local_id: StreamEndpointId) {
152 if !self.reserved_streams.lock().insert(local_id.clone()) {
153 return;
155 }
156 let restart_stream_available_fut = {
157 let self_revoke_fn = Self::make_revocation_fn(&self, &local_id);
158 let reservation = self.permits.reserve_revokable(self_revoke_fn);
159 let open_streams = self.open_streams.clone();
160 let reserved_streams = self.reserved_streams.clone();
161 let label = self.label_for(&local_id);
162 let local_id = local_id.clone();
163 async move {
164 let permit = reservation.await;
165 permit.relabel(label);
166 if open_streams.lock().insert(local_id.clone(), permit).is_some() {
167 warn!("Reservation replaces acquired permit for {}", local_id.clone());
168 }
169 if !reserved_streams.lock().remove(&local_id) {
170 warn!(local_id:%; "Unrecorded reservation resolved");
171 }
172 StreamPermit { local_id, open_streams }
173 }
174 };
175 if let Err(e) = self.sender.unbounded_send(restart_stream_available_fut.boxed()) {
176 warn!(id:% = self.peer_id, local_id:%, e:?; "Couldn't queue reservation to finish");
177 }
178 }
179
180 fn revocation_fn(self, local_id: StreamEndpointId) -> Permit {
185 if let Ok(peer) = PeerInner::upgrade(self.inner.clone()) {
186 {
187 let mut lock = peer.lock();
188 match lock.suspend_local_stream(&local_id) {
189 Ok(remote_id) => drop(lock.peer.suspend(&[remote_id])),
190 Err(e) => warn!("Couldn't stop local stream {local_id:?}: {e:?}"),
191 }
192 }
193 self.setup_reservation_for(local_id.clone());
194 }
195 self.open_streams.lock().remove(&local_id).expect("permit revoked but don't have it")
196 }
197
198 fn make_revocation_fn(&self, local_id: &StreamEndpointId) -> impl FnOnce() -> Permit + use<> {
199 let local_id = local_id.clone();
200 let cloned = self.clone();
201 move || cloned.revocation_fn(local_id)
202 }
203}
204
205impl Peer {
206 pub fn create(
213 id: PeerId,
214 peer: avdtp::Peer,
215 streams: Streams,
216 permits: Option<Permits>,
217 profile: ProfileProxy,
218 avrcp: Option<avrcp::PeerManagerProxy>,
219 metrics: bt_metrics::MetricsLogger,
220 ) -> Self {
221 let inner = Arc::new(Mutex::new(PeerInner::new(peer, id, streams, avrcp, metrics.clone())));
222 inner.lock().self_weak = Arc::downgrade(&inner);
223 let reservations_receiver = if let Some(permits) = permits {
224 let (stream_permits, receiver) =
225 StreamPermits::new(Arc::downgrade(&inner), id, permits);
226 inner.lock().permits = Some(stream_permits);
227 receiver
228 } else {
229 let (_, receiver) = mpsc::unbounded();
230 receiver
231 };
232 let res = Self {
233 id,
234 inner,
235 profile,
236 descriptor: Mutex::new(None),
237 closed_wakers: Arc::new(Mutex::new(Some(Vec::new()))),
238 metrics,
239 };
240 res.start_requests_task(reservations_receiver);
241 res
242 }
243
244 pub fn set_descriptor(&self, descriptor: ProfileDescriptor) -> Option<ProfileDescriptor> {
245 self.descriptor.lock().replace(descriptor)
246 }
247
248 #[cfg(test)]
249 pub fn set_volume_relay_task(&self, task: fasync::Task<()>) {
250 self.inner.lock().volume_relay_task = Some(task);
251 }
252
253 pub fn receive_channel(&self, channel: Channel) -> avdtp::Result<()> {
258 let opened_stream = {
259 let mut lock = self.inner.lock();
260 lock.receive_channel(channel)?
261 };
262 if let Some(stream_id) = opened_stream {
263 PeerInner::schedule_start(Arc::downgrade(&self.inner), stream_id);
264 }
265 PeerInner::maybe_start_volume_relay(&self.inner);
266 Ok(())
267 }
268
269 pub fn avdtp(&self) -> avdtp::Peer {
271 let lock = self.inner.lock();
272 lock.peer.clone()
273 }
274
275 pub fn remote_endpoints(&self) -> Option<Vec<avdtp::StreamEndpoint>> {
277 self.inner.lock().remote_endpoints()
278 }
279
280 pub fn collect_capabilities(
284 &self,
285 ) -> impl Future<Output = avdtp::Result<Vec<avdtp::StreamEndpoint>>> + use<> {
286 let avdtp = self.avdtp();
287 let get_all = self.descriptor.lock().clone().is_some_and(a2dp_version_check);
288 let inner = self.inner.clone();
289 let metrics = self.metrics.clone();
290 let peer_id = self.id;
291 async move {
292 if let Some(caps) = inner.lock().remote_endpoints() {
293 return Ok(caps);
294 }
295 trace!("Discovering peer streams..");
296 let infos = avdtp.discover().await?;
297 trace!("Discovered {} streams", infos.len());
298 let mut remote_streams = Vec::new();
299 for info in infos {
300 let capabilities = if get_all {
301 avdtp.get_all_capabilities(info.id()).await
302 } else {
303 avdtp.get_capabilities(info.id()).await
304 };
305 match capabilities {
306 Ok(capabilities) => {
307 trace!("Stream {:?}", info);
308 for cap in &capabilities {
309 trace!(" - {:?}", cap);
310 }
311 remote_streams.push(avdtp::StreamEndpoint::from_info(&info, capabilities));
312 }
313 Err(e) => {
314 info!(peer_id:%; "Stream {} capabilities failed: {:?}, skipping", info.id(), e);
315 }
316 };
317 }
318 inner.lock().set_remote_endpoints(&remote_streams);
319 Self::record_cobalt_metrics(metrics, &remote_streams);
320 Ok(remote_streams)
321 }
322 }
323
324 fn record_cobalt_metrics(metrics: bt_metrics::MetricsLogger, endpoints: &[StreamEndpoint]) {
325 let codec_metrics: HashSet<_> = endpoints
326 .iter()
327 .filter_map(|endpoint| {
328 endpoint.codec_type().map(|t| codectype_to_availability_metric(t) as u32)
329 })
330 .collect();
331 metrics
332 .log_occurrences(bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID, codec_metrics);
333
334 let cap_metrics: HashSet<_> = endpoints
335 .iter()
336 .flat_map(|endpoint| {
337 endpoint
338 .capabilities()
339 .iter()
340 .filter_map(|t| capability_to_metric(t))
341 .chain(std::iter::once(
342 bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic,
343 ))
344 .map(|t| t as u32)
345 })
346 .collect();
347 metrics.log_occurrences(bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID, cap_metrics);
348 }
349
350 fn transport_channel_params() -> L2capParameters {
351 L2capParameters {
352 psm: Some(PSM_AVDTP),
353 parameters: Some(ChannelParameters {
354 max_rx_packet_size: Some(65535),
355 ..Default::default()
356 }),
357 ..Default::default()
358 }
359 }
360
361 pub fn stream_start(
366 &self,
367 remote_id: StreamEndpointId,
368 capabilities: Vec<ServiceCapability>,
369 ) -> impl Future<Output = avdtp::Result<()>> {
370 let peer = Arc::downgrade(&self.inner);
371 let peer_id = self.id.clone();
372 let avdtp = self.avdtp();
373 let profile = self.profile.clone();
374
375 async move {
376 let codec_params =
377 capabilities.iter().find(|x| x.is_codec()).ok_or(avdtp::Error::InvalidState)?;
378 let (local_id, local_capabilities) = {
379 let peer = PeerInner::upgrade(peer.clone())?;
380 let lock = peer.lock();
381 lock.find_compatible_local_capabilities(codec_params, &remote_id)?
382 };
383
384 let local_by_cat: HashMap<ServiceCategory, ServiceCapability> =
385 local_capabilities.into_iter().map(|i| (i.category(), i)).collect();
386
387 let shared_capabilities: BTreeMap<ServiceCategory, ServiceCapability> = capabilities
390 .into_iter()
391 .filter_map(|cap| {
392 let Some(local_cap) = local_by_cat.get(&cap.category()) else {
393 return None;
394 };
395 if cap.category() == ServiceCategory::MediaCodec {
396 let Ok(a) = MediaCodecConfig::try_from(&cap) else {
397 return None;
398 };
399 let Ok(b) = MediaCodecConfig::try_from(local_cap) else {
400 return None;
401 };
402 let Some(negotiated) = MediaCodecConfig::negotiate(&a, &b) else {
403 return None;
404 };
405 Some((cap.category(), (&negotiated).into()))
406 } else {
407 Some((cap.category(), cap))
408 }
409 })
410 .collect();
411 let shared_capabilities: Vec<_> = shared_capabilities.into_values().collect();
412
413 trace!("Starting stream {local_id} to remote {remote_id} with {shared_capabilities:?}");
414
415 avdtp.set_configuration(&remote_id, &local_id, &shared_capabilities).await?;
416 {
417 let strong = PeerInner::upgrade(peer.clone())?;
418 strong.lock().set_opening(&local_id, &remote_id, shared_capabilities)?;
419 }
420 avdtp.open(&remote_id).await?;
421
422 debug!(peer_id:%; "Connecting transport channel");
423 let channel = profile
424 .connect(
425 &peer_id.into(),
426 &ConnectParameters::L2cap(Self::transport_channel_params()),
427 )
428 .await
429 .context("FIDL error: {}")?
430 .or(Err(avdtp::Error::PeerDisconnected))?;
431
432 trace!(peer_id:%; "Connected transport channel, converting to local Channel");
433
434 let channel = match channel.try_into() {
435 Err(e) => {
436 warn!(peer_id:%, e:?; "Couldn't connect media transport: no channel");
437 return Err(avdtp::Error::PeerDisconnected);
438 }
439 Ok(c) => c,
440 };
441
442 let opened_stream = {
443 let strong = PeerInner::upgrade(peer.clone())?;
444 let mut lock = strong.lock();
445 lock.receive_channel(channel)?
446 };
447 if let Some(stream_id) = opened_stream {
448 PeerInner::schedule_start(peer, stream_id);
449 }
450 Ok(())
451 }
452 }
453
454 pub fn streaming_active(&self) -> bool {
456 self.inner.lock().is_streaming()
457 }
458
459 #[cfg(test)]
461 fn is_streaming_now(&self) -> bool {
462 self.inner.lock().is_streaming_now()
463 }
464
465 pub fn stream_suspend(
470 &self,
471 local_id: StreamEndpointId,
472 ) -> impl Future<Output = avdtp::Result<()>> {
473 let peer = Arc::downgrade(&self.inner);
474 PeerInner::suspend(peer, local_id)
475 }
476
477 fn start_requests_task(
480 &self,
481 mut reservations_receiver: mpsc::UnboundedReceiver<BoxFuture<'static, StreamPermit>>,
482 ) {
483 let lock = self.inner.lock();
484 let mut request_stream = lock.peer.take_request_stream();
485 let id = self.id.clone();
486 let peer = Arc::downgrade(&self.inner);
487 let mut stream_reservations = FuturesUnordered::new();
488 let disconnect_wakers = Arc::downgrade(&self.closed_wakers);
489 fuchsia_async::Task::local(async move {
490 loop {
491 select! {
492 request = request_stream.next() => {
493 match request {
494 None => break,
495 Some(Err(e)) => info!(peer_id:% = id, e:?; "Request stream error"),
496 Some(Ok(request)) => match peer.upgrade() {
497 None => return,
498 Some(p) => {
499 let result_or_future = p.lock().handle_request(request);
500 let result = match result_or_future {
501 Either::Left(result) => result,
502 Either::Right(future) => future.await,
503 };
504 if let Err(e) = result {
505 warn!(peer_id:% = id, e:?; "Error handling request");
506 }
507 PeerInner::maybe_start_volume_relay(&p);
508 }
509 },
510 }
511 },
512 reservation_fut = reservations_receiver.select_next_some() => {
513 stream_reservations.push(reservation_fut)
514 },
515 permit = stream_reservations.select_next_some() => {
516 if let Err(e) = PeerInner::start_permit(peer.clone(), permit).await {
517 warn!(peer_id:% = id, e:?; "Couldn't start stream after unpause");
518 }
519 }
520 complete => break,
521 }
522 }
523 info!(peer_id:% = id; "disconnected");
524 if let Some(wakers) = disconnect_wakers.upgrade() {
525 for waker in wakers.lock().take().unwrap_or_else(Vec::new) {
526 waker.wake();
527 }
528 }
529 })
530 .detach();
531 }
532
533 pub fn closed(&self) -> ClosedPeer {
535 ClosedPeer { inner: Arc::downgrade(&self.closed_wakers) }
536 }
537}
538
539#[must_use = "futures do nothing unless you `.await` or poll them"]
542pub struct ClosedPeer {
543 inner: Weak<Mutex<Option<Vec<Waker>>>>,
544}
545
546impl Future for ClosedPeer {
547 type Output = ();
548
549 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
550 match self.inner.upgrade() {
551 None => Poll::Ready(()),
552 Some(inner) => match inner.lock().as_mut() {
553 None => Poll::Ready(()),
554 Some(wakers) => {
555 wakers.push(cx.waker().clone());
556 Poll::Pending
557 }
558 },
559 }
560 }
561}
562
563fn a2dp_version_check(profile: ProfileDescriptor) -> bool {
565 let (Some(major), Some(minor)) = (profile.major_version, profile.minor_version) else {
566 return false;
567 };
568 (major == 1 && minor >= 3) || major > 1
569}
570
571struct PeerInner {
575 peer: avdtp::Peer,
577 peer_id: PeerId,
579 opening: Option<StreamEndpointId>,
582 local: Streams,
584 permits: Option<StreamPermits>,
586 started: HashMap<StreamEndpointId, WatchedStream>,
588 waiting_start_tasks: HashMap<StreamEndpointId, fasync::Task<()>>,
591 inspect: fuchsia_inspect::Node,
593 remote_endpoints: Option<Vec<StreamEndpoint>>,
595 remote_inspect: fuchsia_inspect::Node,
597 metrics: bt_metrics::MetricsLogger,
599 avrcp: Option<avrcp::PeerManagerProxy>,
601 self_weak: Weak<Mutex<PeerInner>>,
603 volume_relay_task: Option<fasync::Task<()>>,
605}
606
607impl Inspect for &mut PeerInner {
608 fn iattach(self, parent: &inspect::Node, name: impl AsRef<str>) -> Result<(), AttachError> {
611 self.inspect = parent.create_child(name.as_ref());
612 self.inspect.record_string("id", self.peer_id.to_string());
613 self.local.iattach(&self.inspect, "local_streams")
614 }
615}
616
617impl PeerInner {
618 pub fn new(
619 peer: avdtp::Peer,
620 peer_id: PeerId,
621 local: Streams,
622 avrcp: Option<avrcp::PeerManagerProxy>,
623 metrics: bt_metrics::MetricsLogger,
624 ) -> Self {
625 Self {
626 peer,
627 peer_id,
628 opening: None,
629 local,
630 permits: None,
631 started: HashMap::new(),
632 waiting_start_tasks: HashMap::new(),
633 inspect: Default::default(),
634 remote_endpoints: None,
635 remote_inspect: Default::default(),
636 metrics,
637 avrcp,
638 self_weak: Weak::new(),
639 volume_relay_task: None,
640 }
641 }
642
643 pub fn maybe_start_volume_relay(this: &Arc<Mutex<Self>>) {
644 let mut lock = this.lock();
645 let Some(avrcp) = lock.avrcp.clone() else {
646 if lock.volume_relay_task.take().is_some() {
647 let peer_id = lock.peer_id;
648 trace!(peer_id:%; "Stopping volume relay task (AVRCP proxy missing)");
649 }
650 return;
651 };
652
653 let active_source_stream = lock
655 .local
656 .streaming()
657 .any(|s| s.endpoint().endpoint_type() == &avdtp::EndpointType::Source)
658 || lock
659 .local
660 .open()
661 .any(|s| s.endpoint().endpoint_type() == &avdtp::EndpointType::Source);
662
663 if !active_source_stream {
664 if lock.volume_relay_task.take().is_some() {
665 let peer_id = lock.peer_id;
666 trace!(peer_id:%; "Stopping volume relay task (No active local Source streams)");
667 }
668 return;
669 }
670
671 if let Some(task) = lock.volume_relay_task.as_mut() {
672 if let Some(res) = task.now_or_never() {
673 let peer_id = lock.peer_id;
674 trace!(peer_id:%; "Volume relay task completed: {res:?}");
675 lock.volume_relay_task = None;
676 }
677 }
678
679 if lock.volume_relay_task.is_some() {
680 return;
681 }
682
683 let peer_id = lock.peer_id;
684 trace!(peer_id:%; "Spawning volume relay task for A2DP sink peer");
685 let task = fasync::Task::spawn(run_avrcp_volume_relay(peer_id, avrcp));
686 lock.volume_relay_task = Some(task);
687 }
688
689 fn get_mut(&mut self, local_id: &StreamEndpointId) -> Result<&mut Stream, avdtp::ErrorCode> {
691 self.local.get_mut(&local_id).ok_or(avdtp::ErrorCode::BadAcpSeid)
692 }
693
694 fn set_remote_endpoints(&mut self, endpoints: &[StreamEndpoint]) {
695 self.remote_inspect = self.inspect.create_child("remote_endpoints");
696 for endpoint in endpoints {
697 self.remote_inspect.record_child(inspect::unique_name("remote_"), |node| {
698 node.record_string("endpoint_id", endpoint.local_id().debug());
699 node.record_string("capabilities", endpoint.capabilities().debug());
700 node.record_string("type", endpoint.endpoint_type().debug());
701 });
702 }
703 self.remote_endpoints = Some(endpoints.iter().map(StreamEndpoint::as_new).collect());
704 }
705
706 fn remote_endpoints(&self) -> Option<Vec<StreamEndpoint>> {
708 self.remote_endpoints.as_ref().map(|v| v.iter().map(StreamEndpoint::as_new).collect())
709 }
710
711 fn remote_endpoint(&self, id: &StreamEndpointId) -> Option<StreamEndpoint> {
713 self.remote_endpoints
714 .as_ref()
715 .and_then(|v| v.iter().find(|v| v.local_id() == id).map(StreamEndpoint::as_new))
716 }
717
718 fn is_streaming(&self) -> bool {
720 self.is_streaming_now() || self.opening.is_some() || !self.waiting_start_tasks.is_empty()
721 }
722
723 fn is_streaming_now(&self) -> bool {
725 self.local.streaming().next().is_some()
726 }
727
728 fn set_opening(
729 &mut self,
730 local_id: &StreamEndpointId,
731 remote_id: &StreamEndpointId,
732 capabilities: Vec<ServiceCapability>,
733 ) -> avdtp::Result<()> {
734 if self.opening.is_some() {
735 return Err(avdtp::Error::InvalidState);
736 }
737 let peer_id = self.peer_id;
738 let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
739 stream
740 .configure(&peer_id, &remote_id, capabilities)
741 .map_err(|(cat, c)| avdtp::Error::RequestInvalidExtra(c, (&cat).into()))?;
742 stream.endpoint_mut().establish().or(Err(avdtp::Error::InvalidState))?;
743 self.opening = Some(local_id.clone());
744 Ok(())
745 }
746
747 fn upgrade(weak: Weak<Mutex<Self>>) -> avdtp::Result<Arc<Mutex<Self>>> {
748 weak.upgrade().ok_or(avdtp::Error::PeerDisconnected)
749 }
750
751 fn schedule_start(weak: Weak<Mutex<Self>>, local_id: StreamEndpointId) {
755 let Ok(peer) = Self::upgrade(weak.clone()) else {
756 return;
757 };
758 let mut peer_lock = peer.lock();
759 if peer_lock.started.contains_key(&local_id) {
760 return;
761 }
762 if let Some(waiting_start_task) = peer_lock.waiting_start_tasks.get_mut(&local_id) {
763 let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref());
764 if let Poll::Pending = waiting_start_task.poll_unpin(&mut noop_cx) {
765 return;
766 }
767 }
768 let task = fasync::Task::spawn(Self::wait_and_start(weak, local_id.clone()));
769 let _ = peer_lock.waiting_start_tasks.insert(local_id, task);
770 }
771
772 fn forget_waiting_start(weak: &Weak<Mutex<Self>>, local_id: &StreamEndpointId) {
776 let Ok(peer) = Self::upgrade(weak.clone()) else {
777 return;
778 };
779 let waiting_start_task = peer.lock().waiting_start_tasks.remove(local_id);
780 if let Some(task) = waiting_start_task {
781 task.detach();
782 }
783 }
784
785 async fn wait_and_start(weak: Weak<Mutex<Self>>, local_id: StreamEndpointId) {
799 let watch_active = |weak: &Weak<Mutex<Self>>, local_id| {
803 let Ok(peer) = Self::upgrade(weak.clone()) else {
804 return None;
805 };
806 let mut peer = peer.lock();
807 let Ok(stream) = peer.get_mut(local_id) else {
808 return None;
809 };
810 Some(stream.watch_active())
811 };
812
813 let is_source = {
814 let Ok(peer) = Self::upgrade(weak.clone()) else {
815 return;
816 };
817 let mut peer = peer.lock();
818 let Ok(stream) = peer.get_mut(&local_id) else {
819 return;
820 };
821 stream.endpoint().endpoint_type() == &avdtp::EndpointType::Source
822 };
823
824 'wait_start: loop {
825 if is_source {
826 let Some(mut wait_start_fut) = watch_active(&weak, &local_id) else {
827 return;
828 };
829 while !wait_start_fut.await {
830 let Some(new_wait_start_fut) = watch_active(&weak, &local_id) else {
832 return;
833 };
834 wait_start_fut = new_wait_start_fut;
835 }
836 } else {
837 fasync::Timer::new(fasync::MonotonicInstant::after(STREAM_DWELL)).await;
839 }
840
841 for attempt in 1..=START_ATTEMPTS {
842 let Err(e) = Self::start_stream(&weak, &local_id).await else {
843 return;
844 };
845 warn!(local_id:%, attempt, e:?; "Error starting stream");
846 if attempt == START_ATTEMPTS {
847 break;
848 }
849 let mut retry_timer = std::pin::pin!(fasync::Timer::new(
850 fasync::MonotonicInstant::after(START_RETRY_DELAY)
851 ));
852 if !is_source {
853 retry_timer.await;
854 continue;
855 }
856 let Some(wait_stop_fut) = watch_active(&weak, &local_id) else {
857 return;
858 };
859 match futures::future::select(retry_timer.as_mut(), wait_stop_fut).await {
860 Either::Right((false, _)) => continue 'wait_start,
862 Either::Right((true, _)) => retry_timer.await,
865 Either::Left(_) => {}
867 }
868 }
869
870 if !is_source {
871 Self::forget_waiting_start(&weak, &local_id);
874 return;
875 }
876 }
877 }
878
879 async fn start_stream(
883 weak: &Weak<Mutex<Self>>,
884 local_id: &StreamEndpointId,
885 ) -> avdtp::Result<()> {
886 let peer = Self::upgrade(weak.clone())?;
887 let (avdtp, remote_id, permit_result) = {
888 let mut peer = peer.lock();
889 let stream = peer.get_mut(local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
890 let remote_id =
891 stream.endpoint().remote_id().cloned().ok_or(avdtp::Error::InvalidState)?;
892 let avdtp = peer.peer.clone();
893 let permit_result = peer.get_permit_or_reserve(local_id);
894 (avdtp, remote_id, permit_result)
895 };
896 let Ok(permit) = permit_result else {
897 return Ok(());
899 };
900 Self::initiated_start(avdtp, weak.clone(), permit, local_id, &remote_id).await
901 }
902
903 async fn start_permit(weak: Weak<Mutex<Self>>, permit: StreamPermit) -> avdtp::Result<()> {
904 let local_id = permit.local_id().clone();
905 let (avdtp, remote_id) = {
906 let peer = Self::upgrade(weak.clone())?;
907 let mut peer = peer.lock();
908 let remote_id = peer
909 .get_mut(&local_id)
910 .map_err(|e| avdtp::Error::RequestInvalid(e))?
911 .endpoint()
912 .remote_id()
913 .ok_or(avdtp::Error::InvalidState)?
914 .clone();
915 (peer.peer.clone(), remote_id)
916 };
917 Self::initiated_start(avdtp, weak, Some(permit), &local_id, &remote_id).await
918 }
919
920 async fn initiated_start(
922 avdtp: avdtp::Peer,
923 weak: Weak<Mutex<Self>>,
924 permit: Option<StreamPermit>,
925 local_id: &StreamEndpointId,
926 remote_id: &StreamEndpointId,
927 ) -> avdtp::Result<()> {
928 trace!(permit:?, local_id:?, remote_id:?; "Making outgoing start request");
929 let to_start = std::slice::from_ref(remote_id);
930 avdtp.start(to_start).await?;
931 trace!("Start response received: {permit:?}");
932 let peer = Self::upgrade(weak.clone())?;
933 let (peer_id, start_result) = {
934 let mut peer = peer.lock();
935 (peer.peer_id, peer.start_local_stream(permit, &local_id))
936 };
937 if let Err(e) = start_result {
938 warn!(peer_id:%, local_id:%, remote_id:%, e:?; "Failed to start local stream, suspending");
939 avdtp.suspend(to_start).await?;
940 } else {
941 Self::maybe_start_volume_relay(&peer);
942 }
943 Ok(())
944 }
945
946 fn suspend(
948 weak: Weak<Mutex<Self>>,
949 local_id: StreamEndpointId,
950 ) -> impl Future<Output = avdtp::Result<()>> {
951 let res = (move || {
952 let peer = Self::upgrade(weak.clone())?;
953 let mut peer = peer.lock();
954 Ok((peer.peer.clone(), peer.suspend_local_stream(&local_id)?))
955 })();
956 let (avdtp, remote_id) = match res {
957 Err(e) => return futures::future::err(e).left_future(),
958 Ok(r) => r,
959 };
960 let to_suspend = &[remote_id];
961 avdtp.suspend(to_suspend).right_future()
962 }
963
964 pub fn find_compatible_local_capabilities(
967 &self,
968 codec_params: &ServiceCapability,
969 remote_id: &StreamEndpointId,
970 ) -> avdtp::Result<(StreamEndpointId, Vec<ServiceCapability>)> {
971 let config = codec_params.try_into()?;
972 let our_direction = self.remote_endpoint(remote_id).map(|e| e.endpoint_type().opposite());
973 debug!(codec_params:?, local:? = self.local; "Looking for compatible local stream");
974 self.local
975 .compatible(config)
976 .find_map(|s| {
977 let endpoint = s.endpoint();
978 if let Some(d) = our_direction {
979 if &d != endpoint.endpoint_type() {
980 return None;
981 }
982 }
983 Some((endpoint.local_id().clone(), endpoint.capabilities().clone()))
984 })
985 .ok_or(avdtp::Error::OutOfRange)
986 }
987
988 fn get_permit_or_reserve(
992 &self,
993 local_id: &StreamEndpointId,
994 ) -> Result<Option<StreamPermit>, ()> {
995 let Some(permits) = self.permits.as_ref() else {
996 return Ok(None);
997 };
998 if let Some(permit) = permits.get(local_id.clone()) {
999 return Ok(Some(permit));
1000 }
1001 info!(peer_id:% = self.peer_id, local_id:%; "No permit to start stream, adding a reservation");
1002 permits.setup_reservation_for(local_id.clone());
1003 Err(())
1004 }
1005
1006 fn start_local_stream(
1009 &mut self,
1010 permit: Option<StreamPermit>,
1011 local_id: &StreamEndpointId,
1012 ) -> avdtp::Result<()> {
1013 let peer_id = self.peer_id;
1014 let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
1015 if permit.as_ref().is_some_and(|p| !p.is_held()) {
1018 return Err(avdtp::Error::Other(anyhow::format_err!(
1019 "streaming permit revoked during setup"
1020 )));
1021 }
1022
1023 info!(peer_id:%, stream:?; "Starting");
1024 let stream_finished = stream.start().map_err(|c| avdtp::Error::RequestInvalid(c))?;
1025 let watched_stream =
1026 WatchedStream::new(permit, stream_finished, self.self_weak.clone(), local_id.clone());
1027 if self.started.insert(local_id.clone(), watched_stream).is_some() {
1028 warn!(peer_id:%, local_id:%; "Stream that was already started");
1029 }
1030 let _ = self.waiting_start_tasks.remove(local_id);
1031 Ok(())
1032 }
1033
1034 fn suspend_local_stream(
1037 &mut self,
1038 local_id: &StreamEndpointId,
1039 ) -> avdtp::Result<StreamEndpointId> {
1040 let peer_id = self.peer_id;
1041 let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
1042 let remote_id = stream.endpoint().remote_id().ok_or(avdtp::Error::InvalidState)?.clone();
1043 info!(peer_id:%; "Suspend stream local {local_id} <-> {remote_id} remote");
1044 stream.suspend().map_err(|c| avdtp::Error::RequestInvalid(c))?;
1045 let _ = self.started.remove(local_id);
1046 Ok(remote_id)
1047 }
1048
1049 fn receive_channel(&mut self, channel: Channel) -> avdtp::Result<Option<StreamEndpointId>> {
1054 let stream_id = self.opening.as_ref().cloned().ok_or(avdtp::Error::InvalidState)?;
1055 let stream = self.get_mut(&stream_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
1056 let done = !stream.endpoint_mut().receive_channel(channel)?;
1057 if done {
1058 self.opening = None;
1059 }
1060 info!(peer_id:% = self.peer_id, stream_id:%; "Transport connected");
1061 Ok(done.then_some(stream_id))
1062 }
1063
1064 fn handle_request(
1066 &mut self,
1067 request: avdtp::Request,
1068 ) -> Either<avdtp::Result<()>, impl Future<Output = avdtp::Result<()>> + use<>> {
1069 use avdtp::ErrorCode;
1070 use avdtp::Request::*;
1071 trace!("Handling {request:?} from peer..");
1072 let immediate_result = 'result: {
1073 match request {
1074 Discover { responder } => responder.send(&self.local.information()),
1075 GetCapabilities { responder, stream_id }
1076 | GetAllCapabilities { responder, stream_id } => match self.local.get(&stream_id) {
1077 None => responder.reject(ErrorCode::BadAcpSeid),
1078 Some(stream) => responder.send(stream.endpoint().capabilities()),
1079 },
1080 Open { responder, stream_id } => {
1081 if self.opening.is_none() {
1082 break 'result responder.reject(ErrorCode::BadState);
1083 }
1084 let Ok(stream) = self.get_mut(&stream_id) else {
1085 break 'result responder.reject(ErrorCode::BadAcpSeid);
1086 };
1087 match stream.endpoint_mut().establish() {
1088 Ok(()) => responder.send(),
1089 Err(_) => responder.reject(ErrorCode::BadState),
1090 }
1091 }
1092 Close { responder, stream_id } => {
1093 let _ = self.waiting_start_tasks.remove(&stream_id);
1094 let peer = self.peer.clone();
1095 let Ok(stream) = self.get_mut(&stream_id) else {
1096 break 'result responder.reject(ErrorCode::BadAcpSeid);
1097 };
1098 stream.release(responder, &peer)
1099 }
1100 SetConfiguration { responder, local_stream_id, remote_stream_id, capabilities } => {
1101 if self.opening.is_some() {
1102 break 'result responder.reject(ServiceCategory::None, ErrorCode::BadState);
1103 }
1104 let peer_id = self.peer_id;
1105 let Ok(stream) = self.get_mut(&local_stream_id) else {
1106 break 'result responder
1107 .reject(ServiceCategory::None, ErrorCode::BadAcpSeid);
1108 };
1109 match stream.configure(&peer_id, &remote_stream_id, capabilities) {
1110 Ok(_) => {
1111 self.opening = Some(local_stream_id);
1112 responder.send()
1113 }
1114 Err((category, code)) => responder.reject(category, code),
1115 }
1116 }
1117 GetConfiguration { stream_id, responder } => {
1118 let Ok(stream) = self.get_mut(&stream_id) else {
1119 break 'result responder.reject(ErrorCode::BadAcpSeid);
1120 };
1121 let Some(vec_capabilities) = stream.endpoint().get_configuration() else {
1122 break 'result responder.reject(ErrorCode::BadState);
1123 };
1124 responder.send(vec_capabilities.as_slice())
1125 }
1126 Reconfigure { responder, local_stream_id, capabilities } => {
1127 let Ok(stream) = self.get_mut(&local_stream_id) else {
1128 break 'result responder
1129 .reject(ServiceCategory::None, ErrorCode::BadAcpSeid);
1130 };
1131 match stream.reconfigure(capabilities) {
1132 Ok(_) => responder.send(),
1133 Err((cat, code)) => responder.reject(cat, code),
1134 }
1135 }
1136 Start { responder, stream_ids } => {
1137 let mut immediate_suspend = Vec::new();
1138 let result = stream_ids.into_iter().try_for_each(|seid| {
1140 let Some(stream) = self.local.get_mut(&seid) else {
1141 return Err((seid, ErrorCode::BadAcpSeid));
1142 };
1143 let remote_id = stream.endpoint().remote_id().cloned();
1144 let Some(remote_id) = remote_id else {
1145 return Err((seid, ErrorCode::BadState));
1146 };
1147 let Ok(permit) = self.get_permit_or_reserve(&seid) else {
1148 immediate_suspend.push(remote_id);
1152 return Ok(());
1153 };
1154 match self.start_local_stream(permit, &seid) {
1155 Ok(()) => Ok(()),
1156 Err(avdtp::Error::RequestInvalid(code)) => Err((seid, code)),
1157 Err(_) => Err((seid, ErrorCode::BadState)),
1158 }
1159 });
1160 let response_result = match result {
1161 Ok(()) => responder.send(),
1162 Err((seid, code)) => responder.reject(&seid, code),
1163 };
1164 {
1165 let peer = self.peer.clone();
1166 return Either::Right(async move {
1167 if !immediate_suspend.is_empty() {
1168 peer.suspend(immediate_suspend.as_slice()).await?;
1169 }
1170 response_result
1171 });
1172 }
1173 }
1174 Suspend { responder, stream_ids } => {
1175 for seid in stream_ids {
1176 match self.suspend_local_stream(&seid) {
1177 Ok(_remote_id) => {}
1178 Err(avdtp::Error::RequestInvalid(code)) => {
1179 break 'result responder.reject(&seid, code);
1180 }
1181 Err(_e) => break 'result responder.reject(&seid, ErrorCode::BadState),
1182 }
1183 }
1184 responder.send()
1185 }
1186 Abort { responder, stream_id } => {
1187 let _ = self.waiting_start_tasks.remove(&stream_id);
1188 let Ok(stream) = self.get_mut(&stream_id) else {
1189 break 'result Ok(());
1191 };
1192 stream.abort();
1193 self.opening = self.opening.take().filter(|local_id| local_id != &stream_id);
1194 responder.send()
1195 }
1196 DelayReport { responder, delay, stream_id } => {
1197 let delay_ns = delay as u64 * 100000;
1199 self.metrics.log_integer(
1201 bt_metrics::AVDTP_DELAY_REPORT_IN_NANOSECONDS_METRIC_ID,
1202 delay_ns.try_into().unwrap_or(-1),
1203 vec![],
1204 );
1205 let Some(stream) = self.local.get_mut(&stream_id) else {
1207 break 'result responder.reject(avdtp::ErrorCode::BadAcpSeid);
1208 };
1209 let delay_str = format!("delay {}.{} ms", delay / 10, delay % 10);
1210 let peer = self.peer_id;
1211 match stream.set_delay(std::time::Duration::from_nanos(delay_ns)) {
1212 Ok(()) => info!(peer:%, stream_id:%; "reported {delay_str}"),
1213 Err(avdtp::ErrorCode::BadState) => {
1214 info!(peer:%, stream_id:%; "bad state {delay_str}");
1215 break 'result responder.reject(avdtp::ErrorCode::BadState);
1216 }
1217 Err(e) => info!(peer:%, stream_id:%, e:?; "failed {delay_str}"),
1218 };
1219 responder.send()
1221 }
1222 }
1223 };
1224 Either::Left(immediate_result)
1225 }
1226}
1227
1228struct WatchedStream {
1231 _permit_task: fasync::Task<()>,
1232}
1233
1234impl WatchedStream {
1235 fn new(
1236 permit: Option<StreamPermit>,
1237 finish_fut: BoxFuture<'static, Result<MediaTaskStatus, anyhow::Error>>,
1238 weak: Weak<Mutex<PeerInner>>,
1239 local_id: StreamEndpointId,
1240 ) -> Self {
1241 let permit_task = fasync::Task::spawn(async move {
1242 let finish_result = finish_fut.await;
1243 drop(permit);
1244 let Some(peer) = weak.upgrade() else {
1245 return;
1246 };
1247 let finished_task = fasync::Task::spawn(Self::handle_stream_finished(
1250 finish_result,
1251 weak,
1252 local_id.clone(),
1253 ));
1254 let _ = peer.lock().waiting_start_tasks.insert(local_id, finished_task);
1255 });
1256 Self { _permit_task: permit_task }
1257 }
1258
1259 async fn handle_stream_finished(
1260 finish_result: Result<MediaTaskStatus, anyhow::Error>,
1261 weak: Weak<Mutex<PeerInner>>,
1262 local_id: StreamEndpointId,
1263 ) {
1264 let is_audio_disabled = matches!(finish_result, Ok(MediaTaskStatus::AudioDisabled));
1265 let is_stopped = matches!(finish_result, Ok(MediaTaskStatus::Stopped));
1266
1267 let stream_id = local_id.clone();
1268 info!("Audio stopped {finish_result:?} - suspending A2DP stream for {stream_id:?}");
1269 if !is_stopped {
1270 if let Err(e) = PeerInner::suspend(weak.clone(), stream_id.clone()).await {
1271 warn!("Error suspending stream after audio stopped: {e:?}");
1272 }
1273 }
1274 PeerInner::forget_waiting_start(&weak, &local_id);
1276 if is_audio_disabled {
1277 PeerInner::schedule_start(weak, stream_id);
1278 }
1279 }
1280}
1281
1282fn codectype_to_availability_metric(
1283 codec_type: &MediaCodecType,
1284) -> bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec {
1285 match codec_type {
1286 &MediaCodecType::AUDIO_SBC => {
1287 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc
1288 }
1289 &MediaCodecType::AUDIO_MPEG12 => {
1290 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Mpeg12
1291 }
1292 &MediaCodecType::AUDIO_AAC => {
1293 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Aac
1294 }
1295 &MediaCodecType::AUDIO_ATRAC => {
1296 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Atrac
1297 }
1298 &MediaCodecType::AUDIO_NON_A2DP => {
1299 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::VendorSpecific
1300 }
1301 _ => bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Unknown,
1302 }
1303}
1304
1305fn capability_to_metric(
1306 cap: &ServiceCapability,
1307) -> Option<bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability> {
1308 match cap {
1309 ServiceCapability::DelayReporting => {
1310 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::DelayReport)
1311 }
1312 ServiceCapability::Reporting => {
1313 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Reporting)
1314 }
1315 ServiceCapability::Recovery { .. } => {
1316 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Recovery)
1317 }
1318 ServiceCapability::ContentProtection { .. } => {
1319 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::ContentProtection)
1320 }
1321 ServiceCapability::HeaderCompression { .. } => {
1322 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::HeaderCompression)
1323 }
1324 ServiceCapability::Multiplexing { .. } => {
1325 Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Multiplexing)
1326 }
1327 other => {
1329 trace!("untracked remote peer capability: {:?}", other);
1330 None
1331 }
1332 }
1333}
1334
1335#[cfg(test)]
1336mod tests {
1337 use super::*;
1338
1339 use async_utils::PollExt;
1340 use bt_channel_test_support::{Transport, create_test_channels};
1341 use bt_metrics::respond_to_metrics_req_for_test;
1342 use fidl::endpoints::create_proxy_and_stream;
1343 use fidl_fuchsia_bluetooth::ErrorCode;
1344
1345 use fidl_fuchsia_bluetooth_bredr::{
1346 ProfileMarker, ProfileRequest, ProfileRequestStream, ServiceClassProfileIdentifier,
1347 };
1348 use fidl_fuchsia_metrics::{MetricEvent, MetricEventPayload};
1349 use futures::{SinkExt, StreamExt};
1350 use std::pin::pin;
1351 use test_case::test_case;
1352
1353 use crate::media_task::tests::{TestMediaTask, TestMediaTaskBuilder};
1354 use crate::media_types::*;
1355 use crate::stream::tests::{make_sbc_endpoint, sbc_mediacodec_capability};
1356
1357 fn fake_metrics()
1358 -> (bt_metrics::MetricsLogger, fidl_fuchsia_metrics::MetricEventLoggerRequestStream) {
1359 let (c, s) = fidl::endpoints::create_proxy_and_stream::<
1360 fidl_fuchsia_metrics::MetricEventLoggerMarker,
1361 >();
1362 (bt_metrics::MetricsLogger::from_proxy(c), s)
1363 }
1364
1365 fn setup_avdtp_peer(transport: Transport) -> (avdtp::Peer, Channel) {
1366 let (signaling, remote) = create_test_channels(transport);
1367 let peer = avdtp::Peer::new(signaling);
1368 (peer, remote)
1369 }
1370
1371 fn build_test_streams() -> Streams {
1372 let mut streams = Streams::default();
1373 let source = Stream::build(
1374 make_sbc_endpoint(1, avdtp::EndpointType::Source),
1375 TestMediaTaskBuilder::new_delayable().builder(),
1376 );
1377 streams.insert(source);
1378 let sink = Stream::build(
1379 make_sbc_endpoint(2, avdtp::EndpointType::Sink),
1380 TestMediaTaskBuilder::new().builder(),
1381 );
1382 streams.insert(sink);
1383 streams
1384 }
1385
1386 fn build_test_streams_delayable() -> Streams {
1387 fn with_delay(seid: u8, direction: avdtp::EndpointType) -> StreamEndpoint {
1388 StreamEndpoint::new(
1389 seid,
1390 avdtp::MediaType::Audio,
1391 direction,
1392 vec![
1393 avdtp::ServiceCapability::MediaTransport,
1394 avdtp::ServiceCapability::DelayReporting,
1395 sbc_mediacodec_capability(),
1396 ],
1397 )
1398 .expect("endpoint creation should succeed")
1399 }
1400 let mut streams = Streams::default();
1401 let source = Stream::build(
1402 with_delay(1, avdtp::EndpointType::Source),
1403 TestMediaTaskBuilder::new_delayable().builder(),
1404 );
1405 streams.insert(source);
1406 let sink = Stream::build(
1407 with_delay(2, avdtp::EndpointType::Sink),
1408 TestMediaTaskBuilder::new().builder(),
1409 );
1410 streams.insert(sink);
1411 streams
1412 }
1413
1414 #[track_caller]
1415 pub(crate) fn recv_remote(
1416 exec: &mut fasync::TestExecutor,
1417 remote: &mut Channel,
1418 ) -> Result<Vec<u8>, zx::Status> {
1419 let mut fut = remote.next();
1420 match exec.run_until_stalled(&mut fut) {
1421 Poll::Ready(Some(res)) => res,
1422 Poll::Ready(None) => Err(zx::Status::PEER_CLOSED),
1423 Poll::Pending => Err(zx::Status::SHOULD_WAIT),
1424 }
1425 }
1426
1427 fn setup_test_peer(
1430 transport: Transport,
1431 use_cobalt: bool,
1432 streams: Streams,
1433 permits: Option<Permits>,
1434 ) -> (
1435 Channel,
1436 ProfileRequestStream,
1437 Option<fidl_fuchsia_metrics::MetricEventLoggerRequestStream>,
1438 Peer,
1439 ) {
1440 let (avdtp, remote) = setup_avdtp_peer(transport);
1441 let (metrics_logger, cobalt_receiver) = if use_cobalt {
1442 let (l, r) = fake_metrics();
1443 (l, Some(r))
1444 } else {
1445 (bt_metrics::MetricsLogger::default(), None)
1446 };
1447 let (profile_proxy, requests) = create_proxy_and_stream::<ProfileMarker>();
1448 let peer =
1449 Peer::create(PeerId(1), avdtp, streams, permits, profile_proxy, None, metrics_logger);
1450
1451 (remote, requests, cobalt_receiver, peer)
1452 }
1453
1454 #[track_caller]
1455 fn expect_send(exec: &mut fasync::TestExecutor, remote: &mut Channel, data: Vec<u8>) {
1456 exec.run_until_stalled(&mut remote.send(data))
1457 .expect("poll is ready")
1458 .expect("write successful");
1459 }
1460
1461 fn expect_get_capabilities_and_respond(
1462 exec: &mut fasync::TestExecutor,
1463 remote: &mut Channel,
1464 expected_seid: u8,
1465 response_capabilities: &[u8],
1466 ) {
1467 let received = recv_remote(exec, remote).unwrap();
1468 assert_eq!(0x00, received[0] & 0xF);
1470 assert_eq!(0x02, received[1]); assert_eq!(expected_seid << 2, received[2]);
1472
1473 let txlabel_raw = received[0] & 0xF0;
1474
1475 #[rustfmt::skip]
1477 let mut get_capabilities_rsp = vec![
1478 txlabel_raw << 4 | 0x2, 0x02 ];
1481
1482 get_capabilities_rsp.extend_from_slice(response_capabilities);
1483
1484 expect_send(exec, remote, get_capabilities_rsp);
1485 }
1486
1487 fn expect_get_all_capabilities_and_respond(
1488 exec: &mut fasync::TestExecutor,
1489 remote: &mut Channel,
1490 expected_seid: u8,
1491 response_capabilities: &[u8],
1492 ) {
1493 let received = recv_remote(exec, remote).unwrap();
1494 assert_eq!(0x00, received[0] & 0xF);
1496 assert_eq!(0x0C, received[1]); assert_eq!(expected_seid << 2, received[2]);
1498
1499 let txlabel_raw = received[0] & 0xF0;
1500
1501 #[rustfmt::skip]
1503 let mut get_capabilities_rsp = vec![
1504 txlabel_raw << 4 | 0x2, 0x0C ];
1507
1508 get_capabilities_rsp.extend_from_slice(response_capabilities);
1509
1510 expect_send(exec, remote, get_capabilities_rsp);
1511 }
1512
1513 #[test_case(Transport::Socket ; "socket")]
1514 #[test_case(Transport::Fidl ; "fidl")]
1515 fn disconnected(transport: Transport) {
1516 let mut exec = fasync::TestExecutor::new();
1517 let (proxy, _stream) = create_proxy_and_stream::<ProfileMarker>();
1518 let (signaling, remote) = create_test_channels(transport);
1519
1520 let id = PeerId(1);
1521
1522 let avdtp = avdtp::Peer::new(signaling);
1523 let peer = Peer::create(
1524 id,
1525 avdtp,
1526 Streams::default(),
1527 None,
1528 proxy,
1529 None,
1530 bt_metrics::MetricsLogger::default(),
1531 );
1532
1533 let closed_fut = peer.closed();
1534
1535 let mut closed_fut = pin!(closed_fut);
1536
1537 assert!(exec.run_until_stalled(&mut closed_fut).is_pending());
1538
1539 drop(remote);
1541
1542 assert!(exec.run_until_stalled(&mut closed_fut).is_ready());
1543 }
1544
1545 #[test_case(Transport::Socket ; "socket")]
1546 #[test_case(Transport::Fidl ; "fidl")]
1547 fn peer_collect_capabilities_success(transport: Transport) {
1548 let mut exec = fasync::TestExecutor::new();
1549
1550 let (mut remote, _, cobalt_receiver, peer) =
1551 setup_test_peer(transport, true, build_test_streams(), None);
1552
1553 let p: ProfileDescriptor = ProfileDescriptor {
1554 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
1555 major_version: Some(1),
1556 minor_version: Some(2),
1557 ..Default::default()
1558 };
1559 let _ = peer.set_descriptor(p);
1560
1561 let collect_future = peer.collect_capabilities();
1562 let mut collect_future = pin!(collect_future);
1563
1564 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1565
1566 let received = recv_remote(&mut exec, &mut remote).unwrap();
1568 assert_eq!(0x00, received[0] & 0xF);
1570 assert_eq!(0x01, received[1]); let txlabel_raw = received[0] & 0xF0;
1573
1574 let response: &[u8] = &[
1576 txlabel_raw << 4 | 0x0 << 2 | 0x2, 0x01, 0x3E << 2 | 0x0 << 1, 0x00 << 4 | 0x1 << 3, 0x01 << 2 | 0x1 << 1, 0x00 << 4 | 0x1 << 3, ];
1583 expect_send(&mut exec, &mut remote, response.to_vec());
1584
1585 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1586
1587 #[rustfmt::skip]
1589 let capabilities_rsp = &[
1590 0x01, 0x00,
1592 0x07, 0x06, 0x00, 0x04, 0xF0, 0x9F, 0x92, 0x96
1594 ];
1595 expect_get_capabilities_and_respond(&mut exec, &mut remote, 0x3E, capabilities_rsp);
1596
1597 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1598
1599 #[rustfmt::skip]
1601 let capabilities_rsp = &[
1602 0x01, 0x00,
1604 0x07, 0x04, 0x00, 0x00, 0xC0, 0xDE
1606 ];
1607 expect_get_capabilities_and_respond(&mut exec, &mut remote, 0x01, capabilities_rsp);
1608
1609 match exec.run_until_stalled(&mut collect_future) {
1610 Poll::Pending => panic!("collect capabilities should be complete"),
1611 Poll::Ready(Err(e)) => panic!("collect capabilities should have succeeded: {}", e),
1612 Poll::Ready(Ok(endpoints)) => {
1613 let first_seid: StreamEndpointId = 0x3E_u8.try_into().unwrap();
1614 let second_seid: StreamEndpointId = 0x01_u8.try_into().unwrap();
1615 for stream in endpoints {
1616 if stream.local_id() == &first_seid {
1617 let expected_caps = vec![
1618 ServiceCapability::MediaTransport,
1619 ServiceCapability::MediaCodec {
1620 media_type: avdtp::MediaType::Audio,
1621 codec_type: avdtp::MediaCodecType::new(0x04),
1622 codec_extra: vec![0xF0, 0x9F, 0x92, 0x96],
1623 },
1624 ];
1625 assert_eq!(&expected_caps, stream.capabilities());
1626 } else if stream.local_id() == &second_seid {
1627 let expected_codec_type = avdtp::MediaCodecType::new(0x00);
1628 assert_eq!(Some(&expected_codec_type), stream.codec_type());
1629 } else {
1630 panic!("Unexpected endpoint in the streams collected");
1631 }
1632 }
1633 }
1634 }
1635
1636 let mut recv = cobalt_receiver.expect("should have receiver");
1638 let mut log_events = Vec::new();
1639 while let Poll::Ready(Some(Ok(req))) = exec.run_until_stalled(&mut recv.next()) {
1640 log_events.push(respond_to_metrics_req_for_test(req));
1641 }
1642
1643 assert_eq!(3, log_events.len());
1645 assert!(log_events.contains(&MetricEvent {
1646 metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1647 event_codes: vec![
1648 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc as u32
1649 ],
1650 payload: MetricEventPayload::Count(1),
1651 }));
1652 assert!(log_events.contains(&MetricEvent {
1653 metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1654 event_codes: vec![
1655 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Atrac as u32
1656 ],
1657 payload: MetricEventPayload::Count(1),
1658 }));
1659 assert!(log_events.contains(&MetricEvent {
1660 metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1661 event_codes: vec![
1662 bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic as u32
1663 ],
1664 payload: MetricEventPayload::Count(1),
1665 }));
1666
1667 let collect_future = peer.collect_capabilities();
1669 let mut collect_future = pin!(collect_future);
1670
1671 match exec.run_until_stalled(&mut collect_future) {
1672 Poll::Ready(Ok(endpoints)) => assert_eq!(2, endpoints.len()),
1673 x => panic!("Expected get remote capabilities to be done, got {:?}", x),
1674 };
1675 }
1676
1677 #[test_case(Transport::Socket ; "socket")]
1678 #[test_case(Transport::Fidl ; "fidl")]
1679 fn peer_collect_all_capabilities_success(transport: Transport) {
1680 let mut exec = fasync::TestExecutor::new();
1681
1682 let (mut remote, _, cobalt_receiver, peer) =
1683 setup_test_peer(transport, true, build_test_streams(), None);
1684 let p: ProfileDescriptor = ProfileDescriptor {
1685 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
1686 major_version: Some(1),
1687 minor_version: Some(3),
1688 ..Default::default()
1689 };
1690 let _ = peer.set_descriptor(p);
1691
1692 let collect_future = peer.collect_capabilities();
1693 let mut collect_future = pin!(collect_future);
1694
1695 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1696
1697 let received = recv_remote(&mut exec, &mut remote).unwrap();
1699 assert_eq!(0x00, received[0] & 0xF);
1701 assert_eq!(0x01, received[1]); let txlabel_raw = received[0] & 0xF0;
1704
1705 let response: &[u8] = &[
1707 txlabel_raw << 4 | 0x0 << 2 | 0x2, 0x01, 0x3E << 2 | 0x0 << 1, 0x00 << 4 | 0x1 << 3, 0x01 << 2 | 0x1 << 1, 0x00 << 4 | 0x1 << 3, ];
1714 expect_send(&mut exec, &mut remote, response.to_vec());
1715
1716 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1717
1718 #[rustfmt::skip]
1720 let capabilities_rsp = &[
1721 0x01, 0x00,
1723 0x07, 0x06, 0x00, 0x40, 0xF0, 0x9F, 0x92, 0x96,
1725 0x08, 0x00
1727 ];
1728 expect_get_all_capabilities_and_respond(&mut exec, &mut remote, 0x3E, capabilities_rsp);
1729
1730 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1731
1732 #[rustfmt::skip]
1734 let capabilities_rsp = &[
1735 0x01, 0x00,
1737 0x07, 0x04, 0x00, 0x00, 0xC0, 0xDE
1739 ];
1740 expect_get_all_capabilities_and_respond(&mut exec, &mut remote, 0x01, capabilities_rsp);
1741
1742 match exec.run_until_stalled(&mut collect_future) {
1743 Poll::Pending => panic!("collect capabilities should be complete"),
1744 Poll::Ready(Err(e)) => panic!("collect capabilities should have succeeded: {}", e),
1745 Poll::Ready(Ok(endpoints)) => {
1746 let first_seid: StreamEndpointId = 0x3E_u8.try_into().unwrap();
1747 let second_seid: StreamEndpointId = 0x01_u8.try_into().unwrap();
1748 for stream in endpoints {
1749 if stream.local_id() == &first_seid {
1750 let expected_caps = vec![
1751 ServiceCapability::MediaTransport,
1752 ServiceCapability::MediaCodec {
1753 media_type: avdtp::MediaType::Audio,
1754 codec_type: avdtp::MediaCodecType::new(0x40),
1755 codec_extra: vec![0xF0, 0x9F, 0x92, 0x96],
1756 },
1757 ServiceCapability::DelayReporting,
1758 ];
1759 assert_eq!(&expected_caps, stream.capabilities());
1760 } else if stream.local_id() == &second_seid {
1761 let expected_codec_type = avdtp::MediaCodecType::new(0x00);
1762 assert_eq!(Some(&expected_codec_type), stream.codec_type());
1763 } else {
1764 panic!("Unexpected endpoint in the streams collected");
1765 }
1766 }
1767 }
1768 }
1769
1770 let mut recv = cobalt_receiver.expect("should have receiver");
1772 let mut log_events = Vec::new();
1773 while let Poll::Ready(Some(Ok(req))) = exec.run_until_stalled(&mut recv.next()) {
1774 log_events.push(respond_to_metrics_req_for_test(req));
1775 }
1776
1777 assert_eq!(4, log_events.len());
1779 assert!(log_events.contains(&MetricEvent {
1780 metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1781 event_codes: vec![
1782 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Unknown as u32
1783 ],
1784 payload: MetricEventPayload::Count(1),
1785 }));
1786 assert!(log_events.contains(&MetricEvent {
1787 metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1788 event_codes: vec![
1789 bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc as u32
1790 ],
1791 payload: MetricEventPayload::Count(1),
1792 }));
1793 assert!(log_events.contains(&MetricEvent {
1794 metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1795 event_codes: vec![
1796 bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic as u32
1797 ],
1798 payload: MetricEventPayload::Count(1),
1799 }));
1800 assert!(log_events.contains(&MetricEvent {
1801 metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1802 event_codes: vec![
1803 bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::DelayReport as u32
1804 ],
1805 payload: MetricEventPayload::Count(1),
1806 }));
1807
1808 let collect_future = peer.collect_capabilities();
1810 let mut collect_future = pin!(collect_future);
1811
1812 match exec.run_until_stalled(&mut collect_future) {
1813 Poll::Ready(Ok(endpoints)) => assert_eq!(2, endpoints.len()),
1814 x => panic!("Expected get remote capabilities to be done, got {:?}", x),
1815 };
1816 }
1817
1818 #[test_case(Transport::Socket ; "socket")]
1819 #[test_case(Transport::Fidl ; "fidl")]
1820 fn peer_collect_capabilities_discovery_fails(transport: Transport) {
1821 let mut exec = fasync::TestExecutor::new();
1822
1823 let (mut remote, _, _, peer) =
1824 setup_test_peer(transport, false, build_test_streams(), None);
1825
1826 let collect_future = peer.collect_capabilities();
1827 let mut collect_future = pin!(collect_future);
1828
1829 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1831
1832 let received = recv_remote(&mut exec, &mut remote).unwrap();
1834 assert_eq!(0x00, received[0] & 0xF);
1836 assert_eq!(0x01, received[1]); let txlabel_raw = received[0] & 0xF0;
1839
1840 let response: &[u8] = &[
1842 txlabel_raw | 0x0 << 2 | 0x3, 0x01, 0x31, ];
1846 expect_send(&mut exec, &mut remote, response.to_vec());
1847
1848 match exec.run_until_stalled(&mut collect_future) {
1851 Poll::Pending => panic!("Should be ready after discovery failure"),
1852 Poll::Ready(Ok(x)) => panic!("Should be an error but returned {x:?}"),
1853 Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
1854 assert_eq!(Some(Ok(avdtp::ErrorCode::BadState)), e.error_code());
1855 }
1856 Poll::Ready(Err(e)) => panic!("Should have been a RemoteRejected was was {e:?}"),
1857 }
1858 }
1859
1860 #[test_case(Transport::Socket ; "socket")]
1861 #[test_case(Transport::Fidl ; "fidl")]
1862 fn peer_collect_capabilities_get_capability_fails(transport: Transport) {
1863 let mut exec = fasync::TestExecutor::new();
1864
1865 let (mut remote, _, _, peer) = setup_test_peer(transport, true, build_test_streams(), None);
1866
1867 let collect_future = peer.collect_capabilities();
1868 let mut collect_future = pin!(collect_future);
1869
1870 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1872
1873 let received = recv_remote(&mut exec, &mut remote).unwrap();
1875 assert_eq!(0x00, received[0] & 0xF);
1877 assert_eq!(0x01, received[1]); let txlabel_raw = received[0] & 0xF0;
1880
1881 let response: &[u8] = &[
1883 txlabel_raw << 4 | 0x0 << 2 | 0x2, 0x01, 0x3E << 2 | 0x0 << 1, 0x00 << 4 | 0x1 << 3, 0x01 << 2 | 0x1 << 1, 0x00 << 4 | 0x1 << 3, ];
1890 expect_send(&mut exec, &mut remote, response.to_vec());
1891
1892 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1893
1894 let expected_seid = 0x3E;
1896 let received = recv_remote(&mut exec, &mut remote).unwrap();
1897 assert_eq!(0x00, received[0] & 0xF);
1899 assert_eq!(0x02, received[1]); assert_eq!(expected_seid << 2, received[2]);
1901
1902 let txlabel_raw = received[0] & 0xF0;
1903
1904 let response: &[u8] = &[
1905 txlabel_raw | 0x0 << 2 | 0x3, 0x02, 0x12, ];
1909 expect_send(&mut exec, &mut remote, response.to_vec());
1910
1911 assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1912
1913 let expected_seid = 0x01;
1915 let received = recv_remote(&mut exec, &mut remote).unwrap();
1916 assert_eq!(0x00, received[0] & 0xF);
1918 assert_eq!(0x02, received[1]); assert_eq!(expected_seid << 2, received[2]);
1920
1921 let txlabel_raw = received[0] & 0xF0;
1922
1923 let response: &[u8] = &[
1924 txlabel_raw | 0x0 << 2 | 0x3, 0x02, 0x12, ];
1928 expect_send(&mut exec, &mut remote, response.to_vec());
1929
1930 match exec.run_until_stalled(&mut collect_future) {
1932 Poll::Pending => panic!("Should be ready after discovery failure"),
1933 Poll::Ready(Err(e)) => panic!("Shouldn't be an error but returned {:?}", e),
1934 Poll::Ready(Ok(map)) => assert_eq!(0, map.len()),
1935 }
1936 }
1937
1938 fn receive_simple_accept(exec: &mut fasync::TestExecutor, remote: &mut Channel, signal_id: u8) {
1939 let received = recv_remote(exec, remote).expect("expected a packet");
1940 assert_eq!(0x00, received[0] & 0xF);
1942 assert_eq!(signal_id, received[1]);
1943
1944 let txlabel_raw = received[0] & 0xF0;
1945
1946 let response: &[u8] = &[
1947 txlabel_raw | 0x0 << 2 | 0x2, signal_id,
1949 ];
1950 expect_send(exec, remote, response.to_vec());
1951 }
1952
1953 #[test_case(Transport::Socket ; "socket")]
1954 #[test_case(Transport::Fidl ; "fidl")]
1955 fn peer_stream_start_success(transport: Transport) {
1956 let mut exec = fasync::TestExecutor::new();
1957
1958 let (mut remote, mut profile_request_stream, _, peer) =
1959 setup_test_peer(transport, false, build_test_streams(), None);
1960
1961 let remote_seid = 2_u8.try_into().unwrap();
1962
1963 let codec_params = ServiceCapability::MediaCodec {
1964 media_type: avdtp::MediaType::Audio,
1965 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
1966 codec_extra: vec![0x11, 0x45, 51, 51],
1967 };
1968
1969 let remote_endpoint = avdtp::StreamEndpoint::new(
1971 2,
1972 avdtp::MediaType::Audio,
1973 avdtp::EndpointType::Sink,
1974 vec![codec_params.clone()],
1975 )
1976 .expect("valid endpoint");
1977 peer.inner.lock().set_remote_endpoints(&[remote_endpoint]);
1978
1979 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
1980 let mut start_future = pin!(start_future);
1981
1982 match exec.run_until_stalled(&mut start_future) {
1983 Poll::Pending => {}
1984 x => panic!("Expected pending, but got {x:?}"),
1985 };
1986
1987 receive_simple_accept(&mut exec, &mut remote, 0x03); assert!(exec.run_until_stalled(&mut start_future).is_pending());
1990
1991 receive_simple_accept(&mut exec, &mut remote, 0x06); match exec.run_until_stalled(&mut start_future) {
1994 Poll::Pending => {}
1995 Poll::Ready(Err(e)) => panic!("Expected to be pending but error: {:?}", e),
1996 Poll::Ready(Ok(_)) => panic!("Expected to be pending but finished!"),
1997 };
1998
1999 let (transport_chan, _remote_chan) = create_test_channels(transport);
2001
2002 let request = exec.run_until_stalled(&mut profile_request_stream.next());
2003 match request {
2004 Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, connection, responder }))) => {
2005 assert_eq!(PeerId(1), peer_id.into());
2006 assert_eq!(connection, ConnectParameters::L2cap(Peer::transport_channel_params()));
2007 let channel = transport_chan.try_into().unwrap();
2008 responder.send(Ok(channel)).expect("responder sends");
2009 }
2010 x => panic!("Should have sent a open l2cap request, but got {:?}", x),
2011 };
2012
2013 exec.run_until_stalled(&mut start_future)
2016 .expect("start setup finished")
2017 .expect("stream setup is ok");
2018
2019 receive_simple_accept(&mut exec, &mut remote, 0x07); let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
2024 assert!(peer.is_streaming_now());
2025 }
2026
2027 #[test_case(Transport::Socket ; "socket")]
2028 #[test_case(Transport::Fidl ; "fidl")]
2029 fn peer_stream_start_picks_correct_direction(transport: Transport) {
2030 let mut exec = fasync::TestExecutor::new();
2031
2032 let (remote, _, _, peer) = setup_test_peer(transport, false, build_test_streams(), None);
2033 let remote = avdtp::Peer::new(remote);
2034 let mut remote_events = remote.take_request_stream();
2035
2036 fn remote_handle_request(req: avdtp::Request) {
2038 let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2039 let res = match req {
2040 avdtp::Request::Discover { responder } => {
2041 let infos = [avdtp::StreamInformation::new(
2042 expected_stream_id,
2043 false,
2044 avdtp::MediaType::Audio,
2045 avdtp::EndpointType::Source,
2046 )];
2047 responder.send(&infos)
2048 }
2049 avdtp::Request::GetAllCapabilities { stream_id, responder }
2050 | avdtp::Request::GetCapabilities { stream_id, responder } => {
2051 assert_eq!(expected_stream_id, stream_id);
2052 let caps = vec![
2053 ServiceCapability::MediaTransport,
2054 ServiceCapability::MediaCodec {
2055 media_type: avdtp::MediaType::Audio,
2056 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2057 codec_extra: vec![0x11, 0x45, 51, 250],
2058 },
2059 ];
2060 responder.send(&caps[..])
2061 }
2062 avdtp::Request::Open { responder, stream_id } => {
2063 assert_eq!(expected_stream_id, stream_id);
2064 responder.send()
2065 }
2066 avdtp::Request::SetConfiguration {
2067 responder,
2068 local_stream_id,
2069 remote_stream_id,
2070 ..
2071 } => {
2072 assert_eq!(local_stream_id, expected_stream_id);
2073 assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
2075 responder.send()
2076 }
2077 x => panic!("Unexpected request: {:?}", x),
2078 };
2079 res.expect("should be able to respond");
2080 }
2081
2082 let collect_capabilities_fut = peer.collect_capabilities();
2084 let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2085
2086 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2087
2088 let request = exec.run_singlethreaded(&mut remote_events.next());
2089 remote_handle_request(request.expect("should have a discovery request").unwrap());
2090
2091 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2092 let request = exec.run_singlethreaded(&mut remote_events.next());
2093 remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2094
2095 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2096
2097 let remote_seid = 4_u8.try_into().unwrap();
2099
2100 let codec_params = ServiceCapability::MediaCodec {
2101 media_type: avdtp::MediaType::Audio,
2102 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2103 codec_extra: vec![0x11, 0x45, 51, 51],
2104 };
2105 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2106 let mut start_future = pin!(start_future);
2107
2108 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2109 let request = exec.run_singlethreaded(&mut remote_events.next());
2110 remote_handle_request(request.expect("should have a set_capabilities request").unwrap());
2111
2112 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2113 let request = exec.run_singlethreaded(&mut remote_events.next());
2114 remote_handle_request(request.expect("should have an open request").unwrap());
2115 }
2116
2117 #[test_case(Transport::Socket ; "socket")]
2118 #[test_case(Transport::Fidl ; "fidl")]
2119 fn peer_stream_start_strips_unsupported_local_capabilities(transport: Transport) {
2120 let mut exec = fasync::TestExecutor::new();
2121
2122 let (remote, _, _, peer) = setup_test_peer(transport, false, build_test_streams(), None);
2123 let remote = avdtp::Peer::new(remote);
2124 let mut remote_events = remote.take_request_stream();
2125
2126 fn remote_handle_request(req: avdtp::Request) {
2128 let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2129 let res = match req {
2130 avdtp::Request::Discover { responder } => {
2131 let infos = [avdtp::StreamInformation::new(
2132 expected_stream_id,
2133 false,
2134 avdtp::MediaType::Audio,
2135 avdtp::EndpointType::Source,
2136 )];
2137 responder.send(&infos)
2138 }
2139 avdtp::Request::GetAllCapabilities { stream_id, responder }
2140 | avdtp::Request::GetCapabilities { stream_id, responder } => {
2141 assert_eq!(expected_stream_id, stream_id);
2142 let caps = vec![
2143 ServiceCapability::MediaTransport,
2144 ServiceCapability::DelayReporting,
2146 ServiceCapability::MediaCodec {
2147 media_type: avdtp::MediaType::Audio,
2148 codec_type: avdtp::MediaCodecType::AUDIO_AAC,
2149 codec_extra: vec![128, 0, 132, 134, 0, 0],
2150 },
2151 ];
2152 responder.send(&caps[..])
2153 }
2154 avdtp::Request::Open { responder, stream_id } => {
2155 assert_eq!(expected_stream_id, stream_id);
2156 responder.send()
2157 }
2158 avdtp::Request::SetConfiguration {
2159 responder,
2160 local_stream_id,
2161 remote_stream_id,
2162 capabilities,
2163 } => {
2164 assert_eq!(local_stream_id, expected_stream_id);
2165 assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
2167 assert!(!capabilities.contains(&ServiceCapability::DelayReporting));
2170 responder.send()
2171 }
2172 x => panic!("Unexpected request: {:?}", x),
2173 };
2174 res.expect("should be able to respond");
2175 }
2176
2177 let collect_capabilities_fut = peer.collect_capabilities();
2179 let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2180
2181 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2182
2183 let request = exec.run_singlethreaded(&mut remote_events.next());
2184 remote_handle_request(request.expect("should have a discovery request").unwrap());
2185
2186 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2187 let request = exec.run_singlethreaded(&mut remote_events.next());
2188 remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2189
2190 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2191
2192 let remote_seid = 4_u8.try_into().unwrap();
2194
2195 let codec_params = ServiceCapability::MediaCodec {
2196 media_type: avdtp::MediaType::Audio,
2197 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2198 codec_extra: vec![0x11, 0x45, 51, 51],
2199 };
2200 let start_future =
2201 peer.stream_start(remote_seid, vec![codec_params, ServiceCapability::DelayReporting]);
2202 let mut start_future = pin!(start_future);
2203
2204 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2205 let request = exec.run_singlethreaded(&mut remote_events.next());
2206 remote_handle_request(request.expect("should have a set_configuration request").unwrap());
2207
2208 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2209 let request = exec.run_singlethreaded(&mut remote_events.next());
2210 remote_handle_request(request.expect("should have an open request").unwrap());
2211 }
2212
2213 #[test_case(Transport::Socket ; "socket")]
2214 #[test_case(Transport::Fidl ; "fidl")]
2215 fn peer_stream_start_orders_local_capabilities(transport: Transport) {
2216 let mut exec = fasync::TestExecutor::new();
2217
2218 let (remote, _, _, peer) =
2219 setup_test_peer(transport, false, build_test_streams_delayable(), None);
2220 let remote = avdtp::Peer::new(remote);
2221 let mut remote_events = remote.take_request_stream();
2222
2223 fn remote_handle_request(req: avdtp::Request) {
2225 let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2226 let res = match req {
2227 avdtp::Request::Discover { responder } => {
2228 let infos = [avdtp::StreamInformation::new(
2229 expected_stream_id,
2230 false,
2231 avdtp::MediaType::Audio,
2232 avdtp::EndpointType::Source,
2233 )];
2234 responder.send(&infos)
2235 }
2236 avdtp::Request::GetAllCapabilities { stream_id, responder }
2237 | avdtp::Request::GetCapabilities { stream_id, responder } => {
2238 assert_eq!(expected_stream_id, stream_id);
2239 let caps = &[
2240 ServiceCapability::MediaTransport,
2241 ServiceCapability::MediaCodec {
2242 media_type: avdtp::MediaType::Audio,
2243 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2244 codec_extra: vec![0x11, 0x45, 51, 250],
2245 },
2246 ServiceCapability::DelayReporting,
2247 ];
2248 responder.send(caps)
2249 }
2250 avdtp::Request::Open { responder, stream_id } => {
2251 assert_eq!(expected_stream_id, stream_id);
2252 responder.send()
2253 }
2254 avdtp::Request::SetConfiguration {
2255 responder,
2256 local_stream_id,
2257 remote_stream_id,
2258 capabilities,
2259 } => {
2260 assert_eq!(local_stream_id, expected_stream_id);
2261 assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
2263 let mut capabilities_ordered = capabilities.clone();
2265 capabilities_ordered.sort_by_key(ServiceCapability::category);
2266 assert_eq!(capabilities, capabilities_ordered);
2267 responder.send()
2268 }
2269 x => panic!("Unexpected request: {:?}", x),
2270 };
2271 res.expect("should be able to respond");
2272 }
2273
2274 let collect_capabilities_fut = peer.collect_capabilities();
2276 let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2277
2278 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2279
2280 let request = exec.run_singlethreaded(&mut remote_events.next());
2281 remote_handle_request(request.expect("should have a discovery request").unwrap());
2282
2283 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2284 let request = exec.run_singlethreaded(&mut remote_events.next());
2285 remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2286
2287 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2288
2289 let remote_seid = 4_u8.try_into().unwrap();
2291
2292 let codec_params = ServiceCapability::MediaCodec {
2293 media_type: avdtp::MediaType::Audio,
2294 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2295 codec_extra: vec![0x11, 0x45, 51, 51],
2296 };
2297 let start_future = peer.stream_start(
2298 remote_seid,
2299 vec![
2300 ServiceCapability::MediaTransport,
2301 ServiceCapability::DelayReporting,
2302 codec_params,
2303 ],
2304 );
2305 let mut start_future = pin!(start_future);
2306
2307 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2308 let request = exec.run_singlethreaded(&mut remote_events.next());
2309 remote_handle_request(request.expect("should have a set_configuration request").unwrap());
2310
2311 assert!(exec.run_until_stalled(&mut start_future).is_pending());
2312 let request = exec.run_singlethreaded(&mut remote_events.next());
2313 remote_handle_request(request.expect("should have an open request").unwrap());
2314 }
2315
2316 #[test_case(Transport::Socket ; "socket")]
2319 #[test_case(Transport::Fidl ; "fidl")]
2320 fn peer_stream_start_permit_revoked(transport: Transport) {
2321 let mut exec = fasync::TestExecutor::new();
2322
2323 let test_permits = Permits::new(1);
2324 let (mut remote, mut profile_request_stream, _, peer) =
2325 setup_test_peer(transport, false, build_test_streams(), Some(test_permits.clone()));
2326
2327 let remote_seid = 2_u8.try_into().unwrap();
2328
2329 let codec_params = ServiceCapability::MediaCodec {
2330 media_type: avdtp::MediaType::Audio,
2331 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2332 codec_extra: vec![0x11, 0x45, 51, 51],
2333 };
2334
2335 let remote_endpoint = avdtp::StreamEndpoint::new(
2337 2,
2338 avdtp::MediaType::Audio,
2339 avdtp::EndpointType::Sink,
2340 vec![codec_params.clone()],
2341 )
2342 .expect("valid endpoint");
2343 peer.inner.lock().set_remote_endpoints(&[remote_endpoint]);
2344
2345 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2346 let mut start_future = pin!(start_future);
2347
2348 let _ = exec
2349 .run_until_stalled(&mut start_future)
2350 .expect_pending("waiting for set config response");
2351 receive_simple_accept(&mut exec, &mut remote, 0x03); exec.run_until_stalled(&mut start_future).expect_pending("waiting for open response");
2353 receive_simple_accept(&mut exec, &mut remote, 0x06); exec.run_until_stalled(&mut start_future).expect_pending("waiting for media transport");
2355 assert!(!peer.is_streaming_now());
2356
2357 let (transport_chan, _remote_chan) = create_test_channels(transport);
2359
2360 let request = exec.run_until_stalled(&mut profile_request_stream.next());
2361 match request {
2362 Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, connection, responder }))) => {
2363 assert_eq!(PeerId(1), peer_id.into());
2364 assert_eq!(connection, ConnectParameters::L2cap(Peer::transport_channel_params()));
2365 let channel = transport_chan.try_into().unwrap();
2366 responder.send(Ok(channel)).expect("responder sends");
2367 }
2368 x => panic!("Should have sent a open l2cap request, but got {:?}", x),
2369 };
2370
2371 exec.run_until_stalled(&mut start_future)
2374 .expect("start setup finished")
2375 .expect("stream setup is ok");
2376 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
2377 assert!(!peer.is_streaming_now());
2378
2379 let seized_permits = test_permits.seize();
2381 assert_eq!(seized_permits.len(), 1);
2382 receive_simple_accept(&mut exec, &mut remote, 0x07); let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
2387 assert!(!peer.is_streaming_now());
2388 receive_simple_accept(&mut exec, &mut remote, 0x09); let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
2393 assert!(!peer.is_streaming_now());
2394 }
2395
2396 #[test_case(Transport::Socket ; "socket")]
2397 #[test_case(Transport::Fidl ; "fidl")]
2398 fn peer_stream_start_fails_wrong_direction(transport: Transport) {
2399 let mut exec = fasync::TestExecutor::new();
2400
2401 let mut streams = Streams::default();
2403 let source = Stream::build(
2404 make_sbc_endpoint(1, avdtp::EndpointType::Source),
2405 TestMediaTaskBuilder::new().builder(),
2406 );
2407 streams.insert(source);
2408
2409 let (remote, _requests, _, peer) = setup_test_peer(transport, false, streams, None);
2410 let remote = avdtp::Peer::new(remote);
2411 let mut remote_events = remote.take_request_stream();
2412
2413 fn remote_handle_request(req: avdtp::Request) {
2415 let expected_stream_id: StreamEndpointId = 2_u8.try_into().unwrap();
2416 let res = match req {
2417 avdtp::Request::Discover { responder } => {
2418 let infos = [avdtp::StreamInformation::new(
2419 expected_stream_id,
2420 false,
2421 avdtp::MediaType::Audio,
2422 avdtp::EndpointType::Source,
2423 )];
2424 responder.send(&infos)
2425 }
2426 avdtp::Request::GetAllCapabilities { stream_id, responder }
2427 | avdtp::Request::GetCapabilities { stream_id, responder } => {
2428 assert_eq!(expected_stream_id, stream_id);
2429 let caps = vec![
2430 ServiceCapability::MediaTransport,
2431 ServiceCapability::MediaCodec {
2432 media_type: avdtp::MediaType::Audio,
2433 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2434 codec_extra: vec![0x11, 0x45, 51, 250],
2435 },
2436 ];
2437 responder.send(&caps[..])
2438 }
2439 avdtp::Request::Open { responder, .. } => responder.send(),
2440 avdtp::Request::SetConfiguration { responder, .. } => responder.send(),
2441 x => panic!("Unexpected request: {:?}", x),
2442 };
2443 res.expect("should be able to respond");
2444 }
2445
2446 let collect_capabilities_fut = peer.collect_capabilities();
2448 let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2449
2450 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2451
2452 let request = exec.run_singlethreaded(&mut remote_events.next());
2453 remote_handle_request(request.expect("should have a discovery request").unwrap());
2454
2455 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2456 let request = exec.run_singlethreaded(&mut remote_events.next());
2457 remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2458
2459 assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2460
2461 let remote_seid = 2_u8.try_into().unwrap();
2463
2464 let codec_params = ServiceCapability::MediaCodec {
2465 media_type: avdtp::MediaType::Audio,
2466 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2467 codec_extra: vec![0x11, 0x45, 51, 51],
2468 };
2469 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2470 let mut start_future = pin!(start_future);
2471
2472 match exec.run_until_stalled(&mut start_future) {
2473 Poll::Ready(Err(avdtp::Error::OutOfRange)) => {}
2474 x => panic!("Expected a ready OutOfRange error but got {:?}", x),
2475 };
2476 }
2477
2478 #[test_case(Transport::Socket ; "socket")]
2479 #[test_case(Transport::Fidl ; "fidl")]
2480 fn peer_stream_start_fails_to_connect(transport: Transport) {
2481 let mut exec = fasync::TestExecutor::new();
2482
2483 let (mut remote, mut profile_request_stream, _, peer) =
2484 setup_test_peer(transport, false, build_test_streams(), None);
2485
2486 let remote_seid = 2_u8.try_into().unwrap();
2487
2488 let codec_params = ServiceCapability::MediaCodec {
2489 media_type: avdtp::MediaType::Audio,
2490 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2491 codec_extra: vec![0x11, 0x45, 51, 51],
2492 };
2493
2494 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2495 let mut start_future = pin!(start_future);
2496
2497 match exec.run_until_stalled(&mut start_future) {
2498 Poll::Pending => {}
2499 x => panic!("was expecting pending but got {x:?}"),
2500 };
2501
2502 receive_simple_accept(&mut exec, &mut remote, 0x03); assert!(exec.run_until_stalled(&mut start_future).is_pending());
2505
2506 receive_simple_accept(&mut exec, &mut remote, 0x06); match exec.run_until_stalled(&mut start_future) {
2509 Poll::Pending => {}
2510 Poll::Ready(x) => panic!("Expected to be pending but {x:?}"),
2511 };
2512
2513 let request = exec.run_until_stalled(&mut profile_request_stream.next());
2515 match request {
2516 Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, responder, .. }))) => {
2517 assert_eq!(PeerId(1), peer_id.into());
2518 responder.send(Err(ErrorCode::Failed)).expect("responder sends");
2519 }
2520 x => panic!("Should have sent a open l2cap request, but got {:?}", x),
2521 };
2522
2523 match exec.run_until_stalled(&mut start_future) {
2526 Poll::Pending => panic!("Should be ready after start fails"),
2527 Poll::Ready(Ok(_stream)) => panic!("Shouldn't have succeeded stream here"),
2528 Poll::Ready(Err(_)) => {}
2529 }
2530 }
2531
2532 #[test_case(Transport::Socket ; "socket")]
2534 #[test_case(Transport::Fidl ; "fidl")]
2535 #[fuchsia::test]
2536 async fn peer_delay_report(transport: Transport) {
2537 let (remote, _profile_requests, cobalt_recv, peer) =
2538 setup_test_peer(transport, true, build_test_streams(), None);
2539 let remote_peer = avdtp::Peer::new(remote);
2540 let mut remote_events = remote_peer.take_request_stream();
2541
2542 async fn remote_handle_request(req: avdtp::Request, peer: &avdtp::Peer) {
2544 let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2545 let expected_peer_stream_id: StreamEndpointId = 1_u8.try_into().unwrap();
2547 use avdtp::Request::*;
2548 match req {
2549 Discover { responder } => {
2550 let infos = [avdtp::StreamInformation::new(
2551 expected_stream_id,
2552 false,
2553 avdtp::MediaType::Audio,
2554 avdtp::EndpointType::Sink,
2555 )];
2556 responder.send(&infos).expect("response should succeed");
2557 }
2558 GetAllCapabilities { stream_id, responder }
2559 | GetCapabilities { stream_id, responder } => {
2560 assert_eq!(expected_stream_id, stream_id);
2561 let caps = vec![
2562 ServiceCapability::MediaTransport,
2563 ServiceCapability::MediaCodec {
2564 media_type: avdtp::MediaType::Audio,
2565 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2566 codec_extra: vec![0x11, 0x45, 51, 250],
2567 },
2568 ];
2569 responder.send(&caps[..]).expect("response should succeed");
2570 assert!(peer.delay_report(&expected_peer_stream_id, 0xc0de).await.is_err());
2573 }
2574 Open { responder, stream_id } => {
2575 assert!(peer.delay_report(&expected_stream_id, 0xc0de).await.is_err());
2577 peer.delay_report(&expected_peer_stream_id, 0xc0de)
2579 .await
2580 .expect("should get acked correctly");
2581 assert_eq!(expected_stream_id, stream_id);
2582 responder.send().expect("response should succeed");
2583 }
2584 SetConfiguration { responder, local_stream_id, remote_stream_id, .. } => {
2585 assert_eq!(local_stream_id, expected_stream_id);
2586 assert_eq!(remote_stream_id, expected_peer_stream_id);
2587 responder.send().expect("should send back response without issue");
2588 }
2589 x => panic!("Unexpected request: {:?}", x),
2590 };
2591 }
2592
2593 let collect_fut = pin!(peer.collect_capabilities());
2594
2595 let Either::Left((request, collect_fut)) =
2597 futures::future::select(remote_events.next(), collect_fut).await
2598 else {
2599 panic!("Collect future shouldn't finish first");
2600 };
2601 let collect_fut = pin!(collect_fut);
2602 remote_handle_request(request.expect("a request").unwrap(), &remote_peer).await;
2603 let Either::Left((request, collect_fut)) =
2604 futures::future::select(remote_events.next(), collect_fut).await
2605 else {
2606 panic!("Collect future shouldn't finish first");
2607 };
2608 remote_handle_request(request.expect("a request").unwrap(), &remote_peer).await;
2609
2610 assert_eq!(1, collect_fut.await.expect("should get the remote endpoints back").len());
2612
2613 let remote_seid = 4_u8.try_into().unwrap();
2615
2616 let codec_params = ServiceCapability::MediaCodec {
2617 media_type: avdtp::MediaType::Audio,
2618 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2619 codec_extra: vec![0x11, 0x45, 51, 51],
2620 };
2621
2622 let _start_task = fasync::Task::spawn(async move {
2625 let _ = peer.stream_start(remote_seid, vec![codec_params]).await;
2626 panic!("stream start task finished");
2627 });
2628
2629 let request = remote_events.next().await.expect("should have set_config").unwrap();
2630 remote_handle_request(request, &remote_peer).await;
2631
2632 let request = remote_events.next().await.expect("should have open").unwrap();
2633 remote_handle_request(request, &remote_peer).await;
2634
2635 let mut cobalt = cobalt_recv.expect("should have receiver");
2636
2637 let mut got_ids = HashMap::new();
2638 let delay_metric_id = bt_metrics::AVDTP_DELAY_REPORT_IN_NANOSECONDS_METRIC_ID;
2639 while got_ids.len() < 3 || *got_ids.get(&delay_metric_id).unwrap_or(&0) < 3 {
2640 let report = respond_to_metrics_req_for_test(cobalt.next().await.unwrap().unwrap());
2641 let _ = got_ids.entry(report.metric_id).and_modify(|x| *x += 1).or_insert(1);
2642 if report.metric_id == delay_metric_id {
2644 assert_eq!(MetricEventPayload::IntegerValue(0xc0de * 100000), report.payload);
2645 }
2646 }
2647 assert!(got_ids.contains_key(&bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID));
2648 assert!(got_ids.contains_key(&bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID));
2649 assert!(got_ids.contains_key(&delay_metric_id));
2650 assert_eq!(got_ids.get(&delay_metric_id).cloned(), Some(3));
2653 }
2654
2655 fn sbc_capabilities() -> Vec<ServiceCapability> {
2656 let sbc_codec_info = SbcCodecInfo::new(
2657 SbcSamplingFrequency::FREQ48000HZ,
2658 SbcChannelMode::JOINT_STEREO,
2659 SbcBlockCount::SIXTEEN,
2660 SbcSubBands::EIGHT,
2661 SbcAllocation::LOUDNESS,
2662 53,
2663 53,
2664 )
2665 .expect("sbc codec info");
2666
2667 vec![avdtp::ServiceCapability::MediaTransport, sbc_codec_info.into()]
2668 }
2669
2670 #[test_case(Transport::Socket ; "socket")]
2672 #[test_case(Transport::Fidl ; "fidl")]
2673 fn peer_as_acceptor(transport: Transport) {
2674 let mut exec = fasync::TestExecutor::new();
2675
2676 let mut streams = Streams::default();
2677 let mut test_builder = TestMediaTaskBuilder::new_inactive();
2678 streams.insert(Stream::build(
2679 make_sbc_endpoint(1, avdtp::EndpointType::Source),
2680 test_builder.builder(),
2681 ));
2682
2683 let (remote, _requests, _, peer) = setup_test_peer(transport, false, streams, None);
2684 let remote_peer = avdtp::Peer::new(remote);
2685
2686 let discover_fut = remote_peer.discover();
2687 let mut discover_fut = pin!(discover_fut);
2688
2689 let expected = vec![make_sbc_endpoint(1, avdtp::EndpointType::Source).information()];
2690 match exec.run_until_stalled(&mut discover_fut) {
2691 Poll::Ready(Ok(res)) => assert_eq!(res, expected),
2692 x => panic!("Expected discovery to complete and got {:?}", x),
2693 };
2694
2695 let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
2696 let unknown_endpoint_id = 2_u8.try_into().expect("should be able to get sbc endpointid");
2697
2698 let get_caps_fut = remote_peer.get_capabilities(&sbc_endpoint_id);
2699 let mut get_caps_fut = pin!(get_caps_fut);
2700
2701 match exec.run_until_stalled(&mut get_caps_fut) {
2702 Poll::Ready(Ok(caps)) => assert_eq!(2, caps.len()),
2704 x => panic!("Get capabilities should be ready but got {:?}", x),
2705 };
2706
2707 let get_caps_fut = remote_peer.get_capabilities(&unknown_endpoint_id);
2708 let mut get_caps_fut = pin!(get_caps_fut);
2709
2710 match exec.run_until_stalled(&mut get_caps_fut) {
2711 Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
2712 assert_eq!(Some(Ok(avdtp::ErrorCode::BadAcpSeid)), e.error_code())
2713 }
2714 x => panic!("Get capabilities should be a ready error but got {:?}", x),
2715 };
2716
2717 let get_caps_fut = remote_peer.get_all_capabilities(&sbc_endpoint_id);
2718 let mut get_caps_fut = pin!(get_caps_fut);
2719
2720 match exec.run_until_stalled(&mut get_caps_fut) {
2721 Poll::Ready(Ok(caps)) => assert_eq!(2, caps.len()),
2723 x => panic!("Get capabilities should be ready but got {:?}", x),
2724 };
2725
2726 let sbc_caps = sbc_capabilities();
2727 let set_config_fut =
2728 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2729 let mut set_config_fut = pin!(set_config_fut);
2730
2731 match exec.run_until_stalled(&mut set_config_fut) {
2732 Poll::Ready(Ok(())) => {}
2733 x => panic!("Set capabilities should be ready but got {:?}", x),
2734 };
2735
2736 let open_fut = remote_peer.open(&sbc_endpoint_id);
2737 let mut open_fut = pin!(open_fut);
2738 match exec.run_until_stalled(&mut open_fut) {
2739 Poll::Ready(Ok(())) => {}
2740 x => panic!("Open should be ready but got {:?}", x),
2741 };
2742
2743 let (transport_chan, _remote_transport) = create_test_channels(transport);
2745
2746 assert_eq!(Some(()), peer.receive_channel(transport_chan).ok());
2747
2748 let stream_ids = vec![sbc_endpoint_id.clone()];
2749 let start_fut = remote_peer.start(&stream_ids);
2750 let mut start_fut = pin!(start_fut);
2751 match exec.run_until_stalled(&mut start_fut) {
2752 Poll::Ready(Ok(())) => {}
2753 x => panic!("Start should be ready but got {:?}", x),
2754 };
2755
2756 let media_task = test_builder.expect_task();
2758 assert!(media_task.is_started());
2759
2760 let suspend_fut = remote_peer.suspend(&stream_ids);
2761 let mut suspend_fut = pin!(suspend_fut);
2762 match exec.run_until_stalled(&mut suspend_fut) {
2763 Poll::Ready(Ok(())) => {}
2764 x => panic!("Start should be ready but got {:?}", x),
2765 };
2766
2767 assert!(!media_task.is_started());
2769 }
2770
2771 #[test_case(Transport::Socket ; "socket")]
2772 #[test_case(Transport::Fidl ; "fidl")]
2773 fn peer_source_stream_suspends_and_resumes_on_audio_disabled(transport: Transport) {
2774 let mut exec = fasync::TestExecutor::new();
2775
2776 let mut streams = Streams::default();
2777 let mut test_builder = TestMediaTaskBuilder::new_inactive();
2778 let _ = test_builder.with_direction(avdtp::EndpointType::Source);
2779 streams.insert(Stream::build(
2780 make_sbc_endpoint(1, avdtp::EndpointType::Source),
2781 test_builder.builder(),
2782 ));
2783
2784 let (remote, _requests, _, peer) = setup_test_peer(transport, false, streams, None);
2785 let remote_peer = avdtp::Peer::new(remote);
2786
2787 let sbc_endpoint_id = 1_u8.try_into().expect("sbc endpoint id");
2788 let sbc_caps = sbc_capabilities();
2789
2790 let mut set_config_fut =
2792 pin!(remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps));
2793 assert!(exec.run_until_stalled(&mut set_config_fut).is_ready());
2794
2795 let mut open_fut = pin!(remote_peer.open(&sbc_endpoint_id));
2796 assert!(exec.run_until_stalled(&mut open_fut).is_ready());
2797
2798 let (transport_chan, _remote_transport) = create_test_channels(transport);
2800 assert!(peer.receive_channel(transport_chan).is_ok());
2801
2802 let stream_ids = vec![sbc_endpoint_id.clone()];
2804 let mut start_fut = pin!(remote_peer.start(&stream_ids));
2805 assert!(exec.run_until_stalled(&mut start_fut).is_ready());
2806
2807 let media_task = test_builder.expect_task();
2809 assert!(media_task.is_started());
2810
2811 media_task.end_prematurely(Some(Ok(MediaTaskStatus::AudioDisabled)));
2813
2814 let mut remote_events = remote_peer.take_request_stream();
2816 let mut req_fut = remote_events.next();
2817 let Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) =
2818 exec.run_until_stalled(&mut req_fut)
2819 else {
2820 panic!("Expected Suspend request from peer");
2821 };
2822 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
2823 responder.send().expect("suspend response should send");
2824
2825 assert!(!media_task.is_started());
2827
2828 test_builder.set_active(true);
2830
2831 let mut req_fut = remote_events.next();
2833 let Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) =
2834 exec.run_until_stalled(&mut req_fut)
2835 else {
2836 panic!("Expected Start request from peer");
2837 };
2838 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
2839 responder.send().expect("start response should send");
2840
2841 let new_media_task =
2843 exec.run_until_stalled(&mut test_builder.next_task()).expect("ready").unwrap();
2844 assert!(new_media_task.is_started());
2845 }
2846
2847 #[test_case(Transport::Socket ; "socket")]
2848 #[test_case(Transport::Fidl ; "fidl")]
2849 fn peer_set_config_reject_first(transport: Transport) {
2850 let mut exec = fasync::TestExecutor::new();
2851
2852 let mut streams = Streams::default();
2853 let test_builder = TestMediaTaskBuilder::new();
2854 streams.insert(Stream::build(
2855 make_sbc_endpoint(1, avdtp::EndpointType::Source),
2856 test_builder.builder(),
2857 ));
2858
2859 let (remote, _requests, _, _peer) = setup_test_peer(transport, false, streams, None);
2860 let remote_peer = avdtp::Peer::new(remote);
2861
2862 let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
2863
2864 let wrong_freq_sbc = &[SbcCodecInfo::new(
2865 SbcSamplingFrequency::FREQ44100HZ, SbcChannelMode::JOINT_STEREO,
2867 SbcBlockCount::SIXTEEN,
2868 SbcSubBands::EIGHT,
2869 SbcAllocation::LOUDNESS,
2870 53,
2871 53,
2872 )
2873 .expect("sbc codec info")
2874 .into()];
2875
2876 let set_config_fut =
2877 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, wrong_freq_sbc);
2878 let mut set_config_fut = pin!(set_config_fut);
2879
2880 match exec.run_until_stalled(&mut set_config_fut) {
2881 Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
2882 assert!(e.service_category().is_some())
2883 }
2884 x => panic!("Set capabilities should have been rejected but got {:?}", x),
2885 };
2886
2887 let sbc_caps = sbc_capabilities();
2888 let set_config_fut =
2889 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2890 let mut set_config_fut = pin!(set_config_fut);
2891
2892 match exec.run_until_stalled(&mut set_config_fut) {
2893 Poll::Ready(Ok(())) => {}
2894 x => panic!("Set capabilities should be ready but got {:?}", x),
2895 };
2896 }
2897
2898 #[test_case(Transport::Socket ; "socket")]
2901 #[test_case(Transport::Fidl ; "fidl")]
2902 fn peer_retries_failed_start_while_audio_active(transport_mode: Transport) {
2903 let mut exec = fasync::TestExecutor::new_with_fake_time();
2904 exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
2905
2906 let mut streams = Streams::default();
2907 let mut test_builder = TestMediaTaskBuilder::new_inactive();
2908 let _ = test_builder.with_direction(avdtp::EndpointType::Source);
2909 streams.insert(Stream::build(
2910 make_sbc_endpoint(1, avdtp::EndpointType::Source),
2911 test_builder.builder(),
2912 ));
2913
2914 let (remote, _requests, _, peer) = setup_test_peer(transport_mode, false, streams, None);
2915 let remote_peer = avdtp::Peer::new(remote);
2916
2917 let sbc_endpoint_id: StreamEndpointId = 1_u8.try_into().expect("sbc endpoint id");
2918 let sbc_caps = sbc_capabilities();
2919
2920 let mut set_config_fut =
2921 pin!(remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps));
2922 match exec.run_until_stalled(&mut set_config_fut) {
2923 Poll::Ready(Ok(())) => {}
2924 x => panic!("Set capabilities should be ready but got {:?}", x),
2925 };
2926
2927 let mut open_fut = pin!(remote_peer.open(&sbc_endpoint_id));
2928 match exec.run_until_stalled(&mut open_fut) {
2929 Poll::Ready(Ok(())) => {}
2930 x => panic!("Open should be ready but got {:?}", x),
2931 };
2932
2933 let (transport, _remote_transport) = create_test_channels(transport_mode);
2935 assert_eq!(Some(()), peer.receive_channel(transport).ok());
2936
2937 let mut remote_requests = remote_peer.take_request_stream();
2938
2939 test_builder.set_active(true);
2941
2942 let mut next_remote_request_fut = pin!(remote_requests.next());
2944 match exec.run_until_stalled(&mut next_remote_request_fut) {
2945 Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
2946 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
2947 responder
2948 .reject(&sbc_endpoint_id, avdtp::ErrorCode::BadState)
2949 .expect("reject response should send");
2950 }
2951 x => panic!("Expected to receive a start request for the stream, got {:?}", x),
2952 };
2953
2954 assert!(exec.run_until_stalled(&mut test_builder.next_task()).is_pending());
2956 assert!(!peer.is_streaming_now());
2957
2958 let mut next_remote_request_fut = pin!(remote_requests.next());
2960 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
2961
2962 exec.set_fake_time(fasync::MonotonicInstant::after(
2963 START_RETRY_DELAY + zx::MonotonicDuration::from_micros(1),
2964 ));
2965 assert!(exec.wake_expired_timers());
2966
2967 match exec.run_until_stalled(&mut next_remote_request_fut) {
2968 Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
2969 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
2970 responder.send().expect("start response should send");
2971 }
2972 x => panic!("Expected to receive a second start request, got {:?}", x),
2973 };
2974
2975 let media_task =
2977 exec.run_until_stalled(&mut test_builder.next_task()).expect("ready").unwrap();
2978 assert!(media_task.is_started());
2979 }
2980
2981 fn setup_open_sink_stream(
2986 exec: &mut fasync::TestExecutor,
2987 test_builder: &TestMediaTaskBuilder,
2988 seid: u8,
2989 transport_mode: Transport,
2990 ) -> (StreamEndpointId, Peer, avdtp::Peer, Channel) {
2991 let mut streams = Streams::default();
2992 streams.insert(Stream::build(
2993 make_sbc_endpoint(seid, avdtp::EndpointType::Sink),
2994 test_builder.builder(),
2995 ));
2996
2997 let (remote, _requests, _, peer) = setup_test_peer(transport_mode, false, streams, None);
2998 let remote_peer = avdtp::Peer::new(remote);
2999
3000 let local_id: StreamEndpointId = seid.try_into().expect("sbc endpoint id");
3001 let sbc_caps = sbc_capabilities();
3002 let mut set_config_fut =
3003 pin!(remote_peer.set_configuration(&local_id, &local_id, &sbc_caps));
3004 match exec.run_until_stalled(&mut set_config_fut) {
3005 Poll::Ready(Ok(())) => {}
3006 x => panic!("Set configuration should be ready but got {:?}", x),
3007 };
3008
3009 let mut open_fut = pin!(remote_peer.open(&local_id));
3010 match exec.run_until_stalled(&mut open_fut) {
3011 Poll::Ready(Ok(())) => {}
3012 x => panic!("Open should be ready but got {:?}", x),
3013 };
3014
3015 let (transport, remote_transport) = create_test_channels(transport_mode);
3017 assert_eq!(Some(()), peer.receive_channel(transport).ok());
3018
3019 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3021
3022 (local_id, peer, remote_peer, remote_transport)
3023 }
3024
3025 #[test_case(Transport::Socket ; "socket")]
3028 #[test_case(Transport::Fidl ; "fidl")]
3029 fn peer_starts_sink_stream_after_dwell(transport_mode: Transport) {
3030 let mut exec = fasync::TestExecutor::new_with_fake_time();
3031 exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
3032
3033 let mut test_builder = TestMediaTaskBuilder::new();
3034 let (sbc_endpoint_id, _peer, remote_peer, _remote_transport) =
3035 setup_open_sink_stream(&mut exec, &test_builder, 1, transport_mode);
3036 let mut remote_requests = remote_peer.take_request_stream();
3037
3038 let mut next_remote_request_fut = pin!(remote_requests.next());
3040 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
3041
3042 exec.set_fake_time(fasync::MonotonicInstant::after(
3044 STREAM_DWELL + zx::MonotonicDuration::from_micros(1),
3045 ));
3046 assert!(exec.wake_expired_timers());
3047
3048 match exec.run_until_stalled(&mut next_remote_request_fut) {
3049 Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
3050 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3051 responder.send().expect("start response should send");
3052 }
3053 x => panic!("Expected to receive a start request for the stream, got {:?}", x),
3054 };
3055
3056 let media_task =
3057 exec.run_until_stalled(&mut test_builder.next_task()).expect("ready").unwrap();
3058 assert!(media_task.is_started());
3059 }
3060
3061 #[test_case(Transport::Socket ; "socket")]
3064 #[test_case(Transport::Fidl ; "fidl")]
3065 fn peer_does_not_start_sink_stream_started_by_peer(transport_mode: Transport) {
3066 let mut exec = fasync::TestExecutor::new_with_fake_time();
3067 exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
3068
3069 let mut test_builder = TestMediaTaskBuilder::new();
3070 let (sbc_endpoint_id, _peer, remote_peer, _remote_transport) =
3071 setup_open_sink_stream(&mut exec, &test_builder, 1, transport_mode);
3072 let mut remote_requests = remote_peer.take_request_stream();
3073
3074 let mut start_fut = pin!(remote_peer.start(&[sbc_endpoint_id.clone()]));
3076 match exec.run_until_stalled(&mut start_fut) {
3077 Poll::Ready(Ok(())) => {}
3078 x => panic!("Start should be ready but got {:?}", x),
3079 };
3080
3081 let media_task = test_builder.expect_task();
3082 assert!(media_task.is_started());
3083
3084 exec.set_fake_time(fasync::MonotonicInstant::after(
3086 STREAM_DWELL + zx::MonotonicDuration::from_micros(1),
3087 ));
3088 let _ = exec.wake_expired_timers();
3089 let mut next_remote_request_fut = pin!(remote_requests.next());
3090 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
3091 assert!(media_task.is_started());
3092
3093 let mut suspend_fut = pin!(remote_peer.suspend(&[sbc_endpoint_id.clone()]));
3095 match exec.run_until_stalled(&mut suspend_fut) {
3096 Poll::Ready(Ok(())) => {}
3097 x => panic!("Suspend should be ready but got {:?}", x),
3098 };
3099 assert!(!media_task.is_started());
3100
3101 exec.set_fake_time(fasync::MonotonicInstant::after(STREAM_DWELL + START_RETRY_DELAY));
3103 let _ = exec.wake_expired_timers();
3104 let mut next_remote_request_fut = pin!(remote_requests.next());
3105 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
3106 }
3107
3108 #[test_case(Transport::Socket ; "socket")]
3111 #[test_case(Transport::Fidl ; "fidl")]
3112 fn peer_stops_starting_sink_stream_after_attempts(transport_mode: Transport) {
3113 let mut exec = fasync::TestExecutor::new_with_fake_time();
3114 exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
3115
3116 let test_builder = TestMediaTaskBuilder::new();
3117 let (sbc_endpoint_id, peer, remote_peer, _remote_transport) =
3118 setup_open_sink_stream(&mut exec, &test_builder, 1, transport_mode);
3119 let mut remote_requests = remote_peer.take_request_stream();
3120
3121 exec.set_fake_time(fasync::MonotonicInstant::after(
3123 STREAM_DWELL + zx::MonotonicDuration::from_micros(1),
3124 ));
3125 assert!(exec.wake_expired_timers());
3126
3127 for attempt in 1..=START_ATTEMPTS {
3128 let mut next_remote_request_fut = pin!(remote_requests.next());
3129 match exec.run_until_stalled(&mut next_remote_request_fut) {
3130 Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
3131 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3132 responder
3133 .reject(&sbc_endpoint_id, avdtp::ErrorCode::BadState)
3134 .expect("reject response should send");
3135 }
3136 x => panic!("Expected start request number {attempt}, got {:?}", x),
3137 };
3138 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3140 exec.set_fake_time(fasync::MonotonicInstant::after(
3141 START_RETRY_DELAY + zx::MonotonicDuration::from_micros(1),
3142 ));
3143 let retried = exec.wake_expired_timers();
3144 assert_eq!(retried, attempt < START_ATTEMPTS, "attempt {attempt} retry");
3145 }
3146
3147 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3150 exec.set_fake_time(fasync::MonotonicInstant::after(STREAM_DWELL + START_RETRY_DELAY));
3151 assert!(!exec.wake_expired_timers());
3152 let mut next_remote_request_fut = pin!(remote_requests.next());
3153 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
3154 assert!(!peer.streaming_active());
3155 }
3156
3157 #[test_case(Transport::Socket ; "socket")]
3158 #[test_case(Transport::Fidl ; "fidl")]
3159 fn peer_starts_waiting_streams(transport_mode: Transport) {
3160 let mut exec = fasync::TestExecutor::new_with_fake_time();
3161 exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
3162
3163 let mut streams = Streams::default();
3164 let mut test_builder = TestMediaTaskBuilder::new_inactive();
3165 streams.insert(Stream::build(
3166 make_sbc_endpoint(1, avdtp::EndpointType::Source),
3167 test_builder.builder(),
3168 ));
3169
3170 let (remote, _requests, _, peer) = setup_test_peer(transport_mode, false, streams, None);
3171 let remote_peer = avdtp::Peer::new(remote);
3172
3173 let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
3174
3175 let sbc_caps = sbc_capabilities();
3176 let set_config_fut =
3177 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
3178 let mut set_config_fut = pin!(set_config_fut);
3179
3180 match exec.run_until_stalled(&mut set_config_fut) {
3181 Poll::Ready(Ok(())) => {}
3182 x => panic!("Set capabilities should be ready but got {:?}", x),
3183 };
3184
3185 let open_fut = remote_peer.open(&sbc_endpoint_id);
3186 let mut open_fut = pin!(open_fut);
3187 match exec.run_until_stalled(&mut open_fut) {
3188 Poll::Ready(Ok(())) => {}
3189 x => panic!("Open should be ready but got {:?}", x),
3190 };
3191
3192 let (transport, _remote_transport) = create_test_channels(transport_mode);
3194 assert_eq!(Some(()), peer.receive_channel(transport).ok());
3195
3196 let mut remote_requests = remote_peer.take_request_stream();
3198 let next_remote_request_fut = remote_requests.next();
3199 let mut next_remote_request_fut = pin!(next_remote_request_fut);
3200
3201 assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
3203
3204 test_builder.set_active(true);
3206
3207 let stream_ids = match exec.run_until_stalled(&mut next_remote_request_fut) {
3208 Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
3209 responder.send().unwrap();
3210 stream_ids
3211 }
3212 x => panic!("Expected to receive a start request for the stream, got {:?}", x),
3213 };
3214
3215 let media_task =
3217 exec.run_until_stalled(&mut test_builder.next_task()).expect("ready").unwrap();
3218 assert!(media_task.is_started());
3219
3220 let suspend_fut = remote_peer.suspend(&stream_ids);
3222 let mut suspend_fut = pin!(suspend_fut);
3223 match exec.run_until_stalled(&mut suspend_fut) {
3224 Poll::Ready(Ok(())) => {}
3225 x => panic!("Suspend should be ready but got {:?}", x),
3226 };
3227
3228 assert!(!media_task.is_started());
3230 }
3231
3232 #[test_case(Transport::Socket ; "socket")]
3233 #[test_case(Transport::Fidl ; "fidl")]
3234 fn needs_permit_to_start_streams(transport_mode: Transport) {
3235 let mut exec = fasync::TestExecutor::new();
3236
3237 let mut streams = Streams::default();
3238 let mut test_builder = TestMediaTaskBuilder::new();
3239 streams.insert(Stream::build(
3240 make_sbc_endpoint(1, avdtp::EndpointType::Sink),
3241 test_builder.builder(),
3242 ));
3243 streams.insert(Stream::build(
3244 make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3245 test_builder.builder(),
3246 ));
3247 let mut next_task_fut = test_builder.next_task();
3248
3249 let permits = Permits::new(1);
3250 let taken_permit = permits.get().expect("permit taken");
3251 let (remote, _profile_request_stream, _, peer) =
3252 setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3253 let remote_peer = avdtp::Peer::new(remote);
3254
3255 let sbc_endpoint_id = 1_u8.try_into().unwrap();
3256
3257 let sbc_caps = sbc_capabilities();
3258 let mut set_config_fut =
3259 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
3260
3261 match exec.run_until_stalled(&mut set_config_fut) {
3262 Poll::Ready(Ok(())) => {}
3263 x => panic!("Set capabilities should be ready but got {:?}", x),
3264 };
3265
3266 let mut open_fut = remote_peer.open(&sbc_endpoint_id);
3267 match exec.run_until_stalled(&mut open_fut) {
3268 Poll::Ready(Ok(())) => {}
3269 x => panic!("Open should be ready but got {:?}", x),
3270 };
3271
3272 let (transport, _remote_transport) = create_test_channels(transport_mode);
3274 assert_eq!(Some(()), peer.receive_channel(transport).ok());
3275
3276 let sbc_endpoint_two = 2_u8.try_into().unwrap();
3278
3279 let mut set_config_fut =
3280 remote_peer.set_configuration(&sbc_endpoint_two, &sbc_endpoint_two, &sbc_caps);
3281
3282 match exec.run_until_stalled(&mut set_config_fut) {
3283 Poll::Ready(Ok(())) => {}
3284 x => panic!("Set capabilities should be ready but got {:?}", x),
3285 };
3286
3287 let mut open_fut = remote_peer.open(&sbc_endpoint_two);
3288 match exec.run_until_stalled(&mut open_fut) {
3289 Poll::Ready(Ok(())) => {}
3290 x => panic!("Open should be ready but got {:?}", x),
3291 };
3292
3293 let (transport_two, _remote_transport_two) = create_test_channels(transport_mode);
3295 assert_eq!(Some(()), peer.receive_channel(transport_two).ok());
3296
3297 let unknown_endpoint_id: StreamEndpointId = 9_u8.try_into().unwrap();
3300 let stream_ids = [sbc_endpoint_id.clone(), unknown_endpoint_id.clone()];
3301 let mut start_fut = remote_peer.start(&stream_ids);
3302 match exec.run_until_stalled(&mut start_fut) {
3303 Poll::Ready(Err(avdtp::Error::RemoteRejected(rejection))) => {
3304 assert_eq!(avdtp::ErrorCode::BadAcpSeid, rejection.error_code().unwrap().unwrap());
3305 assert_eq!(unknown_endpoint_id, rejection.stream_id().unwrap());
3306 }
3307 x => panic!("Start should be ready but got {:?}", x),
3308 };
3309
3310 let mut remote_requests = remote_peer.take_request_stream();
3312
3313 let suspended_stream_ids = match exec.run_singlethreaded(&mut remote_requests.next()) {
3314 Some(Ok(avdtp::Request::Suspend { responder, stream_ids })) => {
3315 responder.send().unwrap();
3316 stream_ids
3317 }
3318 x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
3319 };
3320
3321 assert!(suspended_stream_ids.contains(&sbc_endpoint_id));
3322 assert_eq!(1, suspended_stream_ids.len());
3323
3324 match exec.run_until_stalled(&mut next_task_fut) {
3326 Poll::Pending => {}
3327 x => panic!("Local task should not have been created at this point: {:?}", x),
3328 };
3329
3330 let mut start_fut = remote_peer.start(&[sbc_endpoint_id.clone()]);
3333 match exec.run_until_stalled(&mut start_fut) {
3334 Poll::Ready(Ok(())) => {}
3335 x => panic!("Start should be ready but got {:?}", x),
3336 }
3337
3338 let suspended_stream_ids = match exec.run_until_stalled(&mut remote_requests.next()) {
3339 Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
3340 responder.send().unwrap();
3341 stream_ids
3342 }
3343 x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
3344 };
3345 assert!(suspended_stream_ids.contains(&sbc_endpoint_id));
3346
3347 drop(taken_permit);
3349
3350 match exec.run_singlethreaded(&mut remote_requests.next()) {
3351 Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
3352 assert_eq!(stream_ids, &[sbc_endpoint_id.clone()]);
3353 responder.send().unwrap();
3354 }
3355 x => panic!("Expected start on permit available but got {x:?}"),
3356 };
3357
3358 let media_task = match exec.run_until_stalled(&mut next_task_fut) {
3360 Poll::Ready(Some(task)) => task,
3361 x => panic!("Local task should be created at this point: {:?}", x),
3362 };
3363
3364 assert!(media_task.is_started());
3365
3366 let mut start_fut = remote_peer.start(&[sbc_endpoint_two.clone()]);
3368 match exec.run_until_stalled(&mut start_fut) {
3369 Poll::Ready(Ok(())) => {}
3370 x => panic!("Start should be ready but got {:?}", x),
3371 }
3372
3373 let suspended_stream_ids = match exec.run_until_stalled(&mut remote_requests.next()) {
3374 Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
3375 responder.send().unwrap();
3376 stream_ids
3377 }
3378 x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
3379 };
3380
3381 assert!(suspended_stream_ids.contains(&sbc_endpoint_two));
3382 assert_eq!(1, suspended_stream_ids.len());
3383
3384 let mut suspend_fut = remote_peer.suspend(&[sbc_endpoint_id.clone()]);
3386 match exec.run_until_stalled(&mut suspend_fut) {
3387 Poll::Ready(Ok(())) => {}
3388 x => panic!("Start should be ready but got {:?}", x),
3389 }
3390
3391 match exec.run_singlethreaded(&mut remote_requests.next()) {
3392 Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
3393 assert_eq!(stream_ids, &[sbc_endpoint_two]);
3394 responder.send().unwrap();
3395 }
3396 x => panic!("Expected start on permit available but got {x:?}"),
3397 };
3398 }
3399
3400 fn start_sbc_stream(
3401 exec: &mut fasync::TestExecutor,
3402 media_test_builder: &mut TestMediaTaskBuilder,
3403 peer: &Peer,
3404 remote_peer: &avdtp::Peer,
3405 local_id: &StreamEndpointId,
3406 remote_id: &StreamEndpointId,
3407 transport_mode: Transport,
3408 ) -> TestMediaTask {
3409 let sbc_caps = sbc_capabilities();
3410 let set_config_fut = remote_peer.set_configuration(&local_id, &remote_id, &sbc_caps);
3411 let mut set_config_fut = pin!(set_config_fut);
3412
3413 match exec.run_until_stalled(&mut set_config_fut) {
3414 Poll::Ready(Ok(())) => {}
3415 x => panic!("Set capabilities should be ready but got {:?}", x),
3416 };
3417
3418 let open_fut = remote_peer.open(&local_id);
3419 let mut open_fut = pin!(open_fut);
3420 match exec.run_until_stalled(&mut open_fut) {
3421 Poll::Ready(Ok(())) => {}
3422 x => panic!("Open should be ready but got {:?}", x),
3423 };
3424
3425 let (transport, _remote_transport) = create_test_channels(transport_mode);
3427 assert_eq!(Some(()), peer.receive_channel(transport).ok());
3428
3429 let stream_ids = [local_id.clone()];
3431 let start_fut = remote_peer.start(&stream_ids);
3432 let mut start_fut = pin!(start_fut);
3433 match exec.run_until_stalled(&mut start_fut) {
3434 Poll::Ready(Ok(())) => {}
3435 x => panic!("Start should be ready but got {:?}", x),
3436 };
3437
3438 let media_task = media_test_builder.expect_task();
3440 assert!(media_task.is_started());
3441
3442 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3443 media_task
3444 }
3445
3446 #[test_case(Transport::Socket ; "socket")]
3447 #[test_case(Transport::Fidl ; "fidl")]
3448 fn permits_can_be_revoked_and_reinstated_all(transport_mode: Transport) {
3449 let mut exec = fasync::TestExecutor::new();
3450
3451 let mut streams = Streams::default();
3452 let mut test_builder = TestMediaTaskBuilder::new();
3453 streams.insert(Stream::build(
3454 make_sbc_endpoint(1, avdtp::EndpointType::Sink),
3455 test_builder.builder(),
3456 ));
3457 let sbc_endpoint_id = 1_u8.try_into().unwrap();
3458 let remote_sbc_endpoint_id = 7_u8.try_into().unwrap();
3459
3460 streams.insert(Stream::build(
3461 make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3462 test_builder.builder(),
3463 ));
3464 let sbc2_endpoint_id = 2_u8.try_into().unwrap();
3465 let remote_sbc2_endpoint_id = 6_u8.try_into().unwrap();
3466
3467 let permits = Permits::new(2);
3468
3469 let (remote, _requests, _, peer) =
3470 setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3471 let remote_peer = avdtp::Peer::new(remote);
3472
3473 let one_media_task = start_sbc_stream(
3474 &mut exec,
3475 &mut test_builder,
3476 &peer,
3477 &remote_peer,
3478 &sbc_endpoint_id,
3479 &remote_sbc_endpoint_id,
3480 transport_mode,
3481 );
3482 let two_media_task = start_sbc_stream(
3483 &mut exec,
3484 &mut test_builder,
3485 &peer,
3486 &remote_peer,
3487 &sbc2_endpoint_id,
3488 &remote_sbc2_endpoint_id,
3489 transport_mode,
3490 );
3491
3492 let taken_permits = permits.seize();
3494
3495 let remote_endpoints: HashSet<_> =
3496 [&remote_sbc_endpoint_id, &remote_sbc2_endpoint_id].iter().cloned().collect();
3497
3498 let mut remote_requests = remote_peer.take_request_stream();
3500 let mut expected_suspends = remote_endpoints.clone();
3501 while !expected_suspends.is_empty() {
3502 match exec.run_until_stalled(&mut remote_requests.next()) {
3503 Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
3504 for stream_id in stream_ids {
3505 assert!(expected_suspends.remove(&stream_id));
3506 }
3507 responder.send().expect("send response okay");
3508 }
3509 x => panic!("Expected suspension and got {:?}", x),
3510 }
3511 }
3512
3513 assert!(!one_media_task.is_started());
3515 assert!(!two_media_task.is_started());
3516
3517 drop(taken_permits);
3519
3520 let mut expected_starts = remote_endpoints.clone();
3521 while !expected_starts.is_empty() {
3522 match exec.run_singlethreaded(&mut remote_requests.next()) {
3523 Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
3524 for stream_id in stream_ids {
3525 assert!(expected_starts.remove(&stream_id));
3526 }
3527 responder.send().expect("send response okay");
3528 }
3529 x => panic!("Expected start and got {:?}", x),
3530 }
3531 }
3532 let one_media_task = test_builder.expect_task();
3535 assert!(one_media_task.is_started());
3536 let two_media_task = match exec.run_until_stalled(&mut test_builder.next_task()) {
3537 Poll::Ready(Some(task)) => task,
3538 x => panic!("Expected another ready task but {x:?}"),
3539 };
3540 assert!(two_media_task.is_started());
3541 }
3542
3543 #[test_case(Transport::Socket ; "socket")]
3544 #[test_case(Transport::Fidl ; "fidl")]
3545 fn permits_can_be_revoked_one_at_a_time(transport_mode: Transport) {
3546 let mut exec = fasync::TestExecutor::new();
3547
3548 let mut streams = Streams::default();
3549 let mut test_builder = TestMediaTaskBuilder::new();
3550 streams.insert(Stream::build(
3551 make_sbc_endpoint(1, avdtp::EndpointType::Sink),
3552 test_builder.builder(),
3553 ));
3554 let sbc_endpoint_id = 1_u8.try_into().unwrap();
3555 let remote_sbc_endpoint_id = 7_u8.try_into().unwrap();
3556
3557 streams.insert(Stream::build(
3558 make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3559 test_builder.builder(),
3560 ));
3561 let sbc2_endpoint_id = 2_u8.try_into().unwrap();
3562 let remote_sbc2_endpoint_id = 6_u8.try_into().unwrap();
3563
3564 let permits = Permits::new(2);
3565
3566 let (remote, _requests, _, peer) =
3567 setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3568 let remote_peer = avdtp::Peer::new(remote);
3569
3570 let one_media_task = start_sbc_stream(
3571 &mut exec,
3572 &mut test_builder,
3573 &peer,
3574 &remote_peer,
3575 &sbc_endpoint_id,
3576 &remote_sbc_endpoint_id,
3577 transport_mode,
3578 );
3579 let two_media_task = start_sbc_stream(
3580 &mut exec,
3581 &mut test_builder,
3582 &peer,
3583 &remote_peer,
3584 &sbc2_endpoint_id,
3585 &remote_sbc2_endpoint_id,
3586 transport_mode,
3587 );
3588
3589 let taken_permit = permits.take();
3591
3592 let remote_endpoints: HashSet<_> =
3593 [&remote_sbc_endpoint_id, &remote_sbc2_endpoint_id].iter().cloned().collect();
3594
3595 let mut remote_requests = remote_peer.take_request_stream();
3597 let suspended_id = match exec.run_until_stalled(&mut remote_requests.next()) {
3598 Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
3599 assert!(stream_ids.len() == 1);
3600 assert!(remote_endpoints.contains(&stream_ids[0]));
3601 responder.send().expect("send response okay");
3602 stream_ids[0].clone()
3603 }
3604 x => panic!("Expected suspension and got {:?}", x),
3605 };
3606
3607 if suspended_id == remote_sbc_endpoint_id {
3609 assert!(!one_media_task.is_started());
3610 assert!(two_media_task.is_started());
3611 } else {
3612 assert!(one_media_task.is_started());
3613 assert!(!two_media_task.is_started());
3614 }
3615
3616 drop(taken_permit);
3618
3619 match exec.run_singlethreaded(&mut remote_requests.next()) {
3620 Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
3621 assert_eq!(stream_ids, &[suspended_id]);
3622 responder.send().expect("send response okay");
3623 }
3624 x => panic!("Expected start and got {:?}", x),
3625 }
3626 let media_task = match exec.run_until_stalled(&mut test_builder.next_task()) {
3628 Poll::Ready(Some(task)) => task,
3629 x => panic!("Expected media task to start: {x:?}"),
3630 };
3631 assert!(media_task.is_started());
3632 }
3633
3634 #[test_case(Transport::Socket ; "socket")]
3637 #[test_case(Transport::Fidl ; "fidl")]
3638 fn permit_suspend_start_while_suspending(transport_mode: Transport) {
3639 let mut exec = fasync::TestExecutor::new();
3640
3641 let mut streams = Streams::default();
3642 let mut test_builder = TestMediaTaskBuilder::new();
3643 streams.insert(Stream::build(
3644 make_sbc_endpoint(1, avdtp::EndpointType::Sink),
3645 test_builder.builder(),
3646 ));
3647 streams.insert(Stream::build(
3648 make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3649 test_builder.builder(),
3650 ));
3651 let mut next_task_fut = test_builder.next_task();
3652
3653 let permits = Permits::new(1);
3654 let (remote, _profile_request_stream, _, peer) =
3655 setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3656
3657 let remote_peer = avdtp::Peer::new(remote);
3658 let mut remote_requests = remote_peer.take_request_stream();
3659
3660 let sbc_endpoint_id = 1_u8.try_into().unwrap();
3661
3662 let sbc_caps = sbc_capabilities();
3663 let mut set_config_fut =
3664 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
3665
3666 match exec.run_until_stalled(&mut set_config_fut) {
3667 Poll::Ready(Ok(())) => {}
3668 x => panic!("Set capabilities should be ready but got {:?}", x),
3669 };
3670
3671 let mut open_fut = remote_peer.open(&sbc_endpoint_id);
3672 match exec.run_until_stalled(&mut open_fut) {
3673 Poll::Ready(Ok(())) => {}
3674 x => panic!("Open should be ready but got {:?}", x),
3675 };
3676
3677 let (_remote_transport, transport) = Channel::create_socket_pair();
3679 assert_eq!(Some(()), peer.receive_channel(transport).ok());
3680
3681 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3683 let Some(_deadline) = exec.wake_next_timer() else {
3684 panic!("Expected a timer to be waiting to run");
3685 };
3686
3687 let start_responder = match exec.run_singlethreaded(&mut remote_requests.next()) {
3689 Some(Ok(avdtp::Request::Start { stream_ids, responder })) => {
3690 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3691 responder
3692 }
3693 x => panic!("Expected a Start request, got {x:?}"),
3694 };
3695
3696 assert!(permits.get().is_none());
3697
3698 let mut start_fut = remote_peer.start(&[sbc_endpoint_id.clone()]);
3700
3701 match exec.run_singlethreaded(&mut start_fut) {
3704 Ok(()) => {}
3705 x => panic!("Expected OK response from start future but got {x:?}"),
3706 }
3707
3708 let suspend_responder = match exec.run_singlethreaded(&mut remote_requests.next()) {
3709 Some(Ok(avdtp::Request::Suspend { stream_ids, responder })) => {
3710 assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3711 responder
3712 }
3713 x => panic!("Expected a suspend got {x:?}"),
3714 };
3715
3716 start_responder.send().unwrap();
3718
3719 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3720
3721 suspend_responder.send().unwrap();
3723
3724 let media_task = match exec.run_until_stalled(&mut next_task_fut) {
3726 Poll::Ready(Some(task)) => task,
3727 x => panic!("Local task should be created at this point: {:?}", x),
3728 };
3729
3730 assert!(media_task.is_started());
3731 }
3732
3733 #[fuchsia::test]
3736 fn version_check() {
3737 let p1: ProfileDescriptor = ProfileDescriptor {
3738 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3739 major_version: Some(1),
3740 minor_version: Some(3),
3741 ..Default::default()
3742 };
3743 assert_eq!(true, a2dp_version_check(p1));
3744
3745 let p1: ProfileDescriptor = ProfileDescriptor {
3746 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3747 major_version: Some(2),
3748 minor_version: Some(10),
3749 ..Default::default()
3750 };
3751 assert_eq!(true, a2dp_version_check(p1));
3752
3753 let p1: ProfileDescriptor = ProfileDescriptor {
3754 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3755 major_version: Some(1),
3756 minor_version: Some(0),
3757 ..Default::default()
3758 };
3759 assert_eq!(false, a2dp_version_check(p1));
3760
3761 let p1: ProfileDescriptor = ProfileDescriptor {
3762 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3763 major_version: None,
3764 minor_version: Some(9),
3765 ..Default::default()
3766 };
3767 assert_eq!(false, a2dp_version_check(p1));
3768
3769 let p1: ProfileDescriptor = ProfileDescriptor {
3770 profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3771 major_version: Some(2),
3772 minor_version: Some(2),
3773 ..Default::default()
3774 };
3775 assert_eq!(true, a2dp_version_check(p1));
3776 }
3777
3778 fn setup_test_peer_with_avrcp(
3779 transport: Transport,
3780 streams: Streams,
3781 avrcp: Option<avrcp::PeerManagerProxy>,
3782 ) -> (Channel, ProfileRequestStream, Peer) {
3783 let (avdtp, remote) = setup_avdtp_peer(transport);
3784 let metrics_logger = bt_metrics::MetricsLogger::default();
3785 let (profile_proxy, requests) = create_proxy_and_stream::<ProfileMarker>();
3786 let peer =
3787 Peer::create(PeerId(1), avdtp, streams, None, profile_proxy, avrcp, metrics_logger);
3788
3789 (remote, requests, peer)
3790 }
3791
3792 #[test_case(Transport::Socket ; "socket")]
3793 #[test_case(Transport::Fidl ; "fidl")]
3794 #[fuchsia::test]
3795 fn test_volume_relay_starts_on_receive_channel_for_source(transport: Transport) {
3796 let mut exec = fasync::TestExecutor::new();
3797 let (avrcp_proxy, mut avrcp_stream) =
3798 fidl::endpoints::create_proxy_and_stream::<avrcp::PeerManagerMarker>();
3799
3800 let mut streams = Streams::default();
3801 let test_builder = TestMediaTaskBuilder::new_inactive();
3802 streams.insert(Stream::build(
3803 make_sbc_endpoint(1, avdtp::EndpointType::Source),
3804 test_builder.builder(),
3805 ));
3806
3807 let (remote, _requests, peer) =
3808 setup_test_peer_with_avrcp(transport, streams, Some(avrcp_proxy));
3809 let remote_peer = avdtp::Peer::new(remote);
3810
3811 let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
3812 let sbc_caps = sbc_capabilities();
3813
3814 let set_config_fut =
3816 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
3817 let mut set_config_fut = pin!(set_config_fut);
3818 match exec.run_until_stalled(&mut set_config_fut) {
3819 Poll::Ready(Ok(())) => {}
3820 x => panic!("Set capabilities should be ready but got {:?}", x),
3821 }
3822
3823 let open_fut = remote_peer.open(&sbc_endpoint_id);
3825 let mut open_fut = pin!(open_fut);
3826 match exec.run_until_stalled(&mut open_fut) {
3827 Poll::Ready(Ok(())) => {}
3828 x => panic!("Open should be ready but got {:?}", x),
3829 }
3830
3831 assert!(peer.inner.lock().volume_relay_task.is_none());
3833
3834 let (transport_chan, _remote_transport) = create_test_channels(transport);
3836 assert_eq!(Some(()), peer.receive_channel(transport_chan).ok());
3837
3838 assert!(peer.inner.lock().volume_relay_task.is_some());
3840
3841 let mut get_controller_fut = avrcp_stream.select_next_some();
3843 let _controller_server = match exec.run_until_stalled(&mut get_controller_fut) {
3844 Poll::Ready(Ok(avrcp::PeerManagerRequest::GetControllerForTarget {
3845 peer_id: req_peer_id,
3846 client,
3847 responder,
3848 })) => {
3849 assert_eq!(req_peer_id, PeerId(1).into());
3850 responder.send(Ok(())).expect("should send response");
3851 client
3852 }
3853 x => panic!("Expected GetControllerForTarget request, got {:?}", x),
3854 };
3855 }
3856
3857 #[test_case(Transport::Socket ; "socket")]
3858 #[test_case(Transport::Fidl ; "fidl")]
3859 fn test_volume_relay_starts_on_stream_start(transport: Transport) {
3860 let mut exec = fasync::TestExecutor::new();
3861 let (avrcp_proxy, mut avrcp_stream) =
3862 fidl::endpoints::create_proxy_and_stream::<avrcp::PeerManagerMarker>();
3863
3864 let (mut remote, mut profile_request_stream, peer) =
3865 setup_test_peer_with_avrcp(transport, build_test_streams(), Some(avrcp_proxy));
3866
3867 let remote_seid: StreamEndpointId = 2_u8.try_into().unwrap();
3868 let codec_params = ServiceCapability::MediaCodec {
3869 media_type: avdtp::MediaType::Audio,
3870 codec_type: avdtp::MediaCodecType::AUDIO_SBC,
3871 codec_extra: vec![0x11, 0x45, 51, 51],
3872 };
3873
3874 let remote_endpoint = avdtp::StreamEndpoint::new(
3875 2,
3876 avdtp::MediaType::Audio,
3877 avdtp::EndpointType::Sink,
3878 vec![codec_params.clone()],
3879 )
3880 .expect("valid endpoint");
3881 peer.inner.lock().set_remote_endpoints(&[remote_endpoint]);
3882
3883 assert!(peer.inner.lock().volume_relay_task.is_none());
3885
3886 let start_future = peer.stream_start(remote_seid, vec![codec_params]);
3887 let mut start_future = pin!(start_future);
3888
3889 assert!(exec.run_until_stalled(&mut start_future).is_pending());
3890 receive_simple_accept(&mut exec, &mut remote, 0x03); assert!(exec.run_until_stalled(&mut start_future).is_pending());
3892 receive_simple_accept(&mut exec, &mut remote, 0x06); assert!(exec.run_until_stalled(&mut start_future).is_pending());
3894
3895 let (transport_chan, _remote_transport) = create_test_channels(transport);
3897 let request = exec.run_until_stalled(&mut profile_request_stream.next());
3898 match request {
3899 Poll::Ready(Some(Ok(ProfileRequest::Connect {
3900 peer_id,
3901 connection: _,
3902 responder,
3903 }))) => {
3904 assert_eq!(PeerId(1), peer_id.into());
3905 let channel = transport_chan.try_into().unwrap();
3906 responder.send(Ok(channel)).expect("responder sends");
3907 }
3908 x => panic!("Expected Connect request, got {:?}", x),
3909 };
3910
3911 exec.run_until_stalled(&mut start_future)
3914 .expect("start setup finished")
3915 .expect("stream setup is ok");
3916 receive_simple_accept(&mut exec, &mut remote, 0x07); let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3918
3919 let mut get_controller_fut = avrcp_stream.select_next_some();
3921 let request = loop {
3922 match exec.run_until_stalled(&mut get_controller_fut) {
3923 Poll::Ready(r) => break r,
3924 Poll::Pending => {
3925 if exec.wake_next_timer().is_some() {
3926 continue;
3927 }
3928 panic!("Expected GetControllerForTarget request, but executor stalled");
3929 }
3930 }
3931 };
3932 let _controller_server = match request {
3933 Ok(avrcp::PeerManagerRequest::GetControllerForTarget {
3934 peer_id: req_peer_id,
3935 client,
3936 responder,
3937 }) => {
3938 assert_eq!(req_peer_id, PeerId(1).into());
3939 responder.send(Ok(())).expect("should send response");
3940 client
3941 }
3942 x => panic!("Expected GetControllerForTarget request, got {:?}", x),
3943 };
3944
3945 assert!(peer.inner.lock().volume_relay_task.is_some());
3947
3948 peer.inner.lock().local.get_mut(&1_u8.try_into().unwrap()).unwrap().abort();
3950 PeerInner::maybe_start_volume_relay(&peer.inner);
3951
3952 assert!(peer.inner.lock().volume_relay_task.is_none());
3954 }
3955
3956 #[test_case(Transport::Socket ; "socket")]
3957 #[test_case(Transport::Fidl ; "fidl")]
3958 #[fuchsia::test]
3959 fn test_volume_relay_does_not_start_when_we_are_sink(transport: Transport) {
3960 let _exec = fasync::TestExecutor::new();
3961 let (avrcp_proxy, _avrcp_stream) =
3962 fidl::endpoints::create_proxy_and_stream::<avrcp::PeerManagerMarker>();
3963
3964 let mut sink_streams = build_test_streams();
3966 let remote_id = 2_u8.try_into().unwrap();
3967 {
3968 let sink_endpoint =
3969 sink_streams.get_mut(&2_u8.try_into().unwrap()).unwrap().endpoint_mut();
3970 sink_endpoint
3971 .configure(&remote_id, vec![avdtp::ServiceCapability::MediaTransport])
3972 .unwrap();
3973 sink_endpoint.establish().unwrap();
3974 let (c1, _c2) = create_test_channels(transport);
3975 let _ = sink_endpoint.receive_channel(c1).unwrap();
3976 }
3977
3978 let (_remote, _profile_stream, peer_sink) =
3979 setup_test_peer_with_avrcp(transport, sink_streams, Some(avrcp_proxy));
3980
3981 let (c1, _c2) = create_test_channels(transport);
3983 let _ = peer_sink.receive_channel(c1);
3984 assert!(peer_sink.inner.lock().volume_relay_task.is_none());
3985 }
3986
3987 #[test_case(Transport::Socket ; "socket")]
3988 #[test_case(Transport::Fidl ; "fidl")]
3989 #[fuchsia::test]
3990 fn test_volume_relay_restarts_if_task_finished(transport: Transport) {
3991 let mut exec = fasync::TestExecutor::new();
3992 let (avrcp_proxy, mut avrcp_stream) =
3993 fidl::endpoints::create_proxy_and_stream::<avrcp::PeerManagerMarker>();
3994
3995 let mut streams = Streams::default();
3996 let test_builder = TestMediaTaskBuilder::new();
3997 streams.insert(Stream::build(
3998 make_sbc_endpoint(1, avdtp::EndpointType::Source),
3999 test_builder.builder(),
4000 ));
4001
4002 let (remote, _requests, peer) =
4003 setup_test_peer_with_avrcp(transport, streams, Some(avrcp_proxy));
4004 let remote_peer = avdtp::Peer::new(remote);
4005
4006 let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
4007 let sbc_caps = sbc_capabilities();
4008
4009 let set_config_fut =
4011 remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
4012 let mut set_config_fut = pin!(set_config_fut);
4013 assert!(exec.run_until_stalled(&mut set_config_fut).is_ready());
4014
4015 let open_fut = remote_peer.open(&sbc_endpoint_id);
4017 let mut open_fut = pin!(open_fut);
4018 assert!(exec.run_until_stalled(&mut open_fut).is_ready());
4019
4020 let (_remote_transport, transport_chan) = create_test_channels(transport);
4022 assert_eq!(Some(()), peer.receive_channel(transport_chan).ok());
4023 assert!(peer.inner.lock().volume_relay_task.is_some());
4024
4025 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
4027 let mut get_controller_fut = avrcp_stream.select_next_some();
4028 match exec.run_until_stalled(&mut get_controller_fut) {
4029 Poll::Ready(Ok(avrcp::PeerManagerRequest::GetControllerForTarget {
4030 peer_id: _,
4031 client: _,
4032 responder,
4033 })) => {
4034 responder.send(Err(zx::Status::INTERNAL.into_raw())).expect("should send response");
4035 }
4036 x => panic!("Expected GetControllerForTarget request, got {:?}", x),
4037 }
4038
4039 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
4041
4042 PeerInner::maybe_start_volume_relay(&peer.inner);
4044 assert!(peer.inner.lock().volume_relay_task.is_some());
4045
4046 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
4048 let mut get_controller_fut = avrcp_stream.select_next_some();
4049 match exec.run_until_stalled(&mut get_controller_fut) {
4050 Poll::Ready(Ok(avrcp::PeerManagerRequest::GetControllerForTarget {
4051 peer_id: req_peer_id,
4052 client: _,
4053 responder,
4054 })) => {
4055 assert_eq!(req_peer_id, PeerId(1).into());
4056 responder.send(Ok(())).expect("should send response");
4057 }
4058 x => panic!("Expected second GetControllerForTarget request, got {:?}", x),
4059 }
4060 }
4061}