blob: ca82faa08f35b9acef743e91398586d89806b1c7 [file]
// 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::{
CharacteristicProperties, CharacteristicProperty, GattError, Handle, ServiceKind,
};
use futures::stream::{FuturesUnordered, Stream};
use pin_project::pin_project;
use std::future::Future;
use std::task::{Context, Poll, Waker};
pub mod event;
use event::*;
pub use event::{
ControlPointWriteResponder, McsServerEvent, SetPlaybackSpeedResponder,
SetPlayingOrderResponder, SetTrackPositionResponder, SetValueResponder,
};
mod state;
use state::*;
use crate::types::*;
use crate::Error;
/// 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),
}
}
}
}
}
}
/// 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 bt_gatt::Characteristic;
use futures::{FutureExt, 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)
}
#[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 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()]);
}
}