| // 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(¬if.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(¬if.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(¬if.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::new(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::new(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(), SetIdentityResolvingKey::new(SirkType::Plaintext, [0xAA; 16])); |
| assert_eq!(client.size(), Some(CoordinatedSetSize::new(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::new(NonZeroU8::new(3).unwrap())) |
| ); |
| assert_eq!(client.size(), Some(CoordinatedSetSize::new(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::new( |
| SirkType::Plaintext, |
| [0xBB; 16] |
| )) |
| ); |
| assert_eq!(client.sirk(), SetIdentityResolvingKey::new(SirkType::Plaintext, [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::new(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()); |
| } |
| } |