| // Copyright 2026 The Fuchsia Authors. All rights reserved. |
| // Use of this source code is governed by a BSD-style license that can be |
| // found in the LICENSE file. |
| |
| //! Implements the Media Control Service (MCS) server. |
| |
| use bt_common::packet_encoding::Decodable; |
| use bt_common::{PeerId, Uuid}; |
| use bt_gatt::server::{ |
| LocalService, ReadResponder, Server as _, ServiceDefinition, ServiceEvent, ServiceId, |
| WriteResponder, |
| }; |
| use bt_gatt::types::{ |
| AttributePermissions, CharacteristicProperties, CharacteristicProperty, GattError, Handle, |
| SecurityLevels, ServiceKind, |
| }; |
| use bt_gatt::Characteristic; |
| use futures::channel::oneshot; |
| use futures::stream::{FuturesUnordered, Stream}; |
| use futures::FutureExt; |
| use pin_project::pin_project; |
| use std::future::Future; |
| use std::pin::Pin; |
| use std::task::{Context, Poll, Waker}; |
| |
| use crate::types::*; |
| use crate::Error; |
| |
| // ============================================================================ |
| // Mandatory Characteristic Handle Definitions (MCS v1.0.1 Section 3) |
| // ============================================================================ |
| |
| /// Handle assigned to the Media Player Name characteristic. |
| const MEDIA_PLAYER_NAME_HANDLE: Handle = Handle(1); |
| /// Handle assigned to the Track Changed characteristic. |
| const TRACK_CHANGED_HANDLE: Handle = Handle(2); |
| /// Handle assigned to the Track Title characteristic. |
| const TRACK_TITLE_HANDLE: Handle = Handle(3); |
| /// Handle assigned to the Track Duration characteristic. |
| const TRACK_DURATION_HANDLE: Handle = Handle(4); |
| /// Handle assigned to the Track Position characteristic. |
| const TRACK_POSITION_HANDLE: Handle = Handle(5); |
| /// Handle assigned to the Media State characteristic. |
| const MEDIA_STATE_HANDLE: Handle = Handle(6); |
| /// Handle assigned to the Content Control ID (CCID) characteristic. |
| const CONTENT_CONTROL_ID_HANDLE: Handle = Handle(7); |
| |
| // ============================================================================ |
| // Optional Characteristic Handle Definitions (MCS v1.0.1 Section 3) |
| // ============================================================================ |
| |
| /// Handle assigned to the Media Player Icon URL characteristic. |
| const MEDIA_PLAYER_ICON_URL_HANDLE: Handle = Handle(8); |
| /// Handle assigned to the Playback Speed characteristic. |
| const PLAYBACK_SPEED_HANDLE: Handle = Handle(9); |
| /// Handle assigned to the Seeking Speed characteristic. |
| const SEEKING_SPEED_HANDLE: Handle = Handle(10); |
| /// Handle assigned to the Playing Order characteristic. |
| const PLAYING_ORDER_HANDLE: Handle = Handle(11); |
| /// Handle assigned to the Playing Orders Supported characteristic. |
| const PLAYING_ORDERS_SUPPORTED_HANDLE: Handle = Handle(12); |
| /// Handle assigned to the Media Control Point characteristic. |
| const MEDIA_CONTROL_POINT_HANDLE: Handle = Handle(13); |
| /// Handle assigned to the Media Control Point Opcodes Supported characteristic. |
| const MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE: Handle = Handle(14); |
| |
| /// The mandatory characteristics defined in MCS v1.0.1 Section 3, Table 3.1. |
| fn mandatory_characteristics() -> [Characteristic; 7] { |
| [ |
| build_characteristic( |
| MEDIA_PLAYER_NAME_HANDLE, |
| MEDIA_PLAYER_NAME_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| build_characteristic( |
| TRACK_CHANGED_HANDLE, |
| TRACK_CHANGED_UUID, |
| CharacteristicProperty::Notify, |
| ), |
| build_characteristic( |
| TRACK_TITLE_HANDLE, |
| TRACK_TITLE_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| build_characteristic( |
| TRACK_DURATION_HANDLE, |
| TRACK_DURATION_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| build_characteristic( |
| TRACK_POSITION_HANDLE, |
| TRACK_POSITION_UUID, |
| CharacteristicProperties::READ_WRITE_NOTIFY, |
| ), |
| build_characteristic( |
| MEDIA_STATE_HANDLE, |
| MEDIA_STATE_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| build_characteristic( |
| CONTENT_CONTROL_ID_HANDLE, |
| CONTENT_CONTROL_ID_UUID, |
| CharacteristicProperty::Read, |
| ), |
| ] |
| } |
| |
| /// Constructs a characteristic definition with the specified handle, UUID, |
| /// properties, and encryption-required permissions conforming to MCS v1.0.1 |
| /// Section 3. |
| fn build_characteristic( |
| handle: Handle, |
| uuid: Uuid, |
| properties: impl Into<CharacteristicProperties>, |
| ) -> Characteristic { |
| let properties = properties.into(); |
| Characteristic { |
| handle, |
| uuid, |
| properties, |
| permissions: AttributePermissions::with_levels( |
| &properties, |
| &SecurityLevels::encryption_required(), |
| ), |
| descriptors: Vec::new(), |
| } |
| } |
| |
| /// Internal state for the MCS server. |
| #[pin_project(project = LocalServiceProj)] |
| enum LocalServiceState<T: bt_gatt::ServerTypes> { |
| /// Service definition has not been registered in the GATT database. |
| NotPublished { |
| waker: Option<Waker>, |
| }, |
| /// Service registration is in progress. |
| Preparing { |
| #[pin] |
| fut: T::LocalServiceFut, |
| }, |
| /// Service registration is complete and active in the GATT database. |
| Published { |
| service: T::LocalService, |
| #[pin] |
| events: T::ServiceEventStream, |
| }, |
| Terminated, |
| } |
| |
| impl<T: bt_gatt::ServerTypes> Default for LocalServiceState<T> { |
| fn default() -> Self { |
| Self::NotPublished { waker: None } |
| } |
| } |
| |
| impl<T: bt_gatt::ServerTypes> LocalServiceState<T> { |
| fn is_published(&self) -> bool { |
| matches!(self, LocalServiceState::Published { .. }) |
| } |
| |
| fn notify(&self, handle: &Handle, data: &[u8], peers: &[PeerId]) { |
| if let LocalServiceState::Published { service, .. } = self { |
| service.notify(handle, data, peers); |
| } |
| } |
| } |
| |
| impl<T: bt_gatt::ServerTypes> Stream for LocalServiceState<T> { |
| type Item = Result<ServiceEvent<T>, Error>; |
| |
| fn poll_next( |
| mut self: std::pin::Pin<&mut Self>, |
| cx: &mut Context<'_>, |
| ) -> Poll<Option<Self::Item>> { |
| loop { |
| match self.as_mut().project() { |
| LocalServiceProj::Terminated => return Poll::Ready(None), |
| LocalServiceProj::NotPublished { waker } => { |
| *waker = Some(cx.waker().clone()); |
| return Poll::Pending; |
| } |
| LocalServiceProj::Preparing { fut } => match futures::ready!(fut.poll(cx)) { |
| Ok(service) => { |
| let events = service.publish(); |
| self.as_mut().set(LocalServiceState::Published { service, events }); |
| } |
| Err(e) => { |
| self.as_mut().set(LocalServiceState::NotPublished { waker: None }); |
| return Poll::Ready(Some(Err(Error::Gatt(e)))); |
| } |
| }, |
| LocalServiceProj::Published { service: _, events } => { |
| match futures::ready!(events.poll_next(cx)) { |
| Some(Ok(event)) => return Poll::Ready(Some(Ok(event))), |
| Some(Err(e)) => { |
| self.as_mut().set(LocalServiceState::Terminated); |
| return Poll::Ready(Some(Err(Error::Gatt(e)))); |
| } |
| // Deferred to McsServer to allow draining in-flight control point |
| // responses. |
| None => return Poll::Ready(None), |
| } |
| } |
| } |
| } |
| } |
| } |
| |
| /// Confirms the updated track position to store in local state and notify to |
| /// clients. |
| /// It is safe to drop without responding to ignore the update or if the value |
| /// has not changed. |
| pub type SetTrackPositionResponder = SetValueResponder<TrackPosition>; |
| |
| /// Confirms the updated playback speed to store in local state and notify to |
| /// clients. |
| /// It is safe to drop without responding to ignore the update or if the value |
| /// has not changed. |
| pub type SetPlaybackSpeedResponder = SetValueResponder<PlaybackSpeed>; |
| |
| /// Confirms the updated playing order to store in local state and notify to |
| /// clients. |
| /// It is safe to drop without responding to ignore the update or if the value |
| /// has not changed. |
| pub type SetPlayingOrderResponder = SetValueResponder<PlayingOrder>; |
| |
| /// Confirms the media control command result code to notify to the client. |
| /// A response is expected to confirm the result. If no response is provided, |
| /// then it is assumed that the Control Point request has failed and |
| /// `ControlPointResultCode::CommandCannotBeCompleted` will be sent to the peer. |
| pub type ControlPointWriteResponder = SetValueResponder<ControlPointResultCode>; |
| |
| /// Responder for an asynchronous characteristic write or control point command. |
| /// |
| /// Calling [`send`](Self::send) confirms the new value or command result to |
| /// update local state and notify clients. It is safe to drop the responder |
| /// without responding to cancel the update (or report command failure). |
| pub struct SetValueResponder<T> { |
| response_tx: oneshot::Sender<T>, |
| } |
| |
| impl<T: std::fmt::Debug> SetValueResponder<T> { |
| /// Confirms the value or result code determined by the media player |
| /// application and dispatches the corresponding GATT notification. |
| pub fn send(self, value: T) { |
| if let Err(result) = self.response_tx.send(value) { |
| log::warn!("Failed to send response: server dropped or closed: {result:?}"); |
| } |
| } |
| } |
| |
| impl<T> std::fmt::Debug for SetValueResponder<T> { |
| fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { |
| f.debug_struct("SetValueResponder").finish_non_exhaustive() |
| } |
| } |
| |
| /// An in-flight handler for a single-value asynchronous write response. |
| struct SetValueResponseFut<T> { |
| rx: oneshot::Receiver<T>, |
| } |
| |
| impl<T> SetValueResponseFut<T> { |
| fn create() -> (Self, SetValueResponder<T>) { |
| let (tx, rx) = oneshot::channel(); |
| (Self { rx }, SetValueResponder { response_tx: tx }) |
| } |
| } |
| |
| impl<T> Future for SetValueResponseFut<T> { |
| type Output = Option<T>; |
| |
| fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { |
| self.rx.poll_unpin(cx).map(Result::ok) |
| } |
| } |
| |
| /// Responses produced when an in-flight asynchronous write response completes. |
| enum PendingWriteResponse { |
| ControlPoint { |
| peer_id: PeerId, |
| opcode: MediaControlOpcode, |
| result_code: ControlPointResultCode, |
| }, |
| TrackPosition(TrackPosition), |
| PlaybackSpeed(PlaybackSpeed), |
| PlayingOrder(PlayingOrder), |
| } |
| |
| /// An in-flight asynchronous write response future. |
| enum PendingWriteResponseFut { |
| ControlPoint { |
| peer_id: PeerId, |
| opcode: MediaControlOpcode, |
| fut: SetValueResponseFut<ControlPointResultCode>, |
| }, |
| TrackPosition(SetValueResponseFut<TrackPosition>), |
| PlaybackSpeed(SetValueResponseFut<PlaybackSpeed>), |
| PlayingOrder(SetValueResponseFut<PlayingOrder>), |
| } |
| |
| impl Future for PendingWriteResponseFut { |
| type Output = Option<PendingWriteResponse>; |
| |
| fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { |
| match self.as_mut().get_mut() { |
| Self::ControlPoint { peer_id, opcode, fut } => fut.poll_unpin(cx).map(|res| { |
| let code = res.unwrap_or(ControlPointResultCode::CommandCannotBeCompleted); |
| Some(PendingWriteResponse::ControlPoint { |
| peer_id: *peer_id, |
| opcode: opcode.clone(), |
| result_code: code, |
| }) |
| }), |
| Self::TrackPosition(fut) => { |
| fut.poll_unpin(cx).map(|res| res.map(PendingWriteResponse::TrackPosition)) |
| } |
| Self::PlaybackSpeed(fut) => { |
| fut.poll_unpin(cx).map(|res| res.map(PendingWriteResponse::PlaybackSpeed)) |
| } |
| Self::PlayingOrder(fut) => { |
| fut.poll_unpin(cx).map(|res| res.map(PendingWriteResponse::PlayingOrder)) |
| } |
| } |
| } |
| } |
| |
| /// Events produced by an [`McsServer`] stream representing client actions that |
| /// require additional upper-layer application processing. |
| #[derive(Debug)] |
| pub enum McsServerEvent { |
| /// Request to set the track playback position. |
| SetTrackPosition { |
| peer_id: PeerId, |
| position: TrackPosition, |
| responder: SetValueResponder<TrackPosition>, |
| }, |
| /// Request to set the playback speed. |
| SetPlaybackSpeed { |
| peer_id: PeerId, |
| speed: PlaybackSpeed, |
| responder: SetValueResponder<PlaybackSpeed>, |
| }, |
| /// Request to set the playing order. |
| SetPlayingOrder { |
| peer_id: PeerId, |
| order: PlayingOrder, |
| responder: SetValueResponder<PlayingOrder>, |
| }, |
| /// Request to execute a media control command. |
| ControlPointCommand { |
| opcode: MediaControlOpcode, |
| responder: SetValueResponder<ControlPointResultCode>, |
| }, |
| } |
| |
| /// State of the playing order characteristics when enabled on the server. |
| #[derive(Debug, Clone, PartialEq, Eq)] |
| struct PlayingOrderState { |
| /// Currently selected playing order. |
| current: PlayingOrder, |
| /// Bitmask of supported playing orders. |
| supported: SupportedPlayingOrders, |
| } |
| |
| impl PlayingOrderState { |
| /// Creates a new `PlayingOrderState` with the default current order |
| /// (`PlayingOrder::InOrderOnce`) and the given supported bitmask. |
| fn new(supported: SupportedPlayingOrders) -> Self { |
| Self { current: PlayingOrder::default(), supported } |
| } |
| |
| /// Sets the playing order if it is supported. |
| /// |
| /// Returns `true` if `order` was supported and applied, or `false` if |
| /// ignored because it was unsupported (MCS v1.0.1 Section 3.15.1). |
| fn set_order(&mut self, order: PlayingOrder) -> bool { |
| if self.supported.contains(order.into()) { |
| self.current = order; |
| true |
| } else { |
| false |
| } |
| } |
| } |
| |
| /// Local state of the characteristics in this server. |
| #[derive(Debug, Clone, PartialEq, Eq)] |
| struct McsLocalState { |
| /// Content Control ID (CCID) identifying this media service instance. |
| ccid: u8, |
| /// Human-readable media player application name. |
| player_name: String, |
| /// Title of the currently selected track (empty if no track loaded). |
| track_title: String, |
| /// Total duration of the current track (MCS v1.0.1 Section 3.6). |
| track_duration: TrackDuration, |
| /// Base playback position of the current track (MCS v1.0.1 Section 3.7). |
| track_position: TrackPosition, |
| /// Timestamp when `track_position` was set or updated. |
| position_updated_at: Option<std::time::Instant>, |
| /// Current player activity state. |
| media_state: MediaState, |
| /// URL pointing to media player icon graphic, if supported. |
| icon_url: Option<String>, |
| /// Playback speed multiplier, if supported. |
| playback_speed: Option<PlaybackSpeed>, |
| /// Seeking speed factor, if supported. |
| seeking_speed: Option<SeekingSpeed>, |
| /// Playing order and supported playing orders, if supported. |
| playing_orders: Option<PlayingOrderState>, |
| /// Supported media control point opcodes, if Media Control Point is |
| /// supported. |
| supported_opcodes: Option<SupportedOpcodes>, |
| } |
| |
| impl McsLocalState { |
| /// Creates a new [`McsLocalState`] with default values for an inactive |
| /// player with no track loaded per MCS v1.0.1 Section 3. |
| fn new(ccid: u8, player_name: impl Into<String>) -> Self { |
| Self { |
| ccid, |
| player_name: player_name.into(), |
| track_title: String::new(), |
| track_duration: TrackDuration::Unknown, |
| track_position: TrackPosition::Unavailable, |
| position_updated_at: None, |
| media_state: MediaState::Inactive, |
| icon_url: None, |
| playback_speed: None, |
| seeking_speed: None, |
| playing_orders: None, |
| supported_opcodes: None, |
| } |
| } |
| |
| /// Calculates the instantaneous track position based on elapsed playback |
| /// time (MCS v1.0.1 Section 3.7). |
| fn current_track_position(&self) -> TrackPosition { |
| // Per MCS Section 3.17, an inactive player has no current track, so it is |
| // unavailable. |
| if self.media_state == MediaState::Inactive { |
| return TrackPosition::Unavailable; |
| } |
| |
| // Playback timing has not started since the track was just initialized. |
| let Some(updated_at) = self.position_updated_at else { |
| return self.track_position; |
| }; |
| |
| // Playback is paused/seeking, so the position remains fixed. |
| if self.media_state != MediaState::Playing { |
| return self.track_position; |
| } |
| |
| let base = match (self.track_position, self.track_duration) { |
| (TrackPosition::FromStart(base), _) => base, |
| (TrackPosition::FromEnd(end), TrackDuration::Duration(total)) => { |
| total.saturating_sub(end) |
| } |
| _ => return self.track_position, |
| }; |
| |
| // TODO(b/540400364): Factor in optional playback speed when supported |
| let mut current = base + updated_at.elapsed(); |
| if let TrackDuration::Duration(total) = self.track_duration { |
| current = current.min(total); |
| } |
| TrackPosition::FromStart(current) |
| } |
| |
| /// Sets the track position and records the update timestamp. |
| fn set_track_position(&mut self, position: TrackPosition) { |
| self.track_position = position; |
| self.position_updated_at = Some(std::time::Instant::now()); |
| } |
| |
| /// Sets the playback speed. |
| fn set_playback_speed(&mut self, speed: PlaybackSpeed) { |
| self.playback_speed = Some(speed); |
| } |
| |
| /// Sets the playing order if playing order is supported. |
| fn set_playing_order(&mut self, order: PlayingOrder) { |
| if let Some(ref mut orders) = self.playing_orders { |
| let _ = orders.set_order(order); |
| } |
| } |
| |
| #[cfg(test)] |
| pub(crate) fn set_media_state(&mut self, state: MediaState) { |
| self.media_state = state; |
| } |
| |
| #[cfg(test)] |
| pub(crate) fn set_supported_opcodes(&mut self, opcodes: SupportedOpcodes) { |
| self.supported_opcodes = Some(opcodes); |
| } |
| |
| /// Reads the characteristic bytes for `handle` at `offset`. |
| fn handle_read(&self, handle: Handle, offset: usize) -> Result<Vec<u8>, GattError> { |
| let read_at_offset = |
| |bytes: &[u8]| bytes.get(offset..).map(Vec::from).ok_or(GattError::InvalidOffset); |
| |
| match handle { |
| // Track Changed and Media Control Point are not readable per MCS v1.0.1 Section 3, |
| // Table 3.1. |
| TRACK_CHANGED_HANDLE | MEDIA_CONTROL_POINT_HANDLE => Err(GattError::ReadNotPermitted), |
| MEDIA_PLAYER_NAME_HANDLE => read_at_offset(self.player_name.as_bytes()), |
| TRACK_TITLE_HANDLE => read_at_offset(self.track_title.as_bytes()), |
| TRACK_DURATION_HANDLE => read_at_offset(&self.track_duration.raw_10ms().to_le_bytes()), |
| TRACK_POSITION_HANDLE => { |
| read_at_offset(&self.current_track_position().raw_10ms().to_le_bytes()) |
| } |
| MEDIA_STATE_HANDLE => read_at_offset(&[self.media_state.into()]), |
| CONTENT_CONTROL_ID_HANDLE => read_at_offset(&[self.ccid]), |
| MEDIA_PLAYER_ICON_URL_HANDLE => { |
| let url = self.icon_url.as_ref().ok_or(GattError::InvalidHandle)?; |
| read_at_offset(url.as_bytes()) |
| } |
| PLAYBACK_SPEED_HANDLE => { |
| let speed = self.playback_speed.ok_or(GattError::InvalidHandle)?; |
| read_at_offset(&[speed.into()]) |
| } |
| SEEKING_SPEED_HANDLE => { |
| let speed = self.seeking_speed.ok_or(GattError::InvalidHandle)?; |
| read_at_offset(&[speed.into()]) |
| } |
| PLAYING_ORDER_HANDLE => { |
| let orders = self.playing_orders.as_ref().ok_or(GattError::InvalidHandle)?; |
| read_at_offset(&[orders.current.into()]) |
| } |
| PLAYING_ORDERS_SUPPORTED_HANDLE => { |
| let orders = self.playing_orders.as_ref().ok_or(GattError::InvalidHandle)?; |
| read_at_offset(&orders.supported.bits().to_le_bytes()) |
| } |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE => { |
| let opcodes = self.supported_opcodes.ok_or(GattError::InvalidHandle)?; |
| read_at_offset(&opcodes.bits().to_le_bytes()) |
| } |
| _ => Err(GattError::InvalidHandle), |
| } |
| } |
| } |
| |
| /// Builder for configuring an MCS or GMCS GATT service. |
| // TODO(b/549911651): Add support for OTS (Object Transfer Service) integration. |
| #[derive(Debug, Clone, PartialEq)] |
| pub struct McsServerBuilder { |
| /// Service UUID assigned to the service. |
| service_uuid: Uuid, |
| /// Local state of the characteristics in this server. |
| state: McsLocalState, |
| } |
| |
| impl McsServerBuilder { |
| /// Creates a builder for the Generic Media Control Service. |
| pub fn generic(ccid: u8, player_name: impl Into<String>) -> Self { |
| Self::new(GENERIC_MEDIA_CONTROL_SERVICE_UUID, ccid, player_name) |
| } |
| |
| /// Creates a builder for an application-specific Media Control Service. |
| pub fn instance(ccid: u8, player_name: impl Into<String>) -> Self { |
| Self::new(MEDIA_CONTROL_SERVICE_UUID, ccid, player_name) |
| } |
| |
| fn new(service_uuid: Uuid, ccid: u8, player_name: impl Into<String>) -> Self { |
| Self { service_uuid, state: McsLocalState::new(ccid, player_name) } |
| } |
| |
| /// Enables media player icon URL support. |
| pub fn with_icon_url(mut self, url: impl Into<String>) -> Self { |
| self.state.icon_url = Some(url.into()); |
| self |
| } |
| |
| /// Enables playback and seeking speed support with default speeds. |
| pub fn with_player_speeds(mut self) -> Self { |
| self.state.playback_speed = Some(PlaybackSpeed::NORMAL); |
| self.state.seeking_speed = Some(SeekingSpeed::NOT_SEEKING); |
| self |
| } |
| |
| /// Enables playing order support with the given supported playing orders. |
| pub fn with_playing_orders(mut self, supported: SupportedPlayingOrders) -> Self { |
| self.state.playing_orders = Some(PlayingOrderState::new(supported)); |
| self |
| } |
| |
| /// Enables media control operations with the given supported opcodes. |
| pub fn with_supported_operations(mut self, supported: SupportedOpcodes) -> Self { |
| self.state.supported_opcodes = Some(supported); |
| self |
| } |
| |
| /// Constructs the GATT [`ServiceDefinition`] containing the mandatory |
| /// characteristics per MCS v1.0.1 Section 3 and any configured optional |
| /// characteristics. |
| pub fn build_service_definition(&self) -> Result<ServiceDefinition, Error> { |
| // The local `ServiceId` is derived from the provided `CCID`. This is valid |
| // because the CCID must be unique across all MCS/GMCS instances on the |
| // host server. Adding a duplicate characteristic will result in an |
| // Error. |
| let mut service_def = ServiceDefinition::new( |
| ServiceId::new(self.state.ccid.into()), |
| self.service_uuid, |
| ServiceKind::Primary, |
| ); |
| |
| for chrc in mandatory_characteristics() { |
| service_def.add_characteristic(chrc)?; |
| } |
| |
| if self.state.icon_url.is_some() { |
| service_def.add_characteristic(build_characteristic( |
| MEDIA_PLAYER_ICON_URL_HANDLE, |
| MEDIA_PLAYER_ICON_URL_UUID, |
| CharacteristicProperty::Read, |
| ))?; |
| } |
| if self.state.playback_speed.is_some() { |
| service_def.add_characteristic(build_characteristic( |
| PLAYBACK_SPEED_HANDLE, |
| PLAYBACK_SPEED_UUID, |
| CharacteristicProperties::READ_WRITE_NOTIFY, |
| ))?; |
| } |
| if self.state.seeking_speed.is_some() { |
| service_def.add_characteristic(build_characteristic( |
| SEEKING_SPEED_HANDLE, |
| SEEKING_SPEED_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ))?; |
| } |
| if self.state.playing_orders.is_some() { |
| service_def.add_characteristic(build_characteristic( |
| PLAYING_ORDER_HANDLE, |
| PLAYING_ORDER_UUID, |
| CharacteristicProperties::READ_WRITE_NOTIFY, |
| ))?; |
| service_def.add_characteristic(build_characteristic( |
| PLAYING_ORDERS_SUPPORTED_HANDLE, |
| PLAYING_ORDERS_SUPPORTED_UUID, |
| CharacteristicProperty::Read, |
| ))?; |
| } |
| if self.state.supported_opcodes.is_some() { |
| service_def.add_characteristic(build_characteristic( |
| MEDIA_CONTROL_POINT_HANDLE, |
| MEDIA_CONTROL_POINT_UUID, |
| CharacteristicProperties::WRITE_NOTIFY, |
| ))?; |
| service_def.add_characteristic(build_characteristic( |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE, |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ))?; |
| } |
| |
| Ok(service_def) |
| } |
| |
| /// Builds an [`McsServer`] configured with this builder. |
| pub fn build<T: bt_gatt::ServerTypes>(self) -> Result<McsServer<T>, Error> { |
| let service_def = self.build_service_definition()?; |
| Ok(McsServer::new(service_def, self.state)) |
| } |
| } |
| |
| /// An instance of a Media Control Service (MCS) or Generic Media Control |
| /// Service (GMCS) GATT server. |
| #[pin_project(project = McsServerProj)] |
| pub struct McsServer<T: bt_gatt::ServerTypes> { |
| service_def: ServiceDefinition, |
| #[pin] |
| local_service: LocalServiceState<T>, |
| /// Local state of the characteristics in this server. |
| state: McsLocalState, |
| #[pin] |
| pending_write_responses: FuturesUnordered<PendingWriteResponseFut>, |
| } |
| |
| impl<T: bt_gatt::ServerTypes> McsServer<T> { |
| fn new(service_def: ServiceDefinition, state: McsLocalState) -> Self { |
| Self { |
| service_def, |
| local_service: Default::default(), |
| state, |
| pending_write_responses: FuturesUnordered::new(), |
| } |
| } |
| |
| /// Returns true if this server is a GMCS server. |
| pub fn is_generic_service(&self) -> bool { |
| self.service_def.uuid() == GENERIC_MEDIA_CONTROL_SERVICE_UUID |
| } |
| |
| /// Returns true if the server has successfully published the GATT service. |
| pub fn is_published(&self) -> bool { |
| self.local_service.is_published() |
| } |
| |
| /// Publishes the service to the GATT database. |
| pub fn publish(&mut self, server: T::Server) -> Result<(), Error> { |
| let LocalServiceState::NotPublished { waker } = &mut self.local_service else { |
| return Err(Error::AlreadyPublished); |
| }; |
| |
| let waker = waker.take(); |
| self.local_service = |
| LocalServiceState::Preparing { fut: server.prepare(self.service_def.clone()) }; |
| |
| if let Some(w) = waker { |
| w.wake(); |
| } |
| |
| Ok(()) |
| } |
| |
| #[cfg(test)] |
| pub(crate) fn set_media_state(&mut self, state: MediaState) { |
| self.state.set_media_state(state); |
| } |
| |
| #[cfg(test)] |
| pub(crate) fn set_supported_opcodes(&mut self, opcodes: SupportedOpcodes) { |
| self.state.set_supported_opcodes(opcodes); |
| } |
| } |
| |
| impl<'a, T: bt_gatt::ServerTypes> McsServerProj<'a, T> { |
| fn handle_read<R: ReadResponder>(&self, handle: Handle, offset: usize, responder: R) { |
| match self.state.handle_read(handle, offset) { |
| Ok(bytes) => responder.respond(&bytes), |
| Err(err) => responder.error(err), |
| } |
| } |
| |
| fn handle_write<W: WriteResponder>( |
| &mut self, |
| peer_id: PeerId, |
| handle: Handle, |
| offset: u32, |
| value: &[u8], |
| responder: W, |
| ) -> Option<McsServerEvent> { |
| // Reject non-zero offsets as MCS writable characteristics do not support |
| // long writes. |
| if offset != 0 { |
| responder.error(GattError::InvalidOffset); |
| return None; |
| } |
| |
| match handle { |
| MEDIA_CONTROL_POINT_HANDLE => { |
| self.handle_control_point_write(peer_id, value, responder) |
| } |
| TRACK_POSITION_HANDLE => { |
| if value.len() != 4 { |
| responder.error(GattError::InvalidAttributeValueLength); |
| return None; |
| } |
| let raw = i32::from_le_bytes(value.try_into().unwrap()); |
| let position = TrackPosition::from_raw_10ms(raw); |
| responder.acknowledge(); |
| let (fut, resp) = SetValueResponseFut::create(); |
| self.pending_write_responses.push(PendingWriteResponseFut::TrackPosition(fut)); |
| Some(McsServerEvent::SetTrackPosition { peer_id, position, responder: resp }) |
| } |
| PLAYBACK_SPEED_HANDLE => { |
| if self.state.playback_speed.is_none() { |
| responder.error(GattError::InvalidHandle); |
| return None; |
| } |
| if value.len() != 1 { |
| responder.error(GattError::InvalidAttributeValueLength); |
| return None; |
| } |
| let speed = PlaybackSpeed::from(value[0] as i8); |
| responder.acknowledge(); |
| let (fut, resp) = SetValueResponseFut::create(); |
| self.pending_write_responses.push(PendingWriteResponseFut::PlaybackSpeed(fut)); |
| Some(McsServerEvent::SetPlaybackSpeed { peer_id, speed, responder: resp }) |
| } |
| PLAYING_ORDER_HANDLE => { |
| let Some(ref playing_orders) = self.state.playing_orders else { |
| responder.error(GattError::InvalidHandle); |
| return None; |
| }; |
| if value.len() != 1 { |
| responder.error(GattError::InvalidAttributeValueLength); |
| return None; |
| } |
| let Ok(order) = PlayingOrder::try_from(value[0]) else { |
| // Invalid playing order values (MCS v1.0.1 Section 3.15.1). |
| responder.acknowledge(); |
| return None; |
| }; |
| if !playing_orders.supported.contains(order.into()) { |
| // Unsupported playing order is ignored per MCS v1.0.1 Section 3.15.1. |
| responder.acknowledge(); |
| return None; |
| } |
| responder.acknowledge(); |
| let (fut, resp) = SetValueResponseFut::create(); |
| self.pending_write_responses.push(PendingWriteResponseFut::PlayingOrder(fut)); |
| Some(McsServerEvent::SetPlayingOrder { peer_id, order, responder: resp }) |
| } |
| // Read-only characteristics cannot be written to. |
| MEDIA_PLAYER_NAME_HANDLE |
| | TRACK_CHANGED_HANDLE |
| | TRACK_TITLE_HANDLE |
| | TRACK_DURATION_HANDLE |
| | MEDIA_PLAYER_ICON_URL_HANDLE |
| | SEEKING_SPEED_HANDLE |
| | PLAYING_ORDERS_SUPPORTED_HANDLE |
| | MEDIA_STATE_HANDLE |
| | MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE |
| | CONTENT_CONTROL_ID_HANDLE => { |
| responder.error(GattError::WriteNotPermitted); |
| None |
| } |
| // TODO(b/549911651): Add support for optional characteristics (e.g. Search Control |
| // Point and OTS object IDs). |
| _ => { |
| responder.error(GattError::InvalidHandle); |
| None |
| } |
| } |
| } |
| |
| fn handle_control_point_write<W: WriteResponder>( |
| &mut self, |
| peer_id: PeerId, |
| value: &[u8], |
| responder: W, |
| ) -> Option<McsServerEvent> { |
| // The Media Control Point characteristic is optional - reject all requests if |
| // it is not enabled on this server. |
| let Some(supported_opcodes) = self.state.supported_opcodes else { |
| responder.error(GattError::InvalidHandle); |
| return None; |
| }; |
| |
| let Some(&raw_opcode) = value.first() else { |
| responder.error(GattError::InvalidAttributeValueLength); |
| return None; |
| }; |
| |
| let notify_error = |responder: W, result_code: ControlPointResultCode| { |
| responder.acknowledge(); |
| self.local_service.notify( |
| &MEDIA_CONTROL_POINT_HANDLE, |
| &[raw_opcode, result_code.into()], |
| &[peer_id], |
| ); |
| None |
| }; |
| |
| // Validate opcode decoding and server opcode support (see MCS v1.0.1 Section |
| // 3.18.2). |
| let opcode = match MediaControlOpcode::decode(value) { |
| (Ok(opcode), _) if supported_opcodes.contains((&opcode).into()) => opcode, |
| _ => return notify_error(responder, ControlPointResultCode::OpcodeNotSupported), |
| }; |
| |
| // Reject commands if the media player is inactive. |
| // TODO(b/540400364): The spec also technically allows the handling of the |
| // command if the media player supports it with no active track. Revisit |
| // this if needed. |
| if self.state.media_state == MediaState::Inactive { |
| return notify_error(responder, ControlPointResultCode::MediaPlayerInactive); |
| } |
| |
| // Acknowledge the GATT write and produce a new event for the request. |
| responder.acknowledge(); |
| let (fut, responder) = SetValueResponseFut::create(); |
| self.pending_write_responses.push(PendingWriteResponseFut::ControlPoint { |
| peer_id, |
| opcode: opcode.clone(), |
| fut, |
| }); |
| Some(McsServerEvent::ControlPointCommand { opcode, responder }) |
| } |
| |
| /// Attempts to handle a pending write response from the upper layer |
| /// application. |
| fn handle_pending_write_response(&mut self, response: PendingWriteResponse) { |
| match response { |
| PendingWriteResponse::ControlPoint { peer_id, opcode, result_code } => { |
| // Dispatches a Media Control Point GATT notification containing |
| // [opcode, result_code] to the peer per MCS v1.0.1 Section 3.18.2. |
| self.local_service.notify( |
| &MEDIA_CONTROL_POINT_HANDLE, |
| &[opcode.raw_opcode(), result_code.into()], |
| &[peer_id], |
| ); |
| } |
| PendingWriteResponse::TrackPosition(position) => { |
| self.state.set_track_position(position); |
| self.local_service.notify( |
| &TRACK_POSITION_HANDLE, |
| &position.raw_10ms().to_le_bytes(), |
| &[], |
| ); |
| } |
| PendingWriteResponse::PlaybackSpeed(speed) => { |
| self.state.set_playback_speed(speed); |
| self.local_service.notify(&PLAYBACK_SPEED_HANDLE, &[speed.into()], &[]); |
| } |
| PendingWriteResponse::PlayingOrder(order) => { |
| self.state.set_playing_order(order); |
| self.local_service.notify(&PLAYING_ORDER_HANDLE, &[order.into()], &[]); |
| } |
| } |
| } |
| } |
| |
| impl<T: bt_gatt::ServerTypes> Stream for McsServer<T> { |
| type Item = Result<McsServerEvent, Error>; |
| |
| fn poll_next( |
| mut self: std::pin::Pin<&mut Self>, |
| cx: &mut Context<'_>, |
| ) -> Poll<Option<Self::Item>> { |
| loop { |
| let mut this = self.as_mut().project(); |
| // Drain any completed asynchronous write responses before processing new |
| // GATT events. |
| if let Poll::Ready(Some(response)) = this.pending_write_responses.as_mut().poll_next(cx) |
| { |
| if let Some(response) = response { |
| this.handle_pending_write_response(response); |
| } |
| continue; |
| } |
| |
| let gatt_event = match futures::ready!(this.local_service.as_mut().poll_next(cx)) { |
| None => { |
| // Continue polling until all in-flight responses have been drained before |
| // terminating the stream. |
| if this.pending_write_responses.is_empty() { |
| this.local_service.as_mut().set(LocalServiceState::Terminated); |
| return Poll::Ready(None); |
| } else { |
| return Poll::Pending; |
| } |
| } |
| Some(Err(e)) => return Poll::Ready(Some(Err(e))), |
| Some(Ok(event)) => event, |
| }; |
| |
| match gatt_event { |
| ServiceEvent::Read { peer_id: _, handle, offset, responder } => { |
| this.handle_read(handle, offset as usize, responder); |
| } |
| ServiceEvent::Write { peer_id, handle, offset, value, responder } => { |
| let value_bytes = value.to_owned(); |
| if let Some(event) = |
| this.handle_write(peer_id, handle, offset, &value_bytes, responder) |
| { |
| return Poll::Ready(Some(Ok(event))); |
| } |
| } |
| _ => continue, |
| } |
| } |
| } |
| } |
| |
| #[cfg(test)] |
| mod tests { |
| use super::*; |
| |
| use assert_matches::assert_matches; |
| use bt_gatt::test_utils::{FakeServer, FakeServerEvent, FakeTypes}; |
| use futures::StreamExt; |
| |
| #[test] |
| fn builder_generic_service_definition() { |
| let builder = McsServerBuilder::generic(0x42, "Test Generic Player"); |
| let service_def = |
| builder.build_service_definition().expect("service definition builds successfully"); |
| assert_eq!(service_def.uuid(), GENERIC_MEDIA_CONTROL_SERVICE_UUID); |
| assert_eq!(service_def.id(), ServiceId::new(0x42)); |
| assert_eq!(service_def.kind(), ServiceKind::Primary); |
| assert_eq!(service_def.characteristics().count(), 7); |
| } |
| |
| #[test] |
| fn builder_instance_service_definition() { |
| let builder = McsServerBuilder::instance(0x07, "Test Instance Player"); |
| let service_def = builder |
| .build_service_definition() |
| .expect("instance service definition builds successfully"); |
| assert_eq!(service_def.uuid(), MEDIA_CONTROL_SERVICE_UUID); |
| assert_eq!(service_def.id(), ServiceId::new(0x07)); |
| assert_eq!(service_def.kind(), ServiceKind::Primary); |
| assert_eq!(service_def.characteristics().count(), 7); |
| } |
| |
| #[test] |
| fn builder_with_optional_characteristics_service_definition() { |
| let builder = McsServerBuilder::generic(0x10, "Full Player") |
| .with_icon_url("https://example.com/icon.png") |
| .with_player_speeds() |
| .with_playing_orders(SupportedPlayingOrders::default()) |
| .with_supported_operations(SupportedOpcodes::default()); |
| |
| let service_def = |
| builder.build_service_definition().expect("service definition builds successfully"); |
| assert_eq!(service_def.characteristics().count(), 14); |
| |
| let characteristics: Vec<&Characteristic> = service_def.characteristics().collect(); |
| let expected_characteristics = [ |
| mandatory_characteristics()[0].clone(), |
| mandatory_characteristics()[1].clone(), |
| mandatory_characteristics()[2].clone(), |
| mandatory_characteristics()[3].clone(), |
| mandatory_characteristics()[4].clone(), |
| mandatory_characteristics()[5].clone(), |
| mandatory_characteristics()[6].clone(), |
| build_characteristic( |
| MEDIA_PLAYER_ICON_URL_HANDLE, |
| MEDIA_PLAYER_ICON_URL_UUID, |
| CharacteristicProperty::Read, |
| ), |
| build_characteristic( |
| PLAYBACK_SPEED_HANDLE, |
| PLAYBACK_SPEED_UUID, |
| CharacteristicProperties::READ_WRITE_NOTIFY, |
| ), |
| build_characteristic( |
| SEEKING_SPEED_HANDLE, |
| SEEKING_SPEED_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| build_characteristic( |
| PLAYING_ORDER_HANDLE, |
| PLAYING_ORDER_UUID, |
| CharacteristicProperties::READ_WRITE_NOTIFY, |
| ), |
| build_characteristic( |
| PLAYING_ORDERS_SUPPORTED_HANDLE, |
| PLAYING_ORDERS_SUPPORTED_UUID, |
| CharacteristicProperty::Read, |
| ), |
| build_characteristic( |
| MEDIA_CONTROL_POINT_HANDLE, |
| MEDIA_CONTROL_POINT_UUID, |
| CharacteristicProperties::WRITE_NOTIFY, |
| ), |
| build_characteristic( |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE, |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_UUID, |
| CharacteristicProperties::READ_NOTIFY, |
| ), |
| ]; |
| |
| assert_eq!(characteristics.len(), expected_characteristics.len()); |
| for (i, expected) in expected_characteristics.iter().enumerate() { |
| let chrc = characteristics[i]; |
| assert_eq!(chrc.handle, expected.handle); |
| assert_eq!(chrc.uuid, expected.uuid); |
| assert_eq!(chrc.properties, expected.properties); |
| } |
| } |
| |
| #[test] |
| fn mandatory_characteristic_handles_uuids_and_properties() { |
| let builder = McsServerBuilder::generic(0x01, "Player"); |
| let service_def = |
| builder.build_service_definition().expect("service definition builds successfully"); |
| let characteristics: Vec<&Characteristic> = service_def.characteristics().collect(); |
| |
| let expected_mandatory = mandatory_characteristics(); |
| assert_eq!(characteristics.len(), expected_mandatory.len()); |
| |
| for (i, expected) in expected_mandatory.iter().enumerate() { |
| let chrc = characteristics[i]; |
| assert_eq!(chrc.handle, expected.handle); |
| assert_eq!(chrc.uuid, expected.uuid); |
| assert_eq!(chrc.properties, expected.properties); |
| |
| assert_eq!( |
| chrc.permissions.read.map(|s| s.encryption), |
| expected.permissions.read.map(|s| s.encryption) |
| ); |
| assert_eq!( |
| chrc.permissions.write.map(|s| s.encryption), |
| expected.permissions.write.map(|s| s.encryption) |
| ); |
| assert_eq!( |
| chrc.permissions.update.map(|s| s.encryption), |
| expected.permissions.update.map(|s| s.encryption) |
| ); |
| } |
| } |
| |
| #[test] |
| fn service_id_tracks_ccid() { |
| let builder1 = McsServerBuilder::instance(0x05, "Player 1"); |
| let builder2 = McsServerBuilder::generic(0x05, "Player 2"); |
| let builder3 = McsServerBuilder::instance(0x06, "Player 3"); |
| |
| let def1 = |
| builder1.build_service_definition().expect("service definition builds successfully"); |
| let def2 = |
| builder2.build_service_definition().expect("service definition builds successfully"); |
| let def3 = |
| builder3.build_service_definition().expect("service definition builds successfully"); |
| |
| assert_eq!(def1.id(), def2.id()); |
| assert_eq!(def1.id(), ServiceId::new(5)); |
| assert_ne!(def1.id(), def3.id()); |
| assert_eq!(def3.id(), ServiceId::new(6)); |
| } |
| |
| #[test] |
| fn build_server_success() { |
| let generic_server: McsServer<FakeTypes> = |
| McsServerBuilder::generic(0x42, "Test Generic Player") |
| .build() |
| .expect("generic server builds successfully"); |
| assert!(!generic_server.is_published()); |
| assert!(generic_server.is_generic_service()); |
| |
| let instance_server: McsServer<FakeTypes> = |
| McsServerBuilder::instance(0x42, "Test Instance Player") |
| .build() |
| .expect("instance server builds successfully"); |
| assert!(!instance_server.is_published()); |
| assert!(!instance_server.is_generic_service()); |
| } |
| |
| #[test] |
| fn publish_server_success() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let mut server: McsServer<FakeTypes> = |
| McsServerBuilder::generic(0x42, "Test Generic Player") |
| .build() |
| .expect("server builds successfully"); |
| assert!(!server.is_published()); |
| |
| let (fake_gatt_server, _event_receiver) = FakeServer::new(); |
| let Poll::Pending = server.next().poll_unpin(&mut noop_cx) else { |
| panic!("Should be pending before publish"); |
| }; |
| |
| server.publish(fake_gatt_server).expect("publish succeeds"); |
| |
| // Advance state: Preparing -> Published |
| let Poll::Pending = server.next().poll_unpin(&mut noop_cx) else { |
| panic!("Should be pending after publish"); |
| }; |
| assert!(server.is_published()); |
| } |
| |
| #[test] |
| fn publish_server_already_published_error() { |
| let (fake_gatt_server, _event_receiver) = FakeServer::new(); |
| let mut server: McsServer<FakeTypes> = McsServerBuilder::generic(0x42, "Test Player") |
| .build() |
| .expect("server builds successfully"); |
| |
| server.publish(fake_gatt_server.clone()).expect("initial publish succeeds"); |
| let err = server.publish(fake_gatt_server); |
| assert_matches!(err, Err(Error::AlreadyPublished)); |
| } |
| |
| #[test] |
| fn duplicate_ccid_publish_error() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (fake_gatt_server, _event_receiver) = FakeServer::new(); |
| |
| let mut server1: McsServer<FakeTypes> = McsServerBuilder::instance(0x05, "Player 1") |
| .build() |
| .expect("server1 builds successfully"); |
| let mut server2: McsServer<FakeTypes> = McsServerBuilder::instance(0x05, "Player 2") |
| .build() |
| .expect("server2 builds successfully"); |
| |
| // The first server publishes successfully. |
| server1.publish(fake_gatt_server.clone()).expect("server1 publish call succeeds"); |
| let _ = server1.next().poll_unpin(&mut noop_cx); |
| assert!(server1.is_published()); |
| |
| // The GATT server rejects the second server attempting to publish with the |
| // duplicate CCID / ServiceId. |
| fake_gatt_server.set_next_prepare_result(Err(bt_gatt::types::Error::AlreadyPublished( |
| ServiceId::new(0x05), |
| ))); |
| server2.publish(fake_gatt_server).expect("server2 publish call succeeds"); |
| let poll_result = server2.next().poll_unpin(&mut noop_cx); |
| assert_matches!( |
| poll_result, |
| Poll::Ready(Some(Err(Error::Gatt(bt_gatt::types::Error::AlreadyPublished(_))))) |
| ); |
| } |
| |
| #[test] |
| fn server_stream_terminates_when_event_stream_closes() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (fake_gatt_server, _event_receiver) = FakeServer::new(); |
| let mut server: McsServer<FakeTypes> = McsServerBuilder::generic(0x42, "Test Player") |
| .build() |
| .expect("server builds successfully"); |
| |
| server.publish(fake_gatt_server).expect("publish succeeds"); |
| |
| // Advance to Published |
| let Poll::Pending = server.next().poll_unpin(&mut noop_cx) else { |
| panic!("Should be pending after publish"); |
| }; |
| assert!(server.is_published()); |
| |
| // Replace local_service events stream with a custom channel that can be |
| // explicitly closed. |
| let (sender, receiver) = futures::channel::mpsc::unbounded(); |
| let LocalServiceState::Published { service, .. } = |
| std::mem::replace(&mut server.local_service, LocalServiceState::Terminated) |
| else { |
| panic!("Expected server to be in Published state"); |
| }; |
| server.local_service = LocalServiceState::Published { service, events: receiver }; |
| |
| // Dropping the sender closes the event stream. |
| drop(sender); |
| |
| // Polling the server returns None indicating the stream has terminated. |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Ready(None)); |
| |
| // Subsequent polls on terminated state also return None. |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Ready(None)); |
| } |
| |
| #[track_caller] |
| fn expect_service_event( |
| events: &mut futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| ) -> FakeServerEvent { |
| let mut cx = Context::from_waker(futures::task::noop_waker_ref()); |
| match events.poll_next_unpin(&mut cx) { |
| Poll::Ready(Some(event)) => event, |
| x => panic!("Expected fake server event, got {x:?}"), |
| } |
| } |
| |
| #[track_caller] |
| fn expect_no_service_event( |
| events: &mut futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| ) { |
| let mut cx = Context::from_waker(futures::task::noop_waker_ref()); |
| assert_matches!(events.poll_next_unpin(&mut cx), Poll::Pending); |
| } |
| |
| fn setup_test_server( |
| builder: McsServerBuilder, |
| ) -> ( |
| McsServer<FakeTypes>, |
| FakeServer, |
| futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| ) { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (fake_gatt_server, mut event_receiver) = FakeServer::new(); |
| let mut server: McsServer<FakeTypes> = builder.build().expect("server builds successfully"); |
| server.publish(fake_gatt_server.clone()).expect("publish succeeds"); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| assert!(matches!( |
| expect_service_event(&mut event_receiver), |
| FakeServerEvent::Published { .. } |
| )); |
| (server, fake_gatt_server, event_receiver) |
| } |
| |
| fn setup_test_server_with_control_point( |
| peer: PeerId, |
| ) -> ( |
| McsServer<FakeTypes>, |
| FakeServer, |
| futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| Context<'static>, |
| ) { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_supported_operations(SupportedOpcodes::all()), |
| ); |
| fake_gatt_server.incoming_client_configuration( |
| peer, |
| server.service_def.id(), |
| MEDIA_CONTROL_POINT_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| (server, fake_gatt_server, event_receiver, noop_cx) |
| } |
| |
| fn assert_read_characteristic( |
| server: &mut McsServer<FakeTypes>, |
| fake_gatt_server: &FakeServer, |
| event_receiver: &mut futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| handle: Handle, |
| expected: &[u8], |
| ) { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| fake_gatt_server.incoming_read(peer, service_id, handle, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), expected); |
| } |
| |
| #[test] |
| fn read_mandatory_characteristics_default_values() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| |
| // Media Player Name |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| MEDIA_PLAYER_NAME_HANDLE, |
| b"Test Player", |
| ); |
| |
| // Track Title |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| TRACK_TITLE_HANDLE, |
| b"", |
| ); |
| |
| // Track Duration |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| TRACK_DURATION_HANDLE, |
| &(-1i32).to_le_bytes(), |
| ); |
| |
| // Track Position |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| TRACK_POSITION_HANDLE, |
| &(-1i32).to_le_bytes(), |
| ); |
| |
| // Media State |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| MEDIA_STATE_HANDLE, |
| &[MediaState::Inactive.into()], |
| ); |
| |
| // Content Control ID |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| CONTENT_CONTROL_ID_HANDLE, |
| &[0x42], |
| ); |
| } |
| |
| #[test] |
| fn read_optional_characteristics_configured_values() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_icon_url("https://example.com/icon.png") |
| .with_player_speeds() |
| .with_playing_orders( |
| SupportedPlayingOrders::IN_ORDER_ONCE | SupportedPlayingOrders::SHUFFLE_ONCE, |
| ) |
| .with_supported_operations(SupportedOpcodes::PLAY | SupportedOpcodes::PAUSE), |
| ); |
| |
| // Media Player Icon URL |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| MEDIA_PLAYER_ICON_URL_HANDLE, |
| b"https://example.com/icon.png", |
| ); |
| |
| // Playback Speed |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| PLAYBACK_SPEED_HANDLE, |
| &[0], |
| ); |
| |
| // Seeking Speed |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| SEEKING_SPEED_HANDLE, |
| &[0], |
| ); |
| |
| // Playing Order |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| PLAYING_ORDER_HANDLE, |
| &[PlayingOrder::InOrderOnce.into()], |
| ); |
| |
| // Playing Orders Supported |
| let expected_orders = |
| SupportedPlayingOrders::IN_ORDER_ONCE | SupportedPlayingOrders::SHUFFLE_ONCE; |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| PLAYING_ORDERS_SUPPORTED_HANDLE, |
| &expected_orders.bits().to_le_bytes(), |
| ); |
| |
| // Media Control Point Opcodes Supported |
| let expected_opcodes = SupportedOpcodes::PLAY | SupportedOpcodes::PAUSE; |
| assert_read_characteristic( |
| &mut server, |
| &fake_gatt_server, |
| &mut event_receiver, |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE, |
| &expected_opcodes.bits().to_le_bytes(), |
| ); |
| } |
| |
| #[test] |
| fn read_unconfigured_optional_characteristics_returns_invalid_handle() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Mandatory Only Player")); |
| |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| |
| let unconfigured_handles = [ |
| MEDIA_PLAYER_ICON_URL_HANDLE, |
| PLAYBACK_SPEED_HANDLE, |
| SEEKING_SPEED_HANDLE, |
| PLAYING_ORDER_HANDLE, |
| PLAYING_ORDERS_SUPPORTED_HANDLE, |
| MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE, |
| ]; |
| |
| for handle in unconfigured_handles { |
| fake_gatt_server.incoming_read(peer, service_id, handle, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded for handle {:?}", handle); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidHandle), |
| "handle {:?} should return InvalidHandle when unconfigured", |
| handle |
| ); |
| } |
| } |
| |
| #[test] |
| fn write_unconfigured_optional_characteristics_returns_invalid_handle() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Mandatory Only Player")); |
| |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| |
| let unconfigured_writes = [ |
| (PLAYBACK_SPEED_HANDLE, vec![0x00]), |
| (PLAYING_ORDER_HANDLE, vec![PlayingOrder::SingleOnce.into()]), |
| (MEDIA_CONTROL_POINT_HANDLE, vec![MediaControlOpcode::Play.raw_opcode()]), |
| ]; |
| |
| for (handle, value) in unconfigured_writes { |
| fake_gatt_server.incoming_write(peer, service_id, handle, 0, value); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| // Must NOT emit any McsServerEvent. |
| assert_matches!(poll_result, Poll::Pending); |
| let bt_gatt::test_utils::FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded for handle {:?}", handle); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidHandle), |
| "handle {:?} should return InvalidHandle on write when unconfigured", |
| handle |
| ); |
| } |
| } |
| |
| #[test] |
| fn read_string_characteristics_with_offset() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Long Player Name") |
| .with_icon_url("https://example.com/icon.png"), |
| ); |
| |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| |
| // Valid offset slice for Player Name |
| fake_gatt_server.incoming_read(peer, service_id, MEDIA_PLAYER_NAME_HANDLE, 5); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), b"Player Name"); |
| |
| // Valid offset slice for Icon URL |
| fake_gatt_server.incoming_read(peer, service_id, MEDIA_PLAYER_ICON_URL_HANDLE, 8); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), b"example.com/icon.png"); |
| |
| // Exact end offset returns empty slice |
| fake_gatt_server.incoming_read( |
| peer, |
| service_id, |
| MEDIA_PLAYER_NAME_HANDLE, |
| "Long Player Name".len() as u32, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), b""); |
| |
| // Out-of-bounds offset returns InvalidOffset |
| fake_gatt_server.incoming_read(peer, service_id, MEDIA_PLAYER_NAME_HANDLE, 100); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidOffset)); |
| |
| // Out-of-bounds offset on Icon URL returns InvalidOffset |
| fake_gatt_server.incoming_read(peer, service_id, MEDIA_PLAYER_ICON_URL_HANDLE, 100); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidOffset)); |
| } |
| |
| #[test] |
| fn read_non_readable_characteristics_returns_error() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_supported_operations(SupportedOpcodes::default()), |
| ); |
| |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| |
| // Track Changed is notify-only |
| fake_gatt_server.incoming_read(peer, service_id, TRACK_CHANGED_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ReadNotPermitted) |
| ); |
| |
| // Media Control Point is write/notify-only |
| fake_gatt_server.incoming_read(peer, service_id, MEDIA_CONTROL_POINT_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ReadNotPermitted) |
| ); |
| |
| // Unknown handle returns InvalidHandle |
| fake_gatt_server.incoming_read(peer, service_id, Handle(999), 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidHandle)); |
| } |
| |
| #[test] |
| fn write_track_position_success() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| server.state.media_state = MediaState::Paused; |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| fake_gatt_server.incoming_client_configuration( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| 0, |
| 4200i32.to_le_bytes().to_vec(), |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetTrackPosition { |
| peer_id, |
| position, |
| responder, |
| }))) = poll_result |
| else { |
| panic!("expected SetTrackPosition event, got {poll_result:?}"); |
| }; |
| assert_eq!(peer_id, peer); |
| assert_eq!(position, TrackPosition::from_raw_10ms(4200)); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Upper layer confirms the set track position. |
| responder.send(position); |
| |
| // Server processes the confirmation, updates state, and sends notification. |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, TRACK_POSITION_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, 4200i32.to_le_bytes().to_vec()); |
| |
| // Subsequent read reflects the updated track position |
| fake_gatt_server.incoming_read(peer, service_id, TRACK_POSITION_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), 4200i32.to_le_bytes()); |
| } |
| |
| #[test] |
| fn write_track_position_invalid_length_or_offset() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| // Invalid offset (non-zero) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| 1, |
| 4200i32.to_le_bytes().to_vec(), |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidOffset)); |
| |
| // Invalid value length (3 bytes instead of 4) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| 0, |
| vec![0x01, 0x02, 0x03], |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidAttributeValueLength) |
| ); |
| } |
| |
| #[test] |
| fn write_playback_speed_success() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player").with_player_speeds()); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| fake_gatt_server.incoming_client_configuration( |
| peer, |
| service_id, |
| PLAYBACK_SPEED_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| PLAYBACK_SPEED_HANDLE, |
| 0, |
| vec![-64i8 as u8], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetPlaybackSpeed { peer_id, speed, responder }))) = |
| poll_result |
| else { |
| panic!("expected SetPlaybackSpeed event, got {poll_result:?}"); |
| }; |
| assert_eq!(peer_id, peer); |
| assert_eq!(speed, PlaybackSpeed::HALF); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Upper layer confirms playback speed |
| responder.send(speed); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, PLAYBACK_SPEED_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![-64i8 as u8]); |
| |
| // Subsequent read reflects the updated playback speed |
| fake_gatt_server.incoming_read(peer, service_id, PLAYBACK_SPEED_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![-64i8 as u8]); |
| } |
| |
| #[test] |
| fn write_playback_speed_invalid_length_or_offset() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player").with_player_speeds()); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| // Invalid offset (non-zero) |
| fake_gatt_server.incoming_write(peer, service_id, PLAYBACK_SPEED_HANDLE, 1, vec![0x00]); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidOffset)); |
| |
| // Invalid value length (2 bytes instead of 1) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| PLAYBACK_SPEED_HANDLE, |
| 0, |
| vec![0x00, 0x01], |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidAttributeValueLength) |
| ); |
| } |
| |
| #[test] |
| fn write_playing_order_supported() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_playing_orders(SupportedPlayingOrders::all()), |
| ); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| fake_gatt_server.incoming_client_configuration( |
| peer, |
| service_id, |
| PLAYING_ORDER_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| PLAYING_ORDER_HANDLE, |
| 0, |
| vec![PlayingOrder::ShuffleOnce.into()], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetPlayingOrder { peer_id, order, responder }))) = |
| poll_result |
| else { |
| panic!("expected SetPlayingOrder event, got {poll_result:?}"); |
| }; |
| assert_eq!(peer_id, peer); |
| assert_eq!(order, PlayingOrder::ShuffleOnce); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Upper layer confirms playing order |
| responder.send(order); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, PLAYING_ORDER_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![PlayingOrder::ShuffleOnce.into()]); |
| |
| // Subsequent read reflects the updated playing order |
| fake_gatt_server.incoming_read(peer, service_id, PLAYING_ORDER_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![PlayingOrder::ShuffleOnce.into()]); |
| } |
| |
| #[test] |
| fn write_playing_order_unsupported_is_ignored() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_playing_orders(SupportedPlayingOrders::IN_ORDER_ONCE), |
| ); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| // Writing ShuffleOnce when only InOrderOnce is supported is ignored per MCS |
| // v1.0.1 Section 3.15.1 |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| PLAYING_ORDER_HANDLE, |
| 0, |
| vec![PlayingOrder::ShuffleOnce.into()], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Value remains unchanged (InOrderOnce) |
| fake_gatt_server.incoming_read(peer, service_id, PLAYING_ORDER_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![PlayingOrder::InOrderOnce.into()]); |
| } |
| |
| #[test] |
| fn write_read_only_characteristic_returns_error() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| let read_only_writes = [ |
| (MEDIA_PLAYER_NAME_HANDLE, b"New Name".to_vec()), |
| (TRACK_CHANGED_HANDLE, vec![]), |
| (TRACK_TITLE_HANDLE, b"New Title".to_vec()), |
| (TRACK_DURATION_HANDLE, 1000i32.to_le_bytes().to_vec()), |
| (MEDIA_PLAYER_ICON_URL_HANDLE, b"https://example.com".to_vec()), |
| (SEEKING_SPEED_HANDLE, vec![0x00]), |
| (PLAYING_ORDERS_SUPPORTED_HANDLE, vec![0x01, 0x00]), |
| (MEDIA_STATE_HANDLE, vec![0x01]), |
| (MEDIA_CONTROL_POINT_OPCODES_SUPPORTED_HANDLE, vec![0x01, 0x00, 0x00, 0x00]), |
| (CONTENT_CONTROL_ID_HANDLE, vec![0x99]), |
| ]; |
| |
| for (handle, value) in read_only_writes { |
| fake_gatt_server.incoming_write(peer, service_id, handle, 0, value); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded for handle {:?}", handle); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::WriteNotPermitted), |
| "handle {:?} should return WriteNotPermitted on write", |
| handle |
| ); |
| } |
| } |
| |
| #[test] |
| fn control_point_unsupported_opcode_immediate_notification() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| |
| // Write an unsupported/RFU opcode (0xEE) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![0xEE], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Immediate Media Control Point Notification sent with OPCODE_NOT_SUPPORTED |
| // (0x02) |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![0xEE, ControlPointResultCode::OpcodeNotSupported.into()]); |
| } |
| |
| #[test] |
| fn control_point_inactive_player_immediate_notification() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| |
| // Write Play opcode (0x01) while player is Inactive |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![0x01], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Immediate Media Control Point Notification sent with MEDIA_PLAYER_INACTIVE |
| // (0x03) |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![0x01, ControlPointResultCode::MediaPlayerInactive.into()]); |
| } |
| |
| #[test] |
| fn control_point_valid_command_success() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| |
| // Write Pause opcode (0x02) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![0x02], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { opcode, responder }))) = |
| poll_result |
| else { |
| panic!("expected ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(opcode, MediaControlOpcode::Pause); |
| |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Upper layer executes command and responds with Success |
| responder.send(ControlPointResultCode::Success); |
| |
| // Polling the server dispatches the notification |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![0x02, ControlPointResultCode::Success.into()]); |
| } |
| |
| #[test] |
| fn control_point_parameterized_command_success() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| |
| // MoveRelative opcode with 4-byte offset parameter |
| let opcode = MediaControlOpcode::MoveRelative(1500); |
| let mut write_val = vec![opcode.raw_opcode()]; |
| write_val.extend_from_slice(&1500i32.to_le_bytes()); |
| |
| fake_gatt_server.incoming_write(peer, service_id, MEDIA_CONTROL_POINT_HANDLE, 0, write_val); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { |
| opcode: rx_opcode, |
| responder, |
| }))) = poll_result |
| else { |
| panic!("expected ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(rx_opcode, opcode); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| responder.send(ControlPointResultCode::Success); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![opcode.raw_opcode(), ControlPointResultCode::Success.into()]); |
| } |
| |
| #[test] |
| fn control_point_truncated_parameter_sends_opcode_not_supported() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| |
| // MoveRelative requires 5 octets; provide only 2 octets |
| let raw_op = MediaControlOpcode::MoveRelative(0).raw_opcode(); |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![raw_op, 0x01], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![raw_op, ControlPointResultCode::OpcodeNotSupported.into()]); |
| } |
| |
| #[test] |
| fn control_point_empty_value_or_invalid_offset_returns_gatt_error() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| |
| // Non-zero offset returns InvalidOffset |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 1, |
| vec![0x01], |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!(value.unwrap_err(), bt_gatt::types::Error::Gatt(GattError::InvalidOffset)); |
| |
| // Empty value returns InvalidAttributeValueLength |
| fake_gatt_server.incoming_write(peer, service_id, MEDIA_CONTROL_POINT_HANDLE, 0, vec![]); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert_matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidAttributeValueLength) |
| ); |
| } |
| |
| #[test] |
| fn control_point_unsupported_opcode_on_active_player_returns_opcode_not_supported() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| server.set_supported_opcodes(SupportedOpcodes::PLAY); |
| |
| // Pause is valid in spec, but not in supported_opcodes for this server instance |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![MediaControlOpcode::Pause.raw_opcode()], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!( |
| value, |
| vec![ |
| MediaControlOpcode::Pause.raw_opcode(), |
| ControlPointResultCode::OpcodeNotSupported.into() |
| ] |
| ); |
| } |
| |
| #[test] |
| fn control_point_dropped_responder_sends_cannot_be_completed() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| |
| // Write Stop opcode (0x05) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![0x05], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { opcode, responder }))) = |
| poll_result |
| else { |
| panic!("expected ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(opcode, MediaControlOpcode::Stop); |
| |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Responder is dropped without calling send |
| drop(responder); |
| |
| // Polling the server dispatches the fallback notification |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![0x05, ControlPointResultCode::CommandCannotBeCompleted.into()]); |
| } |
| |
| #[test] |
| fn control_point_multiple_consecutive_notifications() { |
| let peer1 = PeerId(1); |
| let peer2 = PeerId(2); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer1); |
| server.set_media_state(MediaState::Playing); |
| |
| // Also configure notifications for peer2 |
| fake_gatt_server.incoming_client_configuration( |
| peer2, |
| server.service_def.id(), |
| MEDIA_CONTROL_POINT_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| |
| // Issue 1st command (Play from peer1) |
| fake_gatt_server.incoming_write( |
| peer1, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![MediaControlOpcode::Play.raw_opcode()], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { |
| opcode: op1, |
| responder: resp1, |
| }))) = poll_result |
| else { |
| panic!("expected 1st ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(op1, MediaControlOpcode::Play); |
| |
| // Issue 2nd command (NextTrack from peer2) |
| fake_gatt_server.incoming_write( |
| peer2, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![MediaControlOpcode::NextTrack.raw_opcode()], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { |
| opcode: op2, |
| responder: resp2, |
| }))) = poll_result |
| else { |
| panic!("expected 2nd ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(op2, MediaControlOpcode::NextTrack); |
| |
| // Drain write acknowledgements |
| let FakeServerEvent::WriteResponded { value: val1, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected 1st WriteResponded"); |
| }; |
| assert!(val1.is_ok()); |
| |
| let FakeServerEvent::WriteResponded { value: val2, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected 2nd WriteResponded"); |
| }; |
| assert!(val2.is_ok()); |
| |
| // Respond to both commands in a row before polling the stream |
| resp1.send(ControlPointResultCode::Success); |
| resp2.send(ControlPointResultCode::CommandCannotBeCompleted); |
| |
| // Polling the server stream drains both pending control point responses and |
| // issues two local service notifications in a single poll_next invocation. |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| let n1 = expect_service_event(&mut event_receiver); |
| let n2 = expect_service_event(&mut event_receiver); |
| |
| let mut received = Vec::new(); |
| for event in [n1, n2] { |
| let FakeServerEvent::Notified { handle, value, peers, .. } = event else { |
| panic!("expected Notified event, got {event:?}"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| received.push((peers, value)); |
| } |
| |
| assert!(received.contains(&( |
| vec![peer1], |
| vec![MediaControlOpcode::Play.raw_opcode(), ControlPointResultCode::Success.into()] |
| ))); |
| assert!(received.contains(&( |
| vec![peer2], |
| vec![ |
| MediaControlOpcode::NextTrack.raw_opcode(), |
| ControlPointResultCode::CommandCannotBeCompleted.into() |
| ] |
| ))); |
| |
| // No further events pending |
| expect_no_service_event(&mut event_receiver); |
| } |
| |
| #[test] |
| fn server_stream_drains_pending_control_point_responses_before_terminating() { |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server_with_control_point(peer); |
| server.set_media_state(MediaState::Playing); |
| |
| // Write Pause opcode (0x02) |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| MEDIA_CONTROL_POINT_HANDLE, |
| 0, |
| vec![0x02], |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::ControlPointCommand { opcode, responder }))) = |
| poll_result |
| else { |
| panic!("expected ControlPointCommand event, got {poll_result:?}"); |
| }; |
| assert_eq!(opcode, MediaControlOpcode::Pause); |
| |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Replace local_service events stream with a custom channel and close it. |
| let (sender, receiver) = futures::channel::mpsc::unbounded(); |
| let LocalServiceState::Published { service, .. } = |
| std::mem::replace(&mut server.local_service, LocalServiceState::Terminated) |
| else { |
| panic!("Expected server to be in Published state"); |
| }; |
| server.local_service = LocalServiceState::Published { service, events: receiver }; |
| drop(sender); |
| |
| // Polling the server returns Pending because there is still an in-flight |
| // responder |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| // Upper layer executes command and responds with Success |
| responder.send(ControlPointResultCode::Success); |
| |
| // Polling the server dispatches the notification and drains the in-flight |
| // response, then terminates the stream as the GATT service event stream |
| // has closed. |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Ready(None)); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, MEDIA_CONTROL_POINT_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, vec![0x02, ControlPointResultCode::Success.into()]); |
| } |
| |
| #[test] |
| fn write_negative_track_position_success() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| server.state.media_state = MediaState::Paused; |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| fake_gatt_server.incoming_client_configuration( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| bt_gatt::server::NotificationType::Notify, |
| ); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| expect_no_service_event(&mut event_receiver); |
| |
| // Negative value represents offset from end of track per MCS v1.0.1 Section |
| // 3.7.1 |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| 0, |
| (-500i32).to_le_bytes().to_vec(), |
| ); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetTrackPosition { |
| peer_id, |
| position, |
| responder, |
| }))) = poll_result |
| else { |
| panic!("expected SetTrackPosition event, got {poll_result:?}"); |
| }; |
| assert_eq!(peer_id, peer); |
| assert_eq!(position, TrackPosition::from_raw_10ms(-500)); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Upper layer confirms the updated track position |
| responder.send(position); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| |
| let FakeServerEvent::Notified { handle, value, peers, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected Notified"); |
| }; |
| assert_eq!(handle, TRACK_POSITION_HANDLE); |
| assert_eq!(peers, vec![peer]); |
| assert_eq!(value, (-500i32).to_le_bytes().to_vec()); |
| |
| // Subsequent read reflects the updated track position |
| fake_gatt_server.incoming_read(peer, service_id, TRACK_POSITION_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), (-500i32).to_le_bytes()); |
| } |
| |
| #[test] |
| fn write_responder_dropped_without_send_does_not_notify_or_mutate() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_player_speeds() |
| .with_playing_orders(SupportedPlayingOrders::all()), |
| ); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| // 1. Playback speed write |
| fake_gatt_server.incoming_write(peer, service_id, PLAYBACK_SPEED_HANDLE, 0, vec![0x40]); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetPlaybackSpeed { responder, .. }))) = poll_result |
| else { |
| panic!("expected SetPlaybackSpeed event, got {poll_result:?}"); |
| }; |
| // ATT write is acknowledged |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Drop responder without calling send |
| drop(responder); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| // No notification should be emitted |
| expect_no_service_event(&mut event_receiver); |
| |
| // Playback speed remains unchanged |
| fake_gatt_server.incoming_read(peer, service_id, PLAYBACK_SPEED_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![PlaybackSpeed::NORMAL.into()]); |
| |
| // 2. Track position write |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| TRACK_POSITION_HANDLE, |
| 0, |
| 500i32.to_le_bytes().to_vec(), |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetTrackPosition { responder, .. }))) = poll_result |
| else { |
| panic!("expected SetTrackPosition event, got {poll_result:?}"); |
| }; |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| drop(responder); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| expect_no_service_event(&mut event_receiver); |
| |
| // Track position remains unavailable |
| fake_gatt_server.incoming_read(peer, service_id, TRACK_POSITION_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), TrackPosition::Unavailable.raw_10ms().to_le_bytes()); |
| |
| // 3. Playing order write |
| fake_gatt_server.incoming_write( |
| peer, |
| service_id, |
| PLAYING_ORDER_HANDLE, |
| 0, |
| vec![PlayingOrder::ShuffleOnce.into()], |
| ); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| let Poll::Ready(Some(Ok(McsServerEvent::SetPlayingOrder { responder, .. }))) = poll_result |
| else { |
| panic!("expected SetPlayingOrder event, got {poll_result:?}"); |
| }; |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| drop(responder); |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| expect_no_service_event(&mut event_receiver); |
| |
| // Playing order remains default InOrderOnce |
| fake_gatt_server.incoming_read(peer, service_id, PLAYING_ORDER_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![PlayingOrder::InOrderOnce.into()]); |
| } |
| |
| #[test] |
| fn write_playing_order_invalid_value_is_ignored() { |
| let (mut server, fake_gatt_server, mut event_receiver) = setup_test_server( |
| McsServerBuilder::generic(0x42, "Test Player") |
| .with_playing_orders(SupportedPlayingOrders::all()), |
| ); |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let peer = PeerId(1); |
| let service_id = ServiceId::new(0x42); |
| |
| // Write an out-of-range RFU value (0xFF) |
| fake_gatt_server.incoming_write(peer, service_id, PLAYING_ORDER_HANDLE, 0, vec![0xFF]); |
| |
| let poll_result = server.next().poll_unpin(&mut noop_cx); |
| assert_matches!(poll_result, Poll::Pending); |
| |
| let FakeServerEvent::WriteResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected WriteResponded"); |
| }; |
| assert!(value.is_ok()); |
| |
| // Value remains unchanged (default InOrderOnce) |
| fake_gatt_server.incoming_read(peer, service_id, PLAYING_ORDER_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| assert_eq!(value.unwrap(), vec![PlayingOrder::InOrderOnce.into()]); |
| } |
| |
| #[test] |
| fn local_state_handle_read() { |
| let mut state = McsLocalState::new(0x01, "Player"); |
| state.media_state = MediaState::Paused; |
| state.playback_speed = Some(PlaybackSpeed::NORMAL); |
| state.playing_orders = Some(PlayingOrderState::new(SupportedPlayingOrders::all())); |
| |
| // Track Position |
| assert_eq!( |
| state.handle_read(TRACK_POSITION_HANDLE, 0), |
| Ok(TrackPosition::Unavailable.raw_10ms().to_le_bytes().to_vec()) |
| ); |
| |
| // Playback Speed |
| assert_eq!(state.handle_read(PLAYBACK_SPEED_HANDLE, 0), Ok(vec![0x00])); |
| |
| // Playing Order |
| assert_eq!( |
| state.handle_read(PLAYING_ORDER_HANDLE, 0), |
| Ok(vec![PlayingOrder::InOrderOnce.into()]) |
| ); |
| } |
| |
| #[test] |
| fn playing_order_state_helpers() { |
| let mut state = PlayingOrderState::new(SupportedPlayingOrders::IN_ORDER_ONCE); |
| assert_eq!(state.current, PlayingOrder::InOrderOnce); |
| assert_eq!(state.supported, SupportedPlayingOrders::IN_ORDER_ONCE); |
| assert!(!state.set_order(PlayingOrder::SingleOnce)); |
| assert_eq!(state.current, PlayingOrder::InOrderOnce); |
| assert!(state.set_order(PlayingOrder::InOrderOnce)); |
| |
| let mut all_state = PlayingOrderState::new(SupportedPlayingOrders::all()); |
| assert_eq!(all_state.current, PlayingOrder::InOrderOnce); |
| assert_eq!(all_state.supported, SupportedPlayingOrders::all()); |
| assert!(all_state.set_order(PlayingOrder::ShuffleRepeat)); |
| assert_eq!(all_state.current, PlayingOrder::ShuffleRepeat); |
| } |
| |
| #[test] |
| fn local_state_track_position_calculation() { |
| let mut state = McsLocalState::new(0x01, "Test Player"); |
| |
| // When no track is loaded, position is Unavailable regardless of state. |
| assert_eq!(state.current_track_position(), TrackPosition::Unavailable); |
| state.media_state = MediaState::Playing; |
| assert_eq!(state.current_track_position(), TrackPosition::Unavailable); |
| |
| // When paused or seeking with a loaded track, position does not advance with |
| // time. |
| state.media_state = MediaState::Paused; |
| state.track_position = TrackPosition::from_start(std::time::Duration::from_millis(5000)); |
| state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_secs(10)); |
| assert_eq!( |
| state.current_track_position(), |
| TrackPosition::from_start(std::time::Duration::from_millis(5000)) |
| ); |
| state.media_state = MediaState::Seeking; |
| assert_eq!( |
| state.current_track_position(), |
| TrackPosition::from_start(std::time::Duration::from_millis(5000)) |
| ); |
| |
| // When playing with a loaded track (FromStart), position advances based on |
| // elapsed time. |
| state.media_state = MediaState::Playing; |
| state.track_duration = TrackDuration::from_duration(std::time::Duration::from_secs(20)); |
| state.track_position = TrackPosition::from_start(std::time::Duration::from_millis(5000)); |
| state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_millis(500)); |
| let computed = state.current_track_position(); |
| let TrackPosition::FromStart(dur) = computed else { |
| panic!("expected FromStart"); |
| }; |
| assert!( |
| dur >= std::time::Duration::from_millis(5450) |
| && dur <= std::time::Duration::from_millis(5650), |
| "unexpected computed position: {dur:?}" |
| ); |
| |
| // When elapsed time exceeds track duration, position is clamped to duration. |
| state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_secs(100)); |
| assert_eq!( |
| state.current_track_position(), |
| TrackPosition::from_start(std::time::Duration::from_secs(20)) |
| ); |
| |
| // When playing with a track position relative to end (FromEnd), base is |
| // duration - offset. |
| state.track_duration = TrackDuration::from_duration(std::time::Duration::from_secs(30)); |
| state.track_position = TrackPosition::from_end(std::time::Duration::from_secs(10)); |
| state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_millis(500)); |
| // Base is 30s - 10s = 20s. Elapsed 500ms -> ~20.5s from start. |
| let computed = state.current_track_position(); |
| let TrackPosition::FromStart(dur) = computed else { |
| panic!("expected FromStart"); |
| }; |
| assert!( |
| dur >= std::time::Duration::from_millis(20450) |
| && dur <= std::time::Duration::from_millis(20650), |
| "unexpected computed position: {dur:?}" |
| ); |
| |
| // When FromEnd elapsed time exceeds track duration, position is clamped to |
| // duration. |
| state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_secs(100)); |
| assert_eq!( |
| state.current_track_position(), |
| TrackPosition::from_start(std::time::Duration::from_secs(30)) |
| ); |
| |
| // When FromEnd is used with unknown duration, it cannot resolve to FromStart |
| // and returns FromEnd. |
| state.track_duration = TrackDuration::Unknown; |
| state.track_position = TrackPosition::from_end(std::time::Duration::from_secs(10)); |
| assert_eq!( |
| state.current_track_position(), |
| TrackPosition::from_end(std::time::Duration::from_secs(10)) |
| ); |
| } |
| |
| #[test] |
| fn read_track_position_during_active_playback() { |
| let (mut server, fake_gatt_server, mut event_receiver) = |
| setup_test_server(McsServerBuilder::generic(0x42, "Test Player")); |
| |
| // Set server state to Playing with a loaded track at position 500 (5.0 |
| // seconds). |
| server.state.media_state = MediaState::Playing; |
| server.state.track_duration = |
| TrackDuration::from_duration(std::time::Duration::from_secs(60)); |
| server.state.track_position = |
| TrackPosition::from_start(std::time::Duration::from_millis(5000)); |
| server.state.position_updated_at = |
| Some(std::time::Instant::now() - std::time::Duration::from_millis(1000)); |
| |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let peer = PeerId(1); |
| let service_id = server.service_def.id(); |
| fake_gatt_server.incoming_read(peer, service_id, TRACK_POSITION_HANDLE, 0); |
| let _ = server.next().poll_unpin(&mut noop_cx); |
| let bt_gatt::test_utils::FakeServerEvent::ReadResponded { value, .. } = |
| expect_service_event(&mut event_receiver) |
| else { |
| panic!("expected ReadResponded"); |
| }; |
| |
| let raw_bytes = value.unwrap(); |
| assert_eq!(raw_bytes.len(), 4); |
| let raw_10ms = i32::from_le_bytes(raw_bytes.try_into().unwrap()); |
| // 5000ms + 1000ms elapsed = 6000ms = 600 units of 10ms (allow minor jitter +/- |
| // 30 units) |
| assert!(raw_10ms >= 580 && raw_10ms <= 640, "unexpected raw_10ms: {raw_10ms}"); |
| } |
| } |