blob: 29a20b8e7e78791ca785e3fa51ccc9bb369331ee [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 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 log::warn;
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)]
struct PublishedReceiveStateCharacteristic {
/// Handle assigned to this GATT characteristic.
handle: Handle,
/// Current Broadcast Receive State value.
state: BroadcastReceiveState,
}
/// A pending GATT notification to be dispatched following a state modification.
#[derive(Debug, PartialEq, Eq)]
struct PendingNotification {
/// Handle of the characteristic to notify.
handle: Handle,
/// Encoded GATT characteristic value.
value: Vec<u8>,
}
/// A parsed request to write to the Control Point characteristic.
struct ControlPointWriteResult {
/// The decoded event associated with this request.
event: ServerEvent,
/// An optional notification to be sent to GATT peers.
notification: Option<PendingNotification>,
}
/// 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 { .. })
}
fn service(&self) -> Option<&T::LocalService> {
match self {
Self::Published { service, .. } => Some(service),
_ => None,
}
}
fn notify(&self, notification: &PendingNotification) {
let Some(service) = self.service() else {
warn!("Attempted to notify GATT peers on an unpublished service");
return;
};
service.notify(&notification.handle, &notification.value, &[]);
}
}
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)
}
/// Allocates a new SourceId for the next available Receive State
/// characteristic entry.
///
/// Returns the assigned [`Handle`] and [`SourceId`] on success.
fn allocate_receive_state_chrc(&mut self) -> Result<(Handle, SourceId), GattError> {
let handle = self
.receive_state_characteristics
.values()
.find(|c| c.state.is_empty())
.map(|c| c.handle)
.ok_or(GattError::InsufficientResources)?;
let source_id = self.allocate_source_id(handle)?;
Ok((handle, source_id))
}
/// 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],
) -> Result<ControlPointWriteResult, GattError> {
// 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 {
return Err(GattError::WriteNotPermitted);
}
self.handle_control_point_write(peer_id, value)
}
/// 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<ControlPointWriteResult, GattError> {
let mut event = ServerEvent::decode(peer_id, val_bytes)?;
let notification = match &mut event {
ServerEvent::RemoteScanState { .. } => None,
ServerEvent::AddSource { source_id, operation, .. } => {
let (notification, assigned_source_id) = self.handle_add_source_write(operation)?;
*source_id = assigned_source_id;
Some(notification)
}
ServerEvent::ModifySource { operation, .. } => {
let notification = self.handle_modify_source_write(operation)?;
Some(notification)
}
ServerEvent::SetBroadcastCode { source_id, .. } => {
if !self.is_valid_source_id(*source_id) {
return Err(ERROR_INVALID_SOURCE_ID);
}
None
}
ServerEvent::RemoveSource { source_id, .. } => {
let notification = self.handle_remove_source_write(*source_id)?;
Some(notification)
}
};
Ok(ControlPointWriteResult { event, notification })
}
fn handle_add_source_write(
&mut self,
operation: &AddSourceOperation,
) -> Result<(PendingNotification, SourceId), GattError> {
let (handle, source_id) = self.allocate_receive_state_chrc()?;
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,
);
let new_state = BroadcastReceiveState::NonEmpty(initial_state);
let mut value = vec![0u8; new_state.encoded_len()];
new_state.encode(&mut value).map_err(|_| GattError::UnlikelyError)?;
chrc.state = new_state;
Ok((PendingNotification { handle, value }, source_id))
}
fn handle_modify_source_write(
&mut self,
operation: &ModifySourceOperation,
) -> Result<PendingNotification, GattError> {
let Some(handle) = self.source_id_to_handle.get(&operation.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);
};
let BroadcastReceiveState::NonEmpty(ref 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();
let mut modified_state = state.clone();
modified_state.subgroups = subgroups;
modified_state.pa_sync_state = PaSyncState::from_pa_sync(operation.pa_sync);
let new_state = BroadcastReceiveState::NonEmpty(modified_state);
let mut value = vec![0u8; new_state.encoded_len()];
new_state.encode(&mut value).map_err(|_| GattError::UnlikelyError)?;
chrc.state = new_state;
Ok(PendingNotification { handle, value })
}
fn handle_remove_source_write(
&mut self,
source_id: SourceId,
) -> Result<PendingNotification, 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;
let mut value = vec![0u8; chrc.state.encoded_len()];
chrc.state.encode(&mut value).map_err(|_| GattError::UnlikelyError)?;
Ok(PendingNotification { handle, value })
}
}
/// 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(())
}
/// Adds a new [`ReceiveState`] to the next available characteristic slot
/// and sends a GATT notification to all subscribed peers.
///
/// Characteristic slots are allocated up-front when constructing the server
/// via [`ServerBuilder`]. This method populates an available empty slot
/// with `state` when initiated locally.
///
/// Returns the assigned [`SourceId`] for the characteristic on success.
/// Returns [`Error::ServerFull`] if all characteristic slots are occupied.
pub fn add_receive_state(&mut self, mut state: ReceiveState) -> Result<SourceId, Error> {
let (_handle, source_id) =
self.state.allocate_receive_state_chrc().map_err(|_| Error::ServerFull)?;
state.source_id = source_id;
self.set_receive_state_and_notify(source_id, BroadcastReceiveState::NonEmpty(state))?;
Ok(source_id)
}
/// Updates the existing [`BroadcastReceiveState`] characteristic with the
/// provided `state` and sends a GATT notification to all subscribed peers.
///
/// Returns [`Error::InvalidSourceId`] if there is no such characteristic.
pub fn update_receive_state(&mut self, state: ReceiveState) -> Result<(), Error> {
self.set_receive_state_and_notify(state.source_id(), BroadcastReceiveState::NonEmpty(state))
}
/// Clears the entry for the specified Broadcast Receive State
/// characteristic and sends a GATT notification to all subscribed peers.
///
/// Returns [`Error::InvalidSourceId`] if there is no such characteristic.
pub fn clear_receive_state(&mut self, source_id: SourceId) -> Result<(), Error> {
self.set_receive_state_and_notify(source_id, BroadcastReceiveState::Empty)
}
/// Updates the state for the specified Broadcast Receive Characteristic and
/// sends a GATT notification to all subscribed peers.
///
/// Returns [`Error::InvalidSourceId`] if there is no such characteristic.
fn set_receive_state_and_notify(
&mut self,
source_id: SourceId,
new_state: BroadcastReceiveState,
) -> Result<(), Error> {
let handle = *self
.state
.source_id_to_handle
.get(&source_id)
.ok_or(Error::InvalidSourceId(source_id))?;
let chrc = self
.state
.receive_state_characteristics
.get_mut(&handle)
.ok_or(Error::InvalidSourceId(source_id))?;
let mut value = vec![0u8; new_state.encoded_len()];
new_state.encode(&mut value)?;
if new_state.is_empty() {
self.state.source_id_to_handle.remove(&source_id);
} else {
self.state.source_id_to_handle.insert(source_id, handle);
}
chrc.state = new_state;
self.local_service.notify(&PendingNotification { handle, value });
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();
match this.state.handle_write(peer_id, handle, &val_vec) {
Ok(result) => {
if let Some(notification) = result.notification {
this.local_service.notify(&notification);
}
responder.acknowledge();
return Poll::Ready(Some(Ok(result.event)));
}
Err(err) => {
responder.error(err);
}
}
}
_ => {
// TODO(b/534436439): Track CCC for Broadcast Receive State
// characteristics and `PeerInfo` if needed.
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use assert_matches::assert_matches;
use bt_bap::types::BroadcastId;
use bt_common::core::{Address, AddressType, AdvertisingSetId, PeriodicAdvertisingInterval};
use bt_common::generic_audio::metadata_ltv::Metadata;
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::{BigSubgroup, EncryptionStatus, PaSync, PaSyncState, ReceiveState};
fn make_test_inner_receive_state(source_id: SourceId) -> ReceiveState {
ReceiveState::new(
source_id,
AddressType::Public,
Address::new([0x01, 0x02, 0x03, 0x04, 0x05, 0x06]),
AdvertisingSetId::try_from(1).unwrap(),
BroadcastId::try_from(0x123456).unwrap(),
PaSyncState::NotSynced,
EncryptionStatus::NotEncrypted,
vec![],
)
}
fn make_test_receive_state(source_id: SourceId) -> BroadcastReceiveState {
BroadcastReceiveState::NonEmpty(make_test_inner_receive_state(source_id))
}
#[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();
fake_gatt_server.incoming_client_configuration(
PeerId(1),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
// 1. Remote client writes AddSource to Control Point (Handle 1)
let add_op = AddSourceOperation::new(
AddressType::Public,
Address::new([1, 2, 3, 4, 5, 6]),
AdvertisingSetId::try_from(1).unwrap(),
BroadcastId::try_from(0x123456).unwrap(),
PaSync::SyncPastAvailable,
PeriodicAdvertisingInterval(0x0010),
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);
// Server stream emits AddSource event
let Poll::Ready(Some(Ok(event))) = server.poll_next_unpin(&mut noop_cx) else {
panic!("Expected ServerEvent");
};
let ServerEvent::AddSource { peer_id, source_id, operation } = event else {
panic!("Expected AddSource event, got {event:?}");
};
assert_eq!(peer_id, PeerId(1));
assert_eq!(source_id, 1);
assert_eq!(operation, add_op);
// First notification: Initial minimal state created upon Control Point write
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected first Notified event");
};
assert_eq!(handle, Handle(2));
let (first_notification_state, _) = BroadcastReceiveState::decode(&value);
let BroadcastReceiveState::NonEmpty(initial_state) =
first_notification_state.expect("valid decode")
else {
panic!("Expected NonEmpty state in first notification");
};
assert_eq!(initial_state.source_id, 1);
assert_eq!(initial_state.pa_sync_state, PaSyncState::SyncInfoRequest);
// Write response is sent to client
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());
// 2. Upper layer client processes the event, establishes PA sync, and updates
// state
let mut updated_state = initial_state;
updated_state.pa_sync_state = PaSyncState::Synced;
assert!(server.update_receive_state(updated_state).is_ok());
// Second notification: Upper layer client calls public API method
// (update_receive_state)
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, peers, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected second Notified event");
};
assert_eq!(handle, Handle(2));
assert_eq!(peers, vec![PeerId(1)]);
let (second_notification_state, _) = BroadcastReceiveState::decode(&value);
let BroadcastReceiveState::NonEmpty(final_state) =
second_notification_state.expect("valid decode")
else {
panic!("Expected NonEmpty state in second notification");
};
assert_eq!(final_state.source_id, 1);
assert_eq!(final_state.pa_sync_state, PaSyncState::Synced);
}
#[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();
fake_gatt_server.incoming_client_configuration(
PeerId(1),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
// 1. Remote client writes ModifySource to Control Point (Handle 1)
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);
// Server stream emits ModifySource event
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 }
);
// First notification: State updated upon Control Point write
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected first Notified event");
};
assert_eq!(handle, Handle(2));
let (first_state, _) = BroadcastReceiveState::decode(&value);
let BroadcastReceiveState::NonEmpty(modified_state) = first_state.expect("valid decode")
else {
panic!("Expected NonEmpty state");
};
assert_eq!(modified_state.pa_sync_state, PaSyncState::NotSynced);
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());
// 2. Upper layer client processes the event and calls public API method
// (update_receive_state)
let mut updated_state = modified_state;
updated_state.pa_sync_state = PaSyncState::FailedToSync;
assert!(server.update_receive_state(updated_state).is_ok());
// Second notification: Upper layer client calls public API method
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, peers, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected second Notified event");
};
assert_eq!(handle, Handle(2));
assert_eq!(peers, vec![PeerId(1)]);
let (second_state, _) = BroadcastReceiveState::decode(&value);
let BroadcastReceiveState::NonEmpty(final_state) = second_state.expect("valid decode")
else {
panic!("Expected NonEmpty state");
};
assert_eq!(final_state.pa_sync_state, PaSyncState::FailedToSync);
}
#[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();
fake_gatt_server.incoming_client_configuration(
PeerId(1),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
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::Notified { handle, value, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected Notified event");
};
assert_eq!(handle, Handle(2));
assert!(value.is_empty());
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 add_receive_state_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();
fake_gatt_server.incoming_client_configuration(
PeerId(10),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
let state = make_test_inner_receive_state(99);
let result = server.add_receive_state(state);
assert_eq!(result.expect("add ok"), 1);
// Expect the GATT notification for the updated characteristic.
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, peers, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected Notified event");
};
assert_eq!(handle, Handle(2));
assert_eq!(value[0], 1);
assert_eq!(peers, vec![PeerId(10)]);
}
#[test]
fn add_receive_state_when_full_fails() {
let (mut server, _fake_gatt_server, _event_receiver, _noop_cx) =
setup_test_server(vec![make_test_receive_state(1)]);
let state = make_test_inner_receive_state(2);
let result = server.add_receive_state(state);
assert_matches!(result, Err(Error::ServerFull));
}
#[test]
fn update_receive_state_invalid_source_id_fails() {
let (mut server, _fake_gatt_server, _event_receiver, _noop_cx) =
setup_test_server(vec![BroadcastReceiveState::Empty]);
let state = make_test_inner_receive_state(99);
let result = server.update_receive_state(state);
assert_matches!(result, Err(Error::InvalidSourceId(99)));
}
#[test]
fn update_receive_state_notifies_peers() {
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();
fake_gatt_server.incoming_client_configuration(
PeerId(10),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
let new_inner = make_test_inner_receive_state(1);
let expected_state = BroadcastReceiveState::NonEmpty(new_inner.clone());
let mut expected_bytes = vec![0u8; expected_state.encoded_len()];
expected_state.encode(&mut expected_bytes).expect("encode ok");
let result = server.update_receive_state(new_inner);
assert!(result.is_ok());
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, peers, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected Notified event");
};
assert_eq!(handle, Handle(2));
assert_eq!(value, expected_bytes);
assert_eq!(peers, vec![PeerId(10)]);
}
#[test]
fn clear_receive_state_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();
fake_gatt_server.incoming_client_configuration(
PeerId(10),
BASS_SERVICE_ID,
Handle(2),
bt_gatt::server::NotificationType::Notify,
);
let _ = server.poll_next_unpin(&mut noop_cx);
let result = server.clear_receive_state(1);
assert!(result.is_ok());
let Poll::Ready(Some(FakeServerEvent::Notified { handle, value, peers, .. })) =
event_stream.poll_unpin(&mut noop_cx)
else {
panic!("Expected Notified event");
};
assert_eq!(handle, Handle(2));
assert!(value.is_empty());
assert_eq!(peers, vec![PeerId(10)]);
}
#[test]
fn clear_and_reuse_receive_state_slot() {
let (mut server, _fake_gatt_server, _event_receiver, _noop_cx) =
setup_test_server(vec![make_test_receive_state(1)]);
// Clear existing state in slot 1
assert!(server.clear_receive_state(1).is_ok());
// Adding a new state should successfully reuse the cleared slot and allocate
// the next ID
let new_state = make_test_inner_receive_state(50);
let result = server.add_receive_state(new_state);
assert_eq!(result.expect("add ok"), 2);
}
#[test]
fn add_receive_state_sequential_slots() {
let (mut server, _fake_gatt_server, _event_receiver, _noop_cx) = setup_test_server(vec![
BroadcastReceiveState::Empty,
BroadcastReceiveState::Empty,
BroadcastReceiveState::Empty,
]);
let state1 = make_test_inner_receive_state(10);
let state2 = make_test_inner_receive_state(20);
let state3 = make_test_inner_receive_state(30);
assert_eq!(server.add_receive_state(state1).unwrap(), 1);
assert_eq!(server.add_receive_state(state2).unwrap(), 2);
assert_eq!(server.add_receive_state(state3).unwrap(), 3);
// Fourth add should fail as server is full
let state4 = make_test_inner_receive_state(40);
assert_matches!(server.add_receive_state(state4), Err(Error::ServerFull));
}
#[test]
fn read_cleared_receive_state_returns_empty() {
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();
// Clear the state
assert!(server.clear_receive_state(1).is_ok());
// Perform a GATT Read on Handle(2)
fake_gatt_server.incoming_read(PeerId(10), 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));
let buf = value.expect("read ok");
assert!(buf.is_empty());
}
#[test]
fn update_receive_state_encoding_failure_preserves_state() {
let initial_inner = make_test_inner_receive_state(1);
let (mut server, fake_gatt_server, mut event_receiver, mut noop_cx) =
setup_test_server(vec![BroadcastReceiveState::NonEmpty(initial_inner.clone())]);
let mut event_stream = event_receiver.next();
// Create an invalid ReceiveState whose metadata total length exceeds 255 bytes
// (causing BigSubgroup encode to fail)
let vendor_metadata = vec![Metadata::AudioActiveState(true); 130];
let mut invalid_inner = initial_inner.clone();
invalid_inner.subgroups = vec![BigSubgroup::new(None).with_metadata(vendor_metadata)];
// Attempting to update to invalid_inner should return a Packet encoding error
let result = server.update_receive_state(invalid_inner);
assert_matches!(result, Err(Error::Packet(_)));
// Perform a GATT Read on Handle(2) to verify internal memory state was NOT
// mutated to invalid_state
fake_gatt_server.incoming_read(PeerId(10), 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));
// Read value should successfully decode back to initial_inner state
let buf = value.expect("read ok");
let (decoded_state, _) = BroadcastReceiveState::decode(&buf);
assert_eq!(
decoded_state.expect("decode ok"),
BroadcastReceiveState::NonEmpty(initial_inner)
);
}
#[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));
}
}