1use 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
32struct PeerStats {
35 id: PeerId,
36 inspect_node: inspect::Node,
37 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 descriptors: HashMap<PeerId, (ProfileDescriptor, HashSet<avdtp::EndpointType>)>,
70 stats: HashMap<PeerId, PeerStats>,
72 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 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
122fn 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 None
133 }
134}
135
136async 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
166pub struct ConnectedPeers {
169 connected: DetachableMap<PeerId, Peer>,
171 connection_attempts: Mutex<HashMap<PeerId, fasync::Task<()>>>,
174 discovered: Mutex<DiscoveredPeers>,
176 streams_builder: StreamsBuilder,
178 permits: Permits,
180 profile: ProfileProxy,
182 metrics: bt_metrics::MetricsLogger,
184 avrcp: Option<avrcp::PeerManagerProxy>,
186 inspect: inspect::Node,
188 inspect_peer_direction: inspect::StringProperty,
190 connected_peer_senders: Mutex<Vec<mpsc::Sender<DetachableWeak<PeerId, Peer>>>>,
192 start_stream_tasks: Mutex<HashMap<PeerId, fasync::Task<()>>>,
195 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 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 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 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 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 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 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 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
444pub 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 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 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 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 const AAC_SEID: u8 = 8;
608 const SBC_SINK_SEID: u8 = 9;
610 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, )
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 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 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 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 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 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 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 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 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 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 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 #[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 peers.set_preferred_peer_direction(avdtp::EndpointType::Source);
895
896 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 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
916
917 let mut remote_requests = remote.take_request_stream();
918
919 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 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 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 #[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 peers.set_preferred_peer_direction(avdtp::EndpointType::Source);
986
987 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 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 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 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 let _ = exec.run_until_stalled(&mut futures::future::pending::<()>());
1083
1084 let mut remote_requests = remote.take_request_stream();
1085
1086 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 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 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 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 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 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 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 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 drop(remote);
1235
1236 run_to_stalled(&mut exec);
1237
1238 assert!(peers.get(&id).is_none());
1239
1240 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 assert!(peers.get(&id).is_some());
1248 }
1249}