Skip to main content

test_call_manager/
lib.rs

1// Copyright 2021 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, bail, format_err};
6use async_utils::hanging_get::client::HangingGetStream;
7use derivative::Derivative;
8use fidl_fuchsia_bluetooth::PeerId;
9use fidl_fuchsia_bluetooth_hfp::{
10    CallAction, CallDirection, CallManagerMarker, CallManagerRequest, CallManagerRequestStream,
11    CallMarker, CallRequest, CallRequestStream, CallState as FidlCallState,
12    CallWatchStateResponder, DtmfCode, HeadsetGainProxy, HfpMarker, HfpProxy, NetworkInformation,
13    NextCall, PeerHandlerRequest, PeerHandlerRequestStream,
14    PeerHandlerWatchNetworkInformationResponder, PeerHandlerWatchNextCallResponder, SignalStrength,
15};
16use fidl_fuchsia_bluetooth_hfp_test::{ConnectionBehavior, HfpTestMarker, HfpTestProxy};
17use fuchsia_async as fasync;
18use fuchsia_component::client;
19use futures::lock::Mutex;
20use futures::stream::StreamExt;
21use futures::{FutureExt, TryStreamExt};
22use log::*;
23use serde::Serialize;
24use std::collections::{HashMap, VecDeque};
25use std::sync::Arc;
26
27type CallId = u64;
28type Number = String;
29type Memory = String;
30
31/// Handles call actions initiated by the Hands Free.
32#[derive(Debug, Default, Clone)]
33struct Dialer {
34    /// The last dialed Number if one exists.
35    last_dialed: Option<Number>,
36    /// A map of Memory locations to Numbers.
37    address_book: HashMap<Memory, Number>,
38    /// The result that should be returned from a request to dial a Number.
39    dial_result: HashMap<Number, Result<(), zx::Status>>,
40}
41
42impl Dialer {
43    /// Performs an outgoing call initiation action, simulating a request to the network.
44    /// If the request was a success, the number of the outgoing call is returned.
45    /// If the request failed, the failure status is returned.
46    /// Defaults to failure with `zx::Status::NOT_FOUND` if the number associated with the call
47    /// action has not been explicitly set to return a result.
48    ///
49    /// Panics if `action` is a `CallAction::TransferActive`.
50    pub fn dial(&mut self, action: CallAction) -> Result<Number, zx::Status> {
51        let number = match &action {
52            CallAction::TransferActive(_) => panic!("TransferActive action passed to dial"),
53            CallAction::DialFromNumber(number) => Ok(number),
54            CallAction::DialFromLocation(location) => {
55                self.address_book.get(location).ok_or(zx::Status::NOT_FOUND)
56            }
57            CallAction::RedialLast(_) => self.last_dialed.as_ref().ok_or(zx::Status::NOT_FOUND),
58        }?
59        .to_owned();
60
61        let result = self.dial_result.get(&number).cloned().unwrap_or(Err(zx::Status::NOT_FOUND));
62        info!("Dial action result: {:?} - {:?}", action, result);
63        match result {
64            Ok(()) => {
65                self.last_dialed = Some(number.clone());
66                Ok(number)
67            }
68            Err(e) => Err(e),
69        }
70    }
71}
72
73#[derive(Derivative)]
74#[derivative(Debug)]
75/// State associated with the call manager (client) end of the HFP fidl service.
76struct ManagerState {
77    #[derivative(Debug = "ignore")]
78    peer_watcher: Option<fasync::Task<()>>,
79    network: NetworkInformation,
80    operator: String,
81    subscriber_numbers: Vec<String>,
82    nrec_support: bool,
83    battery_level: Option<u8>,
84    dialer: Dialer,
85}
86
87impl Default for ManagerState {
88    fn default() -> Self {
89        Self {
90            peer_watcher: None,
91            network: NetworkInformation::default(),
92            operator: String::new(),
93            subscriber_numbers: vec![],
94            nrec_support: true,
95            battery_level: None,
96            dialer: Dialer::default(),
97        }
98    }
99}
100
101/// State associated with a single Peer HF device.
102#[derive(Derivative)]
103#[derivative(Debug, Default)]
104struct PeerState {
105    reported_network: Option<NetworkInformation>,
106    network_responder: Option<PeerHandlerWatchNetworkInformationResponder>,
107    // nrec is enabled by default when a peer connects
108    #[derivative(Default(value = "true"))]
109    nrec_enabled: bool,
110    battery_level: u8,
111    speaker_gain: u8,
112    requested_speaker_gain: Option<u8>,
113    microphone_gain: u8,
114    requested_microphone_gain: Option<u8>,
115    #[derivative(Debug = "ignore")]
116    gain_control_watcher: Option<fasync::Task<()>>,
117    gain_control: Option<HeadsetGainProxy>,
118    call_responder: Option<PeerHandlerWatchNextCallResponder>,
119    // The tasks for managing a peer's call actions is owned by the peer.
120    // This task is separate from the manager's view of the call's state.
121    #[derivative(Debug = "ignore")]
122    call_tasks: HashMap<CallId, fasync::Task<()>>,
123}
124
125/// State associated with a single Call.
126#[derive(Derivative)]
127#[derivative(Debug)]
128struct CallState {
129    remote: String,
130    peer_id: Option<PeerId>,
131    responder: Option<CallWatchStateResponder>,
132    state: FidlCallState,
133    direction: CallDirection,
134    reported_state: Option<FidlCallState>,
135    dtmf_codes: Vec<DtmfCode>,
136}
137
138impl CallState {
139    /// Update the `state` and report the state if it is a new state and there is a
140    /// responder to report with.
141    pub fn update_state(&mut self, state: FidlCallState) -> Result<(), Error> {
142        self.state = state;
143        if self.reported_state != Some(state) && self.responder.is_some() {
144            let responder =
145                self.responder.take().expect("responder must be some after checking for presence");
146            responder.send(state)?;
147            self.reported_state = Some(state);
148        }
149        Ok(())
150    }
151}
152
153#[derive(Serialize)]
154pub struct StateSer {
155    manager: ManagerStateSer,
156    peers: HashMap<u64, PeerStateSer>,
157    calls: HashMap<CallId, CallStateSer>,
158}
159
160#[derive(Serialize)]
161pub struct ManagerStateSer {
162    network: NetworkInformationSer,
163    operator: String,
164    subscriber_numbers: Vec<String>,
165    nrec_support: bool,
166    battery_level: Option<u8>,
167    dialer: DialerSer,
168}
169
170#[derive(Serialize)]
171pub struct NetworkInformationSer {
172    service_available: Option<bool>,
173    signal_strength: Option<u8>,
174    roaming: Option<bool>,
175}
176
177impl From<NetworkInformation> for NetworkInformationSer {
178    fn from(info: NetworkInformation) -> Self {
179        let signal_strength = info.signal_strength.map(|strength| match strength {
180            SignalStrength::None => 0,
181            SignalStrength::VeryLow => 1,
182            SignalStrength::Low => 2,
183            SignalStrength::Medium => 3,
184            SignalStrength::High => 4,
185            SignalStrength::VeryHigh => 5,
186        });
187        Self { service_available: info.service_available, signal_strength, roaming: info.roaming }
188    }
189}
190
191impl From<&ManagerState> for ManagerStateSer {
192    fn from(state: &ManagerState) -> Self {
193        Self {
194            network: state.network.clone().into(),
195            operator: state.operator.clone(),
196            subscriber_numbers: state.subscriber_numbers.clone(),
197            nrec_support: state.nrec_support.clone(),
198            battery_level: state.battery_level.clone(),
199            dialer: state.dialer.clone().into(),
200        }
201    }
202}
203
204#[derive(Serialize)]
205struct PeerStateSer {
206    reported_network: Option<NetworkInformationSer>,
207    nrec_enabled: bool,
208    battery_level: u8,
209    speaker_gain: u8,
210    requested_speaker_gain: Option<u8>,
211    microphone_gain: u8,
212    requested_microphone_gain: Option<u8>,
213}
214
215impl From<&PeerState> for PeerStateSer {
216    fn from(state: &PeerState) -> Self {
217        Self {
218            reported_network: state.reported_network.clone().map(Into::into),
219            nrec_enabled: state.nrec_enabled,
220            battery_level: state.battery_level,
221            speaker_gain: state.speaker_gain,
222            requested_speaker_gain: state.requested_speaker_gain,
223            microphone_gain: state.microphone_gain,
224            requested_microphone_gain: state.requested_microphone_gain,
225        }
226    }
227}
228
229#[derive(Serialize)]
230struct CallStateSer {
231    remote: String,
232    direction: String,
233    state: String,
234    reported_state: Option<String>,
235    dtmf_codes: Vec<String>,
236}
237
238impl From<&CallState> for CallStateSer {
239    fn from(state: &CallState) -> Self {
240        Self {
241            remote: state.remote.clone(),
242            direction: format!("{:?}", state.direction),
243            state: format!("{:?}", state.state),
244            reported_state: state.reported_state.clone().map(|s| format!("{:?}", s)),
245            dtmf_codes: state.dtmf_codes.iter().map(|code| format!("{:?}", code)).collect(),
246        }
247    }
248}
249
250#[derive(Serialize)]
251struct DialerSer {
252    /// The last dialed Number if one exists.
253    last_dialed: Option<Number>,
254    /// A map of Memory locations to Numbers.
255    address_book: HashMap<Memory, Number>,
256    /// The result that should be returned from a request to dial a Number.
257    dial_result: HashMap<Number, String>,
258}
259
260impl From<Dialer> for DialerSer {
261    fn from(dialer: Dialer) -> Self {
262        Self {
263            last_dialed: dialer.last_dialed,
264            address_book: dialer.address_book,
265            dial_result: dialer
266                .dial_result
267                .into_iter()
268                .map(|(k, v)| (k, format!("{:?}", v)))
269                .collect(),
270        }
271    }
272}
273
274#[derive(Derivative, Default)]
275#[derivative(Debug)]
276struct TestCallManagerInner {
277    /// Connection to HFP Test interface
278    #[derivative(Debug = "ignore")]
279    test_proxy: Option<HfpTestProxy>,
280    /// Call manager state that is not associated with a particular peer or call.
281    manager: ManagerState,
282    /// Most commands are directed at a single peer. These commands are sent to the `active_peer`.
283    active_peer: Option<PeerId>,
284    /// State for all connected peer devices.
285    peers: HashMap<PeerId, PeerState>,
286    /// The next CallId to be assigned to a new call
287    next_call_id: CallId,
288    /// State for all ongoing calls
289    #[derivative(Debug = "ignore")]
290    calls: HashMap<CallId, CallState>,
291    /// Unreported calls
292    unreported_calls: VecDeque<CallId>,
293}
294
295impl TestCallManagerInner {
296    /// Remove a peer by `id` and all references to that peer.
297    pub fn remove_peer(&mut self, id: PeerId) {
298        let _ = self.peers.remove(&id);
299
300        for call in self.calls.values_mut() {
301            if call.peer_id == Some(id) {
302                call.peer_id = None;
303            }
304        }
305
306        if self.active_peer == Some(id) {
307            self.active_peer = None;
308        }
309    }
310
311    pub fn active_peer_mut(&mut self) -> Option<&mut PeerState> {
312        if let Some(id) = &self.active_peer {
313            Some(self.peers.get_mut(id).expect("Active peer must exist in peers map"))
314        } else {
315            None
316        }
317    }
318}
319
320#[derive(Debug, Clone)]
321pub struct TestCallManager {
322    inner: Arc<Mutex<TestCallManagerInner>>,
323}
324
325/// Perform Bluetooth HFP functions by acting as the call manager (client) side of the
326/// fuchsia.bluetooth.hfp.Hfp protocol.
327impl TestCallManager {
328    pub fn new() -> TestCallManager {
329        TestCallManager { inner: Arc::new(Mutex::new(TestCallManagerInner::default())) }
330    }
331
332    pub async fn set_request_stream(&self, stream: CallManagerRequestStream) {
333        let task =
334            fasync::Task::spawn(self.clone().watch_for_peers(stream).map(|f| f.unwrap_or({})));
335        self.inner.lock().await.manager.peer_watcher = Some(task);
336    }
337
338    pub async fn register_manager(&self, proxy: HfpProxy) -> Result<(), Error> {
339        let (client_end, stream) = fidl::endpoints::create_request_stream::<CallManagerMarker>();
340        proxy.register(client_end)?;
341        self.set_request_stream(stream).await;
342        Ok(())
343    }
344
345    /// Initialize the HFP service.
346    pub async fn init_hfp_service(&self) -> Result<(), Error> {
347        let mut inner = self.inner.lock().await;
348
349        if inner.manager.peer_watcher.is_none() {
350            info!("Connecting to HFP and setting new service proxy");
351
352            let hfp_service_proxy = client::connect_to_protocol::<HfpMarker>()
353                .map_err(|err| format_err!("Failed to create HFP service proxy: {}", err))?;
354
355            inner.test_proxy =
356                Some(client::connect_to_protocol::<HfpTestMarker>().map_err(|err| {
357                    format_err!("Failed to create HFP Test service proxy: {}", err)
358                })?);
359
360            let (client_end, stream) =
361                fidl::endpoints::create_request_stream::<CallManagerMarker>();
362            hfp_service_proxy.register(client_end)?;
363
364            let task =
365                fasync::Task::spawn(self.clone().watch_for_peers(stream).map(|f| f.unwrap_or({})));
366
367            inner.manager.peer_watcher = Some(task);
368        }
369        Ok(())
370    }
371
372    /// Return a list of Calls by call id and remote number.
373    pub async fn list_calls(&self) -> Result<Vec<(u64, String)>, Error> {
374        let inner = self.inner.lock().await;
375        Ok(inner.calls.iter().map(|(&id, state)| (id, state.remote.clone())).collect())
376    }
377
378    /// Notify HFP of an ongoing call. Simulates a new call from the network in the
379    /// "incoming ringing" state.
380    ///
381    /// Arguments:
382    ///     `remote`: The number associated with the remote party. This can be any string formatted
383    ///     number (e.g. +1-555-555-5555).
384    ///     `fidl_state`: The state to assign to the newly created call.
385    pub async fn new_call(
386        &self,
387        remote: &str,
388        fidl_state: FidlCallState,
389        direction: CallDirection,
390    ) -> Result<CallId, Error> {
391        let remote = remote.to_string();
392        let mut inner = self.inner.lock().await;
393        let call_id = inner.next_call_id;
394        inner.next_call_id += 1;
395        let mut state = CallState {
396            remote: remote.clone(),
397            peer_id: None,
398            responder: None,
399            state: fidl_state,
400            direction,
401            reported_state: None,
402            dtmf_codes: vec![],
403        };
404
405        if let Some(peer_id) = inner.active_peer.clone() {
406            let peer = inner.peers.get_mut(&peer_id).expect("Active peer must exist in peers map");
407
408            let (client_end, stream) = fidl::endpoints::create_request_stream::<CallMarker>();
409            // This does not handle the case where there is no peer responder
410            let responder = peer
411                .call_responder
412                .take()
413                .ok_or_else(|| format_err!("No peer call responder for {:?}", peer_id))?;
414            let next_call = NextCall {
415                call: Some(client_end),
416                remote: Some(remote),
417                state: Some(fidl_state),
418                direction: Some(direction),
419                ..Default::default()
420            };
421            if let Ok(()) = responder.send(next_call) {
422                let task = fasync::Task::local(self.clone().manage_call(peer_id, call_id, stream));
423                let _ = peer.call_tasks.insert(call_id, task);
424                state.peer_id = Some(peer_id);
425            }
426        } else {
427            inner.unreported_calls.push_back(call_id);
428        }
429
430        let _ = inner.calls.insert(call_id, state);
431        Ok(call_id)
432    }
433
434    /// Notify HFP of an incoming call. Simulates a new call from the network in the
435    /// "incoming ringing" state.
436    ///
437    /// Arguments:
438    ///     `remote`: The number associated with the remote party. This can be any string formatted
439    ///     number (e.g. +1-555-555-5555).
440    pub async fn incoming_ringing_call(&self, remote: &str) -> Result<CallId, Error> {
441        self.new_call(remote, FidlCallState::IncomingRinging, CallDirection::MobileTerminated).await
442    }
443
444    /// Notify HFP of an incoming waiting call. Simulates a new call from the network in the
445    /// "incoming waiting" state.
446    ///
447    /// Arguments:
448    ///     `remote`: The number associated with the remote party. This can be any string formatted
449    ///     number (e.g. +1-555-555-5555).
450    pub async fn incoming_waiting_call(&self, remote: &str) -> Result<CallId, Error> {
451        self.new_call(remote, FidlCallState::IncomingWaiting, CallDirection::MobileTerminated).await
452    }
453
454    /// Notify HFP of an outgoing call. Simulates a new call to the network in the
455    /// "outgoing notifying" state.
456    ///
457    /// Arguments:
458    ///     `remote`: The number associated with the remote party. This can be any string formatted
459    ///     number (e.g. +1-555-555-5555).
460    pub async fn outgoing_call(&self, remote: &str) -> Result<CallId, Error> {
461        self.new_call(remote, FidlCallState::OutgoingDialing, CallDirection::MobileOriginated).await
462    }
463
464    /// Notify HFP of an update to the state of an ongoing call.
465    ///
466    /// Arguments:
467    ///     `call_id`: The unique id of the call as assigned by the call manager.
468    ///     `fidl_state`: The state to assign to the call.
469    async fn update_call(&self, call_id: CallId, fidl_state: FidlCallState) -> Result<(), Error> {
470        // TODO: do not allow invalid state transitions (e.g. Terminated to Active)
471        self.inner
472            .lock()
473            .await
474            .calls
475            .get_mut(&call_id)
476            .ok_or_else(|| format_err!("Unknown Call Id {}", call_id))
477            .and_then(|call| call.update_state(fidl_state))
478    }
479
480    /// Notify HFP that a call is now active.
481    ///
482    /// Arguments:
483    ///     `call_id`: The unique id of the call as assigned by the call manager.
484    pub async fn set_call_active(&self, call_id: CallId) -> Result<(), Error> {
485        match self.update_call(call_id, FidlCallState::OngoingActive).await {
486            Ok(()) => Ok(()),
487            Err(e) => Err(format_err!("Failed to set call active: {}", e)),
488        }
489    }
490
491    /// Notify HFP that a call is now held.
492    ///
493    /// Arguments:
494    ///     `call_id`: The unique id of the call as assigned by the call manager.
495    pub async fn set_call_held(&self, call_id: CallId) -> Result<(), Error> {
496        match self.update_call(call_id, FidlCallState::OngoingHeld).await {
497            Ok(()) => Ok(()),
498            Err(e) => Err(format_err!("Failed to set call active: {}", e)),
499        }
500    }
501
502    /// Notify HFP that a call is now terminated.
503    ///
504    /// Arguments:
505    ///     `call_id`: The unique id of the call as assigned by the call manager.
506    pub async fn set_call_terminated(&self, call_id: CallId) -> Result<(), Error> {
507        match self.update_call(call_id, FidlCallState::Terminated).await {
508            Ok(()) => Ok(()),
509            Err(e) => Err(format_err!("Failed to terminate call: {}", e)),
510        }
511    }
512
513    /// Notify HFP that a call's audio is now transferred to the Audio Gateway.
514    ///
515    /// Arguments:
516    ///     `call_id`: The unique id of the call as assigned by the call manager.
517    pub async fn set_call_transferred_to_ag(&self, call_id: CallId) -> Result<(), Error> {
518        match self.update_call(call_id, FidlCallState::TransferredToAg).await {
519            Ok(()) => Ok(()),
520            Err(e) => Err(format_err!("Failed to transfer call to ag: {}", e)),
521        }
522    }
523
524    /// Return a list of HFP peers along with a boolean specifying whether each peer is the
525    /// "active" peer.
526    pub async fn list_peers(&self) -> Result<Vec<(u64, bool)>, Error> {
527        let inner = self.inner.lock().await;
528        Ok(inner
529            .peers
530            .keys()
531            .map(|&id| (id.value, inner.active_peer.map(|active| active == id).unwrap_or_default()))
532            .collect())
533    }
534
535    /// Set the active peer with the provided id. All future commands will be directed towards that
536    /// peer.
537    ///
538    /// Arguments:
539    ///     `id`: The unique id for the peer that initiated the request.
540    pub async fn set_active_peer(&self, id: u64) -> Result<(), Error> {
541        let id = PeerId { value: id };
542        let mut inner = self.inner.lock().await;
543        if inner.peers.contains_key(&id) {
544            inner.active_peer = Some(id);
545            Ok(())
546        } else {
547            Err(format_err!("Peer {:?} not connected", id))
548        }
549    }
550
551    // Report all calls to the given peer. Panics if `id` does not point to a valid peer in the
552    // peers map.
553    fn report_calls(
554        self,
555        id: PeerId,
556        mut inner: futures::lock::MutexGuard<'_, TestCallManagerInner>,
557    ) -> Result<(), Error> {
558        // keep popping from unreported until we get one
559        // that is also found in the calls map.
560        while let Some(cid) = inner.unreported_calls.pop_front() {
561            if let Some(call) = inner.calls.get(&cid) {
562                let remote = call.remote.clone();
563                let state = call.state;
564                let direction = call.direction;
565                let (client_end, stream) = fidl::endpoints::create_request_stream::<CallMarker>();
566                let peer = inner.peers.get_mut(&id).expect("peer just added");
567                let next_call = NextCall {
568                    call: Some(client_end),
569                    remote: Some(remote),
570                    state: Some(state),
571                    direction: Some(direction),
572                    ..Default::default()
573                };
574                let res = peer.call_responder.take().expect("just put here").send(next_call);
575                if let Ok(()) = res {
576                    let task = fasync::Task::local(self.manage_call(id, cid, stream));
577                    let _ = peer.call_tasks.insert(cid, task);
578                    inner.calls.get_mut(&cid).expect("still here").peer_id = Some(id);
579                }
580                break;
581            }
582        }
583        Ok(())
584    }
585
586    /// Handle a peer `request`. Most requests are handled by immediately responding with the
587    /// relevant data which is stored in the facade's state.
588    ///
589    /// Arguments:
590    ///     `id`: The unique id for the peer that initiated the request.
591    ///     `request`: A request made by a client of the fuchsia.bluetooth.hfp.PeerHandler protocol.
592    async fn handle_peer_request(
593        &mut self,
594        id: PeerId,
595        request: PeerHandlerRequest,
596    ) -> Result<(), Error> {
597        info!("Received Peer Handler request for {:?}: {:?}", id, request);
598        match request {
599            PeerHandlerRequest::WatchNetworkInformation { responder, .. } => {
600                let mut inner = self.inner.lock().await;
601                let current_network = inner.manager.network.clone();
602                let peer = inner
603                    .peers
604                    .get_mut(&id)
605                    .ok_or_else(|| format_err!("peer removed: {:?}", id))?;
606
607                if Some(&current_network) == peer.reported_network.as_ref() {
608                    peer.network_responder = Some(responder);
609                } else {
610                    responder.send(&current_network)?;
611                    peer.reported_network = Some(current_network);
612                }
613            }
614            PeerHandlerRequest::WatchNextCall { responder, .. } => {
615                let this = self.clone();
616                let mut inner = self.inner.lock().await;
617                let peer = inner
618                    .peers
619                    .get_mut(&id)
620                    .ok_or_else(|| format_err!("peer removed: {:?}", id))?;
621                if peer.call_responder.is_none() {
622                    peer.call_responder = Some(responder);
623                    this.report_calls(id, inner)?;
624                } else {
625                    let err = format_err!("double hanging get call on PeerHandler::WatchNextCall");
626                    error!(err:%; "");
627                    *inner = TestCallManagerInner::default();
628                    return Err(err);
629                }
630            }
631            PeerHandlerRequest::RequestOutgoingCall { action, responder } => {
632                if let CallAction::TransferActive(_) = action {
633                    let inner = self.inner.lock().await;
634                    match inner
635                        .calls
636                        .iter()
637                        .find(|(_, call)| call.state == FidlCallState::OngoingActive)
638                    {
639                        Some((&id, _)) => {
640                            drop(inner);
641                            // result can be ignored because id was just found in the call map.
642                            let _ = self.update_call(id, FidlCallState::TransferredToAg).await;
643                        }
644                        None => drop(responder.send(Err(zx::Status::NOT_FOUND.into_raw()))),
645                    };
646                } else {
647                    // Simulate dialing action and then respond to any outstanding WatchForCall
648                    // requests.
649                    let result = match {
650                        // Only hold onto the lock while using it to "dial" the number.
651                        // Holding the lock past this point would cause a deadlock when
652                        // calling `outgoing_call`.
653                        let mut inner = self.inner.lock().await;
654                        inner.manager.dialer.dial(action)
655                    } {
656                        Ok(number) => match self.outgoing_call(&number).await {
657                            Ok(id) => {
658                                info!(CallId = id; "Initiated outgoing call to {}", number);
659                                Ok(())
660                            }
661                            Err(e) => {
662                                error!("Could not initiate outgoing call action: {}", e);
663                                Err(zx::Status::INTERNAL.into_raw())
664                            }
665                        },
666                        Err(status) => Err(status.into_raw()),
667                    };
668                    info!("sending result to peer: {:?}", result);
669
670                    // Once dialing and hanging gets have been handled, send response.
671                    let _ = responder.send(result);
672                }
673            }
674            PeerHandlerRequest::QueryOperator { responder, .. } => {
675                responder.send(Some(&self.inner.lock().await.manager.operator))?;
676            }
677            PeerHandlerRequest::SubscriberNumberInformation { responder, .. } => {
678                responder.send(&self.inner.lock().await.manager.subscriber_numbers)?;
679            }
680            PeerHandlerRequest::SetNrecMode { enabled, responder, .. } => {
681                let mut inner = self.inner.lock().await;
682                if inner.manager.nrec_support {
683                    let peer = inner
684                        .peers
685                        .get_mut(&id)
686                        .ok_or_else(|| format_err!("peer removed: {:?}", id))?;
687                    peer.nrec_enabled = enabled;
688                    responder.send(Ok(()))?;
689                } else {
690                    responder.send(Err(zx::Status::NOT_SUPPORTED.into_raw()))?;
691                }
692            }
693            PeerHandlerRequest::ReportHeadsetBatteryLevel { level, .. } => {
694                self.inner
695                    .lock()
696                    .await
697                    .peers
698                    .get_mut(&id)
699                    .ok_or_else(|| format_err!("peer removed: {:?}", id))?
700                    .battery_level = level;
701            }
702            PeerHandlerRequest::GainControl { control, .. } => {
703                let this = self.clone();
704                let proxy = control.into_proxy();
705                let proxy_ = proxy.clone();
706                let task = fasync::Task::spawn(async move {
707                    let mut speaker_gain_stream =
708                        HangingGetStream::new(proxy.clone(), HeadsetGainProxy::watch_speaker_gain);
709                    let mut microphone_gain_stream = HangingGetStream::new(
710                        proxy.clone(),
711                        HeadsetGainProxy::watch_microphone_gain,
712                    );
713
714                    loop {
715                        futures::select! {
716                            gain = speaker_gain_stream.next() => {
717                                let mut inner = this.inner.lock().await;
718                                let peer = inner.peers.get_mut(&id);
719                                match (peer, gain) {
720                                    (Some(peer), Some(Ok(gain))) => peer.speaker_gain = gain,
721                                    _ => break,
722                                }
723                            }
724                            gain = microphone_gain_stream.next() => {
725                                let mut inner = this.inner.lock().await;
726                                let peer = inner.peers.get_mut(&id);
727                                match (peer, gain) {
728                                    (Some(peer), Some(Ok(gain))) => peer.microphone_gain = gain,
729                                    _ => break,
730                                }
731                            }
732                        }
733                    }
734                    info!("Headset gain control channel for peer {:?} closed", id);
735                });
736                let mut inner = self.inner.lock().await;
737                let peer = inner
738                    .peers
739                    .get_mut(&id)
740                    .ok_or_else(|| format_err!("peer removed: {:?}", id))?;
741                if let Some(requested) = peer.requested_speaker_gain.take() {
742                    proxy_.set_speaker_gain(requested)?;
743                }
744                if let Some(requested) = peer.requested_microphone_gain.take() {
745                    proxy_.set_microphone_gain(requested)?;
746                }
747                peer.gain_control_watcher = Some(task);
748                peer.gain_control = Some(proxy_);
749            }
750        }
751        Ok(())
752    }
753
754    /// Handle all PeerHandlerRequests for a peer, removing the peer after the request stream
755    /// is closed.
756    ///
757    /// Arguments:
758    ///     `id`: The unique id for the peer that initiated the stream.
759    ///     `stream`: A stream of requests associated with a single peer.
760    async fn manage_peer(mut self, id: PeerId, mut stream: PeerHandlerRequestStream) {
761        while let Some(Ok(request)) = stream.next().await {
762            if let Err(err) = self.handle_peer_request(id, request).await {
763                error!(err:%; "");
764                break;
765            };
766        }
767        self.inner.lock().await.remove_peer(id);
768    }
769
770    /// Watch for new Hands Free peer devices that connect to the DUT.
771    ///
772    /// Spawns a new task to manage each peer that connects.
773    ///
774    /// Arguments:
775    ///     `proxy`: The client end of the CallManager protocol which is used to watch for new
776    ///     Bluetooth peers.
777    async fn watch_for_peers(
778        self,
779        mut stream: CallManagerRequestStream,
780    ) -> Result<(), fidl::Error> {
781        // Entries are only removed from the map when they are replaced, so the size of `peers`
782        // will grow as the total number of unique peers grows. This is acceptable as the number of
783        // unique peers that connect to a DUT is expected to be relatively small.
784        let mut peers = HashMap::new();
785        while let Some(CallManagerRequest::PeerConnected { id, handle, responder }) =
786            stream.try_next().await?
787        {
788            let stream = handle.into_stream();
789            info!("Handling Peer: {:?}", id);
790            {
791                let mut inner = self.inner.lock().await;
792                let _ = inner.peers.insert(id, PeerState::default());
793                inner.active_peer = Some(id);
794            }
795
796            let task = fasync::Task::spawn(self.clone().manage_peer(id, stream));
797            let _ = peers.insert(id, task);
798            let _ = responder.send();
799        }
800        Ok(())
801    }
802
803    /// Manage an ongoing call that is being routed to the Hands Free peer. Handle any requests
804    /// made by the peer that are associated with the individual call.
805    ///
806    /// Arguments:
807    ///     `peer_id`: The unique id of the peer as assigned by the Bluetooth stack.
808    ///     `call_id`: The unique id of the call as assigned by the call manager.
809    ///     `stream`: A stream of requests associated with a single call.
810    async fn manage_call(self, peer_id: PeerId, call_id: CallId, mut stream: CallRequestStream) {
811        while let Some(request) = stream.next().await {
812            info!("Got call request: {:?} {:?} -> {:?}", peer_id, call_id, request);
813            let mut inner = self.inner.lock().await;
814            let state = if let Some(state) = inner.calls.get_mut(&call_id) {
815                state
816            } else {
817                info!("Call management by {:?} ended: {:?}", peer_id, call_id);
818                break;
819            };
820            match request {
821                Ok(CallRequest::WatchState { responder, .. }) => {
822                    if state.responder.is_some() {
823                        warn!("Call client sent multiple WatchState requests. Closing channel");
824                        break;
825                    }
826                    state.responder = Some(responder);
827                    // Trigger an update with the existing state to send it on the responder if
828                    // necessary.
829                    if let Err(e) = state.update_state(state.state) {
830                        info!("Call ended: {}", e);
831                        break;
832                    }
833                }
834                Ok(CallRequest::RequestHold { .. }) => {
835                    if let Err(e) = state.update_state(FidlCallState::OngoingHeld) {
836                        info!("Call ended: {}", e);
837                        break;
838                    }
839                }
840                Ok(CallRequest::RequestActive { .. }) => {
841                    if let Err(e) = state.update_state(FidlCallState::OngoingActive) {
842                        info!("Call ended: {}", e);
843                        break;
844                    }
845                }
846                Ok(CallRequest::RequestTerminate { .. }) => {
847                    if let Err(e) = state.update_state(FidlCallState::Terminated) {
848                        info!("Call ended: {}", e);
849                        break;
850                    }
851                }
852                Ok(CallRequest::RequestTransferAudio { .. }) => {
853                    if let Err(e) = state.update_state(FidlCallState::TransferredToAg) {
854                        info!("Call ended: {}", e);
855                        break;
856                    }
857                }
858                Ok(CallRequest::SendDtmfCode { code, responder, .. }) => {
859                    state.dtmf_codes.push(code);
860                    if let Err(e) = responder.send(Ok(())) {
861                        info!("Call ended: {}", e);
862                        break;
863                    }
864                }
865                Err(e) => {
866                    warn!("Call fidl channel error: {}", e);
867                }
868            }
869        }
870
871        // Cleanup before exiting call task
872        let mut inner = self.inner.lock().await;
873        if let Some(call) = inner.calls.get_mut(&call_id) {
874            call.peer_id = None
875        }
876        if let Some(peer) = inner.peers.get_mut(&peer_id) {
877            let _ = peer.call_tasks.remove(&call_id);
878        }
879    }
880
881    /// Request that the active peer's speaker gain be set to `value`.
882    ///
883    /// Arguments:
884    ///     `value`: must be between 0-15 inclusive.
885    pub async fn set_speaker_gain(&self, value: u64) -> Result<(), Error> {
886        let value = value as u8;
887        let mut inner = self.inner.lock().await;
888        let peer = inner.active_peer_mut().ok_or_else(|| format_err!("No active peer"))?;
889        if peer.speaker_gain != value {
890            if let Some(gain_control) = &peer.gain_control {
891                gain_control.set_speaker_gain(value)?;
892            } else {
893                peer.requested_speaker_gain = Some(value);
894            }
895        }
896        Ok(())
897    }
898
899    /// Request that the active peer's microphone gain be set to `value`.
900    ///
901    /// Arguments:
902    ///     `value`: must be between 0-15 inclusive.
903    pub async fn set_microphone_gain(&self, value: u64) -> Result<(), Error> {
904        let value = value as u8;
905        let mut inner = self.inner.lock().await;
906        let peer = inner.active_peer_mut().ok_or_else(|| format_err!("No active peer"))?;
907        if peer.microphone_gain != value {
908            if let Some(gain_control) = &peer.gain_control {
909                gain_control.set_microphone_gain(value)?;
910            } else {
911                peer.requested_microphone_gain = Some(value);
912            }
913        }
914        Ok(())
915    }
916
917    /// Update the facade's network information with the provided `network`.
918    /// Any fields in `network` that are `None` will not be updated.
919    ///
920    /// Arguments:
921    ///     `network`: The updated network information fields.
922    pub async fn update_network_information(
923        &self,
924        network: NetworkInformation,
925    ) -> Result<(), Error> {
926        let mut inner = self.inner.lock().await;
927
928        // Update network state
929        let last_net = std::mem::replace(&mut inner.manager.network, NetworkInformation::default());
930        inner.manager.network = NetworkInformation {
931            service_available: network.service_available.or(last_net.service_available),
932            signal_strength: network.signal_strength.or(last_net.signal_strength),
933            roaming: network.roaming.or(last_net.roaming),
934            ..last_net
935        };
936        let current_network = inner.manager.network.clone();
937
938        for peer in inner.peers.values_mut() {
939            // Update the client if a responder is present
940            if Some(&current_network) != peer.reported_network.as_ref() {
941                if let Some(responder) = peer.network_responder.take() {
942                    responder.send(&current_network)?;
943                    peer.reported_network = Some(current_network.clone());
944                }
945            }
946        }
947
948        Ok(())
949    }
950
951    pub async fn set_subscriber_number(&self, number: &str) {
952        self.inner.lock().await.manager.subscriber_numbers = vec![number.to_owned()];
953    }
954
955    pub async fn set_operator(&self, value: &str) {
956        self.inner.lock().await.manager.operator = value.to_owned();
957    }
958
959    pub async fn set_nrec_support(&self, value: bool) {
960        self.inner.lock().await.manager.nrec_support = value;
961    }
962
963    pub async fn set_battery_level(&self, value: u64) -> Result<(), Error> {
964        if value > 5 {
965            bail!("Value out of range: {}. Battery level must be 0-5.", value);
966        }
967        let mut inner = self.inner.lock().await;
968        let proxy = inner
969            .test_proxy
970            .as_ref()
971            .ok_or_else(|| format_err!("Cannot set battery without HfpTest proxy"))?;
972        proxy.battery_indicator(value as u8)?;
973        inner.manager.battery_level = Some(value as u8);
974        Ok(())
975    }
976
977    pub async fn get_state(&self) -> StateSer {
978        let inner = self.inner.lock().await;
979        StateSer {
980            manager: (&inner.manager).into(),
981            peers: inner
982                .peers
983                .iter()
984                .map(|(&PeerId { value: id }, peer)| (id, peer.into()))
985                .collect(),
986            calls: inner.calls.iter().map(|(&id, call)| (id, call.into())).collect(),
987        }
988    }
989
990    /// Set the simulated "last dialed" number.
991    ///
992    /// Arguments:
993    ///     `number`: Number to be set. To clear the last dialed number, set `number` to `None`.
994    pub async fn set_last_dialed(&self, number: Option<Number>) {
995        self.inner.lock().await.manager.dialer.last_dialed = number;
996    }
997
998    /// Store a number at a specific location in address book memory.
999    ///
1000    /// Arguments:
1001    ///     `location`: The key used to look up a specific number in address book memory.
1002    ///     `number`: Number to be set. To remove the address book entry, set `number` to `None`.
1003    pub async fn set_memory_location(&self, location: Memory, number: Option<Number>) {
1004        let _ = match number {
1005            Some(number) => {
1006                self.inner.lock().await.manager.dialer.address_book.insert(location, number)
1007            }
1008            None => self.inner.lock().await.manager.dialer.address_book.remove(&location),
1009        };
1010    }
1011
1012    /// Set the simulated result that will be returned after HFP requests an outgoing call.
1013    /// This result is used regardless of whether a number was specified directly through
1014    /// CallAction::dial_from_number or indirectly through either CallAction::dial_from_location or
1015    /// CallAction::redial_last.
1016    ///
1017    /// Arguments:
1018    ///     `number`: Number that maps to a simulated result.
1019    ///     `status`: The simulated result value for `number`.
1020    pub async fn set_dial_result(&self, number: Number, status: Result<(), zx::Status>) {
1021        let _ = self.inner.lock().await.manager.dialer.dial_result.insert(number, status);
1022    }
1023
1024    /// Configure the connection behavior when the component receives new search results from
1025    /// the bredr.Profile protocol.
1026    ///
1027    /// Arguments:
1028    ///     `autoconnect`: determine whether the component should automatically attempt to
1029    ///                    make a new RFCOMM connection.
1030    pub async fn set_connection_behavior(&self, autoconnect: bool) -> Result<(), Error> {
1031        let inner = self.inner.lock().await;
1032        let proxy = inner.test_proxy.as_ref().ok_or_else(|| {
1033            format_err!("Cannot set slc connection behavior on command without HfpTest proxy")
1034        })?;
1035        let () = proxy.set_connection_behavior(&ConnectionBehavior {
1036            autoconnect: Some(autoconnect),
1037            ..Default::default()
1038        })?;
1039        Ok(())
1040    }
1041
1042    /// Cleanup any HFP related objects.
1043    pub async fn cleanup(&self) {
1044        *self.inner.lock().await = TestCallManagerInner::default();
1045    }
1046}
1047
1048#[cfg(test)]
1049mod tests {
1050    use super::*;
1051    use fidl_fuchsia_bluetooth_hfp::PeerHandlerMarker;
1052
1053    #[fuchsia::test]
1054    async fn outgoing_call_does_not_deadlock() {
1055        let manager = TestCallManager::new();
1056
1057        // set up the dial result so that an outgoing call request will be a success.
1058        manager.set_dial_result("123".to_string(), Ok(())).await;
1059
1060        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<PeerHandlerMarker>();
1061
1062        // Create a background task to manage a peer channel.
1063        fasync::Task::local({
1064            let manager = manager.clone();
1065            async move {
1066                manager.manage_peer(PeerId { value: 1 }, stream).await;
1067            }
1068        })
1069        .detach();
1070
1071        // requesting an outgoing call should complete successfully
1072        let result =
1073            proxy.request_outgoing_call(&CallAction::DialFromNumber("123".to_string())).await;
1074        assert!(result.is_ok());
1075    }
1076}