blob: 715a80326585a5cabf128e33676e1a847fa09c3d [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 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));
}
}