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