rust/bt-mcs: Add McsServer and publish GATT service Implement `McsServer` and `LocalServiceState` state machine managing the asynchronous GATT publication lifecycle for Media Control Service (MCS) and Generic Media Control Service (GMCS). Add `McsServerBuilder::build` and `McsServer::publish` methods to support registering the service in the GATT database. Bug: 540400364 Test: cargo test -p bt-mcs Change-Id: I20c1a742f714fe80f57f22e5bca16175c9b45ee7 Reviewed-on: https://bluetooth-review.googlesource.com/c/bluetooth/+/3680
diff --git a/rust/bt-mcs/Cargo.toml b/rust/bt-mcs/Cargo.toml index 1db1b45..a30afdb 100644 --- a/rust/bt-mcs/Cargo.toml +++ b/rust/bt-mcs/Cargo.toml
@@ -8,4 +8,9 @@ bitflags.workspace = true bt-common.workspace = true bt-gatt.workspace = true +futures.workspace = true +pin-project.workspace = true thiserror.workspace = true + +[dev-dependencies] +bt-gatt = { workspace = true, features = ["test-utils"] }
diff --git a/rust/bt-mcs/src/error.rs b/rust/bt-mcs/src/error.rs index eec71eb..d57a6f7 100644 --- a/rust/bt-mcs/src/error.rs +++ b/rust/bt-mcs/src/error.rs
@@ -6,6 +6,9 @@ #[derive(Debug, Error)] pub enum Error { + #[error("Service is already published")] + AlreadyPublished, + #[error("GATT operation error: {0}")] Gatt(#[from] bt_gatt::types::Error),
diff --git a/rust/bt-mcs/src/lib.rs b/rust/bt-mcs/src/lib.rs index da58e1a..2ef3b54 100644 --- a/rust/bt-mcs/src/lib.rs +++ b/rust/bt-mcs/src/lib.rs
@@ -7,4 +7,4 @@ pub mod types; pub use crate::error::Error; -pub use crate::server::McsServerBuilder; +pub use crate::server::{McsServer, McsServerBuilder};
diff --git a/rust/bt-mcs/src/server.rs b/rust/bt-mcs/src/server.rs index 5abb3fa..59e5d68 100644 --- a/rust/bt-mcs/src/server.rs +++ b/rust/bt-mcs/src/server.rs
@@ -5,12 +5,16 @@ //! Implements the Media Control Service (MCS) server. use bt_common::Uuid; -use bt_gatt::server::{ServiceDefinition, ServiceId}; +use bt_gatt::server::{LocalService, Server as _, ServiceDefinition, ServiceEvent, ServiceId}; use bt_gatt::types::{ AttributePermissions, CharacteristicProperties, CharacteristicProperty, Handle, SecurityLevels, ServiceKind, }; use bt_gatt::Characteristic; +use futures::stream::Stream; +use pin_project::pin_project; +use std::future::Future; +use std::task::{Context, Poll, Waker}; use crate::types::*; use crate::Error; @@ -139,6 +143,81 @@ } } +/// 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 { .. }) + } +} + +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)))); + } + None => { + self.as_mut().set(LocalServiceState::Terminated); + return Poll::Ready(None); + } + } + } + } + } + } +} + /// Builder for configuring an MCS or GMCS GATT service. #[derive(Debug, Clone, PartialEq)] pub struct McsServerBuilder { @@ -185,12 +264,88 @@ 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 { + service_def, + local_service: Default::default(), + ccid: self.ccid, + player_name: self.player_name, + }) + } +} + +/// An instance of a Media Control Service (MCS) or Generic Media Control +/// Service (GMCS) GATT server. +#[pin_project] +pub struct McsServer<T: bt_gatt::ServerTypes> { + service_def: ServiceDefinition, + #[pin] + local_service: LocalServiceState<T>, + ccid: u8, + player_name: String, +} + +impl<T: bt_gatt::ServerTypes> McsServer<T> { + /// 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(()) + } +} + +impl<T: bt_gatt::ServerTypes> Stream for McsServer<T> { + type Item = Result<(), 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(); + let gatt_event = match futures::ready!(this.local_service.as_mut().poll_next(cx)) { + None => return Poll::Ready(None), + Some(Err(e)) => return Poll::Ready(Some(Err(e))), + Some(Ok(event)) => event, + }; + match gatt_event { + // TODO(b/540400364): Add support for characteristic reads and writes + _ => continue, + } + } + } } #[cfg(test)] mod tests { use super::*; + use bt_gatt::test_utils::{FakeServer, FakeTypes}; + use futures::{FutureExt, StreamExt}; + #[test] fn builder_generic_service_definition() { let builder = McsServerBuilder::generic(0x42, "Test Generic Player"); @@ -262,7 +417,123 @@ assert_eq!(def3.id(), ServiceId::new(6)); } - // TODO(b/540400364): Add a test verifying that publishing two McsServer - // instances with the same CCID to a GATT server fails with an - // AlreadyPublished error. + #[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))); + } }