| // 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 Broadcast Audio Scan Service (BASS) server role. |
| |
| use bt_bap::types::BroadcastCode; |
| use bt_common::packet_encoding::{Decodable, Encodable}; |
| use bt_common::PeerId; |
| use bt_gatt::server::{ |
| LocalService, ReadResponder, Server as _, ServiceDefinition, ServiceEvent, ServiceId, |
| WriteResponder, |
| }; |
| use bt_gatt::types::{ |
| AttributePermissions, CharacteristicProperty, GattError, Handle, SecurityLevels, ServiceKind, |
| }; |
| use bt_gatt::Characteristic; |
| use futures::stream::Stream; |
| use futures::Future; |
| use pin_project::pin_project; |
| use std::collections::{BTreeMap, HashMap}; |
| use std::num::NonZeroU8; |
| use std::task::{Poll, Waker}; |
| |
| use crate::types::*; |
| |
| pub mod error; |
| use error::Error; |
| |
| /// Service identifier assigned to the published BASS GATT service instance. |
| const BASS_SERVICE_ID: ServiceId = ServiceId::new(1); |
| |
| /// Handle assigned to the Broadcast Audio Scan Control Point characteristic. |
| const CONTROL_POINT_HANDLE: Handle = Handle(1); |
| |
| /// Maximum number of Broadcast Receive State characteristics that can be hosted |
| /// by the server. See BASS v1.0 Section 3.2.1. |
| const MAX_RECEIVE_STATES: usize = 255; |
| |
| /// Internal representation of a Broadcast Receive State characteristic slot. |
| /// See BASS v1.0 Section 3.2. |
| #[derive(Debug)] |
| pub(crate) struct PublishedReceiveStateCharacteristic { |
| /// Handle assigned to this GATT characteristic. |
| handle: Handle, |
| /// Current Broadcast Receive State value. |
| state: BroadcastReceiveState, |
| } |
| |
| /// Internal publication state lifecycle of the BASS GATT service. |
| #[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<bt_gatt::server::ServiceEvent<T>, Error>; |
| |
| fn poll_next( |
| mut self: std::pin::Pin<&mut Self>, |
| cx: &mut std::task::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 } => { |
| let service_result = futures::ready!(fut.poll(cx)); |
| let Ok(service) = service_result else { |
| self.as_mut().set(LocalServiceState::NotPublished { waker: None }); |
| return Poll::Ready(Some(Err(Error::Gatt(service_result.err().unwrap())))); |
| }; |
| let events = service.publish(); |
| self.as_mut().set(LocalServiceState::Published { service, events }); |
| } |
| LocalServiceProj::Published { service: _, events } => { |
| return match futures::ready!(events.poll_next(cx)) { |
| Some(Ok(event)) => Poll::Ready(Some(Ok(event))), |
| Some(Err(e)) => { |
| self.as_mut().set(LocalServiceState::Terminated); |
| Poll::Ready(Some(Err(Error::Gatt(e)))) |
| } |
| None => { |
| self.as_mut().set(LocalServiceState::Terminated); |
| Poll::Ready(None) |
| } |
| }; |
| } |
| } |
| } |
| } |
| } |
| |
| /// Builder for constructing a BASS GATT [`Server`]. |
| #[derive(Default)] |
| pub struct ServerBuilder { |
| /// Staged Broadcast Receive State characteristics. |
| /// There must be at least one such staged characteristic. |
| receive_states: Vec<BroadcastReceiveState>, |
| } |
| |
| impl ServerBuilder { |
| pub fn new() -> Self { |
| Self::default() |
| } |
| |
| /// Adds a Broadcast Receive State characteristic slot to the builder. |
| pub fn add_receive_state_characteristic(mut self, state: BroadcastReceiveState) -> Self { |
| self.receive_states.push(state); |
| self |
| } |
| |
| /// Constructs the Broadcast Audio Scan Control Point characteristic. |
| /// Defined in BASS v1.0 Section 3.1. |
| fn build_control_point() -> Characteristic { |
| let cp_properties = |
| CharacteristicProperty::Write | CharacteristicProperty::WriteWithoutResponse; |
| Characteristic { |
| handle: CONTROL_POINT_HANDLE, |
| uuid: BROADCAST_AUDIO_SCAN_CONTROL_POINT_UUID, |
| properties: cp_properties.clone(), |
| permissions: AttributePermissions::with_levels( |
| &cp_properties, |
| &SecurityLevels::encryption_required(), |
| ), |
| descriptors: Vec::new(), |
| } |
| } |
| |
| /// Constructs the Broadcast Receive State characteristic. |
| /// Defined in BASS v1.0 Section 3.2. |
| fn build_receive_state(handle: Handle) -> Characteristic { |
| let properties = CharacteristicProperty::Read | CharacteristicProperty::Notify; |
| Characteristic { |
| handle, |
| uuid: BROADCAST_RECEIVE_STATE_UUID, |
| properties: properties.clone(), |
| permissions: AttributePermissions::with_levels( |
| &properties, |
| &SecurityLevels::encryption_required(), |
| ), |
| descriptors: Vec::new(), |
| } |
| } |
| |
| /// Builds a [`Server`] instance after verifying required service |
| /// characteristics. |
| pub fn build<T: bt_gatt::ServerTypes>(self) -> Result<Server<T>, Error> { |
| // Per BASS v1.0 Section 3.2, there must be [1,255] Broadcast Receive State |
| // characteristics. |
| if self.receive_states.is_empty() { |
| return Err(Error::MissingReceiveState); |
| } |
| if self.receive_states.len() > MAX_RECEIVE_STATES { |
| return Err(Error::ExceedsMaxReceiveStates); |
| } |
| |
| let mut service_def = ServiceDefinition::new( |
| BASS_SERVICE_ID, |
| BROADCAST_AUDIO_SCAN_SERVICE_UUID, |
| ServiceKind::Primary, |
| ); |
| |
| let _ = service_def.add_characteristic(Self::build_control_point()); |
| |
| // Broadcast Receive State characteristics (Read, Notify; Encryption Required). |
| // Handle(1) is allocated to the Control Point Characteristic. |
| const FIRST_RECEIVE_STATE_HANDLE: Handle = Handle(2); |
| let num_receive_states = self.receive_states.len(); |
| let mut state = ServerState::new(num_receive_states); |
| |
| for (i, mut receive_state) in self.receive_states.into_iter().enumerate() { |
| let handle = Handle(FIRST_RECEIVE_STATE_HANDLE.0 + i as u64); |
| let _ = service_def.add_characteristic(Self::build_receive_state(handle)); |
| |
| if let BroadcastReceiveState::NonEmpty(ref mut r) = receive_state { |
| r.source_id = state |
| .allocate_source_id(handle) |
| .map_err(|e| Error::Gatt(bt_gatt::types::Error::Gatt(e)))?; |
| } |
| |
| state.receive_state_characteristics.insert( |
| handle, |
| PublishedReceiveStateCharacteristic { handle, state: receive_state }, |
| ); |
| } |
| |
| Ok(Server { service_def, local_service: Default::default(), state }) |
| } |
| } |
| |
| /// Events emitted by the [`Server`] when a remote GATT client writes to the |
| /// Control Point. |
| #[derive(Debug, PartialEq)] |
| pub enum ServerEvent { |
| /// A client requested a change in remote scanning state. |
| RemoteScanState { peer_id: PeerId, is_scanning: bool }, |
| /// A client requested adding a new Broadcast Source. |
| AddSource { peer_id: PeerId, source_id: SourceId, operation: AddSourceOperation }, |
| /// A client requested modifying an existing Broadcast Source. |
| ModifySource { peer_id: PeerId, source_id: SourceId, operation: ModifySourceOperation }, |
| /// A client provided a Broadcast Code for an encrypted Broadcast Source. |
| SetBroadcastCode { peer_id: PeerId, source_id: SourceId, broadcast_code: BroadcastCode }, |
| /// A client requested removing a Broadcast Source. |
| RemoveSource { peer_id: PeerId, source_id: SourceId }, |
| } |
| |
| impl ServerEvent { |
| /// Attempts to decode a raw control point write request into a |
| /// [`ServerEvent`]. |
| fn decode(peer_id: PeerId, val_bytes: &[u8]) -> Result<Self, GattError> { |
| if val_bytes.is_empty() { |
| return Err(GattError::WriteRequestRejected); |
| } |
| |
| let raw_opcode = val_bytes[0]; |
| let Ok(opcode) = ControlPointOpcode::try_from(raw_opcode) else { |
| return Err(ERROR_OPCODE_NOT_SUPPORTED); |
| }; |
| |
| match opcode { |
| ControlPointOpcode::RemoteScanStopped => { |
| let _ = decode_control_point_op::<RemoteScanStoppedOperation>(val_bytes)?; |
| Ok(ServerEvent::RemoteScanState { peer_id, is_scanning: false }) |
| } |
| ControlPointOpcode::RemoteScanStarted => { |
| let _ = decode_control_point_op::<RemoteScanStartedOperation>(val_bytes)?; |
| Ok(ServerEvent::RemoteScanState { peer_id, is_scanning: true }) |
| } |
| ControlPointOpcode::AddSource => { |
| let operation = decode_control_point_op::<AddSourceOperation>(val_bytes)?; |
| // `source_id` will be assigned internally by the Server. |
| Ok(ServerEvent::AddSource { peer_id, source_id: 0, operation }) |
| } |
| ControlPointOpcode::ModifySource => { |
| let operation = decode_control_point_op::<ModifySourceOperation>(val_bytes)?; |
| Ok(ServerEvent::ModifySource { peer_id, source_id: operation.source_id, operation }) |
| } |
| ControlPointOpcode::SetBroadcastCode => { |
| let operation = decode_control_point_op::<SetBroadcastCodeOperation>(val_bytes)?; |
| Ok(ServerEvent::SetBroadcastCode { |
| peer_id, |
| source_id: operation.source_id, |
| broadcast_code: operation.broadcast_code, |
| }) |
| } |
| ControlPointOpcode::RemoveSource => { |
| let operation = decode_control_point_op::<RemoveSourceOperation>(val_bytes)?; |
| Ok(ServerEvent::RemoveSource { peer_id, source_id: operation.0 }) |
| } |
| } |
| } |
| } |
| |
| /// Attempts to decode a control point operation payload. |
| fn decode_control_point_op<O: Decodable>(val_bytes: &[u8]) -> Result<O, GattError> { |
| let (Ok(op), consumed) = O::decode(val_bytes) else { |
| return Err(GattError::WriteRequestRejected); |
| }; |
| if consumed != val_bytes.len() { |
| return Err(GattError::WriteRequestRejected); |
| } |
| Ok(op) |
| } |
| |
| /// Internal state of the BASS GATT Server. |
| #[derive(Debug)] |
| struct ServerState { |
| /// Broadcast Receive State characteristics identified by the assigned |
| /// GATT handle. |
| receive_state_characteristics: BTreeMap<Handle, PublishedReceiveStateCharacteristic>, |
| /// Unique IDs tracking each assigned GATT Handle. |
| source_id_to_handle: HashMap<SourceId, Handle>, |
| /// Next available Source ID to be assigned to an empty Receive State |
| /// characteristic. |
| /// Must be in [1, 255]. See BASS spec v1.0 Section 3.2.1. |
| next_source_id: NonZeroU8, |
| } |
| |
| impl ServerState { |
| fn new(capacity: usize) -> Self { |
| Self { |
| receive_state_characteristics: BTreeMap::new(), |
| source_id_to_handle: HashMap::with_capacity(capacity), |
| next_source_id: NonZeroU8::MIN, |
| } |
| } |
| |
| /// Allocates the next available, unused [`SourceId`] and maps it to |
| /// `handle`. |
| /// |
| /// Returns the assigned ID on success. |
| /// Returns `GattError::InsufficientResources` if all 255 Source IDs |
| /// are currently in use. |
| fn allocate_source_id(&mut self, handle: Handle) -> Result<SourceId, GattError> { |
| // Upper bound on the maximum number of unique SourceIDs (u8 max). |
| for _ in 0..NonZeroU8::MAX.get() { |
| let id = self.next_source_id.get(); |
| self.next_source_id = self.next_source_id.checked_add(1).unwrap_or(NonZeroU8::MIN); |
| if !self.source_id_to_handle.contains_key(&id) { |
| self.source_id_to_handle.insert(id, handle); |
| return Ok(id); |
| } |
| } |
| Err(GattError::InsufficientResources) |
| } |
| |
| /// Handles a read request for a characteristic on the BASS server. |
| fn handle_read(&self, handle: Handle, offset: usize, responder: impl ReadResponder) { |
| // Only the Broadcast Receive State characteristics can be read. |
| // See BASS spec v1.0 Section 3.3.2. |
| if handle == CONTROL_POINT_HANDLE { |
| responder.error(GattError::ReadNotPermitted); |
| return; |
| } |
| |
| let Some(chrc) = self.receive_state_characteristics.get(&handle) else { |
| responder.error(GattError::InvalidHandle); |
| return; |
| }; |
| |
| let len = chrc.state.encoded_len(); |
| if offset > len { |
| responder.error(GattError::InvalidOffset); |
| return; |
| } |
| |
| let mut buf = vec![0u8; len]; |
| match chrc.state.encode(&mut buf) { |
| Ok(_) => responder.respond(&buf[offset..]), |
| Err(_) => responder.error(GattError::UnlikelyError), |
| } |
| } |
| |
| /// Handles a write request for a characteristic on the BASS server. |
| fn handle_write( |
| &mut self, |
| peer_id: PeerId, |
| handle: Handle, |
| value: &[u8], |
| responder: impl WriteResponder, |
| ) -> Option<ServerEvent> { |
| // Only the Broadcast Audio Scan Control Point characteristic can be written to |
| // by a client. See BASS v1.0 Section 3.3. |
| if handle != CONTROL_POINT_HANDLE { |
| responder.error(GattError::WriteNotPermitted); |
| return None; |
| } |
| |
| match self.handle_control_point_write(peer_id, value) { |
| Ok(event) => { |
| responder.acknowledge(); |
| Some(event) |
| } |
| Err(err) => { |
| responder.error(err); |
| None |
| } |
| } |
| } |
| |
| /// Returns true if the specified `source_id` is assigned to a non-empty |
| /// Receive State characteristic. |
| fn is_valid_source_id(&self, source_id: SourceId) -> bool { |
| self.source_id_to_handle |
| .get(&source_id) |
| .and_then(|h| self.receive_state_characteristics.get(h)) |
| .map_or(false, |s| !s.state.is_empty()) |
| } |
| |
| fn handle_control_point_write( |
| &mut self, |
| peer_id: PeerId, |
| val_bytes: &[u8], |
| ) -> Result<ServerEvent, GattError> { |
| let mut event = ServerEvent::decode(peer_id, val_bytes)?; |
| |
| match &mut event { |
| ServerEvent::RemoteScanState { .. } => {} |
| ServerEvent::AddSource { source_id, operation, .. } => { |
| *source_id = self.handle_add_source_write(operation)?; |
| } |
| ServerEvent::ModifySource { operation, .. } => { |
| self.handle_modify_source_write(operation)?; |
| } |
| ServerEvent::SetBroadcastCode { source_id, .. } => { |
| if !self.is_valid_source_id(*source_id) { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| } |
| } |
| ServerEvent::RemoveSource { source_id, .. } => { |
| self.handle_remove_source_write(*source_id)?; |
| } |
| } |
| |
| Ok(event) |
| } |
| |
| fn handle_add_source_write( |
| &mut self, |
| operation: &AddSourceOperation, |
| ) -> Result<SourceId, GattError> { |
| // Find the next available empty Receive State characteristic entry. |
| let handle = { |
| let Some(chrc) = |
| self.receive_state_characteristics.values().find(|c| c.state.is_empty()) |
| else { |
| return Err(GattError::InsufficientResources); |
| }; |
| chrc.handle |
| }; |
| |
| let source_id = self.allocate_source_id(handle)?; |
| let chrc = |
| self.receive_state_characteristics.get_mut(&handle).expect("just checked existence"); |
| let pa_sync_state = PaSyncState::from_pa_sync(operation.pa_sync); |
| let subgroups = operation |
| .subgroups |
| .iter() |
| .map(|s| BigSubgroup::new(Some(s.bis_sync.clone())).with_metadata(s.metadata.clone())) |
| .collect(); |
| |
| let initial_state = ReceiveState::new( |
| source_id, |
| operation.advertiser_address_type, |
| operation.advertiser_address, |
| operation.advertising_sid, |
| operation.broadcast_id, |
| pa_sync_state, |
| EncryptionStatus::NotEncrypted, |
| subgroups, |
| ); |
| |
| chrc.state = BroadcastReceiveState::NonEmpty(initial_state); |
| Ok(source_id) |
| } |
| |
| fn handle_modify_source_write( |
| &mut self, |
| operation: &ModifySourceOperation, |
| ) -> Result<(), GattError> { |
| let Some(handle) = self.source_id_to_handle.get(&operation.source_id) else { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| }; |
| let Some(chrc) = self.receive_state_characteristics.get_mut(handle) else { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| }; |
| let BroadcastReceiveState::NonEmpty(ref mut state) = chrc.state else { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| }; |
| |
| let subgroups = operation |
| .subgroups |
| .iter() |
| .map(|s| BigSubgroup::new(Some(s.bis_sync.clone())).with_metadata(s.metadata.clone())) |
| .collect(); |
| |
| state.subgroups = subgroups; |
| state.pa_sync_state = PaSyncState::from_pa_sync(operation.pa_sync); |
| |
| Ok(()) |
| } |
| |
| fn handle_remove_source_write(&mut self, source_id: SourceId) -> Result<(), GattError> { |
| let Some(handle) = self.source_id_to_handle.get(&source_id).copied() else { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| }; |
| let Some(chrc) = self.receive_state_characteristics.get_mut(&handle) else { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| }; |
| if chrc.state.is_empty() { |
| return Err(ERROR_INVALID_SOURCE_ID); |
| } |
| |
| let _ = self.source_id_to_handle.remove(&source_id); |
| chrc.state = BroadcastReceiveState::Empty; |
| Ok(()) |
| } |
| } |
| |
| /// The BASS GATT Server implementation. |
| /// |
| /// Manages the Control Point characteristic and one or more Broadcast Receive |
| /// State characteristics. |
| #[pin_project] |
| pub struct Server<T: bt_gatt::ServerTypes> { |
| service_def: ServiceDefinition, |
| #[pin] |
| local_service: LocalServiceState<T>, |
| state: ServerState, |
| } |
| |
| impl<T: bt_gatt::ServerTypes> Server<T> { |
| /// 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> { |
| if !matches!(self.local_service, LocalServiceState::NotPublished { .. }) { |
| return Err(Error::AlreadyPublished); |
| } |
| |
| let LocalServiceState::NotPublished { waker } = std::mem::replace( |
| &mut self.local_service, |
| LocalServiceState::Preparing { fut: server.prepare(self.service_def.clone()) }, |
| ) else { |
| unreachable!(); |
| }; |
| |
| if let Some(w) = waker { |
| w.wake(); |
| } |
| |
| Ok(()) |
| } |
| } |
| |
| impl<T: bt_gatt::ServerTypes> Stream for Server<T> { |
| type Item = Result<ServerEvent, Error>; |
| |
| fn poll_next( |
| mut self: std::pin::Pin<&mut Self>, |
| cx: &mut std::task::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)) { |
| Some(Ok(event)) => event, |
| Some(Err(e)) => return Poll::Ready(Some(Err(e))), |
| None => return Poll::Ready(None), |
| }; |
| |
| match gatt_event { |
| ServiceEvent::Read { handle, offset, responder, .. } => { |
| this.state.handle_read(handle, offset as usize, responder); |
| } |
| ServiceEvent::Write { handle, peer_id, value, responder, .. } => { |
| let val_vec = value.to_owned(); |
| if let Some(event) = |
| this.state.handle_write(peer_id, handle, &val_vec, responder) |
| { |
| return Poll::Ready(Some(Ok(event))); |
| } |
| } |
| _ => { |
| // TODO(b/534436439): Track CCC for Broadcast Receive State |
| // characteristics and `PeerInfo` if needed. |
| } |
| } |
| } |
| } |
| } |
| |
| #[cfg(test)] |
| mod tests { |
| use super::*; |
| use bt_bap::types::BroadcastId; |
| use bt_common::core::{AddressType, AdvertisingSetId, PeriodicAdvertisingInterval}; |
| use bt_gatt::test_utils::{FakeServer, FakeServerEvent, FakeTypes}; |
| use bt_gatt::types::GattError; |
| use futures::task::Context; |
| use futures::FutureExt; |
| use futures::StreamExt; |
| use std::task::Poll; |
| |
| use crate::types::{EncryptionStatus, PaSync, PaSyncState, ReceiveState}; |
| |
| fn make_test_receive_state(source_id: SourceId) -> BroadcastReceiveState { |
| BroadcastReceiveState::NonEmpty(ReceiveState::new( |
| source_id, |
| AddressType::Public, |
| [0x01, 0x02, 0x03, 0x04, 0x05, 0x06], |
| AdvertisingSetId::try_from(1).unwrap(), |
| BroadcastId::try_from(0x123456).unwrap(), |
| PaSyncState::NotSynced, |
| EncryptionStatus::NotEncrypted, |
| vec![], |
| )) |
| } |
| |
| #[test] |
| fn empty_builder_fails() { |
| let builder = ServerBuilder::new(); |
| assert!(matches!(builder.build::<FakeTypes>(), Err(Error::MissingReceiveState))); |
| } |
| |
| #[test] |
| fn too_many_receive_states_fails() { |
| let mut builder = ServerBuilder::new(); |
| for _ in 0..=MAX_RECEIVE_STATES { |
| builder = builder.add_receive_state_characteristic(BroadcastReceiveState::Empty); |
| } |
| assert!(matches!(builder.build::<FakeTypes>(), Err(Error::ExceedsMaxReceiveStates))); |
| } |
| |
| fn setup_test_server( |
| receive_states: Vec<BroadcastReceiveState>, |
| ) -> ( |
| Server<FakeTypes>, |
| FakeServer, |
| futures::channel::mpsc::UnboundedReceiver<FakeServerEvent>, |
| Context<'static>, |
| ) { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| let (fake_gatt_server, mut event_receiver) = FakeServer::new(); |
| let mut event_stream = event_receiver.next(); |
| |
| let mut builder = ServerBuilder::new(); |
| for state in receive_states { |
| builder = builder.add_receive_state_characteristic(state); |
| } |
| let mut server = builder.build::<FakeTypes>().expect("building server works"); |
| |
| server.publish(fake_gatt_server.clone()).expect("publish ok"); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| assert!(matches!( |
| event_stream.poll_unpin(&mut noop_cx), |
| Poll::Ready(Some(FakeServerEvent::Published { .. })) |
| )); |
| |
| (server, fake_gatt_server, event_receiver, noop_cx) |
| } |
| |
| #[test] |
| fn publish_server() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let mut server = ServerBuilder::new() |
| .add_receive_state_characteristic(BroadcastReceiveState::Empty) |
| .add_receive_state_characteristic(make_test_receive_state(2)) |
| .build::<FakeTypes>() |
| .expect("building server works"); |
| |
| assert!(!server.is_published()); |
| |
| let (fake_gatt_server, mut event_receiver) = FakeServer::new(); |
| let mut event_stream = event_receiver.next(); |
| |
| assert!(server.publish(fake_gatt_server.clone()).is_ok()); |
| assert!(server.publish(fake_gatt_server).is_err()); |
| |
| // Polling stream completes publication preparation |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| assert!(server.is_published()); |
| |
| let Poll::Ready(Some(FakeServerEvent::Published { id, definition })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected published event"); |
| }; |
| |
| assert_eq!(id, BASS_SERVICE_ID); |
| assert_eq!(definition.uuid(), BROADCAST_AUDIO_SCAN_SERVICE_UUID); |
| assert_eq!(definition.characteristics().count(), 3); |
| } |
| |
| #[test] |
| fn read_receive_state() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty, make_test_receive_state(2)]); |
| let mut event_stream = event_receiver.next(); |
| assert!(server.is_published()); |
| |
| // 1. Read empty Receive State slot (Handle 2) |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, Handle(2), 0); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| assert_eq!(handle, Handle(2)); |
| assert_eq!(value.unwrap(), vec![]); |
| |
| // 2. Read populated Receive State slot (Handle 3) |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, Handle(3), 0); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| assert_eq!(handle, Handle(3)); |
| assert!(!value.unwrap().is_empty()); |
| |
| // 3. Read unknown handle (Handle 99) |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, Handle(99), 0); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| assert_eq!(handle, Handle(99)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidHandle) |
| )); |
| } |
| |
| #[test] |
| fn control_point_remote_scan_state_changed() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Simulate RemoteScanStarted request |
| let raw_cmd = vec![ControlPointOpcode::RemoteScanStarted as u8]; |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_cmd); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!(event, ServerEvent::RemoteScanState { peer_id: PeerId(1), is_scanning: true }); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| |
| // Simulate RemoteScanStopped request |
| let raw_cmd = vec![ControlPointOpcode::RemoteScanStopped as u8]; |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_cmd); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!(event, ServerEvent::RemoteScanState { peer_id: PeerId(1), is_scanning: false }); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| } |
| |
| #[test] |
| fn control_point_invalid_opcode_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Simulate invalid opcode request |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, vec![0xFF]); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ApplicationError80) |
| )); |
| } |
| |
| #[test] |
| fn control_point_unallocated_source_id_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Operation targeting SourceId = 1 (slot exists, but is empty) |
| let modify_op = ModifySourceOperation::new( |
| 1, |
| PaSync::DoNotSync, |
| PeriodicAdvertisingInterval(0xFFFF), |
| vec![], |
| ); |
| let mut raw_bytes = vec![0u8; modify_op.encoded_len()]; |
| modify_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ApplicationError81) |
| )); |
| } |
| |
| #[test] |
| fn control_point_nonexistent_source_id_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // ModifySource operation targeting SourceId = 99 (does not exist) |
| let modify_op = ModifySourceOperation::new( |
| 99, |
| PaSync::DoNotSync, |
| PeriodicAdvertisingInterval(0xFFFF), |
| vec![], |
| ); |
| let mut raw_bytes = vec![0u8; modify_op.encoded_len()]; |
| modify_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ApplicationError81) |
| )); |
| } |
| |
| #[test] |
| fn control_point_extra_bytes_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // RemoteScanStarted opcode (1 byte) with extra trailing byte (0xFF) |
| let raw_cmd = vec![ControlPointOpcode::RemoteScanStarted as u8, 0xFF]; |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_cmd); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::WriteRequestRejected) |
| )); |
| } |
| |
| #[test] |
| fn control_point_add_source_success() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| let add_op = AddSourceOperation::new( |
| AddressType::Public, |
| [1, 2, 3, 4, 5, 6], |
| AdvertisingSetId::try_from(1).unwrap(), |
| BroadcastId::try_from(0x123456).unwrap(), |
| PaSync::DoNotSync, |
| PeriodicAdvertisingInterval(0xFFFF), |
| vec![], |
| ); |
| let mut raw_bytes = vec![0u8; add_op.encoded_len()]; |
| add_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!( |
| event, |
| ServerEvent::AddSource { peer_id: PeerId(1), source_id: 1, operation: add_op } |
| ); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| } |
| |
| #[test] |
| fn control_point_modify_source_success() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![make_test_receive_state(1)]); |
| let mut event_stream = event_receiver.next(); |
| |
| let modify_op = ModifySourceOperation::new( |
| 1, |
| PaSync::SyncPastAvailable, |
| PeriodicAdvertisingInterval(0x0010), |
| vec![], |
| ); |
| let mut raw_bytes = vec![0u8; modify_op.encoded_len()]; |
| modify_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!( |
| event, |
| ServerEvent::ModifySource { peer_id: PeerId(1), source_id: 1, operation: modify_op } |
| ); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| } |
| |
| #[test] |
| fn control_point_set_broadcast_code_success() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![make_test_receive_state(1)]); |
| let mut event_stream = event_receiver.next(); |
| |
| let code = BroadcastCode::new([0xAB; 16]); |
| let set_code_op = SetBroadcastCodeOperation::new(1, code); |
| let mut raw_bytes = vec![0u8; set_code_op.encoded_len()]; |
| set_code_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!( |
| event, |
| ServerEvent::SetBroadcastCode { |
| peer_id: PeerId(1), |
| source_id: 1, |
| broadcast_code: code, |
| } |
| ); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| } |
| |
| #[test] |
| fn control_point_remove_source_success() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![make_test_receive_state(1)]); |
| let mut event_stream = event_receiver.next(); |
| |
| let remove_op = RemoveSourceOperation::new(1); |
| let mut raw_bytes = vec![0u8; remove_op.encoded_len()]; |
| remove_op.encode(&mut raw_bytes).expect("encode ok"); |
| |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, raw_bytes); |
| |
| let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else { |
| panic!("Expected ServerEvent"); |
| }; |
| assert_eq!(event, ServerEvent::RemoveSource { peer_id: PeerId(1), source_id: 1 }); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(value.is_ok()); |
| } |
| |
| #[test] |
| fn write_non_control_point_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Writing to BroadcastReceiveState characteristic handle (Handle 2) |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(2), 0, vec![0x01]); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(2)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::WriteNotPermitted) |
| )); |
| } |
| |
| #[test] |
| fn control_point_empty_payload_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Writing empty payload to Control Point |
| fake_gatt_server.incoming_write(PeerId(1), BASS_SERVICE_ID, Handle(1), 0, vec![]); |
| |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| |
| let Poll::Ready(Some(FakeServerEvent::WriteResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected WriteResponded event"); |
| }; |
| assert_eq!(handle, Handle(1)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::WriteRequestRejected) |
| )); |
| } |
| |
| #[test] |
| fn read_control_point_fails() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, CONTROL_POINT_HANDLE, 0); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| assert_eq!(handle, CONTROL_POINT_HANDLE); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::ReadNotPermitted) |
| )); |
| } |
| |
| #[test] |
| fn read_receive_state_invalid_offset() { |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![BroadcastReceiveState::Empty]); |
| let mut event_stream = event_receiver.next(); |
| |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, Handle(2), 10); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| assert_eq!(handle, Handle(2)); |
| assert!(matches!( |
| value.unwrap_err(), |
| bt_gatt::types::Error::Gatt(GattError::InvalidOffset) |
| )); |
| } |
| |
| #[test] |
| fn publish_failure() { |
| let mut noop_cx = Context::from_waker(futures::task::noop_waker_ref()); |
| |
| let mut server = ServerBuilder::new() |
| .add_receive_state_characteristic(BroadcastReceiveState::Empty) |
| .build::<FakeTypes>() |
| .expect("building server works"); |
| |
| let (fake_gatt_server, _event_receiver) = FakeServer::new(); |
| fake_gatt_server |
| .set_next_prepare_result(Err(bt_gatt::types::Error::Gatt(GattError::UnlikelyError))); |
| |
| assert!(server.publish(fake_gatt_server).is_ok()); |
| |
| assert!(matches!( |
| server.poll_next_unpin(&mut noop_cx), |
| Poll::Ready(Some(Err(Error::Gatt(bt_gatt::types::Error::Gatt( |
| GattError::UnlikelyError |
| ))))) |
| )); |
| |
| // Verifying server reset state to NotPublished so it is not published |
| assert!(!server.is_published()); |
| } |
| |
| #[test] |
| fn builder_normalizes_source_id() { |
| // Create state with an arbitrary source_id = 99 |
| let state = make_test_receive_state(99); |
| let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) = |
| setup_test_server(vec![state]); |
| let mut event_stream = event_receiver.next(); |
| |
| // Verify slot's internal state and encoded read value have source_id = 1 |
| fake_gatt_server.incoming_read(PeerId(1), BASS_SERVICE_ID, Handle(2), 0); |
| let _ = server.poll_next_unpin(&mut noop_cx); |
| let Poll::Ready(Some(FakeServerEvent::ReadResponded { handle: _, value, .. })) = |
| event_stream.poll_unpin(&mut noop_cx) |
| else { |
| panic!("Expected ReadResponded event"); |
| }; |
| let buf = value.expect("ok"); |
| assert_eq!(buf[0], 1); |
| } |
| |
| #[test] |
| fn allocate_source_id_skips_in_use_ids() { |
| let mut state = ServerState::new(5); |
| state.source_id_to_handle.insert(1, Handle(10)); |
| // Force the next ID to be 1 to simulate wrap around. |
| state.next_source_id = NonZeroU8::new(1).unwrap(); |
| |
| let allocated_id = state.allocate_source_id(Handle(20)).unwrap(); |
| assert_eq!(allocated_id, 2); |
| assert_eq!(state.source_id_to_handle.get(&2), Some(&Handle(20))); |
| } |
| |
| #[test] |
| fn allocate_source_id_error_when_full() { |
| let mut state = ServerState::new(255); |
| for id in 1..=255 { |
| state.source_id_to_handle.insert(id, Handle(u64::from(id))); |
| } |
| |
| assert_eq!(state.allocate_source_id(Handle(300)), Err(GattError::InsufficientResources)); |
| } |
| } |