Skip to main content

bt_a2dp/
connected_peers.rs

1// Copyright 2019 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::{Error, format_err};
6use bt_avdtp as avdtp;
7use fidl_fuchsia_bluetooth::ChannelParameters;
8use fidl_fuchsia_bluetooth_bredr::{self as bredr, ProfileDescriptor, ProfileProxy};
9use fuchsia_async as fasync;
10use fuchsia_bluetooth::detachable_map::{DetachableMap, DetachableWeak};
11use fuchsia_bluetooth::inspect::DebugExt;
12use fuchsia_bluetooth::types::{Channel, PeerId};
13use fuchsia_inspect::{self as inspect, NumericProperty, Property};
14use fuchsia_inspect_derive::{AttachError, Inspect};
15use fuchsia_sync::Mutex;
16use futures::channel::{mpsc, oneshot};
17use futures::stream::{Stream, StreamExt};
18use futures::task::{Context, Poll};
19use futures::{Future, FutureExt, TryFutureExt};
20use log::{info, warn};
21use std::collections::hash_map::Entry;
22use std::collections::{HashMap, HashSet};
23use std::pin::Pin;
24use std::sync::Arc;
25
26use crate::codec::CodecNegotiation;
27use crate::peer::Peer;
28use crate::permits::Permits;
29use crate::stream::StreamsBuilder;
30
31/// Statistics node for tracking various information about a peer that has been encountered.
32/// Typically used as an inspect tree node.
33struct PeerStats {
34    id: PeerId,
35    inspect_node: inspect::Node,
36    /// The number of times that this peer has been successfully connected to since discovery.
37    connection_count: inspect::UintProperty,
38}
39
40impl PeerStats {
41    fn new(id: PeerId) -> Self {
42        Self { id, inspect_node: Default::default(), connection_count: Default::default() }
43    }
44
45    fn set_descriptor(&mut self, descriptor: &ProfileDescriptor) {
46        self.inspect_node.record_string("descriptor", descriptor.debug());
47    }
48
49    fn record_connected(&mut self) {
50        let _ = self.connection_count.add(1);
51    }
52}
53
54impl Inspect for &mut PeerStats {
55    fn iattach(self, parent: &inspect::Node, name: impl AsRef<str>) -> Result<(), AttachError> {
56        self.inspect_node = parent.create_child(name.as_ref());
57        self.inspect_node.record_string("id", self.id.to_string());
58        self.connection_count = self.inspect_node.create_uint("connection_count", 0);
59        Ok(())
60    }
61}
62
63#[derive(Default)]
64struct DiscoveredPeers {
65    /// The peers that we have discovered, with their descriptors and potential preferred
66    /// endpoint directions. Because the same peer can be discovered multiple times, with
67    /// potentially different endpoints, we maintain a set of advertised directions.
68    descriptors: HashMap<PeerId, (ProfileDescriptor, HashSet<avdtp::EndpointType>)>,
69    /// Holds the child nodes which include the ids and profile descriptors for inspect.
70    stats: HashMap<PeerId, PeerStats>,
71    /// Inspect node, usually at "discovered" in the tree.
72    inspect_node: inspect::Node,
73}
74
75impl DiscoveredPeers {
76    fn insert(
77        &mut self,
78        id: PeerId,
79        descriptor: ProfileDescriptor,
80        directions: HashSet<avdtp::EndpointType>,
81    ) {
82        self.stats
83            .entry(id)
84            .or_insert_with(|| {
85                let mut new_stats = PeerStats::new(id);
86                let _ = new_stats.iattach(&self.inspect_node, inspect::unique_name("peer_"));
87                new_stats
88            })
89            .set_descriptor(&descriptor);
90
91        match self.descriptors.entry(id) {
92            Entry::Occupied(mut entry) => {
93                entry.get_mut().0 = descriptor;
94                entry.get_mut().1.extend(&directions);
95            }
96            Entry::Vacant(entry) => {
97                let _ = entry.insert((descriptor, directions));
98            }
99        };
100    }
101
102    fn connected(&mut self, id: PeerId) {
103        if let Some(stats) = self.stats.get_mut(&id) {
104            stats.record_connected();
105        }
106    }
107
108    /// Returns the descriptor and preferred endpoint direction associated with the peer `id`.
109    fn get(&self, id: &PeerId) -> Option<(ProfileDescriptor, Option<avdtp::EndpointType>)> {
110        self.descriptors.get(id).map(|(desc, dirs)| (desc.clone(), find_preferred_direction(dirs)))
111    }
112}
113
114impl Inspect for &mut DiscoveredPeers {
115    fn iattach(self, parent: &inspect::Node, name: impl AsRef<str>) -> Result<(), AttachError> {
116        self.inspect_node = parent.create_child(name.as_ref());
117        Ok(())
118    }
119}
120
121/// Given a set of endpoint `directions`, returns the preferred direction or None
122/// if both Sink and Source are specified.
123fn find_preferred_direction(
124    directions: &HashSet<avdtp::EndpointType>,
125) -> Option<avdtp::EndpointType> {
126    if directions.len() == 1 {
127        directions.iter().next().cloned()
128    } else {
129        // Otherwise, either there are no A2DP services or both Sink & Source are specified
130        // in which case there is no preferred direction.
131        None
132    }
133}
134
135/// Make an outgoing connection to a peer.
136async fn connect_peer(
137    proxy: ProfileProxy,
138    id: PeerId,
139    channel_params: ChannelParameters,
140) -> Result<Channel, Error> {
141    info!(id:%; "Connecting to peer");
142    let connect_fut = proxy.connect(
143        &id.into(),
144        &bredr::ConnectParameters::L2cap(bredr::L2capParameters {
145            psm: Some(bredr::PSM_AVDTP),
146            parameters: Some(channel_params),
147            ..Default::default()
148        }),
149    );
150    let channel = match connect_fut.await {
151        Err(e) => {
152            warn!(id:%, e:?; "FIDL error on connect");
153            return Err(e.into());
154        }
155        Ok(Err(e)) => return Err(format_err!("Bluetooth connect error: {e:?}")),
156        Ok(Ok(channel)) => channel,
157    };
158
159    let channel = channel
160        .try_into()
161        .map_err(|e| format_err!("Couldn't convert FIDL to BT channel: {e:?}"))?;
162    Ok(channel)
163}
164
165/// ConnectedPeers manages the set of connected peers based on discovery, new connection, and
166/// peer session lifetime.
167pub struct ConnectedPeers {
168    /// The set of connected peers.
169    connected: DetachableMap<PeerId, Peer>,
170    /// Tasks for peers that we are attempting to connect to.
171    /// Used to ensure only one outgoing attempt exists at once.
172    connection_attempts: Mutex<HashMap<PeerId, fasync::Task<()>>>,
173    /// ProfileDescriptors from discovering the peer, stored here even if the peer is disconnected
174    discovered: Mutex<DiscoveredPeers>,
175    /// Streams builder, provides a set of streams and negotiation when a peer is connected
176    streams_builder: StreamsBuilder,
177    /// The permits that each peer uses to validate that we can start a stream.
178    permits: Permits,
179    /// Profile Proxy, used to connect new transport sockets.
180    profile: ProfileProxy,
181    /// Cobalt logger to use and hand out to peers, if we are using one.
182    metrics: bt_metrics::MetricsLogger,
183    /// The 'peers' node of the inspect tree. All connected peers own a child node of this node.
184    inspect: inspect::Node,
185    /// Inspect node for which is the current preferred peer direction.
186    inspect_peer_direction: inspect::StringProperty,
187    /// Listeners for new connected peers
188    connected_peer_senders: Mutex<Vec<mpsc::Sender<DetachableWeak<PeerId, Peer>>>>,
189    /// Task handles for newly connected peer stream starts.
190    // TODO(https://fxbug.dev/42146917): Completed tasks aren't garbage-collected yet.
191    start_stream_tasks: Mutex<HashMap<PeerId, fasync::Task<()>>>,
192    /// Preferred direction for new peers.  This is the direction we prefer the peer's endpoint to
193    /// be, i.e. if we prefer Sink, locally we are Source.
194    preferred_peer_direction: Mutex<avdtp::EndpointType>,
195}
196
197impl ConnectedPeers {
198    pub fn new(
199        streams_builder: StreamsBuilder,
200        permits: Permits,
201        profile: ProfileProxy,
202        metrics: bt_metrics::MetricsLogger,
203    ) -> Self {
204        Self {
205            connected: DetachableMap::new(),
206            connection_attempts: Mutex::new(HashMap::new()),
207            discovered: Default::default(),
208            streams_builder,
209            profile,
210            permits,
211            inspect: inspect::Node::default(),
212            inspect_peer_direction: inspect::StringProperty::default(),
213            metrics,
214            connected_peer_senders: Default::default(),
215            start_stream_tasks: Default::default(),
216            preferred_peer_direction: Mutex::new(avdtp::EndpointType::Sink),
217        }
218    }
219
220    pub(crate) fn get_weak(&self, id: &PeerId) -> Option<DetachableWeak<PeerId, Peer>> {
221        self.connected.get(id)
222    }
223
224    pub(crate) fn get(&self, id: &PeerId) -> Option<Arc<Peer>> {
225        self.get_weak(id).and_then(|p| p.upgrade())
226    }
227
228    pub fn is_connected(&self, id: &PeerId) -> bool {
229        self.connected.contains_key(id)
230    }
231
232    /// Attempts to start streaming on `peer` by collecting the remote streaming endpoint
233    /// information, selecting a compatible peer using `negotiation` and starting the stream.
234    /// Does nothing and returns Ok(()) if the peer is already streaming or will start streaming
235    /// on it's own.
236    async fn start_streaming(
237        peer: &DetachableWeak<PeerId, Peer>,
238        negotiation: CodecNegotiation,
239    ) -> Result<(), anyhow::Error> {
240        let remote_streams = {
241            let strong = peer.upgrade().ok_or_else(|| format_err!("Disconnected"))?;
242            if strong.streaming_active() {
243                return Ok(());
244            }
245            strong.collect_capabilities()
246        }
247        .await?;
248
249        let (negotiated, remote_seid) = negotiation
250            .select(&remote_streams)
251            .ok_or_else(|| format_err!("No compatible stream found"))?;
252
253        let strong = peer.upgrade().ok_or_else(|| format_err!("Disconnected"))?;
254        if strong.streaming_active() {
255            let peer_id = peer.key();
256            info!(peer_id:%; "Not starting streaming, it's already started");
257            return Ok(());
258        }
259        strong.stream_start(remote_seid, negotiated).await.map_err(Into::into)
260    }
261
262    pub fn found(
263        &self,
264        id: PeerId,
265        desc: ProfileDescriptor,
266        preferred_directions: HashSet<avdtp::EndpointType>,
267    ) {
268        self.discovered.lock().insert(id, desc.clone(), preferred_directions);
269        if let Some(peer) = self.get(&id) {
270            let _ = peer.set_descriptor(desc);
271        }
272    }
273
274    pub fn set_preferred_peer_direction(&self, direction: avdtp::EndpointType) {
275        *self.preferred_peer_direction.lock() = direction;
276        self.inspect_peer_direction.set(&format!("{direction:?}"));
277    }
278
279    pub fn preferred_peer_direction(&self) -> avdtp::EndpointType {
280        *self.preferred_peer_direction.lock()
281    }
282
283    pub fn try_connect(
284        &self,
285        id: PeerId,
286        channel_params: ChannelParameters,
287    ) -> impl Future<Output = Result<Option<Channel>, Error>> {
288        let proxy = self.profile.clone();
289        let connected = self.is_connected(&id);
290        let (sender, recv) = oneshot::channel();
291        let recv =
292            recv.map_ok_or_else(|_e| Err(format_err!("Connection task canceled")), Into::into);
293        if connected {
294            if let Err(e) = sender.send(Ok(None)) {
295                warn!(id:%, e:?; "Failed to notify already-connected");
296            }
297            return recv;
298        }
299        let mut attempts = self.connection_attempts.lock();
300        if let Some(previous_connect_task) = attempts.remove(&id) {
301            // We are the only place that can poll the connect task, check if it finished.
302            if previous_connect_task.now_or_never().is_none() {
303                warn!(id:%; "Cancelling previous connect attempt");
304            }
305        }
306        let connect_task = fasync::Task::spawn(async move {
307            if let Err(e) = sender.send(connect_peer(proxy, id, channel_params).await.map(Some)) {
308                warn!(id:%, e:?; "Failed to send channel connect result");
309            }
310        });
311        let _ = attempts.insert(id, connect_task);
312        recv
313    }
314
315    /// Accept a channel that is connected to the peer `id`.
316    /// If `initiator_delay` is set, attempt to start a stream after the specified delay.
317    /// `initiator_delay` has no effect if the peer already has a control channel.
318    /// Returns a weak peer pointer (even if it was previously connected) if successful.
319    pub async fn connected(
320        &self,
321        id: PeerId,
322        channel: Channel,
323        initiator_delay: Option<zx::MonotonicDuration>,
324    ) -> Result<DetachableWeak<PeerId, Peer>, Error> {
325        if let Some(weak) = self.get_weak(&id) {
326            let peer =
327                weak.upgrade().ok_or_else(|| format_err!("Disconnected connecting transport"))?;
328            if let Err(e) = peer.receive_channel(channel) {
329                warn!(id:%, e:%; "failed to connect channel");
330                return Err(e.into());
331            }
332            return Ok(weak);
333        }
334
335        let entry = self.connected.lazy_entry(&id);
336
337        info!(id:%; "peer connected");
338        let audio_offload = channel.audio_offload();
339        let avdtp_peer = avdtp::Peer::new(channel);
340
341        let mut peer = Peer::create(
342            id,
343            avdtp_peer,
344            self.streams_builder.peer_streams(&id, audio_offload.clone()).await?,
345            Some(self.permits.clone()),
346            self.profile.clone(),
347            self.metrics.clone(),
348        );
349
350        self.discovered.lock().connected(id);
351
352        let peer_preferred_direction = if let Some((desc, dir)) = self.discovered.lock().get(&id) {
353            let _ = peer.set_descriptor(desc);
354            dir
355        } else {
356            None
357        };
358
359        if let Err(e) = peer.iattach(&self.inspect, inspect::unique_name("peer_")) {
360            warn!(id:%, e:?; "Couldn't attach inspect");
361        }
362
363        let closed_fut = peer.closed();
364        let peer = match entry.try_insert(peer) {
365            Err(_peer) => {
366                warn!(id:%; "Peer connected while we were setting up");
367                return self.get_weak(&id).ok_or_else(|| format_err!("Peer missing"));
368            }
369            Ok(weak_peer) => weak_peer,
370        };
371
372        if let Some(delay) = initiator_delay {
373            let peer = peer.clone();
374            let peer_id = peer.key().clone();
375
376            // Bias the codec negotiation with the peer's preferred direction that was discovered
377            // from the SDP service search.
378            let negotiation = self
379                .streams_builder
380                .negotiation(
381                    &id,
382                    audio_offload,
383                    peer_preferred_direction.unwrap_or_else(|| self.preferred_peer_direction()),
384                )
385                .await?;
386            let start_stream_task = fuchsia_async::Task::local(async move {
387                let delay_sec = delay.into_millis() as f64 / 1000.0;
388                info!(id:% = peer.key(); "dwelling {delay_sec}s for peer initiation");
389                fasync::Timer::new(fasync::MonotonicInstant::after(delay)).await;
390
391                if let Err(e) = ConnectedPeers::start_streaming(&peer, negotiation).await {
392                    info!(id:% = peer.key(), e:?; "Peer start streaming failed");
393                    peer.detach();
394                }
395            });
396            if self.start_stream_tasks.lock().insert(peer_id, start_stream_task).is_some() {
397                info!(peer_id:%; "Replacing previous start stream dwell");
398            }
399        }
400
401        // Remove the peer when we disconnect.
402        fasync::Task::local(async move {
403            closed_fut.await;
404            peer.detach();
405        })
406        .detach();
407
408        let peer = self.get_weak(&id).ok_or_else(|| format_err!("Peer missing"))?;
409        self.notify_connected(&peer);
410        Ok(peer)
411    }
412
413    /// Notify the listeners that a new peer has been connected to.
414    fn notify_connected(&self, peer: &DetachableWeak<PeerId, Peer>) {
415        let mut senders = self.connected_peer_senders.lock();
416        senders.retain_mut(|sender| sender.try_send(peer.clone()).is_ok());
417    }
418
419    /// Get a stream that produces peers that have been connected.
420    pub fn connected_stream(&self) -> PeerConnections {
421        let (sender, receiver) = mpsc::channel(0);
422        self.connected_peer_senders.lock().push(sender);
423        PeerConnections { stream: receiver }
424    }
425}
426
427impl Inspect for &mut ConnectedPeers {
428    fn iattach(self, parent: &inspect::Node, name: impl AsRef<str>) -> Result<(), AttachError> {
429        self.inspect = parent.create_child(name.as_ref());
430        let peer_dir_str = format!("{:?}", self.preferred_peer_direction());
431        self.inspect_peer_direction =
432            self.inspect.create_string("preferred_peer_direction", peer_dir_str);
433        self.streams_builder.iattach(&self.inspect, "streams_builder")?;
434        self.discovered.lock().iattach(&self.inspect, "discovered")
435    }
436}
437
438/// Provides a stream of peers that have been connected to. This stream produces an item whenever
439/// an A2DP peer has been connected.  It will produce None when no more peers will be connected.
440pub struct PeerConnections {
441    stream: mpsc::Receiver<DetachableWeak<PeerId, Peer>>,
442}
443
444impl Stream for PeerConnections {
445    type Item = DetachableWeak<PeerId, Peer>;
446
447    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
448        self.stream.poll_next_unpin(cx)
449    }
450}
451
452#[cfg(test)]
453mod tests {
454    use super::*;
455
456    use async_utils::PollExt;
457    use bt_avdtp::{Request, ServiceCapability};
458    use bt_channel_test_support::{Transport, create_test_channels};
459    use diagnostics_assertions::assert_data_tree;
460    use fidl::endpoints::create_proxy_and_stream;
461    use fidl_fuchsia_bluetooth_bredr::{
462        AudioOffloadExtProxy, ProfileMarker, ProfileRequestStream, ServiceClassProfileIdentifier,
463    };
464    use futures::future::BoxFuture;
465    use std::pin::pin;
466    use test_case::test_case;
467
468    use crate::codec::MediaCodecConfig;
469    use crate::media_task::{MediaTaskBuilder, MediaTaskError, MediaTaskRunner};
470    use crate::media_types::*;
471
472    fn run_to_stalled(exec: &mut fasync::TestExecutor) {
473        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
474    }
475
476    fn exercise_avdtp(exec: &mut fasync::TestExecutor, remote: Channel, peer: &Peer) {
477        let remote_avdtp = avdtp::Peer::new(remote);
478        let mut remote_requests = remote_avdtp.take_request_stream();
479
480        // Should be able to actually communicate via the peer.
481        let avdtp = peer.avdtp();
482        let discover_fut = avdtp.discover();
483
484        let mut discover_fut = pin!(discover_fut);
485
486        assert!(exec.run_until_stalled(&mut discover_fut).is_pending());
487
488        let responder = match exec.run_until_stalled(&mut remote_requests.next()) {
489            Poll::Ready(Some(Ok(Request::Discover { responder }))) => responder,
490            x => panic!("Expected a Ready Discovery request but got {:?}", x),
491        };
492
493        let endpoint_id = avdtp::StreamEndpointId::try_from(1).expect("endpointid creation");
494
495        let information = avdtp::StreamInformation::new(
496            endpoint_id,
497            false,
498            avdtp::MediaType::Audio,
499            avdtp::EndpointType::Source,
500        );
501
502        responder.send(&[information]).expect("Sending response should have worked");
503
504        let _stream_infos = match exec.run_until_stalled(&mut discover_fut) {
505            Poll::Ready(Ok(infos)) => infos,
506            x => panic!("Expected a Ready response but got {:?}", x),
507        };
508    }
509
510    fn setup_connected_peer_test()
511    -> (fasync::TestExecutor, PeerId, ConnectedPeers, ProfileRequestStream) {
512        let exec = fasync::TestExecutor::new();
513        let (proxy, stream) = create_proxy_and_stream::<ProfileMarker>();
514        let id = PeerId(1);
515
516        let peers = ConnectedPeers::new(
517            StreamsBuilder::default(),
518            Permits::new(1),
519            proxy,
520            bt_metrics::MetricsLogger::default(),
521        );
522
523        (exec, id, peers, stream)
524    }
525
526    #[test_case(Transport::Socket ; "socket")]
527    #[test_case(Transport::Fidl ; "fidl")]
528    #[fuchsia::test]
529    fn connect_creates_peer(transport_mode: Transport) {
530        let (mut exec, id, peers, _stream) = setup_connected_peer_test();
531
532        let (channel, remote) = create_test_channels(transport_mode);
533
534        let peer = exec
535            .run_singlethreaded(peers.connected(id, channel, None))
536            .expect("peer should connect");
537        let peer = peer.upgrade().expect("peer should be connected");
538
539        exercise_avdtp(&mut exec, remote, &peer);
540    }
541
542    #[test_case(Transport::Socket ; "socket")]
543    #[test_case(Transport::Fidl ; "fidl")]
544    #[fuchsia::test]
545    fn connect_notifies_streams(transport_mode: Transport) {
546        let (mut exec, id, peers, _stream) = setup_connected_peer_test();
547
548        let (channel, remote) = create_test_channels(transport_mode);
549
550        let mut peer_stream = peers.connected_stream();
551        let mut peer_stream_two = peers.connected_stream();
552
553        let peer = exec
554            .run_singlethreaded(peers.connected(id, channel, None))
555            .expect("peer should connect");
556        let peer = peer.upgrade().expect("peer should be connected");
557
558        // Peers should have been notified of the new peer
559        let weak = exec.run_singlethreaded(peer_stream.next()).expect("peer stream to produce");
560        assert_eq!(weak.key(), &id);
561        let weak = exec.run_singlethreaded(peer_stream_two.next()).expect("peer stream to produce");
562        assert_eq!(weak.key(), &id);
563
564        exercise_avdtp(&mut exec, remote, &peer);
565
566        // If you drop one stream, the other one should still produce.
567        drop(peer_stream);
568
569        let id2 = PeerId(2);
570        let (channel2, remote2) = create_test_channels(transport_mode);
571        let peer2 = exec
572            .run_singlethreaded(peers.connected(id2, channel2, None))
573            .expect("peer should connect");
574        let peer2 = peer2.upgrade().expect("peer two should be connected");
575
576        let weak = exec.run_singlethreaded(peer_stream_two.next()).expect("peer stream to produce");
577        assert_eq!(weak.key(), &id2);
578
579        exercise_avdtp(&mut exec, remote2, &peer2);
580    }
581
582    #[fuchsia::test]
583    fn find_preferred_direction_returns_correct_endpoints() {
584        let empty = HashSet::new();
585        assert_eq!(find_preferred_direction(&empty), None);
586
587        let sink_only = HashSet::from_iter(vec![avdtp::EndpointType::Sink].into_iter());
588        assert_eq!(find_preferred_direction(&sink_only), Some(avdtp::EndpointType::Sink));
589
590        let source_only = HashSet::from_iter(vec![avdtp::EndpointType::Source].into_iter());
591        assert_eq!(find_preferred_direction(&source_only), Some(avdtp::EndpointType::Source));
592
593        let both = HashSet::from_iter(
594            vec![avdtp::EndpointType::Sink, avdtp::EndpointType::Source].into_iter(),
595        );
596        assert_eq!(find_preferred_direction(&both), None);
597    }
598
599    // Expected chosen ID for the AAC stream endpoint.
600    const AAC_SEID: u8 = 8;
601    // Expected chosen ID for the SBC sink stream endpoint.
602    const SBC_SINK_SEID: u8 = 9;
603    // Expected chosen ID for the SBC source stream endpoint.
604    const SBC_SOURCE_SEID: u8 = 10;
605
606    fn aac_sink_codec() -> avdtp::ServiceCapability {
607        AacCodecInfo::new(
608            AacObjectType::MANDATORY_SNK,
609            AacSamplingFrequency::MANDATORY_SNK,
610            AacChannels::MANDATORY_SNK,
611            true,
612            0, // 0 = Unknown constant bitrate support (A2DP Sec. 4.5.2.4)
613        )
614        .unwrap()
615        .into()
616    }
617
618    fn sbc_sink_codec() -> avdtp::ServiceCapability {
619        SbcCodecInfo::new(
620            SbcSamplingFrequency::MANDATORY_SNK,
621            SbcChannelMode::MANDATORY_SNK,
622            SbcBlockCount::MANDATORY_SNK,
623            SbcSubBands::MANDATORY_SNK,
624            SbcAllocation::MANDATORY_SNK,
625            SbcCodecInfo::BITPOOL_MIN,
626            SbcCodecInfo::BITPOOL_MAX,
627        )
628        .unwrap()
629        .into()
630    }
631
632    fn sbc_source_codec() -> avdtp::ServiceCapability {
633        SbcCodecInfo::new(
634            SbcSamplingFrequency::FREQ48000HZ,
635            SbcChannelMode::JOINT_STEREO,
636            SbcBlockCount::MANDATORY_SRC,
637            SbcSubBands::MANDATORY_SRC,
638            SbcAllocation::MANDATORY_SRC,
639            SbcCodecInfo::BITPOOL_MIN,
640            SbcCodecInfo::BITPOOL_MAX,
641        )
642        .unwrap()
643        .into()
644    }
645
646    #[derive(Clone)]
647    struct FakeBuilder {
648        capability: avdtp::ServiceCapability,
649        direction: avdtp::EndpointType,
650    }
651
652    impl MediaTaskBuilder for FakeBuilder {
653        fn configure(
654            &self,
655            _peer_id: &PeerId,
656            codec_config: &MediaCodecConfig,
657        ) -> Result<Box<dyn MediaTaskRunner>, MediaTaskError> {
658            if self.capability.codec_type() == Some(codec_config.codec_type()) {
659                return Ok(Box::new(FakeRunner {}));
660            }
661            Err(MediaTaskError::Other(String::from("Unsupported configuring")))
662        }
663
664        fn direction(&self) -> bt_avdtp::EndpointType {
665            self.direction
666        }
667
668        fn supported_configs(
669            &self,
670            _peer_id: &PeerId,
671            _offload: Option<AudioOffloadExtProxy>,
672        ) -> BoxFuture<'static, Result<Vec<MediaCodecConfig>, MediaTaskError>> {
673            futures::future::ready(Ok(vec![(&self.capability).try_into().unwrap()])).boxed()
674        }
675    }
676
677    struct FakeRunner {}
678
679    impl MediaTaskRunner for FakeRunner {
680        fn start(
681            &mut self,
682            _stream: avdtp::MediaStream,
683            _offload: Option<AudioOffloadExtProxy>,
684        ) -> Result<Box<dyn crate::media_task::MediaTask>, MediaTaskError> {
685            Err(MediaTaskError::Other(String::from("unimplemented starting")))
686        }
687    }
688
689    /// Sets up a test in which we expect to select a stream and connect to a peer.
690    /// Returns the executor, connected peers (under test), request stream for profile interaction,
691    /// and an SBC and AAC Sink service capability.
692    fn setup_negotiation_test() -> (
693        fasync::TestExecutor,
694        ConnectedPeers,
695        ProfileRequestStream,
696        ServiceCapability,
697        ServiceCapability,
698    ) {
699        let exec = fasync::TestExecutor::new_with_fake_time();
700        exec.set_fake_time(fasync::MonotonicInstant::from_nanos(1_000_000));
701        let (proxy, stream) = create_proxy_and_stream::<ProfileMarker>();
702
703        let aac_sink_codec = aac_sink_codec();
704        let sbc_sink_codec = sbc_sink_codec();
705        let aac_sink_builder = FakeBuilder {
706            capability: aac_sink_codec.clone(),
707            direction: avdtp::EndpointType::Sink,
708        };
709        let sbc_sink_builder = FakeBuilder {
710            capability: sbc_sink_codec.clone(),
711            direction: avdtp::EndpointType::Sink,
712        };
713        let sbc_source_builder =
714            FakeBuilder { capability: sbc_source_codec(), direction: avdtp::EndpointType::Source };
715
716        let mut streams_builder = StreamsBuilder::default();
717        streams_builder.add_builder(aac_sink_builder);
718        streams_builder.add_builder(sbc_sink_builder);
719        streams_builder.add_builder(sbc_source_builder);
720
721        let peers = ConnectedPeers::new(
722            streams_builder,
723            Permits::new(1),
724            proxy,
725            bt_metrics::MetricsLogger::default(),
726        );
727
728        (exec, peers, stream, sbc_sink_codec, aac_sink_codec)
729    }
730
731    #[test_case(Transport::Socket ; "socket")]
732    #[test_case(Transport::Fidl ; "fidl")]
733    #[fuchsia::test]
734    fn streaming_start_with_streaming_peer_is_noop(transport_mode: Transport) {
735        let (mut exec, peers, _stream, sbc_codec, _aac_codec) = setup_negotiation_test();
736        let id = PeerId(1);
737        let (channel, remote) = create_test_channels(transport_mode);
738        let remote = avdtp::Peer::new(remote);
739
740        let delay = zx::MonotonicDuration::from_seconds(1);
741
742        let mut remote_requests = remote.take_request_stream();
743
744        // This starts the task in the background waiting.
745        let mut connected_fut = std::pin::pin!(peers.connected(id, channel, Some(delay)));
746        assert!(exec.run_until_stalled(&mut connected_fut).expect("ready").is_ok());
747        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
748
749        // Before the delay expires, the peer starts the stream.
750
751        let seid: avdtp::StreamEndpointId = SBC_SINK_SEID.try_into().expect("seid to be okay");
752        let config_caps = &[ServiceCapability::MediaTransport, sbc_codec];
753        let set_config_fut = remote.set_configuration(&seid, &seid, config_caps);
754        let mut set_config_fut = pin!(set_config_fut);
755        match exec.run_until_stalled(&mut set_config_fut) {
756            Poll::Ready(Ok(())) => {}
757            x => panic!("Expected set config to be ready and Ok, got {:?}", x),
758        };
759
760        // The remote peer doesn't need to actually open, Set Configuration is enough of a signal.
761        // wait for the delay to expire now.
762
763        exec.set_fake_time(
764            fasync::MonotonicInstant::after(delay) + zx::MonotonicDuration::from_micros(1),
765        );
766        let _ = exec.wake_expired_timers();
767
768        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
769
770        // Shouldn't start a discovery, since the stream is scheduled to start already.
771        assert!(exec.run_until_stalled(&mut remote_requests.next()).is_pending());
772    }
773
774    fn sbc_source_endpoint() -> (avdtp::StreamEndpointId, avdtp::StreamInformation) {
775        let remote_sbc_seid: avdtp::StreamEndpointId = 1u8.try_into().unwrap();
776        let info = avdtp::StreamInformation::new(
777            remote_sbc_seid.clone(),
778            false,
779            avdtp::MediaType::Audio,
780            avdtp::EndpointType::Source,
781        );
782        (remote_sbc_seid, info)
783    }
784
785    fn aac_source_endpoint() -> (avdtp::StreamEndpointId, avdtp::StreamInformation) {
786        let remote_aac_seid: avdtp::StreamEndpointId = 2u8.try_into().unwrap();
787        let info = avdtp::StreamInformation::new(
788            remote_aac_seid.clone(),
789            false,
790            avdtp::MediaType::Audio,
791            avdtp::EndpointType::Source,
792        );
793        (remote_aac_seid, info)
794    }
795
796    fn sbc_sink_endpoint() -> (avdtp::StreamEndpointId, avdtp::StreamInformation) {
797        let remote_sbc_seid: avdtp::StreamEndpointId = 3u8.try_into().unwrap();
798        let info = avdtp::StreamInformation::new(
799            remote_sbc_seid.clone(),
800            false,
801            avdtp::MediaType::Audio,
802            avdtp::EndpointType::Sink,
803        );
804        (remote_sbc_seid, info)
805    }
806
807    /// Expects an AVDTP Discovery request on the `requests` stream. Responds to
808    /// the request with the provided `response` endpoints.
809    fn expect_peer_discovery(
810        exec: &mut fasync::TestExecutor,
811        requests: &mut avdtp::RequestStream,
812        response: Vec<avdtp::StreamInformation>,
813    ) {
814        match exec.run_until_stalled(&mut requests.next()) {
815            Poll::Ready(Some(Ok(avdtp::Request::Discover { responder }))) => {
816                responder.send(&response).expect("response succeeds");
817            }
818            x => panic!("Expected a discovery request to be sent after delay, got {:?}", x),
819        };
820    }
821
822    #[test_case(Transport::Socket ; "socket")]
823    #[test_case(Transport::Fidl ; "fidl")]
824    #[fuchsia::test]
825    fn streaming_start_configure_while_discovery(transport_mode: Transport) {
826        let (mut exec, peers, _stream, sbc_codec, _aac_codec) = setup_negotiation_test();
827        let id = PeerId(1);
828        let (channel, remote) = create_test_channels(transport_mode);
829        let remote = avdtp::Peer::new(remote);
830
831        let delay = zx::MonotonicDuration::from_seconds(1);
832
833        let mut remote_requests = remote.take_request_stream();
834
835        // This starts the task in the background waiting.
836        let mut connected_fut = std::pin::pin!(peers.connected(id, channel, Some(delay)));
837        assert!(exec.run_until_stalled(&mut connected_fut).expect("ready").is_ok());
838        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
839
840        // The delay expires, and the discovery is start!
841        exec.set_fake_time(
842            fasync::MonotonicInstant::after(delay) + zx::MonotonicDuration::from_micros(1),
843        );
844        let _ = exec.wake_expired_timers();
845        expect_peer_discovery(
846            &mut exec,
847            &mut remote_requests,
848            vec![sbc_source_endpoint().1, aac_source_endpoint().1],
849        );
850
851        // The remote peer doesn't need to actually open, Set Configuration is enough of a signal.
852        let seid: avdtp::StreamEndpointId = SBC_SINK_SEID.try_into().expect("seid to be okay");
853        let config_caps = &[ServiceCapability::MediaTransport, sbc_codec.clone()];
854        let set_config_fut = remote.set_configuration(&seid, &seid, config_caps);
855        let mut set_config_fut = pin!(set_config_fut);
856        match exec.run_until_stalled(&mut set_config_fut) {
857            Poll::Ready(Ok(())) => {}
858            x => panic!("Expected set config to be ready and Ok, got {:?}", x),
859        };
860
861        // Can finish the collection process, but not attempt to configure or start a stream.
862        loop {
863            match exec.run_until_stalled(&mut remote_requests.next()) {
864                Poll::Ready(Some(Ok(avdtp::Request::GetCapabilities { responder, .. }))) => {
865                    responder
866                        .send(&[avdtp::ServiceCapability::MediaTransport, sbc_codec.clone()])
867                        .expect("respond succeeds");
868                }
869                Poll::Ready(x) => panic!("Got unexpected request: {:?}", x),
870                Poll::Pending => break,
871            }
872        }
873    }
874
875    /// Tests connection initiation selects the appropriate stream endpoint based
876    /// on a biased codec negotiation that is set from the peer's discovered services.
877    #[test_case(Transport::Socket ; "socket")]
878    #[test_case(Transport::Fidl ; "fidl")]
879    #[fuchsia::test]
880    fn connect_initiation_uses_biased_codec_negotiation_by_peer(transport_mode: Transport) {
881        let (mut exec, peers, _stream, sbc_codec, _aac_codec) = setup_negotiation_test();
882        let id = PeerId(1);
883        let (channel, remote) = create_test_channels(transport_mode);
884
885        // System biases towards the Source direction (called when the AudioMode FIDL changes).
886        peers.set_preferred_peer_direction(avdtp::EndpointType::Source);
887
888        // New fake peer discovered with some descriptor - the peer's SDP entry shows Sink.
889        let remote = avdtp::Peer::new(remote);
890        let desc = ProfileDescriptor {
891            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
892            major_version: Some(1),
893            minor_version: Some(2),
894            ..Default::default()
895        };
896        let preferred_direction = vec![avdtp::EndpointType::Sink];
897        let delay = zx::MonotonicDuration::from_seconds(1);
898        peers.found(id, desc, HashSet::from_iter(preferred_direction.into_iter()));
899
900        let connected_fut = peers.connected(id, channel, Some(delay));
901        let mut connected_fut = std::pin::pin!(connected_fut);
902        let _ = exec
903            .run_until_stalled(&mut connected_fut)
904            .expect("is ready")
905            .expect("connect control channel is ok");
906        // run the start task until it's stalled.
907        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
908
909        let mut remote_requests = remote.take_request_stream();
910
911        // Should wait for the specified amount of time.
912        assert!(exec.run_until_stalled(&mut remote_requests.next()).is_pending());
913
914        exec.set_fake_time(fasync::MonotonicInstant::after(
915            delay + zx::MonotonicDuration::from_micros(1),
916        ));
917        let _ = exec.wake_expired_timers();
918
919        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
920        // Even though the peer supports both SBC Sink and Source, we expect to negotiate and start
921        // on the Sink endpoint since that is the peer's preferred one.
922        let (peer_sbc_source_seid, peer_sbc_source_endpoint) = sbc_source_endpoint();
923        let (peer_sbc_sink_seid, peer_sbc_sink_endpoint) = sbc_sink_endpoint();
924        expect_peer_discovery(
925            &mut exec,
926            &mut remote_requests,
927            vec![peer_sbc_source_endpoint, peer_sbc_sink_endpoint],
928        );
929        for _twice in 1..=2 {
930            match exec.run_until_stalled(&mut remote_requests.next()) {
931                Poll::Ready(Some(Ok(avdtp::Request::GetCapabilities { stream_id, responder }))) => {
932                    let codec = match stream_id {
933                        id if id == peer_sbc_source_seid => sbc_codec.clone(),
934                        id if id == peer_sbc_sink_seid => sbc_codec.clone(),
935                        x => panic!("Got unexpected get_capabilities seid {:?}", x),
936                    };
937                    responder
938                        .send(&[avdtp::ServiceCapability::MediaTransport, codec])
939                        .expect("respond succeeds");
940                }
941                x => panic!("Expected a ready get capabilities request, got {:?}", x),
942            };
943        }
944
945        match exec.run_until_stalled(&mut remote_requests.next()) {
946            Poll::Ready(Some(Ok(avdtp::Request::SetConfiguration {
947                local_stream_id,
948                remote_stream_id,
949                capabilities: _,
950                responder,
951            }))) => {
952                // We expect the set configuration to apply to the remote peer's Sink SEID and the
953                // channel Source SEID.
954                assert_eq!(peer_sbc_sink_seid, local_stream_id);
955                let local_sbc_source_seid: avdtp::StreamEndpointId =
956                    SBC_SOURCE_SEID.try_into().unwrap();
957                assert_eq!(local_sbc_source_seid, remote_stream_id);
958                responder.send().expect("response sends");
959            }
960            x => panic!("Expected a ready set configuration request, got {:?}", x),
961        };
962    }
963
964    /// Tests connection initiation selects the appropriate stream endpoint based
965    /// on a biased codec negotiation that is set from by the system (in practice, the AudioMode
966    /// FIDL). This case typically occurs when a peer advertises both sink and source, and therefore
967    /// has no preference for the endpoint direction.
968    #[test_case(Transport::Socket ; "socket")]
969    #[test_case(Transport::Fidl ; "fidl")]
970    #[fuchsia::test]
971    fn connect_initiation_uses_biased_codec_negotiation_by_system(transport_mode: Transport) {
972        let (mut exec, peers, _stream, sbc_codec, _aac_codec) = setup_negotiation_test();
973        let id = PeerId(1);
974        let (channel, remote) = create_test_channels(transport_mode);
975
976        // System biases towards the Source direction (called when the AudioMode FIDL changes).
977        peers.set_preferred_peer_direction(avdtp::EndpointType::Source);
978
979        // New fake peer discovered with separate Sink and Source entries.
980        let remote = avdtp::Peer::new(remote);
981        let desc = ProfileDescriptor {
982            profile_id: Some(ServiceClassProfileIdentifier::AdvancedAudioDistribution),
983            major_version: Some(1),
984            minor_version: Some(2),
985            ..Default::default()
986        };
987        peers.found(
988            id,
989            desc.clone(),
990            HashSet::from_iter(vec![avdtp::EndpointType::Source].into_iter()),
991        );
992        peers.found(id, desc, HashSet::from_iter(vec![avdtp::EndpointType::Sink].into_iter()));
993
994        let delay = zx::MonotonicDuration::from_seconds(1);
995        let connect_fut = peers.connected(id, channel, Some(delay));
996        let mut connect_fut = std::pin::pin!(connect_fut);
997        let _ = exec
998            .run_until_stalled(&mut connect_fut)
999            .expect("ready")
1000            .expect("connect control channel is ok");
1001        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
1002
1003        let mut remote_requests = remote.take_request_stream();
1004        // Should wait for the specified amount of time.
1005        assert!(exec.run_until_stalled(&mut remote_requests.next()).is_pending());
1006        exec.set_fake_time(fasync::MonotonicInstant::after(
1007            delay + zx::MonotonicDuration::from_micros(1),
1008        ));
1009        let _ = exec.wake_expired_timers();
1010        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
1011
1012        // Because the peer advertises both Sink and Source, we fall back to the system-biased
1013        // direction, which is Source for the peer.
1014        let (peer_sbc_source_seid, peer_sbc_source_endpoint) = sbc_source_endpoint();
1015        let (peer_sbc_sink_seid, peer_sbc_sink_endpoint) = sbc_sink_endpoint();
1016        expect_peer_discovery(
1017            &mut exec,
1018            &mut remote_requests,
1019            vec![peer_sbc_source_endpoint, peer_sbc_sink_endpoint],
1020        );
1021        for _twice in 1..=2 {
1022            match exec.run_until_stalled(&mut remote_requests.next()) {
1023                Poll::Ready(Some(Ok(avdtp::Request::GetCapabilities { stream_id, responder }))) => {
1024                    let codec = match stream_id {
1025                        id if id == peer_sbc_source_seid => sbc_codec.clone(),
1026                        id if id == peer_sbc_sink_seid => sbc_codec.clone(),
1027                        x => panic!("Got unexpected get_capabilities seid {:?}", x),
1028                    };
1029                    responder
1030                        .send(&[avdtp::ServiceCapability::MediaTransport, codec])
1031                        .expect("respond succeeds");
1032                }
1033                x => panic!("Expected a ready get capabilities request, got {:?}", x),
1034            };
1035        }
1036
1037        match exec.run_until_stalled(&mut remote_requests.next()) {
1038            Poll::Ready(Some(Ok(avdtp::Request::SetConfiguration {
1039                local_stream_id,
1040                remote_stream_id,
1041                capabilities: _,
1042                responder,
1043            }))) => {
1044                // We expect the set configuration to apply to the remote peer's Source SEID and the
1045                // channel Sink SEID.
1046                assert_eq!(peer_sbc_source_seid, local_stream_id);
1047                let local_sbc_sink_seid: avdtp::StreamEndpointId =
1048                    SBC_SINK_SEID.try_into().unwrap();
1049                assert_eq!(local_sbc_sink_seid, remote_stream_id);
1050                responder.send().expect("response sends");
1051            }
1052            x => panic!("Expected a ready set configuration request, got {:?}", x),
1053        };
1054    }
1055
1056    #[test_case(Transport::Socket ; "socket")]
1057    #[test_case(Transport::Fidl ; "fidl")]
1058    #[fuchsia::test]
1059    fn connect_initiation_uses_negotiation(transport_mode: Transport) {
1060        let (mut exec, peers, _stream, sbc_codec, aac_codec) = setup_negotiation_test();
1061        let id = PeerId(1);
1062        let (channel, remote) = create_test_channels(transport_mode);
1063        let remote = avdtp::Peer::new(remote);
1064
1065        let delay = zx::MonotonicDuration::from_seconds(1);
1066
1067        let mut connect_fut = std::pin::pin!(peers.connected(id, channel, Some(delay)));
1068        let _ = exec
1069            .run_until_stalled(&mut connect_fut)
1070            .expect("ready")
1071            .expect("connect control channel is ok");
1072
1073        // run the start task until it's stalled.
1074        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
1075
1076        let mut remote_requests = remote.take_request_stream();
1077
1078        // Should wait for the specified amount of time.
1079        assert!(exec.run_until_stalled(&mut remote_requests.next()).is_pending());
1080
1081        exec.set_fake_time(fasync::MonotonicInstant::after(
1082            delay + zx::MonotonicDuration::from_micros(1),
1083        ));
1084        let _ = exec.wake_expired_timers();
1085
1086        let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
1087
1088        // Should discover remote streams, negotiate, and start.
1089        let (peer_sbc_seid, peer_sbc_endpoint) = sbc_source_endpoint();
1090        let (peer_aac_seid, peer_aac_endpoint) = aac_source_endpoint();
1091        expect_peer_discovery(
1092            &mut exec,
1093            &mut remote_requests,
1094            vec![peer_sbc_endpoint, peer_aac_endpoint],
1095        );
1096        for _twice in 1..=2 {
1097            match exec.run_until_stalled(&mut remote_requests.next()) {
1098                Poll::Ready(Some(Ok(avdtp::Request::GetCapabilities { stream_id, responder }))) => {
1099                    let codec = match stream_id {
1100                        id if id == peer_sbc_seid => sbc_codec.clone(),
1101                        id if id == peer_aac_seid => aac_codec.clone(),
1102                        x => panic!("Got unexpected get_capabilities seid {:?}", x),
1103                    };
1104                    responder
1105                        .send(&[avdtp::ServiceCapability::MediaTransport, codec])
1106                        .expect("respond succeeds");
1107                }
1108                x => panic!("Expected a ready get capabilities request, got {:?}", x),
1109            };
1110        }
1111
1112        match exec.run_until_stalled(&mut remote_requests.next()) {
1113            Poll::Ready(Some(Ok(avdtp::Request::SetConfiguration {
1114                local_stream_id,
1115                remote_stream_id,
1116                capabilities: _,
1117                responder,
1118            }))) => {
1119                // Should set the aac stream, matched with channel AAC seid.
1120                assert_eq!(peer_aac_seid, local_stream_id);
1121                let local_aac_seid: avdtp::StreamEndpointId = AAC_SEID.try_into().unwrap();
1122                assert_eq!(local_aac_seid, remote_stream_id);
1123                responder.send().expect("response sends");
1124            }
1125            x => panic!("Expected a ready set configuration request, got {:?}", x),
1126        };
1127    }
1128
1129    #[test_case(Transport::Socket ; "socket")]
1130    #[test_case(Transport::Fidl ; "fidl")]
1131    #[fuchsia::test]
1132    fn connected_peers_inspect(transport_mode: Transport) {
1133        let (mut exec, id, mut peers, _stream) = setup_connected_peer_test();
1134
1135        let inspect = inspect::Inspector::default();
1136        peers.iattach(inspect.root(), "peers").expect("should attach to inspect tree");
1137
1138        assert_data_tree!(@executor exec, inspect, root: {
1139            peers: { streams_builder: contains {}, discovered: contains {}, preferred_peer_direction: "Sink" }});
1140
1141        peers.set_preferred_peer_direction(avdtp::EndpointType::Source);
1142
1143        assert_data_tree!(@executor exec, inspect, root: {
1144            peers: { streams_builder: contains {}, discovered: contains {}, preferred_peer_direction: "Source" }});
1145
1146        // Connect a peer, it should show up in the tree.
1147        let (channel, _remote) = create_test_channels(transport_mode);
1148        assert!(exec.run_singlethreaded(peers.connected(id, channel, None)).is_ok());
1149
1150        assert_data_tree!(@executor exec, inspect, root: {
1151            peers: {
1152                discovered: contains {},
1153                preferred_peer_direction: "Source",
1154                streams_builder: contains {},
1155                peer_0: { id: "0000000000000001", local_streams: contains {} }
1156            }
1157        });
1158    }
1159
1160    #[test_case(Transport::Socket ; "socket")]
1161    #[test_case(Transport::Fidl ; "fidl")]
1162    #[fuchsia::test]
1163    fn try_connect_cancels_previous_attempt(transport_mode: Transport) {
1164        let (mut exec, id, peers, mut profile_stream) = setup_connected_peer_test();
1165
1166        let mut connect_fut = peers.try_connect(id, ChannelParameters::default());
1167
1168        // Should get a request to connect, which we will stall and not respond to.
1169        let responder = match exec.run_singlethreaded(profile_stream.next()) {
1170            Some(Ok(bredr::ProfileRequest::Connect { responder, .. })) => responder,
1171            x => panic!("Expected Profile connect, got {x:?}"),
1172        };
1173
1174        // Trying to connect again should cancel the first try, and send another connect.
1175        let mut connect_again_fut = peers.try_connect(id, ChannelParameters::default());
1176        let responder_two = match exec.run_singlethreaded(profile_stream.next()) {
1177            Some(Ok(bredr::ProfileRequest::Connect { responder, .. })) => responder,
1178            x => panic!("Expected Profile connect, got {x:?}"),
1179        };
1180
1181        let first_result = exec.run_singlethreaded(&mut connect_fut);
1182        let _ = first_result.expect_err("Should have an error from first attempt");
1183
1184        // Responding on the first connect shouldn't do anything at this point.
1185        responder.send(Err(fidl_fuchsia_bluetooth::ErrorCode::Failed)).unwrap();
1186
1187        exec.run_until_stalled(&mut connect_again_fut).expect_pending("shouldn't finish");
1188
1189        let (local, _remote) = create_test_channels(transport_mode);
1190        responder_two.send(Ok(local.try_into().unwrap())).unwrap();
1191
1192        let second_result = exec.run_singlethreaded(&mut connect_again_fut);
1193        let _ = second_result.expect("should receive the channel");
1194    }
1195
1196    #[test_case(Transport::Socket ; "socket")]
1197    #[test_case(Transport::Fidl ; "fidl")]
1198    #[fuchsia::test]
1199    fn connected_peers_peer_disconnect_removes_peer(transport_mode: Transport) {
1200        let (mut exec, id, peers, _stream) = setup_connected_peer_test();
1201
1202        let (channel, remote) = create_test_channels(transport_mode);
1203
1204        assert!(exec.run_singlethreaded(peers.connected(id, channel, None)).is_ok());
1205        run_to_stalled(&mut exec);
1206
1207        // Disconnect the signaling channel, peer should be gone.
1208        drop(remote);
1209
1210        run_to_stalled(&mut exec);
1211
1212        assert!(peers.get(&id).is_none());
1213    }
1214
1215    #[test_case(Transport::Socket ; "socket")]
1216    #[test_case(Transport::Fidl ; "fidl")]
1217    #[fuchsia::test]
1218    fn connected_peers_reconnect_works(transport_mode: Transport) {
1219        let (mut exec, id, peers, _stream) = setup_connected_peer_test();
1220
1221        let (channel, remote) = create_test_channels(transport_mode);
1222        assert!(exec.run_singlethreaded(peers.connected(id, channel, None)).is_ok());
1223        run_to_stalled(&mut exec);
1224
1225        // Disconnect the signaling channel, peer should be gone.
1226        drop(remote);
1227
1228        run_to_stalled(&mut exec);
1229
1230        assert!(peers.get(&id).is_none());
1231
1232        // Connect another peer with the same ID
1233        let (channel, _remote) = create_test_channels(transport_mode);
1234
1235        assert!(exec.run_singlethreaded(peers.connected(id, channel, None)).is_ok());
1236        run_to_stalled(&mut exec);
1237
1238        // Should be connected.
1239        assert!(peers.get(&id).is_some());
1240    }
1241}