Skip to main content

bt_avrcp_vol_control/
lib.rs

1// Copyright 2025 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 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
17/// Manages absolute volume control for an active AVRCP connection to a remote peer.
18pub 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    /// Creates a new `AbsoluteVolumeControl` for the given peer.
41    ///
42    /// This will fail if the controller can't be obtained or the notification filter can't be set.
43    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    /// Sets the absolute volume on the remote peer.
68    ///
69    /// `volume_percentage` is a value between 0 and 100.
70    ///
71    /// This method returns the result of the operation, and also queues the result to be
72    /// delivered as an item from the stream.
73    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        // Push the result to volume results so clients can get volume update from volume change notification stream.
84        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            // First, check for results from any completed `set_absolute_volume` commands.
109            if let Some(result) = self.set_volume_results.pop_front() {
110                return Poll::Ready(Some(result));
111            }
112
113            // Next, poll for incoming notifications from the peer.
114            match self.stream.poll_next_unpin(cx) {
115                Poll::Ready(Some(Ok(ControllerEvent::OnNotification { notification, .. }))) => {
116                    // Fire and forget. If this fails, `is_closed()` will be true on the next poll.
117                    let _ = self.proxy.notify_notification_handled();
118
119                    if let Some(volume) = notification.volume {
120                        return Poll::Ready(Some(Ok(volume)));
121                    }
122                    // Not a volume notification, so we continue the loop to poll again.
123                }
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                    // The underlying stream is pending. Store the waker so that `set_absolute_volume`
134                    // can wake us up if a result becomes available.
135                    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        // The server should receive a `GetControllerForTarget` request.
177        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; // 50 * 127 / 100
211        let expected_response_volume = 64;
212
213        // No volume update yet.
214        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            // The Controller server should have received a `SetAbsoluteVolume` request.
222            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            // The future should complete with the result.
234            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        // No volume change yet.
262        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        // Mock the server sending a volume notification.
267        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        // Should have acknowledged the notification.
279        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        // The future should be pending before any events.
293        assert!(exec.run_until_stalled(&mut vol_change_fut).is_pending());
294
295        // Drop the server-side stream. This will cause the client stream to terminate
296        // gracefully on its next poll.
297        drop(controller_stream);
298
299        // Poll will now see the closed channel and return Ready with None.
300        assert!(exec.run_until_stalled(&mut vol_change_fut).expect("ready").is_none());
301        assert!(avc.is_terminated());
302    }
303}