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