1use anyhow::{Error, format_err};
6use core::fmt;
7use core::pin::Pin;
8use core::task::{Context, Poll, Waker};
9use fidl::endpoints::Proxy;
10use fidl_fuchsia_bluetooth_avrcp::{
11 ControllerEvent, ControllerEventStream, ControllerProxy, Notifications, PeerManagerProxy,
12};
13use fuchsia_bluetooth::types::PeerId;
14use futures::stream::{FusedStream, Stream, StreamExt};
15use std::collections::VecDeque;
16
17pub struct AbsoluteVolumeControl {
19 id: PeerId,
20 peer_manager: PeerManagerProxy,
21 proxy: ControllerProxy,
22 stream: ControllerEventStream,
23 set_volume_results: VecDeque<Result<u8, Error>>,
24 waker: Option<Waker>,
25 is_terminated: bool,
26}
27
28impl fmt::Debug for AbsoluteVolumeControl {
29 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
30 f.debug_struct("AbsoluteVolumeControl")
31 .field("id", &self.id)
32 .field("peer_manager", &self.peer_manager)
33 .field("proxy", &self.proxy)
34 .field("is_terminated", &self.is_terminated)
35 .finish_non_exhaustive()
36 }
37}
38
39impl AbsoluteVolumeControl {
40 pub async fn connect(id: PeerId, peer_manager: PeerManagerProxy) -> Result<Self, Error> {
44 let (proxy, server_end) =
45 fidl::endpoints::create_proxy::<fidl_fuchsia_bluetooth_avrcp::ControllerMarker>();
46 peer_manager
47 .get_controller_for_target(&id.into(), server_end)
48 .await?
49 .map_err(|e| format_err!("failed to get controller: {:?}", e))?;
50 proxy.set_notification_filter(Notifications::VOLUME, 0)?;
51 let stream = proxy.take_event_stream();
52 Ok(Self {
53 id,
54 peer_manager,
55 proxy,
56 stream,
57 set_volume_results: VecDeque::new(),
58 waker: None,
59 is_terminated: false,
60 })
61 }
62
63 pub fn peer_id(&self) -> PeerId {
64 self.id
65 }
66
67 pub async fn set_absolute_volume(&mut self, volume_percentage: u8) -> Result<u8, Error> {
74 if volume_percentage > 100 {
75 return Err(format_err!(
76 "volume_percentage ({volume_percentage:}) must be between 0 and 100",
77 ));
78 }
79
80 let avrcp_volume = (volume_percentage as u16 * 127 / 100) as u8;
81 let result = self.proxy.set_absolute_volume(avrcp_volume).await?;
82
83 self.set_volume_results.push_back(result.clone().map_err(|e| format_err!("{e:?}")));
85
86 if let Some(waker) = self.waker.take() {
87 waker.wake();
88 }
89
90 result.map_err(|e| format_err!("{e:?}"))
91 }
92}
93
94impl Stream for AbsoluteVolumeControl {
95 type Item = Result<u8, Error>;
96
97 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
98 loop {
99 if self.is_terminated {
100 return Poll::Ready(None);
101 }
102
103 if self.proxy.is_closed() {
104 self.is_terminated = true;
105 return Poll::Ready(None);
106 }
107
108 if let Some(result) = self.set_volume_results.pop_front() {
110 return Poll::Ready(Some(result));
111 }
112
113 match self.stream.poll_next_unpin(cx) {
115 Poll::Ready(Some(Ok(ControllerEvent::OnNotification { notification, .. }))) => {
116 let _ = self.proxy.notify_notification_handled();
118
119 if let Some(volume) = notification.volume {
120 return Poll::Ready(Some(Ok(volume)));
121 }
122 }
124 Poll::Ready(Some(Err(e))) => {
125 self.is_terminated = true;
126 return Poll::Ready(Some(Err(e.into())));
127 }
128 Poll::Ready(None) => {
129 self.is_terminated = true;
130 return Poll::Ready(None);
131 }
132 Poll::Pending => {
133 self.waker = Some(cx.waker().clone());
136 return Poll::Pending;
137 }
138 }
139 }
140 }
141}
142
143impl FusedStream for AbsoluteVolumeControl {
144 fn is_terminated(&self) -> bool {
145 self.is_terminated
146 }
147}
148
149#[cfg(test)]
150mod tests {
151 use super::*;
152 use assert_matches::assert_matches;
153 use async_utils::PollExt;
154 use fidl::endpoints::{RequestStream, create_proxy_and_stream};
155 use fidl_fuchsia_bluetooth_avrcp::{
156 ControllerRequest, ControllerRequestStream, PeerManagerMarker, PeerManagerRequest,
157 PeerManagerRequestStream,
158 };
159 use fuchsia_async as fasync;
160 use futures::{StreamExt, pin_mut};
161
162 const PEER_ID: PeerId = PeerId(1);
163
164 #[track_caller]
165 fn set_up(
166 exec: &mut fasync::TestExecutor,
167 ) -> (AbsoluteVolumeControl, PeerManagerRequestStream, ControllerRequestStream) {
168 let (peer_manager_proxy, mut peer_manager_stream) =
169 create_proxy_and_stream::<PeerManagerMarker>();
170
171 let connect_fut = AbsoluteVolumeControl::connect(PEER_ID, peer_manager_proxy);
172 pin_mut!(connect_fut);
173
174 exec.run_until_stalled(&mut connect_fut).expect_pending("should be pending");
175
176 let mut request_fut = peer_manager_stream.next();
178 let request = exec.run_until_stalled(&mut request_fut).expect("stream item");
179 let mut controller_stream = match request {
180 Some(Ok(PeerManagerRequest::GetControllerForTarget { peer_id, client, responder })) => {
181 assert_eq!(peer_id, PEER_ID.into());
182 responder.send(Ok(())).expect("should succeed");
183 client.into_stream()
184 }
185 r => panic!("Expected GetControllerForTarget request, got {:?}", r),
186 };
187
188 let avc =
189 exec.run_until_stalled(&mut connect_fut).expect("should have succeeded").expect("ok");
190
191 let mut request_fut = controller_stream.next();
192 let request = exec.run_until_stalled(&mut request_fut).expect("stream item");
193 match request {
194 Some(Ok(ControllerRequest::SetNotificationFilter { notifications, .. })) => {
195 assert_eq!(notifications, Notifications::VOLUME);
196 }
197 r => panic!("Expected SetNotificationFilter request, got {:?}", r),
198 }
199 assert_eq!(avc.peer_id(), PEER_ID);
200
201 (avc, peer_manager_stream, controller_stream)
202 }
203
204 #[fuchsia::test]
205 fn test_set_absolute_volume() {
206 let mut exec = fasync::TestExecutor::new();
207 let (mut avc, _peer_manager_stream, mut controller_stream) = set_up(&mut exec);
208
209 let requested_vol_percent = 50;
210 let expected_avrcp_volume = 63; let expected_response_volume = 64;
212
213 exec.run_until_stalled(&mut avc.next()).expect_pending("pending");
215
216 {
217 let set_vol_fut = avc.set_absolute_volume(requested_vol_percent);
218 pin_mut!(set_vol_fut);
219 exec.run_until_stalled(&mut set_vol_fut).expect_pending("pending");
220
221 let request_fut = controller_stream.next();
223 pin_mut!(request_fut);
224 let request = exec.run_until_stalled(&mut request_fut).expect("stream item");
225 match request {
226 Some(Ok(ControllerRequest::SetAbsoluteVolume { requested_volume, responder })) => {
227 assert_eq!(requested_volume, expected_avrcp_volume);
228 responder.send(Ok(expected_response_volume)).unwrap();
229 }
230 r => panic!("Expected SetAbsoluteVolume request, got {:?}", r),
231 }
232
233 let result = exec.run_until_stalled(&mut set_vol_fut).expect("future should complete");
235 assert_eq!(result.unwrap(), expected_response_volume);
236 }
237
238 let stream_result =
239 exec.run_until_stalled(&mut avc.next()).expect("stream should have item");
240 assert_eq!(stream_result.unwrap().unwrap(), expected_response_volume);
241 }
242
243 #[fuchsia::test]
244 fn test_set_absolute_volume_out_of_range() {
245 let mut exec = fasync::TestExecutor::new();
246 let (mut avc, _peer_manager_stream, _controller_stream) = set_up(&mut exec);
247
248 let requested_vol_percent = 101;
249
250 let set_vol_fut = avc.set_absolute_volume(requested_vol_percent);
251 pin_mut!(set_vol_fut);
252 let result = exec.run_until_stalled(&mut set_vol_fut).expect("future should complete");
253 assert!(result.is_err());
254 }
255
256 #[fuchsia::test]
257 fn test_volume_notification_stream() {
258 let mut exec = fasync::TestExecutor::new();
259 let (mut avc, _peer_manager_stream, mut controller_stream) = set_up(&mut exec);
260
261 let vol_change_fut = avc.next();
263 pin_mut!(vol_change_fut);
264 exec.run_until_stalled(&mut vol_change_fut).expect_pending("pending");
265
266 let event =
268 fidl_fuchsia_bluetooth_avrcp::Notification { volume: Some(100), ..Default::default() };
269 controller_stream.control_handle().send_on_notification(1000, &event).unwrap();
270
271 let new_vol = exec
272 .run_until_stalled(&mut vol_change_fut)
273 .expect("stream item")
274 .expect("some")
275 .expect("ok");
276 assert_eq!(new_vol, 100);
277
278 let mut request_fut = controller_stream.next();
280 let request = exec.run_until_stalled(&mut request_fut).expect("stream item").expect("some");
281 assert_matches!(request, Ok(ControllerRequest::NotifyNotificationHandled { .. }));
282 }
283
284 #[fuchsia::test]
285 fn test_stream_terminates() {
286 let mut exec = fasync::TestExecutor::new();
287 let (avc, _peer_manager_stream, controller_stream) = set_up(&mut exec);
288
289 pin_mut!(avc);
290 let mut vol_change_fut = avc.next();
291
292 assert!(exec.run_until_stalled(&mut vol_change_fut).is_pending());
294
295 drop(controller_stream);
298
299 assert!(exec.run_until_stalled(&mut vol_change_fut).expect("ready").is_none());
301 assert!(avc.is_terminated());
302 }
303}