use std::collections::HashMap;
use std::time::Duration;
use futures_util::FutureExt;
use futures_util::stream::{StreamExt, select};
use log::warn;
use tokio::time::{sleep, timeout};
use zbus::fdo::DBusProxy;
use zbus::message::Type;
use zbus::zvariant::{OwnedObjectPath, OwnedValue};
use zbus::{Connection, MatchRule, Message, MessageStream, proxy};
use crate::widget::{Bluetooth, BluetoothState, Msg};
use crate::producer::{MsgSender, Producer, ProducerFuture, ProducerResult};
const DEFAULT_INTERVAL: Duration = Duration::from_secs(60);
const SETTLE: Duration = Duration::from_millis(150);
const BLUEZ: &str = "org.bluez";
const MAX_QUEUED_SIGNALS: usize = 64;
const ADAPTER_IFACE: &str = "org.bluez.Adapter1";
const DEVICE_IFACE: &str = "org.bluez.Device1";
pub fn bluetooth_from_bluez(
adapter_count: usize,
powered: Option<bool>,
connected: u32,
) -> Bluetooth {
if adapter_count == 0 {
return Bluetooth::new(BluetoothState::Unavailable, 0);
}
match powered {
Some(true) => Bluetooth::new(BluetoothState::On, connected),
Some(false) => Bluetooth::new(BluetoothState::Off, 0),
None => Bluetooth::new(BluetoothState::Unavailable, 0),
}
}
fn read_bool(props: &HashMap<String, OwnedValue>, key: &str) -> Option<bool> {
props.get(key).and_then(|v| bool::try_from(v).ok())
}
fn summarize(
objects: &HashMap<OwnedObjectPath, HashMap<String, HashMap<String, OwnedValue>>>,
) -> (usize, Option<bool>, u32) {
let mut adapter_count = 0usize;
let mut any_on = false;
let mut any_off = false;
for interfaces in objects.values() {
let Some(adapter_props) = interfaces.get(ADAPTER_IFACE) else {
continue;
};
adapter_count += 1;
match read_bool(adapter_props, "Powered") {
Some(true) => any_on = true,
Some(false) => any_off = true,
None => {}
}
}
let powered = if any_on {
Some(true)
} else if any_off {
Some(false)
} else {
None
};
let mut connected: u32 = 0;
for interfaces in objects.values() {
if let Some(device_props) = interfaces.get(DEVICE_IFACE) {
if read_bool(device_props, "Connected") == Some(true) {
connected = connected.saturating_add(1);
}
}
}
(adapter_count, powered, connected)
}
#[proxy(
interface = "org.freedesktop.DBus.ObjectManager",
default_service = "org.bluez",
default_path = "/"
)]
trait ObjectManager {
fn get_managed_objects(
&self,
) -> zbus::Result<HashMap<OwnedObjectPath, HashMap<String, HashMap<String, OwnedValue>>>>;
}
async fn read_snapshot(om: &ObjectManagerProxy<'_>) -> Bluetooth {
match om.get_managed_objects().await {
Ok(objects) => {
let (adapter_count, powered, connected) = summarize(&objects);
bluetooth_from_bluez(adapter_count, powered, connected)
}
Err(e) => {
warn!("bluetooth: reading BlueZ state failed: {e}");
Bluetooth::new(BluetoothState::Unavailable, 0)
}
}
}
fn signal_changes_state(message: &Message) -> bool {
let header = message.header();
let (Some(interface), Some(member)) = (header.interface(), header.member()) else {
return false;
};
match (interface.as_str(), member.as_str()) {
("org.freedesktop.DBus.ObjectManager", "InterfacesAdded" | "InterfacesRemoved") => true,
("org.freedesktop.DBus.Properties", "PropertiesChanged") => message
.body()
.deserialize::<(String, HashMap<String, OwnedValue>, Vec<String>)>()
.is_ok_and(|(interface, changed, invalidated)| {
let names = changed.keys().chain(&invalidated).map(String::as_str);
property_change_is_displayed(&interface, names)
}),
_ => false,
}
}
fn property_change_is_displayed<'a>(
interface: &str,
mut names: impl Iterator<Item = &'a str>,
) -> bool {
match interface {
ADAPTER_IFACE => names.any(|name| name == "Powered"),
DEVICE_IFACE => names.any(|name| name == "Connected"),
_ => false,
}
}
pub struct BluetoothProducer {
interval: Duration,
}
impl BluetoothProducer {
pub fn new() -> Self {
Self {
interval: DEFAULT_INTERVAL,
}
}
pub fn with_interval(interval: Duration) -> Self {
Self { interval }
}
}
impl Default for BluetoothProducer {
fn default() -> Self {
Self::new()
}
}
impl Producer for BluetoothProducer {
fn name(&self) -> String {
"bluetooth".to_string()
}
fn run(self: Box<Self>, tx: MsgSender) -> ProducerFuture {
Box::pin(run(tx, self.interval))
}
}
async fn run(tx: MsgSender, resync: Duration) -> ProducerResult {
let conn = Connection::system().await?;
let om = ObjectManagerProxy::new(&conn).await?;
let rule = MatchRule::builder()
.msg_type(Type::Signal)
.sender(BLUEZ)?
.build();
let signals = MessageStream::for_match_rule(rule, &conn, Some(MAX_QUEUED_SIGNALS))
.await?
.filter_map(|message| async move {
message
.is_ok_and(|message| signal_changes_state(&message))
.then_some(())
});
let owners = DBusProxy::new(&conn)
.await?
.receive_name_owner_changed_with_args(&[(0, BLUEZ)])
.await?
.map(|_| ());
let events = select(signals, owners);
futures_util::pin_mut!(events);
let mut sent = None;
loop {
if !send_if_changed(&tx, &mut sent, read_snapshot(&om).await) {
return Ok(());
}
if let Ok(None) = timeout(resync, events.next()).await {
return Ok(());
}
sleep(SETTLE).await;
while let Some(Some(())) = events.next().now_or_never() {}
}
}
fn send_if_changed(tx: &MsgSender, sent: &mut Option<Bluetooth>, snapshot: Bluetooth) -> bool {
if *sent == Some(snapshot) {
return true;
}
*sent = Some(snapshot);
tx.send(Msg::Bluetooth(snapshot)).is_ok()
}
#[cfg(test)]
mod tests {
use super::*;
fn empty_objects() -> HashMap<OwnedObjectPath, HashMap<String, HashMap<String, OwnedValue>>> {
HashMap::new()
}
fn adapter_with(powered: Option<bool>) -> HashMap<String, OwnedValue> {
let mut props = HashMap::new();
if let Some(p) = powered {
props.insert("Powered".to_string(), OwnedValue::from(p));
}
props
}
fn device_with(connected: Option<bool>) -> HashMap<String, OwnedValue> {
let mut props = HashMap::new();
if let Some(c) = connected {
props.insert("Connected".to_string(), OwnedValue::from(c));
}
props
}
fn insert_adapter(
objects: &mut HashMap<OwnedObjectPath, HashMap<String, HashMap<String, OwnedValue>>>,
path: &str,
props: HashMap<String, OwnedValue>,
) {
let mut ifaces = HashMap::new();
ifaces.insert(ADAPTER_IFACE.to_string(), props);
objects.insert(OwnedObjectPath::try_from(path).unwrap(), ifaces);
}
fn insert_device(
objects: &mut HashMap<OwnedObjectPath, HashMap<String, HashMap<String, OwnedValue>>>,
path: &str,
props: HashMap<String, OwnedValue>,
) {
let mut ifaces = HashMap::new();
ifaces.insert(DEVICE_IFACE.to_string(), props);
objects.insert(OwnedObjectPath::try_from(path).unwrap(), ifaces);
}
#[test]
fn no_adapter_normalizes_to_unavailable_regardless_of_powered() {
for powered in [None, Some(true), Some(false)] {
assert_eq!(
bluetooth_from_bluez(0, powered, 0).state(),
BluetoothState::Unavailable
);
}
}
#[test]
fn powered_adapter_normalizes_to_on_with_the_connected_count() {
assert_eq!(
bluetooth_from_bluez(1, Some(true), 0),
Bluetooth::new(BluetoothState::On, 0)
);
assert_eq!(
bluetooth_from_bluez(1, Some(true), 2),
Bluetooth::new(BluetoothState::On, 2)
);
}
#[test]
fn unpowered_adapter_normalizes_to_off_with_zero_connected() {
assert_eq!(
bluetooth_from_bluez(1, Some(false), 5),
Bluetooth::new(BluetoothState::Off, 0)
);
}
#[test]
fn missing_powered_with_an_adapter_normalizes_to_unavailable() {
assert_eq!(
bluetooth_from_bluez(1, None, 0).state(),
BluetoothState::Unavailable
);
}
#[test]
fn summarize_with_no_objects_reports_zero_adapters_and_unknown_power() {
let (count, powered, connected) = summarize(&empty_objects());
assert_eq!(count, 0);
assert_eq!(powered, None);
assert_eq!(connected, 0);
}
#[test]
fn summarize_aggregates_power_across_all_adapters() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(Some(true)));
insert_adapter(&mut objects, "/org/bluez/hci1", adapter_with(Some(false)));
let (count, powered, connected) = summarize(&objects);
assert_eq!(count, 2);
assert_eq!(powered, Some(true));
assert_eq!(connected, 0);
}
#[test]
fn summarize_power_is_off_when_no_adapter_is_on_and_at_least_one_is() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(None));
insert_adapter(&mut objects, "/org/bluez/hci1", adapter_with(Some(false)));
let (_, powered, _) = summarize(&objects);
assert_eq!(powered, Some(false));
}
#[test]
fn summarize_power_is_unknown_when_every_adapter_omits_powered() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(None));
insert_adapter(&mut objects, "/org/bluez/hci1", adapter_with(None));
let (_, powered, _) = summarize(&objects);
assert_eq!(powered, None);
}
#[test]
fn summarize_aggregation_is_independent_of_path_iteration_order() {
let mut objects_a = empty_objects();
insert_adapter(&mut objects_a, "/org/bluez/hci0", adapter_with(Some(true)));
insert_adapter(&mut objects_a, "/org/bluez/hci1", adapter_with(Some(false)));
let mut objects_b = empty_objects();
insert_adapter(&mut objects_b, "/org/bluez/hci1", adapter_with(Some(false)));
insert_adapter(&mut objects_b, "/org/bluez/hci0", adapter_with(Some(true)));
let (count_a, powered_a, connected_a) = summarize(&objects_a);
let (count_b, powered_b, connected_b) = summarize(&objects_b);
assert_eq!(
(count_a, powered_a, connected_a),
(count_b, powered_b, connected_b)
);
}
#[test]
fn summarize_power_state_and_connected_count_share_the_same_scope() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(Some(false)));
insert_adapter(&mut objects, "/org/bluez/hci1", adapter_with(Some(true)));
insert_device(
&mut objects,
"/org/bluez/hci1/dev_AA",
device_with(Some(true)),
);
insert_device(
&mut objects,
"/org/bluez/hci1/dev_BB",
device_with(Some(true)),
);
let (count, powered, connected) = summarize(&objects);
assert_eq!(count, 2);
assert_eq!(powered, Some(true));
assert_eq!(connected, 2);
assert_eq!(
bluetooth_from_bluez(count, powered, connected),
Bluetooth::new(BluetoothState::On, 2)
);
}
#[test]
fn summarize_counts_connected_devices_across_adapters() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(Some(true)));
insert_device(
&mut objects,
"/org/bluez/hci0/dev_AA_BB_CC_DD_EE_FF",
device_with(Some(true)),
);
insert_device(
&mut objects,
"/org/bluez/hci0/dev_11_22_33_44_55_66",
device_with(Some(false)),
);
insert_device(
&mut objects,
"/org/bluez/hci0/dev_DE_AD_BE_EF",
device_with(None),
);
insert_adapter(&mut objects, "/org/bluez/hci1", adapter_with(Some(true)));
insert_device(
&mut objects,
"/org/bluez/hci1/dev_AB_CD_EF_00_11_22",
device_with(Some(true)),
);
let (count, powered, connected) = summarize(&objects);
assert_eq!(count, 2);
assert_eq!(powered, Some(true));
assert_eq!(connected, 2);
}
#[test]
fn summarize_ignores_devices_with_unknown_connected() {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(Some(true)));
insert_device(&mut objects, "/org/bluez/hci0/dev_AA", device_with(None));
let (_, _, connected) = summarize(&objects);
assert_eq!(connected, 0);
}
#[test]
fn bluetooth_bluez_round_trip_drives_visible_labels() {
let powered_on = {
let mut objects = empty_objects();
insert_adapter(&mut objects, "/org/bluez/hci0", adapter_with(Some(true)));
insert_device(
&mut objects,
"/org/bluez/hci0/dev_AA",
device_with(Some(true)),
);
insert_device(
&mut objects,
"/org/bluez/hci0/dev_BB",
device_with(Some(true)),
);
let (count, powered, connected) = summarize(&objects);
bluetooth_from_bluez(count, powered, connected)
};
assert_eq!(powered_on.state(), BluetoothState::On);
assert_eq!(powered_on.connected(), 2);
assert_eq!(powered_on.label(), "2 connected");
let powered_off = bluetooth_from_bluez(1, Some(false), 0);
assert_eq!(powered_off.state(), BluetoothState::Off);
assert_eq!(powered_off.label(), "off");
let no_adapter = bluetooth_from_bluez(0, None, 0);
assert_eq!(no_adapter.state(), BluetoothState::Unavailable);
assert_eq!(no_adapter.label(), "unavailable");
}
fn properties_changed(interface: &str, changed: &[&str], invalidated: &[&str]) -> Message {
let changed: HashMap<&str, OwnedValue> = changed
.iter()
.map(|name| (*name, OwnedValue::from(true)))
.collect();
Message::signal(
"/org/bluez/hci0",
"org.freedesktop.DBus.Properties",
"PropertiesChanged",
)
.unwrap()
.build(&(interface, changed, invalidated))
.unwrap()
}
#[test]
fn displayed_property_changes_cause_a_re_read() {
assert!(signal_changes_state(&properties_changed(
ADAPTER_IFACE,
&["Powered"],
&[]
)));
assert!(signal_changes_state(&properties_changed(
DEVICE_IFACE,
&["ServicesResolved", "Connected"],
&[]
)));
assert!(signal_changes_state(&properties_changed(
DEVICE_IFACE,
&[],
&["Connected"]
)));
}
#[test]
fn scan_noise_and_foreign_interfaces_cause_no_re_read() {
for (interface, name) in [
(DEVICE_IFACE, "RSSI"),
(ADAPTER_IFACE, "Discovering"),
("org.bluez.MediaTransport1", "Connected"),
("org.bluez.Battery1", "Percentage"),
] {
assert!(
!signal_changes_state(&properties_changed(interface, &[name], &[])),
"{interface} {name}"
);
}
}
#[test]
fn objects_appearing_and_vanishing_cause_a_re_read() {
for member in ["InterfacesAdded", "InterfacesRemoved"] {
let message = Message::signal("/", "org.freedesktop.DBus.ObjectManager", member)
.unwrap()
.build(&())
.unwrap();
assert!(signal_changes_state(&message), "{member}");
}
}
#[test]
fn unrelated_and_malformed_signals_cause_no_re_read() {
let unrelated = Message::signal("/", "org.bluez.Custom", "Something")
.unwrap()
.build(&())
.unwrap();
assert!(!signal_changes_state(&unrelated));
let malformed =
Message::signal("/", "org.freedesktop.DBus.Properties", "PropertiesChanged")
.unwrap()
.build(&("not the PropertiesChanged body",))
.unwrap();
assert!(!signal_changes_state(&malformed));
}
#[test]
fn an_unchanged_snapshot_is_not_sent_again() {
let (bridge, channel) = crate::producer::ProducerBridge::new().unwrap();
let tx = bridge.sender();
let on = Bluetooth::new(BluetoothState::On, 0);
let connected = Bluetooth::new(BluetoothState::On, 1);
let mut sent = None;
for snapshot in [on, on, connected, connected] {
assert!(send_if_changed(&tx, &mut sent, snapshot));
}
let received: Vec<_> = std::iter::from_fn(|| channel.try_recv().ok())
.map(|msg| match msg {
Msg::Bluetooth(snapshot) => snapshot,
other => panic!("unexpected message: {other:?}"),
})
.collect();
assert_eq!(received, [on, connected]);
}
}