use std::{
collections::HashMap,
process::Stdio,
str::FromStr,
time::{Duration, Instant},
};
use anyhow::Result;
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
process::Command,
sync::{mpsc, oneshot},
task::JoinHandle,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct Address(pub [u8; 6]);
impl std::fmt::Display for Address {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let [a, b, c, d, e, g] = self.0;
write!(f, "{a:02X}:{b:02X}:{c:02X}:{d:02X}:{e:02X}:{g:02X}")
}
}
impl FromStr for Address {
type Err = anyhow::Error;
fn from_str(s: &str) -> Result<Self> {
let bytes = s
.split(':')
.map(|byte| u8::from_str_radix(byte, 16))
.collect::<Result<Vec<_>, _>>()?;
Ok(Self(
bytes
.try_into()
.map_err(|_| anyhow::anyhow!("Invalid address"))?,
))
}
}
#[derive(Debug, Clone)]
pub struct DeviceInfo {
pub address: Address,
pub name: Option<String>,
pub paired: Option<bool>,
pub connected: Option<bool>,
}
#[derive(Debug)]
pub struct PairingRequest {
pub address: Address,
pub respond: oneshot::Sender<bool>,
}
async fn bluetoothctl(args: &[&str]) -> Result<String> {
let output = Command::new("bluetoothctl").args(args).output().await?;
let stdout = String::from_utf8_lossy(&output.stdout).into_owned();
if let Some(line) = stdout
.lines()
.find(|line| line.starts_with("Failed") || line.contains("not available"))
{
anyhow::bail!("{line}");
}
Ok(stdout)
}
fn list(output: &str) -> HashMap<Address, String> {
output
.lines()
.filter_map(|line| line.strip_prefix("Device ")?.split_once(' '))
.filter_map(|(address, name)| Some((address.parse().ok()?, name.to_owned())))
.collect()
}
#[derive(Debug)]
pub struct BluetoothManager {
pub devices: HashMap<Address, DeviceInfo>,
pub init_task: Option<JoinHandle<Result<()>>>,
pub scan_task: Option<JoinHandle<Result<()>>>,
pub action_task: Option<JoinHandle<Result<DeviceInfo>>>,
pub action_address: Option<Address>,
pub action_label: Option<&'static str>,
refresh_task: Option<JoinHandle<Result<Vec<DeviceInfo>>>>,
last_refresh: Option<Instant>,
pairing_tx: mpsc::UnboundedSender<PairingRequest>,
pairing_rx: mpsc::UnboundedReceiver<PairingRequest>,
pub pairing_request: Option<PairingRequest>,
}
impl BluetoothManager {
pub fn new() -> Self {
let (pairing_tx, pairing_rx) = mpsc::unbounded_channel();
Self {
devices: HashMap::new(),
init_task: None,
scan_task: None,
action_task: None,
action_address: None,
action_label: None,
refresh_task: None,
last_refresh: None,
pairing_tx,
pairing_rx,
pairing_request: None,
}
}
pub fn sorted_devices(&self) -> Vec<&DeviceInfo> {
let mut devices = self.devices.values().collect::<Vec<_>>();
devices.sort_by_key(|d| (d.name.is_none(), d.name.clone(), d.address));
devices
}
pub fn init(&mut self) {
if self.init_task.is_some() {
return;
}
self.init_task = Some(tokio::spawn(async {
bluetoothctl(&["power", "on"]).await?;
Ok(())
}));
}
pub fn scan(&mut self) {
if self.scan_task.is_some() {
return;
}
self.scan_task = Some(tokio::spawn(async {
bluetoothctl(&["--timeout", "10", "scan", "on"]).await?;
Ok(())
}));
}
pub fn pair(&mut self, address: Address) {
if self.action_task.is_some() {
return;
}
let Some(device) = self.devices.get(&address).cloned() else {
return;
};
let pairing_tx = self.pairing_tx.clone();
self.action_address = Some(address);
self.action_label = Some("Pairing");
self.action_task = Some(tokio::spawn(async move {
let target = address.to_string();
if device.paired != Some(true) {
let mut child = Command::new("bluetoothctl")
.args(["--agent", "DisplayYesNo"])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.kill_on_drop(true)
.spawn()?;
let mut stdin = child.stdin.take().unwrap();
let mut stdout = child.stdout.take().unwrap();
stdin
.write_all(format!("default-agent\npair {target}\n").as_bytes())
.await?;
let mut output = String::new();
let mut chunk = [0u8; 1024];
while !output.contains("Pairing successful") {
let read = stdout.read(&mut chunk).await?;
if read == 0 {
anyhow::bail!("bluetoothctl exited");
}
output.push_str(&String::from_utf8_lossy(&chunk[..read]));
if let Some(line) = output.lines().find(|l| l.contains("Failed to pair")) {
anyhow::bail!("{}", line.trim());
}
if output.contains("(yes/no)") {
let (respond, response) = oneshot::channel();
pairing_tx.send(PairingRequest { address, respond })?;
let answer = if response.await.unwrap_or(false) {
"yes"
} else {
"no"
};
stdin.write_all(format!("{answer}\n").as_bytes()).await?;
output.clear();
}
}
}
bluetoothctl(&["trust", &target]).await?;
bluetoothctl(&["connect", &target]).await?;
Ok(DeviceInfo {
paired: Some(true),
connected: Some(true),
..device
})
}));
}
pub fn unpair(&mut self, address: Address) {
if self.action_task.is_some() {
return;
}
let Some(device) = self.devices.get(&address).cloned() else {
return;
};
self.action_address = Some(address);
self.action_label = Some("Unpairing");
self.action_task = Some(tokio::spawn(async move {
bluetoothctl(&["remove", &address.to_string()]).await?;
Ok(DeviceInfo {
paired: Some(false),
connected: Some(false),
..device
})
}));
}
pub fn connect(&mut self, address: Address) {
if self.action_task.is_some() {
return;
}
let Some(device) = self.devices.get(&address).cloned() else {
return;
};
self.action_address = Some(address);
self.action_label = Some("Connecting");
self.action_task = Some(tokio::spawn(async move {
bluetoothctl(&["connect", &address.to_string()]).await?;
Ok(DeviceInfo {
connected: Some(true),
..device
})
}));
}
pub fn disconnect(&mut self, address: Address) {
if self.action_task.is_some() {
return;
}
let Some(device) = self.devices.get(&address).cloned() else {
return;
};
self.action_address = Some(address);
self.action_label = Some("Disconnecting");
self.action_task = Some(tokio::spawn(async move {
bluetoothctl(&["disconnect", &address.to_string()]).await?;
Ok(DeviceInfo {
connected: Some(false),
..device
})
}));
}
pub fn accept_pairing(&mut self) {
if let Some(request) = self.pairing_request.take() {
let _ = request.respond.send(true);
}
}
pub fn reject_pairing(&mut self) {
if let Some(request) = self.pairing_request.take() {
let _ = request.respond.send(false);
}
}
pub async fn run(&mut self) -> Result<()> {
while let Ok(request) = self.pairing_rx.try_recv() {
if self.pairing_request.is_none() {
self.pairing_request = Some(request);
} else {
let _ = request.respond.send(false);
}
}
if self
.pairing_request
.as_ref()
.is_some_and(|request| request.respond.is_closed())
{
self.pairing_request = None;
}
if self.init_task.as_ref().is_some_and(|t| t.is_finished()) {
self.init_task.take().unwrap().await??;
self.scan();
}
if self.scan_task.as_ref().is_some_and(|t| t.is_finished()) {
let _ = self.scan_task.take().unwrap().await;
}
if self.action_task.as_ref().is_some_and(|t| t.is_finished()) {
self.action_address = None;
self.action_label = None;
if let Ok(Ok(info)) = self.action_task.take().unwrap().await {
self.devices.insert(info.address, info);
}
}
if self.refresh_task.as_ref().is_some_and(|t| t.is_finished()) {
if let Ok(Ok(devices)) = self.refresh_task.take().unwrap().await {
if self.action_task.is_none() {
self.devices = devices
.into_iter()
.map(|device| (device.address, device))
.collect();
}
}
}
if self.refresh_task.is_none()
&& self.action_task.is_none()
&& self
.last_refresh
.is_none_or(|at| at.elapsed() > Duration::from_secs(2))
{
self.last_refresh = Some(Instant::now());
self.refresh_task = Some(tokio::spawn(async {
let all = list(&bluetoothctl(&["devices"]).await?);
let paired = list(&bluetoothctl(&["devices", "Paired"]).await?);
let connected = list(&bluetoothctl(&["devices", "Connected"]).await?);
Ok(all
.into_iter()
.map(|(address, name)| DeviceInfo {
address,
name: (name != address.to_string().replace(':', "-")).then_some(name),
paired: Some(paired.contains_key(&address)),
connected: Some(connected.contains_key(&address)),
})
.collect())
}));
}
Ok(())
}
}