blob: 0f78dcf611f3a2111e590dcc7da26d1a75634cd4 [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.
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};
use futures::stream::{BoxStream, FusedStream, SelectAll, Stream, StreamExt};
use parking_lot::Mutex;
use bt_common::packet_encoding::Decodable;
use bt_gatt::client::CharacteristicNotification;
use bt_gatt::types::Handle;
use crate::client::{CsisClientState, Error};
use crate::types::*;
/// A change to a CSIS characteristic value signaled by the Set Member via a
/// GATT notification.
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum CsisNotification {
/// The Set Identity Resolving Key characteristic value changed. The value
/// is reported as exposed by the server, so it may be an encrypted SIRK
/// (CSIS v1.1 Section 4.5).
SirkChanged(SetIdentityResolvingKey),
/// The Coordinated Set Size characteristic value changed.
SizeChanged(CoordinatedSetSize),
/// The Set Member Lock characteristic value changed (CSIS v1.1 Section
/// 5.3.1). A Set Member does not notify the client whose own write caused
/// the change.
LockChanged(SetMemberLock),
}
/// Handles of the CSIS characteristics whose notifications are recognized.
#[derive(Debug, Clone, Copy)]
pub(crate) struct NotifiedHandles {
pub(crate) sirk: Handle,
pub(crate) size: Option<Handle>,
pub(crate) lock: Option<Handle>,
}
/// A stream of CSIS characteristic value changes notified by a Set Member.
///
/// The notified value is written back to the client's cached state before the
/// event is yielded.
pub struct NotificationStream {
notification_streams:
SelectAll<BoxStream<'static, Result<CharacteristicNotification, bt_gatt::types::Error>>>,
handles: NotifiedHandles,
state: Arc<Mutex<CsisClientState>>,
terminated: bool,
}
impl NotificationStream {
pub(crate) fn new(
notification_streams: SelectAll<
BoxStream<'static, Result<CharacteristicNotification, bt_gatt::types::Error>>,
>,
handles: NotifiedHandles,
state: Arc<Mutex<CsisClientState>>,
) -> Self {
Self { notification_streams, handles, state, terminated: false }
}
/// Decodes a notification and updates the cached state.
fn handle_notification(
&mut self,
notif: CharacteristicNotification,
) -> Result<CsisNotification, Error> {
if notif.handle == self.handles.sirk {
let (res, _) = SetIdentityResolvingKey::decode(&notif.value);
let sirk = res?;
self.state.lock().set_sirk(sirk);
return Ok(CsisNotification::SirkChanged(sirk));
}
if Some(notif.handle) == self.handles.size {
let (res, _) = CoordinatedSetSize::decode(&notif.value);
let size = res?;
self.state.lock().set_size(size);
return Ok(CsisNotification::SizeChanged(size));
}
if Some(notif.handle) == self.handles.lock {
let (res, _) = SetMemberLock::decode(&notif.value);
let lock = res?;
self.state.lock().set_lock_state(lock);
return Ok(CsisNotification::LockChanged(lock));
}
Err(Error::UnexpectedNotification(notif.handle))
}
}
impl Stream for NotificationStream {
type Item = Result<CsisNotification, Error>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
if self.terminated {
return Poll::Ready(None);
}
match futures::ready!(self.notification_streams.poll_next_unpin(cx)) {
Some(Ok(notif)) => Poll::Ready(Some(self.handle_notification(notif))),
Some(Err(e)) => Poll::Ready(Some(Err(Error::Gatt(e)))),
None => {
self.terminated = true;
Poll::Ready(None)
}
}
}
}
impl FusedStream for NotificationStream {
fn is_terminated(&self) -> bool {
self.terminated
}
}
#[cfg(test)]
mod tests {
use super::*;
use assert_matches::assert_matches;
use bt_common::Uuid;
use bt_gatt::test_utils::{FakePeerService, FakeTypes};
use bt_gatt::types::{AttributePermissions, Characteristic, CharacteristicProperty};
use core::num::NonZeroU8;
use futures::task::noop_waker_ref;
use futures::FutureExt;
use crate::client::CoordinatedSetIdentificationServiceClient;
const SIRK_HANDLE: Handle = Handle(1);
const SIZE_HANDLE: Handle = Handle(2);
const LOCK_HANDLE: Handle = Handle(3);
const RANK_HANDLE: Handle = Handle(4);
fn add_char(
service: &mut FakePeerService,
handle: Handle,
uuid: Uuid,
value: Vec<u8>,
notifiable: bool,
) {
let properties = if notifiable {
(&[CharacteristicProperty::Read, CharacteristicProperty::Notify]).into()
} else {
CharacteristicProperty::Read.into()
};
service.add_characteristic(
Characteristic {
handle,
uuid,
properties,
permissions: AttributePermissions::default(),
descriptors: vec![],
},
value,
);
}
fn sirk_value(key: [u8; 16]) -> Vec<u8> {
let mut value = vec![0x01]; // Plaintext
value.extend_from_slice(&key);
value
}
/// Builds a Set Member exposing SIRK, Size, Lock and Rank, with SIRK, Size
/// and Lock notifiable.
fn fake_service() -> FakePeerService {
let mut service = FakePeerService::new();
add_char(
&mut service,
SIRK_HANDLE,
SET_IDENTITY_RESOLVING_KEY_UUID,
sirk_value([0xAA; 16]),
true,
);
add_char(&mut service, SIZE_HANDLE, COORDINATED_SET_SIZE_UUID, vec![0x02], true);
add_char(&mut service, LOCK_HANDLE, SET_MEMBER_LOCK_UUID, vec![0x01], true);
// CSIS v1.1 Table 5.1 note C.1: Rank is mandatory when Lock is
// present, and is never notifiable.
add_char(&mut service, RANK_HANDLE, SET_MEMBER_RANK_UUID, vec![0x01], false);
service
}
fn build_client(
service: FakePeerService,
) -> CoordinatedSetIdentificationServiceClient<FakeTypes> {
let mut noop_cx = Context::from_waker(noop_waker_ref());
let create_fut = CoordinatedSetIdentificationServiceClient::<FakeTypes>::create(service);
let mut create_fut = std::pin::pin!(create_fut);
let Poll::Ready(client) = create_fut.poll_unpin(&mut noop_cx) else {
panic!("Expected create to be ready");
};
client.expect("Expected create to succeed")
}
#[test]
fn create_does_not_subscribe() {
let service = fake_service();
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service.clone());
// Dropped, nothing is subscribed yet. Differs from the value below so
// a client that subscribed in `create` would fail this test.
service.notify(
&SIZE_HANDLE,
Ok(CharacteristicNotification {
handle: SIZE_HANDLE,
value: vec![0x05],
maybe_truncated: false,
}),
);
assert_eq!(client.size(), Some(CoordinatedSetSize(NonZeroU8::new(2).unwrap())));
let mut stream = client.take_notification_stream().expect("stream available");
service.notify(
&SIZE_HANDLE,
Ok(CharacteristicNotification {
handle: SIZE_HANDLE,
value: vec![0x03],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a size notification");
};
assert_eq!(
event,
CsisNotification::SizeChanged(CoordinatedSetSize(NonZeroU8::new(3).unwrap()))
);
}
#[test]
fn notifications_update_cached_state() {
let service = fake_service();
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service.clone());
let mut stream = client.take_notification_stream().expect("stream available");
// Values read during discovery are visible before any notification.
assert_eq!(client.sirk().value, [0xAA; 16]);
assert_eq!(client.size(), Some(CoordinatedSetSize(NonZeroU8::new(2).unwrap())));
assert_eq!(client.lock_state(), Some(SetMemberLock::Unlocked));
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Pending);
// A lock timeout expiring, or another client releasing the lock
// (CSIS v1.1 Section 5.3.1.1).
service.notify(
&LOCK_HANDLE,
Ok(CharacteristicNotification {
handle: LOCK_HANDLE,
value: vec![0x02],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a lock notification");
};
assert_eq!(event, CsisNotification::LockChanged(SetMemberLock::Locked));
assert_eq!(client.lock_state(), Some(SetMemberLock::Locked));
service.notify(
&SIZE_HANDLE,
Ok(CharacteristicNotification {
handle: SIZE_HANDLE,
value: vec![0x03],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a size notification");
};
assert_eq!(
event,
CsisNotification::SizeChanged(CoordinatedSetSize(NonZeroU8::new(3).unwrap()))
);
assert_eq!(client.size(), Some(CoordinatedSetSize(NonZeroU8::new(3).unwrap())));
service.notify(
&SIRK_HANDLE,
Ok(CharacteristicNotification {
handle: SIRK_HANDLE,
value: sirk_value([0xBB; 16]),
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a SIRK notification");
};
assert_eq!(
event,
CsisNotification::SirkChanged(SetIdentityResolvingKey {
sirk_type: SirkType::Plaintext,
value: [0xBB; 16],
})
);
assert_eq!(client.sirk().value, [0xBB; 16]);
}
#[test]
fn unrecognized_notifications_are_reported_as_errors() {
let mut service = FakePeerService::new();
add_char(
&mut service,
SIRK_HANDLE,
SET_IDENTITY_RESOLVING_KEY_UUID,
sirk_value([0xAA; 16]),
true,
);
// Not notifiable, so it is never subscribed.
add_char(&mut service, SIZE_HANDLE, COORDINATED_SET_SIZE_UUID, vec![0x02], false);
// CSIS v1.1 Table 5.1 excludes Notify for Set Member Rank, so this
// Set Member is not conformant. It is subscribed but never recognized.
add_char(&mut service, RANK_HANDLE, SET_MEMBER_RANK_UUID, vec![0x01], true);
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service.clone());
let mut stream = client.take_notification_stream().expect("stream available");
// Rank is subscribed but its notifications are not recognized.
service.notify(
&RANK_HANDLE,
Ok(CharacteristicNotification {
handle: RANK_HANDLE,
value: vec![0x02],
maybe_truncated: false,
}),
);
assert_matches!(
stream.poll_next_unpin(&mut noop_cx),
Poll::Ready(Some(Err(Error::UnexpectedNotification(RANK_HANDLE))))
);
assert_eq!(client.rank(), Some(SetMemberRank(NonZeroU8::new(1).unwrap())));
// The SIRK characteristic is still subscribed.
service.notify(
&SIRK_HANDLE,
Ok(CharacteristicNotification {
handle: SIRK_HANDLE,
value: sirk_value([0xBB; 16]),
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a SIRK notification");
};
assert_matches!(event, CsisNotification::SirkChanged(_));
// The Rank subscription is still open, so the stream stays pending.
service.clear_notifier(&SIRK_HANDLE);
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Pending);
}
#[test]
fn stream_terminates_when_nothing_is_notifiable() {
let mut service = FakePeerService::new();
add_char(
&mut service,
SIRK_HANDLE,
SET_IDENTITY_RESOLVING_KEY_UUID,
sirk_value([0xAA; 16]),
false,
);
add_char(&mut service, SIZE_HANDLE, COORDINATED_SET_SIZE_UUID, vec![0x02], false);
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service);
let mut stream = client.take_notification_stream().expect("stream available");
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Ready(None));
assert!(stream.is_terminated());
}
#[test]
fn stream_terminates_once_every_subscription_closes() {
let service = fake_service();
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service.clone());
let mut stream = client.take_notification_stream().expect("stream available");
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Pending);
// Dropping one subscription is not enough: the rest can still notify.
service.clear_notifier(&SIRK_HANDLE);
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Pending);
service.notify(
&LOCK_HANDLE,
Ok(CharacteristicNotification {
handle: LOCK_HANDLE,
value: vec![0x02],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a lock notification");
};
assert_eq!(event, CsisNotification::LockChanged(SetMemberLock::Locked));
// Once they have all closed, the stream ends.
service.clear_notifier(&SIZE_HANDLE);
service.clear_notifier(&LOCK_HANDLE);
assert_matches!(stream.poll_next_unpin(&mut noop_cx), Poll::Ready(None));
assert!(stream.is_terminated());
}
#[test]
fn errors_are_reported_without_ending_the_stream() {
let service = fake_service();
let mut noop_cx = Context::from_waker(noop_waker_ref());
let mut client = build_client(service.clone());
let mut stream = client.take_notification_stream().expect("stream available");
// 0x00 is RFU for the Set Member Lock characteristic.
service.notify(
&LOCK_HANDLE,
Ok(CharacteristicNotification {
handle: LOCK_HANDLE,
value: vec![0x00],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Err(e))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a decoding error");
};
assert_matches!(e, Error::Packet(_));
assert!(!stream.is_terminated());
// The subscription is still live after a GATT error.
service.notify(
&SIZE_HANDLE,
Err(bt_gatt::types::Error::PeerDisconnected(bt_common::PeerId(1))),
);
let Poll::Ready(Some(Err(e))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a GATT error");
};
assert_matches!(e, Error::Gatt(_));
assert!(!stream.is_terminated());
// The cached value is left alone and later notifications still arrive.
assert_eq!(client.lock_state(), Some(SetMemberLock::Unlocked));
service.notify(
&LOCK_HANDLE,
Ok(CharacteristicNotification {
handle: LOCK_HANDLE,
value: vec![0x02],
maybe_truncated: false,
}),
);
let Poll::Ready(Some(Ok(event))) = stream.poll_next_unpin(&mut noop_cx) else {
panic!("Expected a lock notification");
};
assert_eq!(event, CsisNotification::LockChanged(SetMemberLock::Locked));
}
#[test]
fn notification_stream_is_only_available_once() {
let service = fake_service();
let mut client = build_client(service);
assert!(client.take_notification_stream().is_some());
assert!(client.take_notification_stream().is_none());
}
}