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