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)));
+    }
 }