Skip to main content

bt_a2dp/peer/
mod.rs

1// Copyright 2020 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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_bredr::{
12    ConnectParameters, L2capParameters, PSM_AVDTP, ProfileDescriptor, ProfileProxy,
13};
14use fuchsia_async::{self as fasync, DurationExt};
15use fuchsia_bluetooth::inspect::DebugExt;
16use fuchsia_bluetooth::types::{Channel, PeerId};
17use fuchsia_inspect as inspect;
18use fuchsia_inspect_derive::{AttachError, Inspect};
19use fuchsia_sync::Mutex;
20use futures::channel::mpsc;
21use futures::future::{BoxFuture, Either};
22use futures::stream::FuturesUnordered;
23use futures::task::{Context, Poll, Waker};
24use futures::{Future, FutureExt, StreamExt, select};
25use log::{debug, info, trace, warn};
26use std::collections::{BTreeMap, HashMap, HashSet};
27use std::pin::Pin;
28use std::sync::{Arc, Weak};
29
30/// For sending out-of-band commands over the A2DP peer.
31mod controller;
32pub use controller::ControllerPool;
33
34use crate::codec::MediaCodecConfig;
35use crate::permits::{Permit, Permits};
36use crate::stream::{Stream, Streams};
37
38/// A Peer represents an A2DP peer which may be connected to this device.
39/// Only one A2DP peer should exist for each Bluetooth peer.
40#[derive(Inspect)]
41pub struct Peer {
42    /// The id of the peer we are connected to.
43    id: PeerId,
44    /// Inner keeps track of the peer and the streams.
45    #[inspect(forward)]
46    inner: Arc<Mutex<PeerInner>>,
47    /// Profile Proxy to connect new transport channels
48    profile: ProfileProxy,
49    /// The profile descriptor for this peer, if it has been discovered.
50    descriptor: Mutex<Option<ProfileDescriptor>>,
51    /// Wakers that are to be woken when the peer disconnects.  If None, the peers have been woken
52    /// and this peer is disconnected.  Shared weakly with ClosedPeer future objects that complete
53    /// when the peer disconnects.
54    closed_wakers: Arc<Mutex<Option<Vec<Waker>>>>,
55    /// Used to report peer metrics to Cobalt.
56    metrics: bt_metrics::MetricsLogger,
57    /// A task waiting to start a stream if it hasn't been started yet.
58    start_stream_task: Mutex<Option<fasync::Task<avdtp::Result<()>>>>,
59}
60
61/// StreamPermits handles reserving and retrieving permits for streaming audio.
62/// Reservations are automatically retrieved for streams that are revoked, and when the
63/// reservation is completed, the permit is stored and a StreamPermit is sent so it can be started.
64#[derive(Clone)]
65struct StreamPermits {
66    permits: Permits,
67    open_streams: Arc<Mutex<HashMap<StreamEndpointId, Permit>>>,
68    reserved_streams: Arc<Mutex<HashSet<StreamEndpointId>>>,
69    inner: Weak<Mutex<PeerInner>>,
70    peer_id: PeerId,
71    sender: mpsc::UnboundedSender<BoxFuture<'static, StreamPermit>>,
72}
73
74#[derive(Debug)]
75struct StreamPermit {
76    local_id: StreamEndpointId,
77    open_streams: Arc<Mutex<HashMap<StreamEndpointId, Permit>>>,
78}
79
80impl StreamPermit {
81    fn local_id(&self) -> &StreamEndpointId {
82        &self.local_id
83    }
84
85    /// Returns true if a Permit is held for this stream endpoint.
86    fn is_held(&self) -> bool {
87        self.open_streams.lock().contains_key(&self.local_id)
88    }
89}
90
91impl Drop for StreamPermit {
92    fn drop(&mut self) {
93        let _ = self.open_streams.lock().remove(&self.local_id);
94    }
95}
96
97impl StreamPermits {
98    fn new(
99        inner: Weak<Mutex<PeerInner>>,
100        peer_id: PeerId,
101        permits: Permits,
102    ) -> (Self, mpsc::UnboundedReceiver<BoxFuture<'static, StreamPermit>>) {
103        let (sender, reservations_receiver) = futures::channel::mpsc::unbounded();
104        (
105            Self {
106                inner,
107                permits,
108                peer_id,
109                sender,
110                open_streams: Default::default(),
111                reserved_streams: Default::default(),
112            },
113            reservations_receiver,
114        )
115    }
116
117    fn label_for(&self, local_id: &StreamEndpointId) -> String {
118        format!("{} {}", self.peer_id, local_id)
119    }
120
121    /// Get a permit to stream on the stream with id `local_id`.
122    /// Returns Some() if there is a permit available.
123    fn get(&self, local_id: StreamEndpointId) -> Option<StreamPermit> {
124        let revoke_fn = self.make_revocation_fn(&local_id);
125        let Some(permit) = self.permits.get_revokable(revoke_fn) else {
126            info!("No permits available: {:?}", self.permits);
127            return None;
128        };
129        permit.relabel(self.label_for(&local_id));
130        if let Some(_) = self.open_streams.lock().insert(local_id.clone(), permit) {
131            warn!(id:% = self.peer_id; "Started stream {local_id:?} twice, dropping previous permit");
132        }
133        Some(StreamPermit { local_id, open_streams: self.open_streams.clone() })
134    }
135
136    /// Get a reservation that will resolve to a StreamPermit to start a stream with the id
137    /// `local_id`
138    fn setup_reservation_for(&self, local_id: StreamEndpointId) {
139        if !self.reserved_streams.lock().insert(local_id.clone()) {
140            // Already reserved.
141            return;
142        }
143        let restart_stream_available_fut = {
144            let self_revoke_fn = Self::make_revocation_fn(&self, &local_id);
145            let reservation = self.permits.reserve_revokable(self_revoke_fn);
146            let open_streams = self.open_streams.clone();
147            let reserved_streams = self.reserved_streams.clone();
148            let label = self.label_for(&local_id);
149            let local_id = local_id.clone();
150            async move {
151                let permit = reservation.await;
152                permit.relabel(label);
153                if open_streams.lock().insert(local_id.clone(), permit).is_some() {
154                    warn!("Reservation replaces acquired permit for {}", local_id.clone());
155                }
156                if !reserved_streams.lock().remove(&local_id) {
157                    warn!(local_id:%; "Unrecorded reservation resolved");
158                }
159                StreamPermit { local_id, open_streams }
160            }
161        };
162        if let Err(e) = self.sender.unbounded_send(restart_stream_available_fut.boxed()) {
163            warn!(id:% = self.peer_id, local_id:%, e:?; "Couldn't queue reservation to finish");
164        }
165    }
166
167    /// Revokes a permit that was previously delivered, suspending the local stream and signaling
168    /// the peer.
169    /// The permit must have been previously received through StreamPermits::get or a have been
170    /// restarted after being revoked, otherwise will panic.
171    fn revocation_fn(self, local_id: StreamEndpointId) -> Permit {
172        if let Ok(peer) = PeerInner::upgrade(self.inner.clone()) {
173            {
174                let mut lock = peer.lock();
175                match lock.suspend_local_stream(&local_id) {
176                    Ok(remote_id) => drop(lock.peer.suspend(&[remote_id])),
177                    Err(e) => warn!("Couldn't stop local stream {local_id:?}: {e:?}"),
178                }
179            }
180            self.setup_reservation_for(local_id.clone());
181        }
182        self.open_streams.lock().remove(&local_id).expect("permit revoked but don't have it")
183    }
184
185    fn make_revocation_fn(&self, local_id: &StreamEndpointId) -> impl FnOnce() -> Permit + use<> {
186        let local_id = local_id.clone();
187        let cloned = self.clone();
188        move || cloned.revocation_fn(local_id)
189    }
190}
191
192impl Peer {
193    /// Make a new Peer which is connected to the peer `id` using the AVDTP `peer`.
194    /// The `streams` are the local endpoints available to the peer.
195    /// `profile` will be used to initiate connections for Media Transport.
196    /// The `permits`, if provided, will acquire a permit before starting streams on this peer.
197    /// If `metrics` is included, metrics for codec availability will be reported.
198    /// This also starts a task on the executor to handle incoming events from the peer.
199    pub fn create(
200        id: PeerId,
201        peer: avdtp::Peer,
202        streams: Streams,
203        permits: Option<Permits>,
204        profile: ProfileProxy,
205        metrics: bt_metrics::MetricsLogger,
206    ) -> Self {
207        let inner = Arc::new(Mutex::new(PeerInner::new(peer, id, streams, metrics.clone())));
208        let reservations_receiver = if let Some(permits) = permits {
209            let (stream_permits, receiver) =
210                StreamPermits::new(Arc::downgrade(&inner), id, permits);
211            inner.lock().permits = Some(stream_permits);
212            receiver
213        } else {
214            let (_, receiver) = mpsc::unbounded();
215            receiver
216        };
217        let res = Self {
218            id,
219            inner,
220            profile,
221            descriptor: Mutex::new(None),
222            closed_wakers: Arc::new(Mutex::new(Some(Vec::new()))),
223            metrics,
224            start_stream_task: Mutex::new(None),
225        };
226        res.start_requests_task(reservations_receiver);
227        res
228    }
229
230    pub fn set_descriptor(&self, descriptor: ProfileDescriptor) -> Option<ProfileDescriptor> {
231        self.descriptor.lock().replace(descriptor)
232    }
233
234    /// How long to wait after a non-local establishment of a stream to start the stream.
235    /// Chosen to produce reasonably quick startup while allowing for peer start.
236    const STREAM_DWELL: zx::MonotonicDuration = zx::MonotonicDuration::from_millis(500);
237
238    /// Receive a channel from the peer that was initiated remotely.
239    /// This function should be called whenever the peer associated with this opens an L2CAP channel.
240    /// If this completes opening a stream, streams that are suspended will be scheduled to start.
241    pub fn receive_channel(&self, channel: Channel) -> avdtp::Result<()> {
242        let mut lock = self.inner.lock();
243        if lock.receive_channel(channel)? {
244            let weak = Arc::downgrade(&self.inner);
245            let mut task_lock = self.start_stream_task.lock();
246            *task_lock = Some(fasync::Task::local(async move {
247                trace!("Dwelling to start remotely-opened stream..");
248                fasync::Timer::new(Self::STREAM_DWELL.after_now()).await;
249                PeerInner::start_opened(weak).await
250            }));
251        }
252        Ok(())
253    }
254
255    /// Return a handle to the AVDTP peer, to use as initiator of commands.
256    pub fn avdtp(&self) -> avdtp::Peer {
257        let lock = self.inner.lock();
258        lock.peer.clone()
259    }
260
261    /// Returns the stream endpoints discovered by this peer.
262    pub fn remote_endpoints(&self) -> Option<Vec<avdtp::StreamEndpoint>> {
263        self.inner.lock().remote_endpoints()
264    }
265
266    /// Perform Discovery and Collect Capabilities to enumerate the endpoints and capabilities of
267    /// the connected peer.
268    /// Returns a future which performs the work and resolves to a vector of peer stream endpoints.
269    pub fn collect_capabilities(
270        &self,
271    ) -> impl Future<Output = avdtp::Result<Vec<avdtp::StreamEndpoint>>> + use<> {
272        let avdtp = self.avdtp();
273        let get_all = self.descriptor.lock().clone().is_some_and(a2dp_version_check);
274        let inner = self.inner.clone();
275        let metrics = self.metrics.clone();
276        let peer_id = self.id;
277        async move {
278            if let Some(caps) = inner.lock().remote_endpoints() {
279                return Ok(caps);
280            }
281            trace!("Discovering peer streams..");
282            let infos = avdtp.discover().await?;
283            trace!("Discovered {} streams", infos.len());
284            let mut remote_streams = Vec::new();
285            for info in infos {
286                let capabilities = if get_all {
287                    avdtp.get_all_capabilities(info.id()).await
288                } else {
289                    avdtp.get_capabilities(info.id()).await
290                };
291                match capabilities {
292                    Ok(capabilities) => {
293                        trace!("Stream {:?}", info);
294                        for cap in &capabilities {
295                            trace!("  - {:?}", cap);
296                        }
297                        remote_streams.push(avdtp::StreamEndpoint::from_info(&info, capabilities));
298                    }
299                    Err(e) => {
300                        info!(peer_id:%; "Stream {} capabilities failed: {:?}, skipping", info.id(), e);
301                    }
302                };
303            }
304            inner.lock().set_remote_endpoints(&remote_streams);
305            Self::record_cobalt_metrics(metrics, &remote_streams);
306            Ok(remote_streams)
307        }
308    }
309
310    fn record_cobalt_metrics(metrics: bt_metrics::MetricsLogger, endpoints: &[StreamEndpoint]) {
311        let codec_metrics: HashSet<_> = endpoints
312            .iter()
313            .filter_map(|endpoint| {
314                endpoint.codec_type().map(|t| codectype_to_availability_metric(t) as u32)
315            })
316            .collect();
317        metrics
318            .log_occurrences(bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID, codec_metrics);
319
320        let cap_metrics: HashSet<_> = endpoints
321            .iter()
322            .flat_map(|endpoint| {
323                endpoint
324                    .capabilities()
325                    .iter()
326                    .filter_map(|t| capability_to_metric(t))
327                    .chain(std::iter::once(
328                        bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic,
329                    ))
330                    .map(|t| t as u32)
331            })
332            .collect();
333        metrics.log_occurrences(bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID, cap_metrics);
334    }
335
336    fn transport_channel_params() -> L2capParameters {
337        L2capParameters {
338            psm: Some(PSM_AVDTP),
339            parameters: Some(ChannelParameters {
340                max_rx_packet_size: Some(65535),
341                ..Default::default()
342            }),
343            ..Default::default()
344        }
345    }
346
347    /// Open and start a media transport stream, connecting a compatible local stream to the remote
348    /// stream `remote_id`, configuring it with the `capabilities` provided.
349    /// Returns a future which should be awaited on.
350    /// The future returns Ok(()) if successfully started, and an appropriate error otherwise.
351    pub fn stream_start(
352        &self,
353        remote_id: StreamEndpointId,
354        capabilities: Vec<ServiceCapability>,
355    ) -> impl Future<Output = avdtp::Result<()>> {
356        let peer = Arc::downgrade(&self.inner);
357        let peer_id = self.id.clone();
358        let avdtp = self.avdtp();
359        let profile = self.profile.clone();
360
361        async move {
362            let codec_params =
363                capabilities.iter().find(|x| x.is_codec()).ok_or(avdtp::Error::InvalidState)?;
364            let (local_id, local_capabilities) = {
365                let peer = PeerInner::upgrade(peer.clone())?;
366                let lock = peer.lock();
367                lock.find_compatible_local_capabilities(codec_params, &remote_id)?
368            };
369
370            let local_by_cat: HashMap<ServiceCategory, ServiceCapability> =
371                local_capabilities.into_iter().map(|i| (i.category(), i)).collect();
372
373            // Filter things out if they don't have a match in the local capabilities.
374            // Order them by the ServiceCategory ordinal - some noncompliant devices care about it.
375            let shared_capabilities: BTreeMap<ServiceCategory, ServiceCapability> = capabilities
376                .into_iter()
377                .filter_map(|cap| {
378                    let Some(local_cap) = local_by_cat.get(&cap.category()) else {
379                        return None;
380                    };
381                    if cap.category() == ServiceCategory::MediaCodec {
382                        let Ok(a) = MediaCodecConfig::try_from(&cap) else {
383                            return None;
384                        };
385                        let Ok(b) = MediaCodecConfig::try_from(local_cap) else {
386                            return None;
387                        };
388                        let Some(negotiated) = MediaCodecConfig::negotiate(&a, &b) else {
389                            return None;
390                        };
391                        Some((cap.category(), (&negotiated).into()))
392                    } else {
393                        Some((cap.category(), cap))
394                    }
395                })
396                .collect();
397            let shared_capabilities: Vec<_> = shared_capabilities.into_values().collect();
398
399            trace!("Starting stream {local_id} to remote {remote_id} with {shared_capabilities:?}");
400
401            avdtp.set_configuration(&remote_id, &local_id, &shared_capabilities).await?;
402            {
403                let strong = PeerInner::upgrade(peer.clone())?;
404                strong.lock().set_opening(&local_id, &remote_id, shared_capabilities)?;
405            }
406            avdtp.open(&remote_id).await?;
407
408            debug!(peer_id:%; "Connecting transport channel");
409            let channel = profile
410                .connect(
411                    &peer_id.into(),
412                    &ConnectParameters::L2cap(Self::transport_channel_params()),
413                )
414                .await
415                .context("FIDL error: {}")?
416                .or(Err(avdtp::Error::PeerDisconnected))?;
417            trace!(peer_id:%; "Connected transport channel, converting to local Channel");
418            let channel = match channel.try_into() {
419                Err(e) => {
420                    warn!(peer_id:%, e:?; "Couldn't connect media transport: no channel");
421                    return Err(avdtp::Error::PeerDisconnected);
422                }
423                Ok(c) => c,
424            };
425
426            trace!(peer_id:%; "Connected transport channel, passing to Peer..");
427
428            {
429                let strong = PeerInner::upgrade(peer.clone())?;
430                let _ = strong.lock().receive_channel(channel)?;
431            }
432            // Start streams immediately if the channel is locally initiated.
433            PeerInner::start_opened(peer).await
434        }
435    }
436
437    /// Query whether any streams are currently started or scheduled to start.
438    pub fn streaming_active(&self) -> bool {
439        self.inner.lock().is_streaming() || self.will_start_streaming()
440    }
441
442    /// Returns true if there are any streams that are currently started.
443    #[cfg(test)]
444    fn is_streaming_now(&self) -> bool {
445        self.inner.lock().is_streaming_now()
446    }
447
448    /// Polls the task scheduled to start streaming, returning true if the task is still scheduled
449    /// to start streaming.
450    fn will_start_streaming(&self) -> bool {
451        let mut task_lock = self.start_stream_task.lock();
452        if task_lock.is_none() {
453            return false;
454        }
455        // This is the only thing that can poll the start task, so it is okay to ignore the wakeup.
456        let mut cx = Context::from_waker(&std::task::Waker::noop());
457        if let Poll::Pending = task_lock.as_mut().unwrap().poll_unpin(&mut cx) {
458            return true;
459        }
460        // Reset the task to None so that we don't try to re-poll it.
461        let _ = task_lock.take();
462        false
463    }
464
465    /// Suspend a media transport stream `local_id`.
466    /// It's possible that the stream is not active - a suspend will be attempted, but an
467    /// error from the command will be returned.
468    /// Returns the result of the suspend command.
469    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    /// Start an asynchronous task to handle any requests from the AVDTP peer.
478    /// This task completes when the remote end closes the signaling connection.
479    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                                }
508                            },
509                        }
510                    },
511                    reservation_fut = reservations_receiver.select_next_some() => {
512                        stream_reservations.push(reservation_fut)
513                    },
514                    permit = stream_reservations.select_next_some() => {
515                        if let Err(e) = PeerInner::start_permit(peer.clone(), permit).await {
516                            warn!(peer_id:% = id, e:?; "Couldn't start stream after unpause");
517                        }
518                    }
519                    complete => break,
520                }
521            }
522            info!(peer_id:% = id; "disconnected");
523            if let Some(wakers) = disconnect_wakers.upgrade() {
524                for waker in wakers.lock().take().unwrap_or_else(Vec::new) {
525                    waker.wake();
526                }
527            }
528        })
529        .detach();
530    }
531
532    /// Returns a future that will complete when the peer disconnects.
533    pub fn closed(&self) -> ClosedPeer {
534        ClosedPeer { inner: Arc::downgrade(&self.closed_wakers) }
535    }
536}
537
538/// Future which completes when the A2DP peer has closed the control connection.
539/// See `Peer::closed`
540#[must_use = "futures do nothing unless you `.await` or poll them"]
541pub struct ClosedPeer {
542    inner: Weak<Mutex<Option<Vec<Waker>>>>,
543}
544
545impl Future for ClosedPeer {
546    type Output = ();
547
548    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
549        match self.inner.upgrade() {
550            None => Poll::Ready(()),
551            Some(inner) => match inner.lock().as_mut() {
552                None => Poll::Ready(()),
553                Some(wakers) => {
554                    wakers.push(cx.waker().clone());
555                    Poll::Pending
556                }
557            },
558        }
559    }
560}
561
562/// Determines if Peer profile version is newer (>= 1.3) or older (< 1.3)
563fn a2dp_version_check(profile: ProfileDescriptor) -> bool {
564    let (Some(major), Some(minor)) = (profile.major_version, profile.minor_version) else {
565        return false;
566    };
567    (major == 1 && minor >= 3) || major > 1
568}
569
570/// Peer handles the communication with the AVDTP layer, and provides responses as appropriate
571/// based on the current state of local streams available.
572/// Each peer has its own set of local stream endpoints, and tracks a set of remote peer endpoints.
573struct PeerInner {
574    /// AVDTP peer communicating to this.
575    peer: avdtp::Peer,
576    /// The PeerId that this peer is representing
577    peer_id: PeerId,
578    /// Some(local_id) if an endpoint has been configured but hasn't finished opening.
579    /// Per AVDTP Sec 6.11 only up to one stream can be in this state.
580    opening: Option<StreamEndpointId>,
581    /// The local stream endpoint collection
582    local: Streams,
583    /// The permits that are available for this peer.
584    permits: Option<StreamPermits>,
585    /// Tasks watching for the end of a started stream. Key is the local stream id.
586    started: HashMap<StreamEndpointId, WatchedStream>,
587    /// The inspect node for this peer
588    inspect: fuchsia_inspect::Node,
589    /// The set of discovered remote endpoints. None until set.
590    remote_endpoints: Option<Vec<StreamEndpoint>>,
591    /// The inspect node representing the remote endpoints.
592    remote_inspect: fuchsia_inspect::Node,
593    /// Cobalt logger used to report peer metrics.
594    metrics: bt_metrics::MetricsLogger,
595}
596
597impl Inspect for &mut PeerInner {
598    // Set up the StreamEndpoint to update the state
599    // The MediaTask node will be created when the media task is started.
600    fn iattach(self, parent: &inspect::Node, name: impl AsRef<str>) -> Result<(), AttachError> {
601        self.inspect = parent.create_child(name.as_ref());
602        self.inspect.record_string("id", self.peer_id.to_string());
603        self.local.iattach(&self.inspect, "local_streams")
604    }
605}
606
607impl PeerInner {
608    pub fn new(
609        peer: avdtp::Peer,
610        peer_id: PeerId,
611        local: Streams,
612        metrics: bt_metrics::MetricsLogger,
613    ) -> Self {
614        Self {
615            peer,
616            peer_id,
617            opening: None,
618            local,
619            permits: None,
620            started: HashMap::new(),
621            inspect: Default::default(),
622            remote_endpoints: None,
623            remote_inspect: Default::default(),
624            metrics,
625        }
626    }
627
628    /// Returns an endpoint from the local set or a BadAcpSeid error if it doesn't exist.
629    fn get_mut(&mut self, local_id: &StreamEndpointId) -> Result<&mut Stream, avdtp::ErrorCode> {
630        self.local.get_mut(&local_id).ok_or(avdtp::ErrorCode::BadAcpSeid)
631    }
632
633    fn set_remote_endpoints(&mut self, endpoints: &[StreamEndpoint]) {
634        self.remote_inspect = self.inspect.create_child("remote_endpoints");
635        for endpoint in endpoints {
636            self.remote_inspect.record_child(inspect::unique_name("remote_"), |node| {
637                node.record_string("endpoint_id", endpoint.local_id().debug());
638                node.record_string("capabilities", endpoint.capabilities().debug());
639                node.record_string("type", endpoint.endpoint_type().debug());
640            });
641        }
642        self.remote_endpoints = Some(endpoints.iter().map(StreamEndpoint::as_new).collect());
643    }
644
645    /// If the remote endpoints have been set, returns a copy of the endpoints.
646    fn remote_endpoints(&self) -> Option<Vec<StreamEndpoint>> {
647        self.remote_endpoints.as_ref().map(|v| v.iter().map(StreamEndpoint::as_new).collect())
648    }
649
650    /// If the remote endpoint with endpoint `id` exists, return a copy of the endpoint.
651    fn remote_endpoint(&self, id: &StreamEndpointId) -> Option<StreamEndpoint> {
652        self.remote_endpoints
653            .as_ref()
654            .and_then(|v| v.iter().find(|v| v.local_id() == id).map(StreamEndpoint::as_new))
655    }
656
657    /// Returns true if there is at least one stream that has started or is starting for this peer.
658    fn is_streaming(&self) -> bool {
659        self.is_streaming_now() || self.opening.is_some()
660    }
661
662    /// Returns true if there is at least one stream in the started state for this peer.
663    fn is_streaming_now(&self) -> bool {
664        self.local.streaming().next().is_some()
665    }
666
667    fn set_opening(
668        &mut self,
669        local_id: &StreamEndpointId,
670        remote_id: &StreamEndpointId,
671        capabilities: Vec<ServiceCapability>,
672    ) -> avdtp::Result<()> {
673        if self.opening.is_some() {
674            return Err(avdtp::Error::InvalidState);
675        }
676        let peer_id = self.peer_id;
677        let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
678        stream
679            .configure(&peer_id, &remote_id, capabilities)
680            .map_err(|(cat, c)| avdtp::Error::RequestInvalidExtra(c, (&cat).into()))?;
681        stream.endpoint_mut().establish().or(Err(avdtp::Error::InvalidState))?;
682        self.opening = Some(local_id.clone());
683        Ok(())
684    }
685
686    fn upgrade(weak: Weak<Mutex<Self>>) -> avdtp::Result<Arc<Mutex<Self>>> {
687        weak.upgrade().ok_or(avdtp::Error::PeerDisconnected)
688    }
689
690    /// Start the stream that is opening, completing the opened procedure.
691    async fn start_opened(weak: Weak<Mutex<Self>>) -> avdtp::Result<()> {
692        let (avdtp, stream_pairs) = {
693            let peer = Self::upgrade(weak.clone())?;
694            let peer = peer.lock();
695            let stream_pairs: Vec<(StreamEndpointId, StreamEndpointId)> = peer
696                .local
697                .open()
698                .filter_map(|stream| {
699                    let endpoint = stream.endpoint();
700                    endpoint.remote_id().map(|id| (endpoint.local_id().clone(), id.clone()))
701                })
702                .collect();
703            (peer.peer.clone(), stream_pairs)
704        };
705        for (local_id, remote_id) in stream_pairs {
706            let permit_result =
707                Self::upgrade(weak.clone())?.lock().get_permit_or_reserve(&local_id);
708            if let Ok(permit) = permit_result {
709                Self::initiated_start(avdtp.clone(), weak.clone(), permit, &local_id, &remote_id)
710                    .await?;
711            }
712        }
713        Ok(())
714    }
715
716    async fn start_permit(weak: Weak<Mutex<Self>>, permit: StreamPermit) -> avdtp::Result<()> {
717        let local_id = permit.local_id().clone();
718        let (avdtp, remote_id) = {
719            let peer = Self::upgrade(weak.clone())?;
720            let mut peer = peer.lock();
721            let remote_id = peer
722                .get_mut(&local_id)
723                .map_err(|e| avdtp::Error::RequestInvalid(e))?
724                .endpoint()
725                .remote_id()
726                .ok_or(avdtp::Error::InvalidState)?
727                .clone();
728            (peer.peer.clone(), remote_id)
729        };
730        Self::initiated_start(avdtp, weak, Some(permit), &local_id, &remote_id).await
731    }
732
733    /// Start a stream for a local reason.  Requires a Permit to start streaming for the local stream.
734    async fn initiated_start(
735        avdtp: avdtp::Peer,
736        weak: Weak<Mutex<Self>>,
737        permit: Option<StreamPermit>,
738        local_id: &StreamEndpointId,
739        remote_id: &StreamEndpointId,
740    ) -> avdtp::Result<()> {
741        trace!(permit:?, local_id:?, remote_id:?; "Making outgoing start request");
742        let to_start = std::slice::from_ref(remote_id);
743        avdtp.start(to_start).await?;
744        trace!("Start response received: {permit:?}");
745        let peer = Self::upgrade(weak.clone())?;
746        let (peer_id, start_result) = {
747            let mut peer = peer.lock();
748            (peer.peer_id, peer.start_local_stream(permit, &local_id))
749        };
750        if let Err(e) = start_result {
751            warn!(peer_id:%, local_id:%, remote_id:%, e:?; "Failed to start local stream, suspending");
752            avdtp.suspend(to_start).await?;
753        }
754        Ok(())
755    }
756
757    /// Suspend a stream locally, returning a future to get the result from the peer.
758    fn suspend(
759        weak: Weak<Mutex<Self>>,
760        local_id: StreamEndpointId,
761    ) -> impl Future<Output = avdtp::Result<()>> {
762        let res = (move || {
763            let peer = Self::upgrade(weak.clone())?;
764            let mut peer = peer.lock();
765            Ok((peer.peer.clone(), peer.suspend_local_stream(&local_id)?))
766        })();
767        let (avdtp, remote_id) = match res {
768            Err(e) => return futures::future::err(e).left_future(),
769            Ok(r) => r,
770        };
771        let to_suspend = &[remote_id];
772        avdtp.suspend(to_suspend).right_future()
773    }
774
775    /// Finds a stream in the local stream set which is compatible with the remote_id given the codec config.
776    /// Returns the local stream ID and capabilities if found, or OutOfRange if one could not be found.
777    pub fn find_compatible_local_capabilities(
778        &self,
779        codec_params: &ServiceCapability,
780        remote_id: &StreamEndpointId,
781    ) -> avdtp::Result<(StreamEndpointId, Vec<ServiceCapability>)> {
782        let config = codec_params.try_into()?;
783        let our_direction = self.remote_endpoint(remote_id).map(|e| e.endpoint_type().opposite());
784        debug!(codec_params:?, local:? = self.local; "Looking for compatible local stream");
785        self.local
786            .compatible(config)
787            .find_map(|s| {
788                let endpoint = s.endpoint();
789                if let Some(d) = our_direction {
790                    if &d != endpoint.endpoint_type() {
791                        return None;
792                    }
793                }
794                Some((endpoint.local_id().clone(), endpoint.capabilities().clone()))
795            })
796            .ok_or(avdtp::Error::OutOfRange)
797    }
798
799    /// Attempts to acquire a permit for streaming, if the permits are set.
800    /// Returns Ok if is is okay to stream, and Err if the permit was not available and a
801    /// reservation was made.
802    fn get_permit_or_reserve(
803        &self,
804        local_id: &StreamEndpointId,
805    ) -> Result<Option<StreamPermit>, ()> {
806        let Some(permits) = self.permits.as_ref() else {
807            return Ok(None);
808        };
809        if let Some(permit) = permits.get(local_id.clone()) {
810            return Ok(Some(permit));
811        }
812        info!(peer_id:% = self.peer_id, local_id:%; "No permit to start stream, adding a reservation");
813        permits.setup_reservation_for(local_id.clone());
814        Err(())
815    }
816
817    /// Starts the stream which is in the local Streams with `local_id`.
818    /// Requires a permit to stream.
819    fn start_local_stream(
820        &mut self,
821        permit: Option<StreamPermit>,
822        local_id: &StreamEndpointId,
823    ) -> avdtp::Result<()> {
824        let peer_id = self.peer_id;
825        let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
826        // The streaming permit can be revoked while stream setup is in progress. If so, return
827        // without starting the local stream.
828        if permit.as_ref().is_some_and(|p| !p.is_held()) {
829            return Err(avdtp::Error::Other(anyhow::format_err!(
830                "streaming permit revoked during setup"
831            )));
832        }
833
834        info!(peer_id:%, stream:?; "Starting");
835        let stream_finished = stream.start().map_err(|c| avdtp::Error::RequestInvalid(c))?;
836        // TODO(https://fxbug.dev/42147239): if streaming stops unexpectedly, send a suspend to match to peer
837        let watched_stream = WatchedStream::new(permit, stream_finished);
838        if self.started.insert(local_id.clone(), watched_stream).is_some() {
839            warn!(peer_id:%, local_id:%; "Stream that was already started");
840        }
841        Ok(())
842    }
843
844    /// Suspend a stream on the local side. Returns the remote StreamEndpointId if the stream was suspended,
845    /// or a RequestInvalid error with the error code otherwise.
846    fn suspend_local_stream(
847        &mut self,
848        local_id: &StreamEndpointId,
849    ) -> avdtp::Result<StreamEndpointId> {
850        let peer_id = self.peer_id;
851        let stream = self.get_mut(&local_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
852        let remote_id = stream.endpoint().remote_id().ok_or(avdtp::Error::InvalidState)?.clone();
853        info!(peer_id:%; "Suspend stream local {local_id} <-> {remote_id} remote");
854        stream.suspend().map_err(|c| avdtp::Error::RequestInvalid(c))?;
855        let _ = self.started.remove(local_id);
856        Ok(remote_id)
857    }
858
859    /// Provide a new established L2CAP channel to this remote peer.
860    /// This function should be called whenever the remote associated with this peer opens an
861    /// L2CAP channel after the first.
862    /// Returns true if this channel completed the opening sequence.
863    fn receive_channel(&mut self, channel: Channel) -> avdtp::Result<bool> {
864        let stream_id = self.opening.as_ref().cloned().ok_or(avdtp::Error::InvalidState)?;
865        let stream = self.get_mut(&stream_id).map_err(|e| avdtp::Error::RequestInvalid(e))?;
866        let done = !stream.endpoint_mut().receive_channel(channel)?;
867        if done {
868            self.opening = None;
869        }
870        info!(peer_id:% = self.peer_id, stream_id:%; "Transport connected");
871        Ok(done)
872    }
873
874    /// Handle a single request event from the avdtp peer.
875    fn handle_request(
876        &mut self,
877        request: avdtp::Request,
878    ) -> Either<avdtp::Result<()>, impl Future<Output = avdtp::Result<()>> + use<>> {
879        use avdtp::ErrorCode;
880        use avdtp::Request::*;
881        trace!("Handling {request:?} from peer..");
882        let immediate_result = 'result: {
883            match request {
884                Discover { responder } => responder.send(&self.local.information()),
885                GetCapabilities { responder, stream_id }
886                | GetAllCapabilities { responder, stream_id } => match self.local.get(&stream_id) {
887                    None => responder.reject(ErrorCode::BadAcpSeid),
888                    Some(stream) => responder.send(stream.endpoint().capabilities()),
889                },
890                Open { responder, stream_id } => {
891                    if self.opening.is_none() {
892                        break 'result responder.reject(ErrorCode::BadState);
893                    }
894                    let Ok(stream) = self.get_mut(&stream_id) else {
895                        break 'result responder.reject(ErrorCode::BadAcpSeid);
896                    };
897                    match stream.endpoint_mut().establish() {
898                        Ok(()) => responder.send(),
899                        Err(_) => responder.reject(ErrorCode::BadState),
900                    }
901                }
902                Close { responder, stream_id } => {
903                    let peer = self.peer.clone();
904                    let Ok(stream) = self.get_mut(&stream_id) else {
905                        break 'result responder.reject(ErrorCode::BadAcpSeid);
906                    };
907                    stream.release(responder, &peer)
908                }
909                SetConfiguration { responder, local_stream_id, remote_stream_id, capabilities } => {
910                    if self.opening.is_some() {
911                        break 'result responder.reject(ServiceCategory::None, ErrorCode::BadState);
912                    }
913                    let peer_id = self.peer_id;
914                    let Ok(stream) = self.get_mut(&local_stream_id) else {
915                        break 'result responder
916                            .reject(ServiceCategory::None, ErrorCode::BadAcpSeid);
917                    };
918                    match stream.configure(&peer_id, &remote_stream_id, capabilities) {
919                        Ok(_) => {
920                            self.opening = Some(local_stream_id);
921                            responder.send()
922                        }
923                        Err((category, code)) => responder.reject(category, code),
924                    }
925                }
926                GetConfiguration { stream_id, responder } => {
927                    let Ok(stream) = self.get_mut(&stream_id) else {
928                        break 'result responder.reject(ErrorCode::BadAcpSeid);
929                    };
930                    let Some(vec_capabilities) = stream.endpoint().get_configuration() else {
931                        break 'result responder.reject(ErrorCode::BadState);
932                    };
933                    responder.send(vec_capabilities.as_slice())
934                }
935                Reconfigure { responder, local_stream_id, capabilities } => {
936                    let Ok(stream) = self.get_mut(&local_stream_id) else {
937                        break 'result responder
938                            .reject(ServiceCategory::None, ErrorCode::BadAcpSeid);
939                    };
940                    match stream.reconfigure(capabilities) {
941                        Ok(_) => responder.send(),
942                        Err((cat, code)) => responder.reject(cat, code),
943                    }
944                }
945                Start { responder, stream_ids } => {
946                    let mut immediate_suspend = Vec::new();
947                    // Fail on the first failed endpoint, as per the AVDTP spec 8.13 Note 5
948                    let result = stream_ids.into_iter().try_for_each(|seid| {
949                        let Some(stream) = self.local.get_mut(&seid) else {
950                            return Err((seid, ErrorCode::BadAcpSeid));
951                        };
952                        let remote_id = stream.endpoint().remote_id().cloned();
953                        let Some(remote_id) = remote_id else {
954                            return Err((seid, ErrorCode::BadState));
955                        };
956                        let Ok(permit) = self.get_permit_or_reserve(&seid) else {
957                            // Happens when we cannot start because of permits.
958                            // Accept this one, then queue up for suspend.
959                            // We are already reserved for a permit.
960                            immediate_suspend.push(remote_id);
961                            return Ok(());
962                        };
963                        match self.start_local_stream(permit, &seid) {
964                            Ok(()) => Ok(()),
965                            Err(avdtp::Error::RequestInvalid(code)) => Err((seid, code)),
966                            Err(_) => Err((seid, ErrorCode::BadState)),
967                        }
968                    });
969                    let response_result = match result {
970                        Ok(()) => responder.send(),
971                        Err((seid, code)) => responder.reject(&seid, code),
972                    };
973                    {
974                        let peer = self.peer.clone();
975                        return Either::Right(async move {
976                            if !immediate_suspend.is_empty() {
977                                peer.suspend(immediate_suspend.as_slice()).await?;
978                            }
979                            response_result
980                        });
981                    }
982                }
983                Suspend { responder, stream_ids } => {
984                    for seid in stream_ids {
985                        match self.suspend_local_stream(&seid) {
986                            Ok(_remote_id) => {}
987                            Err(avdtp::Error::RequestInvalid(code)) => {
988                                break 'result responder.reject(&seid, code);
989                            }
990                            Err(_e) => break 'result responder.reject(&seid, ErrorCode::BadState),
991                        }
992                    }
993                    responder.send()
994                }
995                Abort { responder, stream_id } => {
996                    let Ok(stream) = self.get_mut(&stream_id) else {
997                        // No response is sent on an invalid ID for an Abort
998                        break 'result Ok(());
999                    };
1000                    stream.abort();
1001                    self.opening = self.opening.take().filter(|local_id| local_id != &stream_id);
1002                    responder.send()
1003                }
1004                DelayReport { responder, delay, stream_id } => {
1005                    // Delay is in 1/10 ms
1006                    let delay_ns = delay as u64 * 100000;
1007                    // Record delay to cobalt.
1008                    self.metrics.log_integer(
1009                        bt_metrics::AVDTP_DELAY_REPORT_IN_NANOSECONDS_METRIC_ID,
1010                        delay_ns.try_into().unwrap_or(-1),
1011                        vec![],
1012                    );
1013                    // Report should only come after a stream is configured
1014                    let Some(stream) = self.local.get_mut(&stream_id) else {
1015                        break 'result responder.reject(avdtp::ErrorCode::BadAcpSeid);
1016                    };
1017                    let delay_str = format!("delay {}.{} ms", delay / 10, delay % 10);
1018                    let peer = self.peer_id;
1019                    match stream.set_delay(std::time::Duration::from_nanos(delay_ns)) {
1020                        Ok(()) => info!(peer:%, stream_id:%; "reported {delay_str}"),
1021                        Err(avdtp::ErrorCode::BadState) => {
1022                            info!(peer:%, stream_id:%; "bad state {delay_str}");
1023                            break 'result responder.reject(avdtp::ErrorCode::BadState);
1024                        }
1025                        Err(e) => info!(peer:%, stream_id:%, e:?; "failed {delay_str}"),
1026                    };
1027                    // Can't really respond with an Error
1028                    responder.send()
1029                }
1030            }
1031        };
1032        Either::Left(immediate_result)
1033    }
1034}
1035
1036/// A WatchedStream holds a task tracking a started stream and ensures actions are performed when
1037/// the stream media task finishes.
1038struct WatchedStream {
1039    _permit_task: fasync::Task<()>,
1040}
1041
1042impl WatchedStream {
1043    fn new(
1044        permit: Option<StreamPermit>,
1045        finish_fut: BoxFuture<'static, Result<(), anyhow::Error>>,
1046    ) -> Self {
1047        let permit_task = fasync::Task::spawn(async move {
1048            let _ = finish_fut.await;
1049            drop(permit);
1050        });
1051        Self { _permit_task: permit_task }
1052    }
1053}
1054
1055fn codectype_to_availability_metric(
1056    codec_type: &MediaCodecType,
1057) -> bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec {
1058    match codec_type {
1059        &MediaCodecType::AUDIO_SBC => {
1060            bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc
1061        }
1062        &MediaCodecType::AUDIO_MPEG12 => {
1063            bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Mpeg12
1064        }
1065        &MediaCodecType::AUDIO_AAC => {
1066            bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Aac
1067        }
1068        &MediaCodecType::AUDIO_ATRAC => {
1069            bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Atrac
1070        }
1071        &MediaCodecType::AUDIO_NON_A2DP => {
1072            bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::VendorSpecific
1073        }
1074        _ => bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Unknown,
1075    }
1076}
1077
1078fn capability_to_metric(
1079    cap: &ServiceCapability,
1080) -> Option<bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability> {
1081    match cap {
1082        ServiceCapability::DelayReporting => {
1083            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::DelayReport)
1084        }
1085        ServiceCapability::Reporting => {
1086            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Reporting)
1087        }
1088        ServiceCapability::Recovery { .. } => {
1089            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Recovery)
1090        }
1091        ServiceCapability::ContentProtection { .. } => {
1092            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::ContentProtection)
1093        }
1094        ServiceCapability::HeaderCompression { .. } => {
1095            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::HeaderCompression)
1096        }
1097        ServiceCapability::Multiplexing { .. } => {
1098            Some(bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Multiplexing)
1099        }
1100        // We ignore capabilities that we don't care to track.
1101        other => {
1102            trace!("untracked remote peer capability: {:?}", other);
1103            None
1104        }
1105    }
1106}
1107
1108#[cfg(test)]
1109mod tests {
1110    use super::*;
1111
1112    use async_utils::PollExt;
1113    use bt_channel_test_support::{Transport, create_test_channels};
1114    use bt_metrics::respond_to_metrics_req_for_test;
1115    use fidl::endpoints::create_proxy_and_stream;
1116    use fidl_fuchsia_bluetooth::ErrorCode;
1117
1118    use fidl_fuchsia_bluetooth_bredr::{
1119        ProfileMarker, ProfileRequest, ProfileRequestStream, ServiceClassProfileIdentifier,
1120    };
1121    use fidl_fuchsia_metrics::{MetricEvent, MetricEventPayload};
1122    use futures::{SinkExt, StreamExt};
1123    use std::pin::pin;
1124    use test_case::test_case;
1125
1126    use crate::media_task::tests::{TestMediaTask, TestMediaTaskBuilder};
1127    use crate::media_types::*;
1128    use crate::stream::tests::{make_sbc_endpoint, sbc_mediacodec_capability};
1129
1130    fn fake_metrics()
1131    -> (bt_metrics::MetricsLogger, fidl_fuchsia_metrics::MetricEventLoggerRequestStream) {
1132        let (c, s) = fidl::endpoints::create_proxy_and_stream::<
1133            fidl_fuchsia_metrics::MetricEventLoggerMarker,
1134        >();
1135        (bt_metrics::MetricsLogger::from_proxy(c), s)
1136    }
1137
1138    fn setup_avdtp_peer(transport: Transport) -> (avdtp::Peer, Channel) {
1139        let (signaling, remote) = create_test_channels(transport);
1140        let peer = avdtp::Peer::new(signaling);
1141        (peer, remote)
1142    }
1143
1144    fn build_test_streams() -> Streams {
1145        let mut streams = Streams::default();
1146        let source = Stream::build(
1147            make_sbc_endpoint(1, avdtp::EndpointType::Source),
1148            TestMediaTaskBuilder::new_delayable().builder(),
1149        );
1150        streams.insert(source);
1151        let sink = Stream::build(
1152            make_sbc_endpoint(2, avdtp::EndpointType::Sink),
1153            TestMediaTaskBuilder::new().builder(),
1154        );
1155        streams.insert(sink);
1156        streams
1157    }
1158
1159    fn build_test_streams_delayable() -> Streams {
1160        fn with_delay(seid: u8, direction: avdtp::EndpointType) -> StreamEndpoint {
1161            StreamEndpoint::new(
1162                seid,
1163                avdtp::MediaType::Audio,
1164                direction,
1165                vec![
1166                    avdtp::ServiceCapability::MediaTransport,
1167                    avdtp::ServiceCapability::DelayReporting,
1168                    sbc_mediacodec_capability(),
1169                ],
1170            )
1171            .expect("endpoint creation should succeed")
1172        }
1173        let mut streams = Streams::default();
1174        let source = Stream::build(
1175            with_delay(1, avdtp::EndpointType::Source),
1176            TestMediaTaskBuilder::new_delayable().builder(),
1177        );
1178        streams.insert(source);
1179        let sink = Stream::build(
1180            with_delay(2, avdtp::EndpointType::Sink),
1181            TestMediaTaskBuilder::new().builder(),
1182        );
1183        streams.insert(sink);
1184        streams
1185    }
1186
1187    #[track_caller]
1188    pub(crate) fn recv_remote(
1189        exec: &mut fasync::TestExecutor,
1190        remote: &mut Channel,
1191    ) -> Result<Vec<u8>, zx::Status> {
1192        let mut fut = remote.next();
1193        match exec.run_until_stalled(&mut fut) {
1194            Poll::Ready(Some(res)) => res,
1195            Poll::Ready(None) => Err(zx::Status::PEER_CLOSED),
1196            Poll::Pending => Err(zx::Status::SHOULD_WAIT),
1197        }
1198    }
1199
1200    /// Creates a Peer object, returning a channel connected ot the remote end, a
1201    /// ProfileRequestStream connected to the profile_proxy, and the Peer object.
1202    fn setup_test_peer(
1203        transport: Transport,
1204        use_cobalt: bool,
1205        streams: Streams,
1206        permits: Option<Permits>,
1207    ) -> (
1208        Channel,
1209        ProfileRequestStream,
1210        Option<fidl_fuchsia_metrics::MetricEventLoggerRequestStream>,
1211        Peer,
1212    ) {
1213        let (avdtp, remote) = setup_avdtp_peer(transport);
1214        let (metrics_logger, cobalt_receiver) = if use_cobalt {
1215            let (l, r) = fake_metrics();
1216            (l, Some(r))
1217        } else {
1218            (bt_metrics::MetricsLogger::default(), None)
1219        };
1220        let (profile_proxy, requests) = create_proxy_and_stream::<ProfileMarker>();
1221        let peer = Peer::create(PeerId(1), avdtp, streams, permits, profile_proxy, metrics_logger);
1222
1223        (remote, requests, cobalt_receiver, peer)
1224    }
1225
1226    #[track_caller]
1227    fn expect_send(exec: &mut fasync::TestExecutor, remote: &mut Channel, data: Vec<u8>) {
1228        exec.run_until_stalled(&mut remote.send(data))
1229            .expect("poll is ready")
1230            .expect("write successful");
1231    }
1232
1233    fn expect_get_capabilities_and_respond(
1234        exec: &mut fasync::TestExecutor,
1235        remote: &mut Channel,
1236        expected_seid: u8,
1237        response_capabilities: &[u8],
1238    ) {
1239        let received = recv_remote(exec, remote).unwrap();
1240        // Last half of header must be Single (0b00) and Command (0b00)
1241        assert_eq!(0x00, received[0] & 0xF);
1242        assert_eq!(0x02, received[1]); // 0x02 = Get Capabilities
1243        assert_eq!(expected_seid << 2, received[2]);
1244
1245        let txlabel_raw = received[0] & 0xF0;
1246
1247        // Expect a get capabilities and respond.
1248        #[rustfmt::skip]
1249        let mut get_capabilities_rsp = vec![
1250            txlabel_raw << 4 | 0x2, // TxLabel (same) + ResponseAccept (0x02)
1251            0x02 // GetCapabilities
1252        ];
1253
1254        get_capabilities_rsp.extend_from_slice(response_capabilities);
1255
1256        expect_send(exec, remote, get_capabilities_rsp);
1257    }
1258
1259    fn expect_get_all_capabilities_and_respond(
1260        exec: &mut fasync::TestExecutor,
1261        remote: &mut Channel,
1262        expected_seid: u8,
1263        response_capabilities: &[u8],
1264    ) {
1265        let received = recv_remote(exec, remote).unwrap();
1266        // Last half of header must be Single (0b00) and Command (0b00)
1267        assert_eq!(0x00, received[0] & 0xF);
1268        assert_eq!(0x0C, received[1]); // 0x0C = Get All Capabilities
1269        assert_eq!(expected_seid << 2, received[2]);
1270
1271        let txlabel_raw = received[0] & 0xF0;
1272
1273        // Expect a get capabilities and respond.
1274        #[rustfmt::skip]
1275        let mut get_capabilities_rsp = vec![
1276            txlabel_raw << 4 | 0x2, // TxLabel (same) + ResponseAccept (0x02)
1277            0x0C // GetAllCapabilities
1278        ];
1279
1280        get_capabilities_rsp.extend_from_slice(response_capabilities);
1281
1282        expect_send(exec, remote, get_capabilities_rsp);
1283    }
1284
1285    #[test_case(Transport::Socket ; "socket")]
1286    #[test_case(Transport::Fidl ; "fidl")]
1287    #[fuchsia::test]
1288    fn disconnected(transport: Transport) {
1289        let mut exec = fasync::TestExecutor::new();
1290        let (proxy, _stream) = create_proxy_and_stream::<ProfileMarker>();
1291        let (signaling, remote) = create_test_channels(transport);
1292
1293        let id = PeerId(1);
1294
1295        let avdtp = avdtp::Peer::new(signaling);
1296        let peer = Peer::create(
1297            id,
1298            avdtp,
1299            Streams::default(),
1300            None,
1301            proxy,
1302            bt_metrics::MetricsLogger::default(),
1303        );
1304
1305        let closed_fut = peer.closed();
1306
1307        let mut closed_fut = pin!(closed_fut);
1308
1309        assert!(exec.run_until_stalled(&mut closed_fut).is_pending());
1310
1311        // Close the remote channel
1312        drop(remote);
1313
1314        assert!(exec.run_until_stalled(&mut closed_fut).is_ready());
1315    }
1316
1317    #[test_case(Transport::Socket ; "socket")]
1318    #[test_case(Transport::Fidl ; "fidl")]
1319    #[fuchsia::test]
1320    fn peer_collect_capabilities_success(transport: Transport) {
1321        let mut exec = fasync::TestExecutor::new();
1322
1323        let (mut remote, _, cobalt_receiver, peer) =
1324            setup_test_peer(transport, true, build_test_streams(), None);
1325
1326        let p: ProfileDescriptor = ProfileDescriptor {
1327            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
1328            major_version: Some(1),
1329            minor_version: Some(2),
1330            ..Default::default()
1331        };
1332        let _ = peer.set_descriptor(p);
1333
1334        let collect_future = peer.collect_capabilities();
1335        let mut collect_future = pin!(collect_future);
1336
1337        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1338
1339        // Expect a discover command.
1340        let received = recv_remote(&mut exec, &mut remote).unwrap();
1341        // Last half of header must be Single (0b00) and Command (0b00)
1342        assert_eq!(0x00, received[0] & 0xF);
1343        assert_eq!(0x01, received[1]); // 0x01 = Discover
1344
1345        let txlabel_raw = received[0] & 0xF0;
1346
1347        // Respond with a set of streams.
1348        let response: &[u8] = &[
1349            txlabel_raw << 4 | 0x0 << 2 | 0x2, // txlabel (same), Single (0b00), Response Accept (0b10)
1350            0x01,                              // Discover
1351            0x3E << 2 | 0x0 << 1,              // SEID (3E), Not In Use (0b0)
1352            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1353            0x01 << 2 | 0x1 << 1,              // SEID (1), In Use (0b1)
1354            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1355        ];
1356        expect_send(&mut exec, &mut remote, response.to_vec());
1357
1358        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1359
1360        // Expect a get capabilities and respond.
1361        #[rustfmt::skip]
1362        let capabilities_rsp = &[
1363            // MediaTransport (Length of Service Capability = 0)
1364            0x01, 0x00,
1365            // Media Codec (LOSC = 2 + 4), Media Type Audio (0x00), Codec type (0x04), Codec specific 0xF09F9296
1366            0x07, 0x06, 0x00, 0x04, 0xF0, 0x9F, 0x92, 0x96
1367        ];
1368        expect_get_capabilities_and_respond(&mut exec, &mut remote, 0x3E, capabilities_rsp);
1369
1370        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1371
1372        // Expect a get capabilities and respond.
1373        #[rustfmt::skip]
1374        let capabilities_rsp = &[
1375            // MediaTransport (Length of Service Capability = 0)
1376            0x01, 0x00,
1377            // Media Codec (LOSC = 2 + 2), Media Type Audio (0x00), Codec type (0x00), Codec specific 0xC0DE
1378            0x07, 0x04, 0x00, 0x00, 0xC0, 0xDE
1379        ];
1380        expect_get_capabilities_and_respond(&mut exec, &mut remote, 0x01, capabilities_rsp);
1381
1382        match exec.run_until_stalled(&mut collect_future) {
1383            Poll::Pending => panic!("collect capabilities should be complete"),
1384            Poll::Ready(Err(e)) => panic!("collect capabilities should have succeeded: {}", e),
1385            Poll::Ready(Ok(endpoints)) => {
1386                let first_seid: StreamEndpointId = 0x3E_u8.try_into().unwrap();
1387                let second_seid: StreamEndpointId = 0x01_u8.try_into().unwrap();
1388                for stream in endpoints {
1389                    if stream.local_id() == &first_seid {
1390                        let expected_caps = vec![
1391                            ServiceCapability::MediaTransport,
1392                            ServiceCapability::MediaCodec {
1393                                media_type: avdtp::MediaType::Audio,
1394                                codec_type: avdtp::MediaCodecType::new(0x04),
1395                                codec_extra: vec![0xF0, 0x9F, 0x92, 0x96],
1396                            },
1397                        ];
1398                        assert_eq!(&expected_caps, stream.capabilities());
1399                    } else if stream.local_id() == &second_seid {
1400                        let expected_codec_type = avdtp::MediaCodecType::new(0x00);
1401                        assert_eq!(Some(&expected_codec_type), stream.codec_type());
1402                    } else {
1403                        panic!("Unexpected endpoint in the streams collected");
1404                    }
1405                }
1406            }
1407        }
1408
1409        // Collect reported cobalt logs.
1410        let mut recv = cobalt_receiver.expect("should have receiver");
1411        let mut log_events = Vec::new();
1412        while let Poll::Ready(Some(Ok(req))) = exec.run_until_stalled(&mut recv.next()) {
1413            log_events.push(respond_to_metrics_req_for_test(req));
1414        }
1415
1416        // Should have sent two metric events for codec and one for capability.
1417        assert_eq!(3, log_events.len());
1418        assert!(log_events.contains(&MetricEvent {
1419            metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1420            event_codes: vec![
1421                bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc as u32
1422            ],
1423            payload: MetricEventPayload::Count(1),
1424        }));
1425        assert!(log_events.contains(&MetricEvent {
1426            metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1427            event_codes: vec![
1428                bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Atrac as u32
1429            ],
1430            payload: MetricEventPayload::Count(1),
1431        }));
1432        assert!(log_events.contains(&MetricEvent {
1433            metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1434            event_codes: vec![
1435                bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic as u32
1436            ],
1437            payload: MetricEventPayload::Count(1),
1438        }));
1439
1440        // The second time, we don't expect to ask the peer again.
1441        let collect_future = peer.collect_capabilities();
1442        let mut collect_future = pin!(collect_future);
1443
1444        match exec.run_until_stalled(&mut collect_future) {
1445            Poll::Ready(Ok(endpoints)) => assert_eq!(2, endpoints.len()),
1446            x => panic!("Expected get remote capabilities to be done, got {:?}", x),
1447        };
1448    }
1449
1450    #[test_case(Transport::Socket ; "socket")]
1451    #[test_case(Transport::Fidl ; "fidl")]
1452    #[fuchsia::test]
1453    fn peer_collect_all_capabilities_success(transport: Transport) {
1454        let mut exec = fasync::TestExecutor::new();
1455
1456        let (mut remote, _, cobalt_receiver, peer) =
1457            setup_test_peer(transport, true, build_test_streams(), None);
1458        let p: ProfileDescriptor = ProfileDescriptor {
1459            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
1460            major_version: Some(1),
1461            minor_version: Some(3),
1462            ..Default::default()
1463        };
1464        let _ = peer.set_descriptor(p);
1465
1466        let collect_future = peer.collect_capabilities();
1467        let mut collect_future = pin!(collect_future);
1468
1469        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1470
1471        // Expect a discover command.
1472        let received = recv_remote(&mut exec, &mut remote).unwrap();
1473        // Last half of header must be Single (0b00) and Command (0b00)
1474        assert_eq!(0x00, received[0] & 0xF);
1475        assert_eq!(0x01, received[1]); // 0x01 = Discover
1476
1477        let txlabel_raw = received[0] & 0xF0;
1478
1479        // Respond with a set of streams.
1480        let response: &[u8] = &[
1481            txlabel_raw << 4 | 0x0 << 2 | 0x2, // txlabel (same), Single (0b00), Response Accept (0b10)
1482            0x01,                              // Discover
1483            0x3E << 2 | 0x0 << 1,              // SEID (3E), Not In Use (0b0)
1484            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1485            0x01 << 2 | 0x1 << 1,              // SEID (1), In Use (0b1)
1486            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1487        ];
1488        expect_send(&mut exec, &mut remote, response.to_vec());
1489
1490        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1491
1492        // Expect a get all capabilities and respond.
1493        #[rustfmt::skip]
1494        let capabilities_rsp = &[
1495            // MediaTransport (Length of Service Capability = 0)
1496            0x01, 0x00,
1497            // Media Codec (LOSC = 2 + 4), Media Type Audio (0x00), Codec type (0x40), Codec specific 0xF09F9296
1498            0x07, 0x06, 0x00, 0x40, 0xF0, 0x9F, 0x92, 0x96,
1499            // Delay Reporting (LOSC = 0)
1500            0x08, 0x00
1501        ];
1502        expect_get_all_capabilities_and_respond(&mut exec, &mut remote, 0x3E, capabilities_rsp);
1503
1504        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1505
1506        // Expect a get all capabilities and respond.
1507        #[rustfmt::skip]
1508        let capabilities_rsp = &[
1509            // MediaTransport (Length of Service Capability = 0)
1510            0x01, 0x00,
1511            // Media Codec (LOSC = 2 + 2), Media Type Audio (0x00), Codec type (0x00), Codec specific 0xC0DE
1512            0x07, 0x04, 0x00, 0x00, 0xC0, 0xDE
1513        ];
1514        expect_get_all_capabilities_and_respond(&mut exec, &mut remote, 0x01, capabilities_rsp);
1515
1516        match exec.run_until_stalled(&mut collect_future) {
1517            Poll::Pending => panic!("collect capabilities should be complete"),
1518            Poll::Ready(Err(e)) => panic!("collect capabilities should have succeeded: {}", e),
1519            Poll::Ready(Ok(endpoints)) => {
1520                let first_seid: StreamEndpointId = 0x3E_u8.try_into().unwrap();
1521                let second_seid: StreamEndpointId = 0x01_u8.try_into().unwrap();
1522                for stream in endpoints {
1523                    if stream.local_id() == &first_seid {
1524                        let expected_caps = vec![
1525                            ServiceCapability::MediaTransport,
1526                            ServiceCapability::MediaCodec {
1527                                media_type: avdtp::MediaType::Audio,
1528                                codec_type: avdtp::MediaCodecType::new(0x40),
1529                                codec_extra: vec![0xF0, 0x9F, 0x92, 0x96],
1530                            },
1531                            ServiceCapability::DelayReporting,
1532                        ];
1533                        assert_eq!(&expected_caps, stream.capabilities());
1534                    } else if stream.local_id() == &second_seid {
1535                        let expected_codec_type = avdtp::MediaCodecType::new(0x00);
1536                        assert_eq!(Some(&expected_codec_type), stream.codec_type());
1537                    } else {
1538                        panic!("Unexpected endpoint in the streams collected");
1539                    }
1540                }
1541            }
1542        }
1543
1544        // Collect reported cobalt logs.
1545        let mut recv = cobalt_receiver.expect("should have receiver");
1546        let mut log_events = Vec::new();
1547        while let Poll::Ready(Some(Ok(req))) = exec.run_until_stalled(&mut recv.next()) {
1548            log_events.push(respond_to_metrics_req_for_test(req));
1549        }
1550
1551        // Should have sent two metric events for codec and two for capability.
1552        assert_eq!(4, log_events.len());
1553        assert!(log_events.contains(&MetricEvent {
1554            metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1555            event_codes: vec![
1556                bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Unknown as u32
1557            ],
1558            payload: MetricEventPayload::Count(1),
1559        }));
1560        assert!(log_events.contains(&MetricEvent {
1561            metric_id: bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID,
1562            event_codes: vec![
1563                bt_metrics::A2dpCodecAvailabilityMigratedMetricDimensionCodec::Sbc as u32
1564            ],
1565            payload: MetricEventPayload::Count(1),
1566        }));
1567        assert!(log_events.contains(&MetricEvent {
1568            metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1569            event_codes: vec![
1570                bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::Basic as u32
1571            ],
1572            payload: MetricEventPayload::Count(1),
1573        }));
1574        assert!(log_events.contains(&MetricEvent {
1575            metric_id: bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID,
1576            event_codes: vec![
1577                bt_metrics::A2dpRemotePeerCapabilitiesMetricDimensionCapability::DelayReport as u32
1578            ],
1579            payload: MetricEventPayload::Count(1),
1580        }));
1581
1582        // The second time, we don't expect to ask the peer again.
1583        let collect_future = peer.collect_capabilities();
1584        let mut collect_future = pin!(collect_future);
1585
1586        match exec.run_until_stalled(&mut collect_future) {
1587            Poll::Ready(Ok(endpoints)) => assert_eq!(2, endpoints.len()),
1588            x => panic!("Expected get remote capabilities to be done, got {:?}", x),
1589        };
1590    }
1591
1592    #[test_case(Transport::Socket ; "socket")]
1593    #[test_case(Transport::Fidl ; "fidl")]
1594    #[fuchsia::test]
1595    fn peer_collect_capabilities_discovery_fails(transport: Transport) {
1596        let mut exec = fasync::TestExecutor::new();
1597
1598        let (mut remote, _, _, peer) =
1599            setup_test_peer(transport, false, build_test_streams(), None);
1600
1601        let collect_future = peer.collect_capabilities();
1602        let mut collect_future = pin!(collect_future);
1603
1604        // Shouldn't finish yet.
1605        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1606
1607        // Expect a discover command.
1608        let received = recv_remote(&mut exec, &mut remote).unwrap();
1609        // Last half of header must be Single (0b00) and Command (0b00)
1610        assert_eq!(0x00, received[0] & 0xF);
1611        assert_eq!(0x01, received[1]); // 0x01 = Discover
1612
1613        let txlabel_raw = received[0] & 0xF0;
1614
1615        // Respond with an error.
1616        let response: &[u8] = &[
1617            txlabel_raw | 0x0 << 2 | 0x3, // txlabel (same), Single (0b00), Response Reject (0b11)
1618            0x01,                         // Discover
1619            0x31,                         // BAD_STATE
1620        ];
1621        expect_send(&mut exec, &mut remote, response.to_vec());
1622
1623        // Should be done with an error.
1624        // Should finish!
1625        match exec.run_until_stalled(&mut collect_future) {
1626            Poll::Pending => panic!("Should be ready after discovery failure"),
1627            Poll::Ready(Ok(x)) => panic!("Should be an error but returned {x:?}"),
1628            Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
1629                assert_eq!(Some(Ok(avdtp::ErrorCode::BadState)), e.error_code());
1630            }
1631            Poll::Ready(Err(e)) => panic!("Should have been a RemoteRejected was was {e:?}"),
1632        }
1633    }
1634
1635    #[test_case(Transport::Socket ; "socket")]
1636    #[test_case(Transport::Fidl ; "fidl")]
1637    #[fuchsia::test]
1638    fn peer_collect_capabilities_get_capability_fails(transport: Transport) {
1639        let mut exec = fasync::TestExecutor::new();
1640
1641        let (mut remote, _, _, peer) = setup_test_peer(transport, true, build_test_streams(), None);
1642
1643        let collect_future = peer.collect_capabilities();
1644        let mut collect_future = pin!(collect_future);
1645
1646        // Shouldn't finish yet.
1647        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1648
1649        // Expect a discover command.
1650        let received = recv_remote(&mut exec, &mut remote).unwrap();
1651        // Last half of header must be Single (0b00) and Command (0b00)
1652        assert_eq!(0x00, received[0] & 0xF);
1653        assert_eq!(0x01, received[1]); // 0x01 = Discover
1654
1655        let txlabel_raw = received[0] & 0xF0;
1656
1657        // Respond with a set of streams.
1658        let response: &[u8] = &[
1659            txlabel_raw << 4 | 0x0 << 2 | 0x2, // txlabel (same), Single (0b00), Response Accept (0b10)
1660            0x01,                              // Discover
1661            0x3E << 2 | 0x0 << 1,              // SEID (3E), Not In Use (0b0)
1662            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1663            0x01 << 2 | 0x1 << 1,              // SEID (1), In Use (0b1)
1664            0x00 << 4 | 0x1 << 3,              // Audio (0x00), Sink (0x01)
1665        ];
1666        expect_send(&mut exec, &mut remote, response.to_vec());
1667
1668        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1669
1670        // Expect a get capabilities request
1671        let expected_seid = 0x3E;
1672        let received = recv_remote(&mut exec, &mut remote).unwrap();
1673        // Last half of header must be Single (0b00) and Command (0b00)
1674        assert_eq!(0x00, received[0] & 0xF);
1675        assert_eq!(0x02, received[1]); // 0x02 = Get Capabilities
1676        assert_eq!(expected_seid << 2, received[2]);
1677
1678        let txlabel_raw = received[0] & 0xF0;
1679
1680        let response: &[u8] = &[
1681            txlabel_raw | 0x0 << 2 | 0x3, // txlabel (same), Single (0b00), Response Reject (0b11)
1682            0x02,                         // Get Capabilities
1683            0x12,                         // BAD_ACP_SEID
1684        ];
1685        expect_send(&mut exec, &mut remote, response.to_vec());
1686
1687        assert!(exec.run_until_stalled(&mut collect_future).is_pending());
1688
1689        // Expect a get capabilities request (skipped the last one)
1690        let expected_seid = 0x01;
1691        let received = recv_remote(&mut exec, &mut remote).unwrap();
1692        // Last half of header must be Single (0b00) and Command (0b00)
1693        assert_eq!(0x00, received[0] & 0xF);
1694        assert_eq!(0x02, received[1]); // 0x02 = Get Capabilities
1695        assert_eq!(expected_seid << 2, received[2]);
1696
1697        let txlabel_raw = received[0] & 0xF0;
1698
1699        let response: &[u8] = &[
1700            txlabel_raw | 0x0 << 2 | 0x3, // txlabel (same), Single (0b00), Response Reject (0b11)
1701            0x02,                         // Get Capabilities
1702            0x12,                         // BAD_ACP_SEID
1703        ];
1704        expect_send(&mut exec, &mut remote, response.to_vec());
1705
1706        // Should be done without an error, but with no streams.
1707        match exec.run_until_stalled(&mut collect_future) {
1708            Poll::Pending => panic!("Should be ready after discovery failure"),
1709            Poll::Ready(Err(e)) => panic!("Shouldn't be an error but returned {:?}", e),
1710            Poll::Ready(Ok(map)) => assert_eq!(0, map.len()),
1711        }
1712    }
1713
1714    fn receive_simple_accept(exec: &mut fasync::TestExecutor, remote: &mut Channel, signal_id: u8) {
1715        let received = recv_remote(exec, remote).expect("expected a packet");
1716        // Last half of header must be Single (0b00) and Command (0b00)
1717        assert_eq!(0x00, received[0] & 0xF);
1718        assert_eq!(signal_id, received[1]);
1719
1720        let txlabel_raw = received[0] & 0xF0;
1721
1722        let response: &[u8] = &[
1723            txlabel_raw | 0x0 << 2 | 0x2, // txlabel (same), Single (0b00), Response Accept (0b10)
1724            signal_id,
1725        ];
1726        expect_send(exec, remote, response.to_vec());
1727    }
1728
1729    #[test_case(Transport::Socket ; "socket")]
1730    #[test_case(Transport::Fidl ; "fidl")]
1731    #[fuchsia::test]
1732    fn peer_stream_start_success(transport: Transport) {
1733        let mut exec = fasync::TestExecutor::new();
1734
1735        let (mut remote, mut profile_request_stream, _, peer) =
1736            setup_test_peer(transport, false, build_test_streams(), None);
1737
1738        let remote_seid = 2_u8.try_into().unwrap();
1739
1740        let codec_params = ServiceCapability::MediaCodec {
1741            media_type: avdtp::MediaType::Audio,
1742            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
1743            codec_extra: vec![0x11, 0x45, 51, 51],
1744        };
1745
1746        let start_future = peer.stream_start(remote_seid, vec![codec_params]);
1747        let mut start_future = pin!(start_future);
1748
1749        match exec.run_until_stalled(&mut start_future) {
1750            Poll::Pending => {}
1751            x => panic!("Expected pending, but got {x:?}"),
1752        };
1753
1754        receive_simple_accept(&mut exec, &mut remote, 0x03); // Set Configuration
1755
1756        assert!(exec.run_until_stalled(&mut start_future).is_pending());
1757
1758        receive_simple_accept(&mut exec, &mut remote, 0x06); // Open
1759
1760        match exec.run_until_stalled(&mut start_future) {
1761            Poll::Pending => {}
1762            Poll::Ready(Err(e)) => panic!("Expected to be pending but error: {:?}", e),
1763            Poll::Ready(Ok(_)) => panic!("Expected to be pending but finished!"),
1764        };
1765
1766        // Should connect the media channel after open.
1767        let (transport_chan, _remote_chan) = create_test_channels(transport);
1768
1769        let request = exec.run_until_stalled(&mut profile_request_stream.next());
1770        match request {
1771            Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, connection, responder }))) => {
1772                assert_eq!(PeerId(1), peer_id.into());
1773                assert_eq!(connection, ConnectParameters::L2cap(Peer::transport_channel_params()));
1774                let channel = transport_chan.try_into().unwrap();
1775                responder.send(Ok(channel)).expect("responder sends");
1776            }
1777            x => panic!("Should have sent a open l2cap request, but got {:?}", x),
1778        };
1779
1780        match exec.run_until_stalled(&mut start_future) {
1781            Poll::Pending => {}
1782            Poll::Ready(Err(e)) => panic!("Expected to be pending but error: {:?}", e),
1783            Poll::Ready(Ok(_)) => panic!("Expected to be pending but finished!"),
1784        };
1785
1786        receive_simple_accept(&mut exec, &mut remote, 0x07); // Start
1787
1788        // Should return the media stream (which should be connected)
1789        // Should be done without an error, but with no streams.
1790        match exec.run_until_stalled(&mut start_future) {
1791            Poll::Pending => panic!("Should be ready after start succeeds"),
1792            Poll::Ready(Err(e)) => panic!("Shouldn't be an error but returned {:?}", e),
1793            // TODO: confirm the stream is usable
1794            Poll::Ready(Ok(())) => {
1795                assert!(peer.is_streaming_now());
1796            }
1797        }
1798    }
1799
1800    #[test_case(Transport::Socket ; "socket")]
1801    #[test_case(Transport::Fidl ; "fidl")]
1802    #[fuchsia::test]
1803    fn peer_stream_start_picks_correct_direction(transport: Transport) {
1804        let mut exec = fasync::TestExecutor::new();
1805
1806        let (remote, _, _, peer) = setup_test_peer(transport, false, build_test_streams(), None);
1807        let remote = avdtp::Peer::new(remote);
1808        let mut remote_events = remote.take_request_stream();
1809
1810        // Respond as if we have a single SBC Source Stream
1811        fn remote_handle_request(req: avdtp::Request) {
1812            let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
1813            let res = match req {
1814                avdtp::Request::Discover { responder } => {
1815                    let infos = [avdtp::StreamInformation::new(
1816                        expected_stream_id,
1817                        false,
1818                        avdtp::MediaType::Audio,
1819                        avdtp::EndpointType::Source,
1820                    )];
1821                    responder.send(&infos)
1822                }
1823                avdtp::Request::GetAllCapabilities { stream_id, responder }
1824                | avdtp::Request::GetCapabilities { stream_id, responder } => {
1825                    assert_eq!(expected_stream_id, stream_id);
1826                    let caps = vec![
1827                        ServiceCapability::MediaTransport,
1828                        ServiceCapability::MediaCodec {
1829                            media_type: avdtp::MediaType::Audio,
1830                            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
1831                            codec_extra: vec![0x11, 0x45, 51, 250],
1832                        },
1833                    ];
1834                    responder.send(&caps[..])
1835                }
1836                avdtp::Request::Open { responder, stream_id } => {
1837                    assert_eq!(expected_stream_id, stream_id);
1838                    responder.send()
1839                }
1840                avdtp::Request::SetConfiguration {
1841                    responder,
1842                    local_stream_id,
1843                    remote_stream_id,
1844                    ..
1845                } => {
1846                    assert_eq!(local_stream_id, expected_stream_id);
1847                    // This is the "sink" local stream id.
1848                    assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
1849                    responder.send()
1850                }
1851                x => panic!("Unexpected request: {:?}", x),
1852            };
1853            res.expect("should be able to respond");
1854        }
1855
1856        // Need to discover the remote streams first, or the stream start will not work.
1857        let collect_capabilities_fut = peer.collect_capabilities();
1858        let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
1859
1860        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
1861
1862        let request = exec.run_singlethreaded(&mut remote_events.next());
1863        remote_handle_request(request.expect("should have a discovery request").unwrap());
1864
1865        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
1866        let request = exec.run_singlethreaded(&mut remote_events.next());
1867        remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
1868
1869        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
1870
1871        // Try to start the stream.  It should continue to configure and connect.
1872        let remote_seid = 4_u8.try_into().unwrap();
1873
1874        let codec_params = ServiceCapability::MediaCodec {
1875            media_type: avdtp::MediaType::Audio,
1876            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
1877            codec_extra: vec![0x11, 0x45, 51, 51],
1878        };
1879        let start_future = peer.stream_start(remote_seid, vec![codec_params]);
1880        let mut start_future = pin!(start_future);
1881
1882        assert!(exec.run_until_stalled(&mut start_future).is_pending());
1883        let request = exec.run_singlethreaded(&mut remote_events.next());
1884        remote_handle_request(request.expect("should have a set_capabilities request").unwrap());
1885
1886        assert!(exec.run_until_stalled(&mut start_future).is_pending());
1887        let request = exec.run_singlethreaded(&mut remote_events.next());
1888        remote_handle_request(request.expect("should have an open request").unwrap());
1889    }
1890
1891    #[test_case(Transport::Socket ; "socket")]
1892    #[test_case(Transport::Fidl ; "fidl")]
1893    #[fuchsia::test]
1894    fn peer_stream_start_strips_unsupported_local_capabilities(transport: Transport) {
1895        let mut exec = fasync::TestExecutor::new();
1896
1897        let (remote, _, _, peer) = setup_test_peer(transport, false, build_test_streams(), None);
1898        let remote = avdtp::Peer::new(remote);
1899        let mut remote_events = remote.take_request_stream();
1900
1901        // Respond as if we have a single SBC Source Stream
1902        fn remote_handle_request(req: avdtp::Request) {
1903            let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
1904            let res = match req {
1905                avdtp::Request::Discover { responder } => {
1906                    let infos = [avdtp::StreamInformation::new(
1907                        expected_stream_id,
1908                        false,
1909                        avdtp::MediaType::Audio,
1910                        avdtp::EndpointType::Source,
1911                    )];
1912                    responder.send(&infos)
1913                }
1914                avdtp::Request::GetAllCapabilities { stream_id, responder }
1915                | avdtp::Request::GetCapabilities { stream_id, responder } => {
1916                    assert_eq!(expected_stream_id, stream_id);
1917                    let caps = vec![
1918                        ServiceCapability::MediaTransport,
1919                        // We don't have a local delay-reporting, so this shouldn't be requested.
1920                        ServiceCapability::DelayReporting,
1921                        ServiceCapability::MediaCodec {
1922                            media_type: avdtp::MediaType::Audio,
1923                            codec_type: avdtp::MediaCodecType::AUDIO_AAC,
1924                            codec_extra: vec![128, 0, 132, 134, 0, 0],
1925                        },
1926                    ];
1927                    responder.send(&caps[..])
1928                }
1929                avdtp::Request::Open { responder, stream_id } => {
1930                    assert_eq!(expected_stream_id, stream_id);
1931                    responder.send()
1932                }
1933                avdtp::Request::SetConfiguration {
1934                    responder,
1935                    local_stream_id,
1936                    remote_stream_id,
1937                    capabilities,
1938                } => {
1939                    assert_eq!(local_stream_id, expected_stream_id);
1940                    // This is the "sink" local stream id.
1941                    assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
1942                    // Make sure we didn't request a DelayReport since the local Sink doesn't
1943                    // support it.
1944                    assert!(!capabilities.contains(&ServiceCapability::DelayReporting));
1945                    responder.send()
1946                }
1947                x => panic!("Unexpected request: {:?}", x),
1948            };
1949            res.expect("should be able to respond");
1950        }
1951
1952        // Need to discover the remote streams first, or the stream start will not work.
1953        let collect_capabilities_fut = peer.collect_capabilities();
1954        let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
1955
1956        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
1957
1958        let request = exec.run_singlethreaded(&mut remote_events.next());
1959        remote_handle_request(request.expect("should have a discovery request").unwrap());
1960
1961        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
1962        let request = exec.run_singlethreaded(&mut remote_events.next());
1963        remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
1964
1965        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
1966
1967        // Try to start the stream.  It should continue to configure and connect.
1968        let remote_seid = 4_u8.try_into().unwrap();
1969
1970        let codec_params = ServiceCapability::MediaCodec {
1971            media_type: avdtp::MediaType::Audio,
1972            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
1973            codec_extra: vec![0x11, 0x45, 51, 51],
1974        };
1975        let start_future =
1976            peer.stream_start(remote_seid, vec![codec_params, ServiceCapability::DelayReporting]);
1977        let mut start_future = pin!(start_future);
1978
1979        assert!(exec.run_until_stalled(&mut start_future).is_pending());
1980        let request = exec.run_singlethreaded(&mut remote_events.next());
1981        remote_handle_request(request.expect("should have a set_configuration request").unwrap());
1982
1983        assert!(exec.run_until_stalled(&mut start_future).is_pending());
1984        let request = exec.run_singlethreaded(&mut remote_events.next());
1985        remote_handle_request(request.expect("should have an open request").unwrap());
1986    }
1987
1988    #[test_case(Transport::Socket ; "socket")]
1989    #[test_case(Transport::Fidl ; "fidl")]
1990    #[fuchsia::test]
1991    fn peer_stream_start_orders_local_capabilities(transport: Transport) {
1992        let mut exec = fasync::TestExecutor::new();
1993
1994        let (remote, _, _, peer) =
1995            setup_test_peer(transport, false, build_test_streams_delayable(), None);
1996        let remote = avdtp::Peer::new(remote);
1997        let mut remote_events = remote.take_request_stream();
1998
1999        // Respond as if we have a single SBC Source Stream
2000        fn remote_handle_request(req: avdtp::Request) {
2001            let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2002            let res = match req {
2003                avdtp::Request::Discover { responder } => {
2004                    let infos = [avdtp::StreamInformation::new(
2005                        expected_stream_id,
2006                        false,
2007                        avdtp::MediaType::Audio,
2008                        avdtp::EndpointType::Source,
2009                    )];
2010                    responder.send(&infos)
2011                }
2012                avdtp::Request::GetAllCapabilities { stream_id, responder }
2013                | avdtp::Request::GetCapabilities { stream_id, responder } => {
2014                    assert_eq!(expected_stream_id, stream_id);
2015                    let caps = &[
2016                        ServiceCapability::MediaTransport,
2017                        ServiceCapability::MediaCodec {
2018                            media_type: avdtp::MediaType::Audio,
2019                            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2020                            codec_extra: vec![0x11, 0x45, 51, 250],
2021                        },
2022                        ServiceCapability::DelayReporting,
2023                    ];
2024                    responder.send(caps)
2025                }
2026                avdtp::Request::Open { responder, stream_id } => {
2027                    assert_eq!(expected_stream_id, stream_id);
2028                    responder.send()
2029                }
2030                avdtp::Request::SetConfiguration {
2031                    responder,
2032                    local_stream_id,
2033                    remote_stream_id,
2034                    capabilities,
2035                } => {
2036                    assert_eq!(local_stream_id, expected_stream_id);
2037                    // This is the "sink" local stream id.
2038                    assert_eq!(remote_stream_id, 2_u8.try_into().unwrap());
2039                    // The capabilities should be in order.
2040                    let mut capabilities_ordered = capabilities.clone();
2041                    capabilities_ordered.sort_by_key(ServiceCapability::category);
2042                    assert_eq!(capabilities, capabilities_ordered);
2043                    responder.send()
2044                }
2045                x => panic!("Unexpected request: {:?}", x),
2046            };
2047            res.expect("should be able to respond");
2048        }
2049
2050        // Need to discover the remote streams first, or the stream start will not work.
2051        let collect_capabilities_fut = peer.collect_capabilities();
2052        let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2053
2054        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2055
2056        let request = exec.run_singlethreaded(&mut remote_events.next());
2057        remote_handle_request(request.expect("should have a discovery request").unwrap());
2058
2059        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2060        let request = exec.run_singlethreaded(&mut remote_events.next());
2061        remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2062
2063        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2064
2065        // Try to start the stream.  It should continue to configure and connect.
2066        let remote_seid = 4_u8.try_into().unwrap();
2067
2068        let codec_params = ServiceCapability::MediaCodec {
2069            media_type: avdtp::MediaType::Audio,
2070            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2071            codec_extra: vec![0x11, 0x45, 51, 51],
2072        };
2073        let start_future = peer.stream_start(
2074            remote_seid,
2075            vec![
2076                ServiceCapability::MediaTransport,
2077                ServiceCapability::DelayReporting,
2078                codec_params,
2079            ],
2080        );
2081        let mut start_future = pin!(start_future);
2082
2083        assert!(exec.run_until_stalled(&mut start_future).is_pending());
2084        let request = exec.run_singlethreaded(&mut remote_events.next());
2085        remote_handle_request(request.expect("should have a set_configuration request").unwrap());
2086
2087        assert!(exec.run_until_stalled(&mut start_future).is_pending());
2088        let request = exec.run_singlethreaded(&mut remote_events.next());
2089        remote_handle_request(request.expect("should have an open request").unwrap());
2090    }
2091
2092    /// Tests that A2DP streaming does not start if the streaming permit is revoked during streaming
2093    /// setup.
2094    #[test_case(Transport::Socket ; "socket")]
2095    #[test_case(Transport::Fidl ; "fidl")]
2096    #[fuchsia::test]
2097    fn peer_stream_start_permit_revoked(transport: Transport) {
2098        let mut exec = fasync::TestExecutor::new();
2099
2100        let test_permits = Permits::new(1);
2101        let (mut remote, mut profile_request_stream, _, peer) =
2102            setup_test_peer(transport, false, build_test_streams(), Some(test_permits.clone()));
2103
2104        let remote_seid = 2_u8.try_into().unwrap();
2105
2106        let codec_params = ServiceCapability::MediaCodec {
2107            media_type: avdtp::MediaType::Audio,
2108            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2109            codec_extra: vec![0x11, 0x45, 51, 51],
2110        };
2111
2112        let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2113        let mut start_future = pin!(start_future);
2114
2115        let _ = exec
2116            .run_until_stalled(&mut start_future)
2117            .expect_pending("waiting for set config response");
2118        receive_simple_accept(&mut exec, &mut remote, 0x03); // Set Configuration
2119        exec.run_until_stalled(&mut start_future).expect_pending("waiting for open response");
2120        receive_simple_accept(&mut exec, &mut remote, 0x06); // Open
2121        exec.run_until_stalled(&mut start_future).expect_pending("waiting for media transport");
2122        assert!(!peer.is_streaming_now());
2123
2124        // Should connect the media channel after open.
2125        let (transport_chan, _remote_chan) = create_test_channels(transport);
2126
2127        let request = exec.run_until_stalled(&mut profile_request_stream.next());
2128        match request {
2129            Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, connection, responder }))) => {
2130                assert_eq!(PeerId(1), peer_id.into());
2131                assert_eq!(connection, ConnectParameters::L2cap(Peer::transport_channel_params()));
2132                let channel = transport_chan.try_into().unwrap();
2133                responder.send(Ok(channel)).expect("responder sends");
2134            }
2135            x => panic!("Should have sent a open l2cap request, but got {:?}", x),
2136        };
2137
2138        exec.run_until_stalled(&mut start_future).expect_pending("waiting for media transport");
2139        assert!(!peer.is_streaming_now());
2140
2141        // Before peer responds to start, the permit gets taken.
2142        let seized_permits = test_permits.seize();
2143        assert_eq!(seized_permits.len(), 1);
2144        receive_simple_accept(&mut exec, &mut remote, 0x07); // Start
2145
2146        // Streaming should not locally begin because there is no available permit. The Start
2147        // response is handled gracefully.
2148        exec.run_until_stalled(&mut start_future)
2149            .expect_pending("waiting to send outgoing suspend");
2150        assert!(!peer.is_streaming_now());
2151        // We should issue an outgoing suspend request to synchronize state with the remote peer.
2152        receive_simple_accept(&mut exec, &mut remote, 0x09); // Suspend
2153
2154        // The start future should resolve without Error, and A2DP should not have started
2155        // streaming.
2156        let () = exec
2157            .run_until_stalled(&mut start_future)
2158            .expect("start finished")
2159            .expect("suspended stream is ok");
2160        assert!(!peer.is_streaming_now());
2161    }
2162
2163    #[test_case(Transport::Socket ; "socket")]
2164    #[test_case(Transport::Fidl ; "fidl")]
2165    #[fuchsia::test]
2166    fn peer_stream_start_fails_wrong_direction(transport: Transport) {
2167        let mut exec = fasync::TestExecutor::new();
2168
2169        // Setup peers with only one Source Stream.
2170        let mut streams = Streams::default();
2171        let source = Stream::build(
2172            make_sbc_endpoint(1, avdtp::EndpointType::Source),
2173            TestMediaTaskBuilder::new().builder(),
2174        );
2175        streams.insert(source);
2176
2177        let (remote, _requests, _, peer) = setup_test_peer(transport, false, streams, None);
2178        let remote = avdtp::Peer::new(remote);
2179        let mut remote_events = remote.take_request_stream();
2180
2181        // Respond as if we have a single SBC Source Stream
2182        fn remote_handle_request(req: avdtp::Request) {
2183            let expected_stream_id: StreamEndpointId = 2_u8.try_into().unwrap();
2184            let res = match req {
2185                avdtp::Request::Discover { responder } => {
2186                    let infos = [avdtp::StreamInformation::new(
2187                        expected_stream_id,
2188                        false,
2189                        avdtp::MediaType::Audio,
2190                        avdtp::EndpointType::Source,
2191                    )];
2192                    responder.send(&infos)
2193                }
2194                avdtp::Request::GetAllCapabilities { stream_id, responder }
2195                | avdtp::Request::GetCapabilities { stream_id, responder } => {
2196                    assert_eq!(expected_stream_id, stream_id);
2197                    let caps = vec![
2198                        ServiceCapability::MediaTransport,
2199                        ServiceCapability::MediaCodec {
2200                            media_type: avdtp::MediaType::Audio,
2201                            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2202                            codec_extra: vec![0x11, 0x45, 51, 250],
2203                        },
2204                    ];
2205                    responder.send(&caps[..])
2206                }
2207                avdtp::Request::Open { responder, .. } => responder.send(),
2208                avdtp::Request::SetConfiguration { responder, .. } => responder.send(),
2209                x => panic!("Unexpected request: {:?}", x),
2210            };
2211            res.expect("should be able to respond");
2212        }
2213
2214        // Need to discover the remote streams first, or the stream start will always work.
2215        let collect_capabilities_fut = peer.collect_capabilities();
2216        let mut collect_capabilities_fut = pin!(collect_capabilities_fut);
2217
2218        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2219
2220        let request = exec.run_singlethreaded(&mut remote_events.next());
2221        remote_handle_request(request.expect("should have a discovery request").unwrap());
2222
2223        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_pending());
2224        let request = exec.run_singlethreaded(&mut remote_events.next());
2225        remote_handle_request(request.expect("should have a get_capabilities request").unwrap());
2226
2227        assert!(exec.run_until_stalled(&mut collect_capabilities_fut).is_ready());
2228
2229        // Try to start the stream.  It should fail with OutOfRange because we can't connect a Source to a Source.
2230        let remote_seid = 2_u8.try_into().unwrap();
2231
2232        let codec_params = ServiceCapability::MediaCodec {
2233            media_type: avdtp::MediaType::Audio,
2234            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2235            codec_extra: vec![0x11, 0x45, 51, 51],
2236        };
2237        let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2238        let mut start_future = pin!(start_future);
2239
2240        match exec.run_until_stalled(&mut start_future) {
2241            Poll::Ready(Err(avdtp::Error::OutOfRange)) => {}
2242            x => panic!("Expected a ready OutOfRange error but got {:?}", x),
2243        };
2244    }
2245
2246    #[test_case(Transport::Socket ; "socket")]
2247    #[test_case(Transport::Fidl ; "fidl")]
2248    #[fuchsia::test]
2249    fn peer_stream_start_fails_to_connect(transport: Transport) {
2250        let mut exec = fasync::TestExecutor::new();
2251
2252        let (mut remote, mut profile_request_stream, _, peer) =
2253            setup_test_peer(transport, false, build_test_streams(), None);
2254
2255        let remote_seid = 2_u8.try_into().unwrap();
2256
2257        let codec_params = ServiceCapability::MediaCodec {
2258            media_type: avdtp::MediaType::Audio,
2259            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2260            codec_extra: vec![0x11, 0x45, 51, 51],
2261        };
2262
2263        let start_future = peer.stream_start(remote_seid, vec![codec_params]);
2264        let mut start_future = pin!(start_future);
2265
2266        match exec.run_until_stalled(&mut start_future) {
2267            Poll::Pending => {}
2268            x => panic!("was expecting pending but got {x:?}"),
2269        };
2270
2271        receive_simple_accept(&mut exec, &mut remote, 0x03); // Set Configuration
2272
2273        assert!(exec.run_until_stalled(&mut start_future).is_pending());
2274
2275        receive_simple_accept(&mut exec, &mut remote, 0x06); // Open
2276
2277        match exec.run_until_stalled(&mut start_future) {
2278            Poll::Pending => {}
2279            Poll::Ready(x) => panic!("Expected to be pending but {x:?}"),
2280        };
2281
2282        // Should connect the media channel after open.
2283        let request = exec.run_until_stalled(&mut profile_request_stream.next());
2284        match request {
2285            Poll::Ready(Some(Ok(ProfileRequest::Connect { peer_id, responder, .. }))) => {
2286                assert_eq!(PeerId(1), peer_id.into());
2287                responder.send(Err(ErrorCode::Failed)).expect("responder sends");
2288            }
2289            x => panic!("Should have sent a open l2cap request, but got {:?}", x),
2290        };
2291
2292        // Should return an error.
2293        // Should be done without an error, but with no streams.
2294        match exec.run_until_stalled(&mut start_future) {
2295            Poll::Pending => panic!("Should be ready after start fails"),
2296            Poll::Ready(Ok(_stream)) => panic!("Shouldn't have succeeded stream here"),
2297            Poll::Ready(Err(_)) => {}
2298        }
2299    }
2300
2301    /// Test that the delay reports get acknowledged and they are sent to cobalt.
2302    #[test_case(Transport::Socket ; "socket")]
2303    #[test_case(Transport::Fidl ; "fidl")]
2304    #[fuchsia::test]
2305    async fn peer_delay_report(transport: Transport) {
2306        let (remote, _profile_requests, cobalt_recv, peer) =
2307            setup_test_peer(transport, true, build_test_streams(), None);
2308        let remote_peer = avdtp::Peer::new(remote);
2309        let mut remote_events = remote_peer.take_request_stream();
2310
2311        // Respond as if we have a single SBC Sink Stream
2312        async fn remote_handle_request(req: avdtp::Request, peer: &avdtp::Peer) {
2313            let expected_stream_id: StreamEndpointId = 4_u8.try_into().unwrap();
2314            // "peer" in this case is the test code Peer stream
2315            let expected_peer_stream_id: StreamEndpointId = 1_u8.try_into().unwrap();
2316            use avdtp::Request::*;
2317            match req {
2318                Discover { responder } => {
2319                    let infos = [avdtp::StreamInformation::new(
2320                        expected_stream_id,
2321                        false,
2322                        avdtp::MediaType::Audio,
2323                        avdtp::EndpointType::Sink,
2324                    )];
2325                    responder.send(&infos).expect("response should succeed");
2326                }
2327                GetAllCapabilities { stream_id, responder }
2328                | GetCapabilities { stream_id, responder } => {
2329                    assert_eq!(expected_stream_id, stream_id);
2330                    let caps = vec![
2331                        ServiceCapability::MediaTransport,
2332                        ServiceCapability::MediaCodec {
2333                            media_type: avdtp::MediaType::Audio,
2334                            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2335                            codec_extra: vec![0x11, 0x45, 51, 250],
2336                        },
2337                    ];
2338                    responder.send(&caps[..]).expect("response should succeed");
2339                    // Sending a delayreport before the stream is configured is not allowed, it's a
2340                    // bad state.
2341                    assert!(peer.delay_report(&expected_peer_stream_id, 0xc0de).await.is_err());
2342                }
2343                Open { responder, stream_id } => {
2344                    // Configuration has happened but open not succeeded yet, send delay reports.
2345                    assert!(peer.delay_report(&expected_stream_id, 0xc0de).await.is_err());
2346                    // Send a delay report to the peer.
2347                    peer.delay_report(&expected_peer_stream_id, 0xc0de)
2348                        .await
2349                        .expect("should get acked correctly");
2350                    assert_eq!(expected_stream_id, stream_id);
2351                    responder.send().expect("response should succeed");
2352                }
2353                SetConfiguration { responder, local_stream_id, remote_stream_id, .. } => {
2354                    assert_eq!(local_stream_id, expected_stream_id);
2355                    assert_eq!(remote_stream_id, expected_peer_stream_id);
2356                    responder.send().expect("should send back response without issue");
2357                }
2358                x => panic!("Unexpected request: {:?}", x),
2359            };
2360        }
2361
2362        let collect_fut = pin!(peer.collect_capabilities());
2363
2364        // Discover then a GetCapabilities.
2365        let Either::Left((request, collect_fut)) =
2366            futures::future::select(remote_events.next(), collect_fut).await
2367        else {
2368            panic!("Collect future shouldn't finish first");
2369        };
2370        let collect_fut = pin!(collect_fut);
2371        remote_handle_request(request.expect("a request").unwrap(), &remote_peer).await;
2372        let Either::Left((request, collect_fut)) =
2373            futures::future::select(remote_events.next(), collect_fut).await
2374        else {
2375            panic!("Collect future shouldn't finish first");
2376        };
2377        remote_handle_request(request.expect("a request").unwrap(), &remote_peer).await;
2378
2379        // Collect future should be able to finish now.
2380        assert_eq!(1, collect_fut.await.expect("should get the remote endpoints back").len());
2381
2382        // Try to start the stream.  It should go through the normal motions,
2383        let remote_seid = 4_u8.try_into().unwrap();
2384
2385        let codec_params = ServiceCapability::MediaCodec {
2386            media_type: avdtp::MediaType::Audio,
2387            codec_type: avdtp::MediaCodecType::AUDIO_SBC,
2388            codec_extra: vec![0x11, 0x45, 51, 51],
2389        };
2390
2391        // We don't expect this task to finish before being dropped, since we never respond to the
2392        // request to open the transport channel.
2393        let _start_task = fasync::Task::spawn(async move {
2394            let _ = peer.stream_start(remote_seid, vec![codec_params]).await;
2395            panic!("stream start task finished");
2396        });
2397
2398        let request = remote_events.next().await.expect("should have set_config").unwrap();
2399        remote_handle_request(request, &remote_peer).await;
2400
2401        let request = remote_events.next().await.expect("should have open").unwrap();
2402        remote_handle_request(request, &remote_peer).await;
2403
2404        let mut cobalt = cobalt_recv.expect("should have receiver");
2405
2406        let mut got_ids = HashMap::new();
2407        let delay_metric_id = bt_metrics::AVDTP_DELAY_REPORT_IN_NANOSECONDS_METRIC_ID;
2408        while got_ids.len() < 3 || *got_ids.get(&delay_metric_id).unwrap_or(&0) < 3 {
2409            let report = respond_to_metrics_req_for_test(cobalt.next().await.unwrap().unwrap());
2410            let _ = got_ids.entry(report.metric_id).and_modify(|x| *x += 1).or_insert(1);
2411            // All the delay reports should report the same value correctly.
2412            if report.metric_id == delay_metric_id {
2413                assert_eq!(MetricEventPayload::IntegerValue(0xc0de * 100000), report.payload);
2414            }
2415        }
2416        assert!(got_ids.contains_key(&bt_metrics::A2DP_CODEC_AVAILABILITY_MIGRATED_METRIC_ID));
2417        assert!(got_ids.contains_key(&bt_metrics::A2DP_REMOTE_PEER_CAPABILITIES_METRIC_ID));
2418        assert!(got_ids.contains_key(&delay_metric_id));
2419        // There should have been three reports.
2420        // We report the delay amount even if it fails to work.
2421        assert_eq!(got_ids.get(&delay_metric_id).cloned(), Some(3));
2422    }
2423
2424    fn sbc_capabilities() -> Vec<ServiceCapability> {
2425        let sbc_codec_info = SbcCodecInfo::new(
2426            SbcSamplingFrequency::FREQ48000HZ,
2427            SbcChannelMode::JOINT_STEREO,
2428            SbcBlockCount::SIXTEEN,
2429            SbcSubBands::EIGHT,
2430            SbcAllocation::LOUDNESS,
2431            /* min_bpv= */ 53,
2432            /* max_bpv= */ 53,
2433        )
2434        .expect("sbc codec info");
2435
2436        vec![avdtp::ServiceCapability::MediaTransport, sbc_codec_info.into()]
2437    }
2438
2439    /// Test that the remote end can configure and start a stream.
2440    #[test_case(Transport::Socket ; "socket")]
2441    #[test_case(Transport::Fidl ; "fidl")]
2442    #[fuchsia::test]
2443    fn peer_as_acceptor(transport: Transport) {
2444        let mut exec = fasync::TestExecutor::new();
2445
2446        let mut streams = Streams::default();
2447        let mut test_builder = TestMediaTaskBuilder::new();
2448        streams.insert(Stream::build(
2449            make_sbc_endpoint(1, avdtp::EndpointType::Source),
2450            test_builder.builder(),
2451        ));
2452
2453        let (remote, _requests, _, peer) = setup_test_peer(transport, false, streams, None);
2454        let remote_peer = avdtp::Peer::new(remote);
2455
2456        let discover_fut = remote_peer.discover();
2457        let mut discover_fut = pin!(discover_fut);
2458
2459        let expected = vec![make_sbc_endpoint(1, avdtp::EndpointType::Source).information()];
2460        match exec.run_until_stalled(&mut discover_fut) {
2461            Poll::Ready(Ok(res)) => assert_eq!(res, expected),
2462            x => panic!("Expected discovery to complete and got {:?}", x),
2463        };
2464
2465        let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
2466        let unknown_endpoint_id = 2_u8.try_into().expect("should be able to get sbc endpointid");
2467
2468        let get_caps_fut = remote_peer.get_capabilities(&sbc_endpoint_id);
2469        let mut get_caps_fut = pin!(get_caps_fut);
2470
2471        match exec.run_until_stalled(&mut get_caps_fut) {
2472            // There are two caps (mediatransport, mediacodec) in the sbc endpoint.
2473            Poll::Ready(Ok(caps)) => assert_eq!(2, caps.len()),
2474            x => panic!("Get capabilities should be ready but got {:?}", x),
2475        };
2476
2477        let get_caps_fut = remote_peer.get_capabilities(&unknown_endpoint_id);
2478        let mut get_caps_fut = pin!(get_caps_fut);
2479
2480        match exec.run_until_stalled(&mut get_caps_fut) {
2481            Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
2482                assert_eq!(Some(Ok(avdtp::ErrorCode::BadAcpSeid)), e.error_code())
2483            }
2484            x => panic!("Get capabilities should be a ready error but got {:?}", x),
2485        };
2486
2487        let get_caps_fut = remote_peer.get_all_capabilities(&sbc_endpoint_id);
2488        let mut get_caps_fut = pin!(get_caps_fut);
2489
2490        match exec.run_until_stalled(&mut get_caps_fut) {
2491            // There are two caps (mediatransport, mediacodec) in the sbc endpoint.
2492            Poll::Ready(Ok(caps)) => assert_eq!(2, caps.len()),
2493            x => panic!("Get capabilities should be ready but got {:?}", x),
2494        };
2495
2496        let sbc_caps = sbc_capabilities();
2497        let set_config_fut =
2498            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2499        let mut set_config_fut = pin!(set_config_fut);
2500
2501        match exec.run_until_stalled(&mut set_config_fut) {
2502            Poll::Ready(Ok(())) => {}
2503            x => panic!("Set capabilities should be ready but got {:?}", x),
2504        };
2505
2506        let open_fut = remote_peer.open(&sbc_endpoint_id);
2507        let mut open_fut = pin!(open_fut);
2508        match exec.run_until_stalled(&mut open_fut) {
2509            Poll::Ready(Ok(())) => {}
2510            x => panic!("Open should be ready but got {:?}", x),
2511        };
2512
2513        // Establish a media transport stream
2514        let (transport_chan, _remote_transport) = create_test_channels(transport);
2515
2516        assert_eq!(Some(()), peer.receive_channel(transport_chan).ok());
2517
2518        let stream_ids = vec![sbc_endpoint_id.clone()];
2519        let start_fut = remote_peer.start(&stream_ids);
2520        let mut start_fut = pin!(start_fut);
2521        match exec.run_until_stalled(&mut start_fut) {
2522            Poll::Ready(Ok(())) => {}
2523            x => panic!("Start should be ready but got {:?}", x),
2524        };
2525
2526        // The task should be created locally and started.
2527        let media_task = test_builder.expect_task();
2528        assert!(media_task.is_started());
2529
2530        let suspend_fut = remote_peer.suspend(&stream_ids);
2531        let mut suspend_fut = pin!(suspend_fut);
2532        match exec.run_until_stalled(&mut suspend_fut) {
2533            Poll::Ready(Ok(())) => {}
2534            x => panic!("Start should be ready but got {:?}", x),
2535        };
2536
2537        // Should have stopped the media task on suspend.
2538        assert!(!media_task.is_started());
2539    }
2540
2541    #[test_case(Transport::Socket ; "socket")]
2542    #[test_case(Transport::Fidl ; "fidl")]
2543    #[fuchsia::test]
2544    fn peer_set_config_reject_first(transport: Transport) {
2545        let mut exec = fasync::TestExecutor::new();
2546
2547        let mut streams = Streams::default();
2548        let test_builder = TestMediaTaskBuilder::new();
2549        streams.insert(Stream::build(
2550            make_sbc_endpoint(1, avdtp::EndpointType::Source),
2551            test_builder.builder(),
2552        ));
2553
2554        let (remote, _requests, _, _peer) = setup_test_peer(transport, false, streams, None);
2555        let remote_peer = avdtp::Peer::new(remote);
2556
2557        let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
2558
2559        let wrong_freq_sbc = &[SbcCodecInfo::new(
2560            SbcSamplingFrequency::FREQ44100HZ, // 44.1 is not supported by the caps from above.
2561            SbcChannelMode::JOINT_STEREO,
2562            SbcBlockCount::SIXTEEN,
2563            SbcSubBands::EIGHT,
2564            SbcAllocation::LOUDNESS,
2565            /* min_bpv= */ 53,
2566            /* max_bpv= */ 53,
2567        )
2568        .expect("sbc codec info")
2569        .into()];
2570
2571        let set_config_fut =
2572            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, wrong_freq_sbc);
2573        let mut set_config_fut = pin!(set_config_fut);
2574
2575        match exec.run_until_stalled(&mut set_config_fut) {
2576            Poll::Ready(Err(avdtp::Error::RemoteRejected(e))) => {
2577                assert!(e.service_category().is_some())
2578            }
2579            x => panic!("Set capabilities should have been rejected but got {:?}", x),
2580        };
2581
2582        let sbc_caps = sbc_capabilities();
2583        let set_config_fut =
2584            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2585        let mut set_config_fut = pin!(set_config_fut);
2586
2587        match exec.run_until_stalled(&mut set_config_fut) {
2588            Poll::Ready(Ok(())) => {}
2589            x => panic!("Set capabilities should be ready but got {:?}", x),
2590        };
2591    }
2592
2593    #[test_case(Transport::Socket ; "socket")]
2594    #[test_case(Transport::Fidl ; "fidl")]
2595    #[fuchsia::test]
2596    fn peer_starts_waiting_streams(transport_mode: Transport) {
2597        let mut exec = fasync::TestExecutor::new_with_fake_time();
2598        exec.set_fake_time(fasync::MonotonicInstant::from_nanos(5_000_000_000));
2599
2600        let mut streams = Streams::default();
2601        let mut test_builder = TestMediaTaskBuilder::new();
2602        streams.insert(Stream::build(
2603            make_sbc_endpoint(1, avdtp::EndpointType::Source),
2604            test_builder.builder(),
2605        ));
2606
2607        let (remote, _requests, _, peer) = setup_test_peer(transport_mode, false, streams, None);
2608        let remote_peer = avdtp::Peer::new(remote);
2609
2610        let sbc_endpoint_id = 1_u8.try_into().expect("should be able to get sbc endpointid");
2611
2612        let sbc_caps = sbc_capabilities();
2613        let set_config_fut =
2614            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2615        let mut set_config_fut = pin!(set_config_fut);
2616
2617        match exec.run_until_stalled(&mut set_config_fut) {
2618            Poll::Ready(Ok(())) => {}
2619            x => panic!("Set capabilities should be ready but got {:?}", x),
2620        };
2621
2622        let open_fut = remote_peer.open(&sbc_endpoint_id);
2623        let mut open_fut = pin!(open_fut);
2624        match exec.run_until_stalled(&mut open_fut) {
2625            Poll::Ready(Ok(())) => {}
2626            x => panic!("Open should be ready but got {:?}", x),
2627        };
2628
2629        // Establish a media transport stream
2630        let (transport, _remote_transport) = create_test_channels(transport_mode);
2631        assert_eq!(Some(()), peer.receive_channel(transport).ok());
2632
2633        // The remote end should get a start request after the timeout.
2634        let mut remote_requests = remote_peer.take_request_stream();
2635        let next_remote_request_fut = remote_requests.next();
2636        let mut next_remote_request_fut = pin!(next_remote_request_fut);
2637
2638        // Nothing should happen immediately.
2639        assert!(exec.run_until_stalled(&mut next_remote_request_fut).is_pending());
2640
2641        // After the timeout has passed..
2642        exec.set_fake_time(zx::MonotonicDuration::from_seconds(3).after_now());
2643        let _ = exec.wake_expired_timers();
2644
2645        let stream_ids = match exec.run_until_stalled(&mut next_remote_request_fut) {
2646            Poll::Ready(Some(Ok(avdtp::Request::Start { responder, stream_ids }))) => {
2647                responder.send().unwrap();
2648                stream_ids
2649            }
2650            x => panic!("Expected to receive a start request for the stream, got {:?}", x),
2651        };
2652
2653        // We should start the media task, so the task should be created locally
2654        let media_task =
2655            exec.run_until_stalled(&mut test_builder.next_task()).expect("ready").unwrap();
2656        assert!(media_task.is_started());
2657
2658        // Remote peer should still be able to suspend the stream.
2659        let suspend_fut = remote_peer.suspend(&stream_ids);
2660        let mut suspend_fut = pin!(suspend_fut);
2661        match exec.run_until_stalled(&mut suspend_fut) {
2662            Poll::Ready(Ok(())) => {}
2663            x => panic!("Suspend should be ready but got {:?}", x),
2664        };
2665
2666        // Should have stopped the media task on suspend.
2667        assert!(!media_task.is_started());
2668    }
2669
2670    #[test_case(Transport::Socket ; "socket")]
2671    #[test_case(Transport::Fidl ; "fidl")]
2672    #[fuchsia::test]
2673    fn needs_permit_to_start_streams(transport_mode: Transport) {
2674        let mut exec = fasync::TestExecutor::new();
2675
2676        let mut streams = Streams::default();
2677        let mut test_builder = TestMediaTaskBuilder::new();
2678        streams.insert(Stream::build(
2679            make_sbc_endpoint(1, avdtp::EndpointType::Sink),
2680            test_builder.builder(),
2681        ));
2682        streams.insert(Stream::build(
2683            make_sbc_endpoint(2, avdtp::EndpointType::Sink),
2684            test_builder.builder(),
2685        ));
2686        let mut next_task_fut = test_builder.next_task();
2687
2688        let permits = Permits::new(1);
2689        let taken_permit = permits.get().expect("permit taken");
2690        let (remote, _profile_request_stream, _, peer) =
2691            setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
2692        let remote_peer = avdtp::Peer::new(remote);
2693
2694        let sbc_endpoint_id = 1_u8.try_into().unwrap();
2695
2696        let sbc_caps = sbc_capabilities();
2697        let mut set_config_fut =
2698            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
2699
2700        match exec.run_until_stalled(&mut set_config_fut) {
2701            Poll::Ready(Ok(())) => {}
2702            x => panic!("Set capabilities should be ready but got {:?}", x),
2703        };
2704
2705        let mut open_fut = remote_peer.open(&sbc_endpoint_id);
2706        match exec.run_until_stalled(&mut open_fut) {
2707            Poll::Ready(Ok(())) => {}
2708            x => panic!("Open should be ready but got {:?}", x),
2709        };
2710
2711        // Establish a media transport stream
2712        let (transport, _remote_transport) = create_test_channels(transport_mode);
2713        assert_eq!(Some(()), peer.receive_channel(transport).ok());
2714
2715        // Do the same, but for the OTHER stream.
2716        let sbc_endpoint_two = 2_u8.try_into().unwrap();
2717
2718        let mut set_config_fut =
2719            remote_peer.set_configuration(&sbc_endpoint_two, &sbc_endpoint_two, &sbc_caps);
2720
2721        match exec.run_until_stalled(&mut set_config_fut) {
2722            Poll::Ready(Ok(())) => {}
2723            x => panic!("Set capabilities should be ready but got {:?}", x),
2724        };
2725
2726        let mut open_fut = remote_peer.open(&sbc_endpoint_two);
2727        match exec.run_until_stalled(&mut open_fut) {
2728            Poll::Ready(Ok(())) => {}
2729            x => panic!("Open should be ready but got {:?}", x),
2730        };
2731
2732        // Establish a media transport stream
2733        let (transport_two, _remote_transport_two) = create_test_channels(transport_mode);
2734        assert_eq!(Some(()), peer.receive_channel(transport_two).ok());
2735
2736        // Remote peer should still be able to try to start the stream, and we will say yes, but
2737        // that last seid looks wonky.
2738        let unknown_endpoint_id: StreamEndpointId = 9_u8.try_into().unwrap();
2739        let stream_ids = [sbc_endpoint_id.clone(), unknown_endpoint_id.clone()];
2740        let mut start_fut = remote_peer.start(&stream_ids);
2741        match exec.run_until_stalled(&mut start_fut) {
2742            Poll::Ready(Err(avdtp::Error::RemoteRejected(rejection))) => {
2743                assert_eq!(avdtp::ErrorCode::BadAcpSeid, rejection.error_code().unwrap().unwrap());
2744                assert_eq!(unknown_endpoint_id, rejection.stream_id().unwrap());
2745            }
2746            x => panic!("Start should be ready but got {:?}", x),
2747        };
2748
2749        // We can't get a permit (none are available) so we suspend the one we didn't error on.
2750        let mut remote_requests = remote_peer.take_request_stream();
2751
2752        let suspended_stream_ids = match exec.run_singlethreaded(&mut remote_requests.next()) {
2753            Some(Ok(avdtp::Request::Suspend { responder, stream_ids })) => {
2754                responder.send().unwrap();
2755                stream_ids
2756            }
2757            x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
2758        };
2759
2760        assert!(suspended_stream_ids.contains(&sbc_endpoint_id));
2761        assert_eq!(1, suspended_stream_ids.len());
2762
2763        // And we should have not tried to start a task.
2764        match exec.run_until_stalled(&mut next_task_fut) {
2765            Poll::Pending => {}
2766            x => panic!("Local task should not have been created at this point: {:?}", x),
2767        };
2768
2769        // No matter how many times they ask to start, we will still suspend (but not queue another
2770        // reservation for the same id)
2771        let mut start_fut = remote_peer.start(&[sbc_endpoint_id.clone()]);
2772        match exec.run_until_stalled(&mut start_fut) {
2773            Poll::Ready(Ok(())) => {}
2774            x => panic!("Start should be ready but got {:?}", x),
2775        }
2776
2777        let suspended_stream_ids = match exec.run_until_stalled(&mut remote_requests.next()) {
2778            Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
2779                responder.send().unwrap();
2780                stream_ids
2781            }
2782            x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
2783        };
2784        assert!(suspended_stream_ids.contains(&sbc_endpoint_id));
2785
2786        // After a permit is available, should try to start the first endpoint that failed.
2787        drop(taken_permit);
2788
2789        match exec.run_singlethreaded(&mut remote_requests.next()) {
2790            Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
2791                assert_eq!(stream_ids, &[sbc_endpoint_id.clone()]);
2792                responder.send().unwrap();
2793            }
2794            x => panic!("Expected start on permit available but got {x:?}"),
2795        };
2796
2797        // And we should start a task.
2798        let media_task = match exec.run_until_stalled(&mut next_task_fut) {
2799            Poll::Ready(Some(task)) => task,
2800            x => panic!("Local task should be created at this point: {:?}", x),
2801        };
2802
2803        assert!(media_task.is_started());
2804
2805        // If the remote asks to start another one, we still suspend it immediately.
2806        let mut start_fut = remote_peer.start(&[sbc_endpoint_two.clone()]);
2807        match exec.run_until_stalled(&mut start_fut) {
2808            Poll::Ready(Ok(())) => {}
2809            x => panic!("Start should be ready but got {:?}", x),
2810        }
2811
2812        let suspended_stream_ids = match exec.run_until_stalled(&mut remote_requests.next()) {
2813            Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
2814                responder.send().unwrap();
2815                stream_ids
2816            }
2817            x => panic!("Expected to receive a suspend request for the stream, got {:?}", x),
2818        };
2819
2820        assert!(suspended_stream_ids.contains(&sbc_endpoint_two));
2821        assert_eq!(1, suspended_stream_ids.len());
2822
2823        // Once the first one is done, the second can start.
2824        let mut suspend_fut = remote_peer.suspend(&[sbc_endpoint_id.clone()]);
2825        match exec.run_until_stalled(&mut suspend_fut) {
2826            Poll::Ready(Ok(())) => {}
2827            x => panic!("Start should be ready but got {:?}", x),
2828        }
2829
2830        match exec.run_singlethreaded(&mut remote_requests.next()) {
2831            Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
2832                assert_eq!(stream_ids, &[sbc_endpoint_two]);
2833                responder.send().unwrap();
2834            }
2835            x => panic!("Expected start on permit available but got {x:?}"),
2836        };
2837    }
2838
2839    fn start_sbc_stream(
2840        exec: &mut fasync::TestExecutor,
2841        media_test_builder: &mut TestMediaTaskBuilder,
2842        peer: &Peer,
2843        remote_peer: &avdtp::Peer,
2844        local_id: &StreamEndpointId,
2845        remote_id: &StreamEndpointId,
2846        transport_mode: Transport,
2847    ) -> TestMediaTask {
2848        let sbc_caps = sbc_capabilities();
2849        let set_config_fut = remote_peer.set_configuration(&local_id, &remote_id, &sbc_caps);
2850        let mut set_config_fut = pin!(set_config_fut);
2851
2852        match exec.run_until_stalled(&mut set_config_fut) {
2853            Poll::Ready(Ok(())) => {}
2854            x => panic!("Set capabilities should be ready but got {:?}", x),
2855        };
2856
2857        let open_fut = remote_peer.open(&local_id);
2858        let mut open_fut = pin!(open_fut);
2859        match exec.run_until_stalled(&mut open_fut) {
2860            Poll::Ready(Ok(())) => {}
2861            x => panic!("Open should be ready but got {:?}", x),
2862        };
2863
2864        // Establish a media transport stream
2865        let (transport, _remote_transport) = create_test_channels(transport_mode);
2866        assert_eq!(Some(()), peer.receive_channel(transport).ok());
2867
2868        // Remote peer should still be able to try to start the stream, and we will say yes.
2869        let stream_ids = [local_id.clone()];
2870        let start_fut = remote_peer.start(&stream_ids);
2871        let mut start_fut = pin!(start_fut);
2872        match exec.run_until_stalled(&mut start_fut) {
2873            Poll::Ready(Ok(())) => {}
2874            x => panic!("Start should be ready but got {:?}", x),
2875        };
2876
2877        // And we should start a media task.
2878        let media_task = media_test_builder.expect_task();
2879        assert!(media_task.is_started());
2880
2881        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
2882        media_task
2883    }
2884
2885    #[test_case(Transport::Socket ; "socket")]
2886    #[test_case(Transport::Fidl ; "fidl")]
2887    #[fuchsia::test]
2888    fn permits_can_be_revoked_and_reinstated_all(transport_mode: Transport) {
2889        let mut exec = fasync::TestExecutor::new();
2890
2891        let mut streams = Streams::default();
2892        let mut test_builder = TestMediaTaskBuilder::new();
2893        streams.insert(Stream::build(
2894            make_sbc_endpoint(1, avdtp::EndpointType::Sink),
2895            test_builder.builder(),
2896        ));
2897        let sbc_endpoint_id = 1_u8.try_into().unwrap();
2898        let remote_sbc_endpoint_id = 7_u8.try_into().unwrap();
2899
2900        streams.insert(Stream::build(
2901            make_sbc_endpoint(2, avdtp::EndpointType::Sink),
2902            test_builder.builder(),
2903        ));
2904        let sbc2_endpoint_id = 2_u8.try_into().unwrap();
2905        let remote_sbc2_endpoint_id = 6_u8.try_into().unwrap();
2906
2907        let permits = Permits::new(2);
2908
2909        let (remote, _requests, _, peer) =
2910            setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
2911        let remote_peer = avdtp::Peer::new(remote);
2912
2913        let one_media_task = start_sbc_stream(
2914            &mut exec,
2915            &mut test_builder,
2916            &peer,
2917            &remote_peer,
2918            &sbc_endpoint_id,
2919            &remote_sbc_endpoint_id,
2920            transport_mode,
2921        );
2922        let two_media_task = start_sbc_stream(
2923            &mut exec,
2924            &mut test_builder,
2925            &peer,
2926            &remote_peer,
2927            &sbc2_endpoint_id,
2928            &remote_sbc2_endpoint_id,
2929            transport_mode,
2930        );
2931
2932        // Someone comes along and revokes our permits.
2933        let taken_permits = permits.seize();
2934
2935        let remote_endpoints: HashSet<_> =
2936            [&remote_sbc_endpoint_id, &remote_sbc2_endpoint_id].iter().cloned().collect();
2937
2938        // We should send a suspend to the other end, for both of them.
2939        let mut remote_requests = remote_peer.take_request_stream();
2940        let mut expected_suspends = remote_endpoints.clone();
2941        while !expected_suspends.is_empty() {
2942            match exec.run_until_stalled(&mut remote_requests.next()) {
2943                Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
2944                    for stream_id in stream_ids {
2945                        assert!(expected_suspends.remove(&stream_id));
2946                    }
2947                    responder.send().expect("send response okay");
2948                }
2949                x => panic!("Expected suspension and got {:?}", x),
2950            }
2951        }
2952
2953        // And the media tasks should be stopped.
2954        assert!(!one_media_task.is_started());
2955        assert!(!two_media_task.is_started());
2956
2957        // After the permits are available again, we send a start, and start the media stream.
2958        drop(taken_permits);
2959
2960        let mut expected_starts = remote_endpoints.clone();
2961        while !expected_starts.is_empty() {
2962            match exec.run_singlethreaded(&mut remote_requests.next()) {
2963                Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
2964                    for stream_id in stream_ids {
2965                        assert!(expected_starts.remove(&stream_id));
2966                    }
2967                    responder.send().expect("send response okay");
2968                }
2969                x => panic!("Expected start and got {:?}", x),
2970            }
2971        }
2972        // And we should start two media tasks.
2973
2974        let one_media_task = test_builder.expect_task();
2975        assert!(one_media_task.is_started());
2976        let two_media_task = match exec.run_until_stalled(&mut test_builder.next_task()) {
2977            Poll::Ready(Some(task)) => task,
2978            x => panic!("Expected another ready task but {x:?}"),
2979        };
2980        assert!(two_media_task.is_started());
2981    }
2982
2983    #[test_case(Transport::Socket ; "socket")]
2984    #[test_case(Transport::Fidl ; "fidl")]
2985    #[fuchsia::test]
2986    fn permits_can_be_revoked_one_at_a_time(transport_mode: Transport) {
2987        let mut exec = fasync::TestExecutor::new();
2988
2989        let mut streams = Streams::default();
2990        let mut test_builder = TestMediaTaskBuilder::new();
2991        streams.insert(Stream::build(
2992            make_sbc_endpoint(1, avdtp::EndpointType::Sink),
2993            test_builder.builder(),
2994        ));
2995        let sbc_endpoint_id = 1_u8.try_into().unwrap();
2996        let remote_sbc_endpoint_id = 7_u8.try_into().unwrap();
2997
2998        streams.insert(Stream::build(
2999            make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3000            test_builder.builder(),
3001        ));
3002        let sbc2_endpoint_id = 2_u8.try_into().unwrap();
3003        let remote_sbc2_endpoint_id = 6_u8.try_into().unwrap();
3004
3005        let permits = Permits::new(2);
3006
3007        let (remote, _requests, _, peer) =
3008            setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3009        let remote_peer = avdtp::Peer::new(remote);
3010
3011        let one_media_task = start_sbc_stream(
3012            &mut exec,
3013            &mut test_builder,
3014            &peer,
3015            &remote_peer,
3016            &sbc_endpoint_id,
3017            &remote_sbc_endpoint_id,
3018            transport_mode,
3019        );
3020        let two_media_task = start_sbc_stream(
3021            &mut exec,
3022            &mut test_builder,
3023            &peer,
3024            &remote_peer,
3025            &sbc2_endpoint_id,
3026            &remote_sbc2_endpoint_id,
3027            transport_mode,
3028        );
3029
3030        // Someone comes along and revokes one of our permits.
3031        let taken_permit = permits.take();
3032
3033        let remote_endpoints: HashSet<_> =
3034            [&remote_sbc_endpoint_id, &remote_sbc2_endpoint_id].iter().cloned().collect();
3035
3036        // We should send a suspend to the other end, for both of them.
3037        let mut remote_requests = remote_peer.take_request_stream();
3038        let suspended_id = match exec.run_until_stalled(&mut remote_requests.next()) {
3039            Poll::Ready(Some(Ok(avdtp::Request::Suspend { responder, stream_ids }))) => {
3040                assert!(stream_ids.len() == 1);
3041                assert!(remote_endpoints.contains(&stream_ids[0]));
3042                responder.send().expect("send response okay");
3043                stream_ids[0].clone()
3044            }
3045            x => panic!("Expected suspension and got {:?}", x),
3046        };
3047
3048        // And the correct one of the media tasks should be stopped.
3049        if suspended_id == remote_sbc_endpoint_id {
3050            assert!(!one_media_task.is_started());
3051            assert!(two_media_task.is_started());
3052        } else {
3053            assert!(one_media_task.is_started());
3054            assert!(!two_media_task.is_started());
3055        }
3056
3057        // After the permits are available again, we send a start, and start the media stream.
3058        drop(taken_permit);
3059
3060        match exec.run_singlethreaded(&mut remote_requests.next()) {
3061            Some(Ok(avdtp::Request::Start { responder, stream_ids })) => {
3062                assert_eq!(stream_ids, &[suspended_id]);
3063                responder.send().expect("send response okay");
3064            }
3065            x => panic!("Expected start and got {:?}", x),
3066        }
3067        // And we should start another media task.
3068        let media_task = match exec.run_until_stalled(&mut test_builder.next_task()) {
3069            Poll::Ready(Some(task)) => task,
3070            x => panic!("Expected media task to start: {x:?}"),
3071        };
3072        assert!(media_task.is_started());
3073    }
3074
3075    // Scenario: when we are waiting for a suspend response from the peer after a permit was not
3076    // available, we try to start the peer (because a dwell has expired)
3077    #[test_case(Transport::Socket ; "socket")]
3078    #[test_case(Transport::Fidl ; "fidl")]
3079    #[fuchsia::test]
3080    fn permit_suspend_start_while_suspending(transport_mode: Transport) {
3081        let mut exec = fasync::TestExecutor::new();
3082
3083        let mut streams = Streams::default();
3084        let mut test_builder = TestMediaTaskBuilder::new();
3085        streams.insert(Stream::build(
3086            make_sbc_endpoint(1, avdtp::EndpointType::Sink),
3087            test_builder.builder(),
3088        ));
3089        streams.insert(Stream::build(
3090            make_sbc_endpoint(2, avdtp::EndpointType::Sink),
3091            test_builder.builder(),
3092        ));
3093        let mut next_task_fut = test_builder.next_task();
3094
3095        let permits = Permits::new(1);
3096        let (remote, _profile_request_stream, _, peer) =
3097            setup_test_peer(transport_mode, false, streams, Some(permits.clone()));
3098
3099        let remote_peer = avdtp::Peer::new(remote);
3100        let mut remote_requests = remote_peer.take_request_stream();
3101
3102        let sbc_endpoint_id = 1_u8.try_into().unwrap();
3103
3104        let sbc_caps = sbc_capabilities();
3105        let mut set_config_fut =
3106            remote_peer.set_configuration(&sbc_endpoint_id, &sbc_endpoint_id, &sbc_caps);
3107
3108        match exec.run_until_stalled(&mut set_config_fut) {
3109            Poll::Ready(Ok(())) => {}
3110            x => panic!("Set capabilities should be ready but got {:?}", x),
3111        };
3112
3113        let mut open_fut = remote_peer.open(&sbc_endpoint_id);
3114        match exec.run_until_stalled(&mut open_fut) {
3115            Poll::Ready(Ok(())) => {}
3116            x => panic!("Open should be ready but got {:?}", x),
3117        };
3118
3119        // Establish a media transport stream
3120        let (_remote_transport, transport) = Channel::create_socket_pair();
3121        assert_eq!(Some(()), peer.receive_channel(transport).ok());
3122
3123        // At this point, we are dwelling, waiting for the peer to start the stream. Skip the timer.
3124        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3125        let Some(_deadline) = exec.wake_next_timer() else {
3126            panic!("Expected a timer to be waiting to run");
3127        };
3128
3129        // We will try to start it ourselves, which will take the only permit and send a start.
3130        let start_responder = match exec.run_singlethreaded(&mut remote_requests.next()) {
3131            Some(Ok(avdtp::Request::Start { stream_ids, responder })) => {
3132                assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3133                responder
3134            }
3135            x => panic!("Expected a Start request, got {x:?}"),
3136        };
3137
3138        assert!(permits.get().is_none());
3139
3140        // The peer doesn't notice. Instead try to start it from the peer side (bad timing)
3141        let mut start_fut = remote_peer.start(&[sbc_endpoint_id.clone()]);
3142
3143        // We get an OK, and then immediately a suspend request because there are no
3144        // permits available.
3145        match exec.run_singlethreaded(&mut start_fut) {
3146            Ok(()) => {}
3147            x => panic!("Expected OK response from start future but got {x:?}"),
3148        }
3149
3150        let suspend_responder = match exec.run_singlethreaded(&mut remote_requests.next()) {
3151            Some(Ok(avdtp::Request::Suspend { stream_ids, responder })) => {
3152                assert_eq!(stream_ids, vec![sbc_endpoint_id.clone()]);
3153                responder
3154            }
3155            x => panic!("Expected a suspend got {x:?}"),
3156        };
3157
3158        // At this point, the peer notices the start request and responds.
3159        start_responder.send().unwrap();
3160
3161        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
3162
3163        // Okay I guess..
3164        suspend_responder.send().unwrap();
3165
3166        // And we should start a task.
3167        let media_task = match exec.run_until_stalled(&mut next_task_fut) {
3168            Poll::Ready(Some(task)) => task,
3169            x => panic!("Local task should be created at this point: {:?}", x),
3170        };
3171
3172        assert!(media_task.is_started());
3173    }
3174
3175    /// Test that the version check method correctly differentiates between newer
3176    /// and older A2DP versions.
3177    #[fuchsia::test]
3178    fn version_check() {
3179        let p1: ProfileDescriptor = ProfileDescriptor {
3180            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3181            major_version: Some(1),
3182            minor_version: Some(3),
3183            ..Default::default()
3184        };
3185        assert_eq!(true, a2dp_version_check(p1));
3186
3187        let p1: ProfileDescriptor = ProfileDescriptor {
3188            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3189            major_version: Some(2),
3190            minor_version: Some(10),
3191            ..Default::default()
3192        };
3193        assert_eq!(true, a2dp_version_check(p1));
3194
3195        let p1: ProfileDescriptor = ProfileDescriptor {
3196            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3197            major_version: Some(1),
3198            minor_version: Some(0),
3199            ..Default::default()
3200        };
3201        assert_eq!(false, a2dp_version_check(p1));
3202
3203        let p1: ProfileDescriptor = ProfileDescriptor {
3204            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3205            major_version: None,
3206            minor_version: Some(9),
3207            ..Default::default()
3208        };
3209        assert_eq!(false, a2dp_version_check(p1));
3210
3211        let p1: ProfileDescriptor = ProfileDescriptor {
3212            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
3213            major_version: Some(2),
3214            minor_version: Some(2),
3215            ..Default::default()
3216        };
3217        assert_eq!(true, a2dp_version_check(p1));
3218    }
3219}