use core::iter::once;
use core::pin::pin;
use bluer::adv::Advertisement;
use bluer::agent::Agent;
use bluer::gatt::local::{
characteristic_control, Application, Characteristic, CharacteristicControl,
CharacteristicControlEvent, CharacteristicNotify, CharacteristicNotifyMethod,
CharacteristicWrite, CharacteristicWriteMethod, CharacteristicWriteRequest, Service,
};
use bluer::gatt::remote::Characteristic as RemoteCharacteristic;
use bluer::gatt::{CharacteristicWriter, WriteOp};
use bluer::{Adapter, AdapterEvent, Address, DiscoveryFilter, DiscoveryTransport, Uuid};
use embassy_futures::select::{select, select3, select4, Either};
use embassy_time::{Duration, Timer};
use tokio::sync::mpsc::Receiver;
use tokio_stream::StreamExt;
use crate::error::{Error, ErrorCode};
use crate::transport::network::btp::Btp;
use crate::transport::network::mdns::CommissionableFilter;
use crate::transport::network::BtAddr;
use crate::utils::select::Coalesce;
use super::{AdvData, C1_CHARACTERISTIC_UUID, C2_CHARACTERISTIC_UUID, MATTER_BLE_SERVICE_UUID};
pub const DEFAULT_SCAN_TIMEOUT_SECS: u16 = 60;
const CONNECT_ATTEMPTS: u8 = 4;
const CONNECT_RETRY_DELAY_MS: u64 = 500;
pub async fn run_peripheral(
adapter_name: Option<&str>,
service_name: &str,
service_adv_data: &AdvData,
btp: &Btp,
) -> Result<(), Error> {
let session = bluer::Session::new().await?;
let _handle = session.register_agent(Agent::default()).await?;
let adapter = if let Some(adapter_name) = adapter_name {
session.adapter(adapter_name)?
} else {
session.default_adapter().await?
};
adapter.set_powered(true).await?;
let le_advertisement = Advertisement {
discoverable: Some(true),
local_name: Some(service_name.into()),
service_uuids: once(Uuid::from_u128(MATTER_BLE_SERVICE_UUID)).collect(),
service_data: once((
Uuid::from_u128(MATTER_BLE_SERVICE_UUID),
service_adv_data.service_payload_iter().collect(),
))
.collect(),
..Default::default()
};
let (write_sender, mut write_receiver) = tokio::sync::mpsc::channel(1);
let (mut notify_cc, notify_cc_handle) = characteristic_control();
let app = Application {
services: vec![Service {
uuid: Uuid::from_u128(MATTER_BLE_SERVICE_UUID),
primary: true,
characteristics: vec![
Characteristic {
uuid: Uuid::from_u128(C1_CHARACTERISTIC_UUID),
write: Some(CharacteristicWrite {
write: true,
method: CharacteristicWriteMethod::Fun(Box::new(move |new_value, req| {
let sender = write_sender.clone();
Box::pin(async move {
sender.send((new_value, req)).await.unwrap();
Ok(())
})
})),
..Default::default()
}),
..Default::default()
},
Characteristic {
uuid: Uuid::from_u128(C2_CHARACTERISTIC_UUID),
notify: Some(CharacteristicNotify {
indicate: true,
method: CharacteristicNotifyMethod::Io,
..Default::default()
}),
control_handle: notify_cc_handle,
..Default::default()
},
],
..Default::default()
}],
..Default::default()
};
let _app_handle = adapter.serve_gatt_application(app).await?;
info!(
"Serving Matter GATT BTP service on Bluetooth adapter {}",
adapter.name()
);
loop {
let notifier = {
let _adv_handle = adapter.advertise(le_advertisement.clone()).await?;
info!(
"Advertising Matter GATT BTP service on Bluetooth adapter {}",
adapter.name(),
);
notifier(&mut notify_cc).await
};
btp.reset();
select4(
wait_complete(btp, ¬ifier),
process_write(btp, &mut write_receiver),
process_indicate(btp, None, ¬ifier, &mut [0; 512]),
process_cc_events(&mut notify_cc),
)
.coalesce()
.await?;
}
}
pub async fn scan<F, R>(
adapter_name: Option<&str>,
filter: &CommissionableFilter,
scan_timeout: Option<u16>,
mut on_found: F,
) -> Result<R, Error>
where
F: FnMut(BtAddr, &AdvData) -> Option<R>,
{
let session = bluer::Session::new().await?;
let adapter = open_adapter(&session, adapter_name).await?;
adapter.set_powered(true).await?;
info!(
"Scanning for a commissionable Matter device on Bluetooth adapter {} (filter: {:?})",
adapter.name(),
filter
);
let discovery_filter = DiscoveryFilter {
uuids: once(Uuid::from_u128(MATTER_BLE_SERVICE_UUID)).collect(),
transport: DiscoveryTransport::Le,
..Default::default()
};
adapter.set_discovery_filter(discovery_filter).await?;
let device_events = adapter.discover_devices_with_changes().await?;
let mut device_events = pin!(device_events);
let mut reported: heapless::Vec<BtAddr, 16> = heapless::Vec::new();
let scan_fut = async {
while let Some(event) = device_events.next().await {
if let AdapterEvent::DeviceAdded(addr) = event {
if let Some(result) =
try_match_device(&adapter, addr, filter, &mut reported, &mut on_found).await?
{
return Ok(Some(result));
}
}
}
Ok::<Option<R>, Error>(None)
};
let timeout = Timer::after(Duration::from_secs(
scan_timeout.unwrap_or(DEFAULT_SCAN_TIMEOUT_SECS) as u64,
));
let outcome = match select(scan_fut, timeout).await {
Either::First(result) => result?,
Either::Second(_) => None,
};
outcome.ok_or_else(|| {
warn!(
"No commissionable Matter device matching the filter was found within the scan timeout"
);
ErrorCode::NoNetworkInterface.into()
})
}
async fn try_match_device<F, R>(
adapter: &Adapter,
addr: Address,
filter: &CommissionableFilter,
reported: &mut heapless::Vec<BtAddr, 16>,
on_found: &mut F,
) -> Result<Option<R>, Error>
where
F: FnMut(BtAddr, &AdvData) -> Option<R>,
{
let bt_addr = BtAddr(addr.0);
if reported.contains(&bt_addr) {
return Ok(None);
}
let Ok(device) = adapter.device(addr) else {
return Ok(None);
};
let Ok(Some(service_data)) = device.service_data().await else {
return Ok(None);
};
let Some(bytes) = service_data.get(&Uuid::from_u128(MATTER_BLE_SERVICE_UUID)) else {
return Ok(None);
};
let Some(adv) = AdvData::parse_service_data(bytes) else {
return Ok(None);
};
if !adv.matches(filter) {
return Ok(None);
}
let _ = reported.push(bt_addr);
debug!("Matched commissionable device {} (adv: {:?})", bt_addr, adv);
Ok(on_found(bt_addr, &adv))
}
pub async fn run_central(adapter_name: Option<&str>, addr: BtAddr, btp: &Btp) -> Result<(), Error> {
let session = bluer::Session::new().await?;
let _handle = session.register_agent(Agent::default()).await?;
let adapter = open_adapter(&session, adapter_name).await?;
adapter.set_powered(true).await?;
let device = adapter.device(Address(addr.0))?;
info!("Connecting to commissionable device {}", addr);
connect_with_retry(&device).await?;
let (c1, c2) = discover_matter_characteristics(&device).await?;
debug!("Discovered Matter GATT characteristics C1/C2");
let c2_notify = c2.notify().await?;
let mut c2_notify = pin!(c2_notify);
btp.set_initiator(true);
let gatt_mtu = None;
select3(
wait_central_complete(btp, &device),
process_c2_indications(btp, addr, gatt_mtu, &mut c2_notify),
process_c1_writes(btp, gatt_mtu, &c1, &mut [0; 512]),
)
.coalesce()
.await
}
async fn process_c2_indications(
btp: &Btp,
peer_addr: BtAddr,
gatt_mtu: Option<u16>,
c2_notify: &mut (impl StreamExt<Item = Vec<u8>> + Unpin),
) -> Result<(), Error> {
while let Some(value) = c2_notify.next().await {
if value.is_empty() {
continue;
}
trace!(
"Received C2 indication from peer {}: {:?}",
peer_addr,
value
);
btp.process_incoming(gatt_mtu, peer_addr, &value)?;
}
Ok(())
}
async fn process_c1_writes(
btp: &Btp,
gatt_mtu: Option<u16>,
c1: &RemoteCharacteristic,
buf: &mut [u8],
) -> Result<(), Error> {
let req = bluer::gatt::remote::CharacteristicWriteRequest {
op_type: WriteOp::Request,
..Default::default()
};
loop {
let len = btp.process_outgoing(gatt_mtu, buf)?;
if len > 0 {
trace!("Writing to C1: {:?}", &buf[..len]);
c1.write_ext(&buf[..len], &req).await?;
} else {
btp.wait_outgoing().await;
}
}
}
async fn wait_central_complete(btp: &Btp, _device: &bluer::Device) -> Result<(), Error> {
btp.wait_timeout().await;
info!("Timeout while waiting for data from the peer");
Ok(())
}
async fn connect_with_retry(device: &bluer::Device) -> Result<(), Error> {
for attempt in 1..=CONNECT_ATTEMPTS {
match device.connect().await {
Ok(()) => return Ok(()),
Err(e) if attempt < CONNECT_ATTEMPTS => {
warn!(
"Connect attempt {}/{} failed ({:?}); retrying",
attempt, CONNECT_ATTEMPTS, e
);
Timer::after(Duration::from_millis(CONNECT_RETRY_DELAY_MS)).await;
}
Err(e) => return Err(e.into()),
}
}
Err(ErrorCode::NoNetworkInterface.into())
}
async fn discover_matter_characteristics(
device: &bluer::Device,
) -> Result<(RemoteCharacteristic, RemoteCharacteristic), Error> {
let matter_service_uuid = Uuid::from_u128(MATTER_BLE_SERVICE_UUID);
let c1_uuid = Uuid::from_u128(C1_CHARACTERISTIC_UUID);
let c2_uuid = Uuid::from_u128(C2_CHARACTERISTIC_UUID);
for service in device.services().await? {
if service.uuid().await? != matter_service_uuid {
continue;
}
let mut c1 = None;
let mut c2 = None;
for chr in service.characteristics().await? {
let uuid = chr.uuid().await?;
if uuid == c1_uuid {
c1 = Some(chr);
} else if uuid == c2_uuid {
c2 = Some(chr);
}
}
let c1 = c1.ok_or_else(|| {
warn!("The Matter GATT service is missing the C1 characteristic");
Error::from(ErrorCode::NoNetworkInterface)
})?;
let c2 = c2.ok_or_else(|| {
warn!("The Matter GATT service is missing the C2 characteristic");
Error::from(ErrorCode::NoNetworkInterface)
})?;
return Ok((c1, c2));
}
warn!("The connected device does not expose the Matter GATT service");
Err(ErrorCode::NoNetworkInterface.into())
}
async fn open_adapter(
session: &bluer::Session,
adapter_name: Option<&str>,
) -> Result<Adapter, Error> {
let adapter = if let Some(adapter_name) = adapter_name {
session.adapter(adapter_name)?
} else {
session.default_adapter().await?
};
Ok(adapter)
}
async fn process_write(
btp: &Btp,
receiver: &mut Receiver<(Vec<u8>, CharacteristicWriteRequest)>,
) -> Result<(), Error> {
while let Some((value, req)) = receiver.recv().await {
btp.process_incoming(Some(req.mtu), BtAddr(req.device_address.0), &value)?;
}
Ok(())
}
async fn process_indicate(
btp: &Btp,
gatt_mtu: Option<u16>,
notifier: &CharacteristicWriter,
buf: &mut [u8],
) -> Result<(), Error> {
loop {
let len = btp.process_outgoing(gatt_mtu, buf)?;
if len > 0 {
notifier.send(&buf[..len]).await?;
} else {
btp.wait_outgoing().await;
}
}
}
async fn process_cc_events(cc: &mut CharacteristicControl) -> Result<(), Error> {
loop {
let _ = notifier(cc).await;
}
}
async fn wait_complete(btp: &Btp, notifier: &CharacteristicWriter) -> Result<(), Error> {
let result = select(notifier.closed(), btp.wait_timeout()).await;
match result {
Either::First(_) => info!("Peer unsubscribed"),
Either::Second(_) => info!("Timeout while waiting for data from the peer"),
}
Ok(())
}
async fn notifier(cc: &mut CharacteristicControl) -> CharacteristicWriter {
loop {
if let Some(notifier) = cc.next().await.map(|event| {
let CharacteristicControlEvent::Notify(notifier) = event else {
unreachable!();
};
notifier
}) {
break notifier;
}
}
}