Files
lightbar/src/modules/audio.rs
2026-09-29 13:01:10 +02:00

813 lines
30 KiB
Rust

use std::{
collections::BTreeMap,
process::Stdio,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
thread::{self, JoinHandle},
time::Duration,
};
use anyhow::{Context, Result, bail};
use calloop::channel::Sender;
use serde_json::Value as Json;
use tokio::{io::AsyncReadExt, process::Command, sync::mpsc};
use crate::{
config::{AudioModule, UnavailableBehavior},
format::{Value, expand},
model::{
AudioDevice, AudioModel, ModuleEvent, ModuleSnapshot, PopupContent, PopupModel, Segment,
},
};
const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
const MAX_OBJECTS: usize = 4096;
const COMMAND_TIMEOUT: Duration = Duration::from_secs(2);
#[derive(Debug, Clone, Copy)]
pub enum AudioTarget {
DefaultOutput,
DefaultInput,
Device(u32),
}
impl AudioTarget {
fn argument(self) -> String {
match self {
Self::DefaultOutput => "@DEFAULT_AUDIO_SINK@".to_owned(),
Self::DefaultInput => "@DEFAULT_AUDIO_SOURCE@".to_owned(),
Self::Device(id) => id.to_string(),
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum AudioAction {
SetVolume(u32, f64),
ChangeVolume(AudioTarget, i16),
ToggleMute(AudioTarget),
SetDefault(u32),
}
impl AudioAction {
fn arguments(self) -> Vec<String> {
match self {
Self::SetVolume(id, percent) => vec![
"set-volume".into(),
id.to_string(),
format!("{:.3}", percent.clamp(0.0, 150.0) / 100.0),
],
Self::ChangeVolume(target, step) => vec![
"set-volume".into(),
"-l".into(),
"1.5".into(),
target.argument(),
format!(
"{}%{}",
step.unsigned_abs(),
if step < 0 { "-" } else { "+" }
),
],
Self::ToggleMute(target) => vec!["set-mute".into(), target.argument(), "toggle".into()],
Self::SetDefault(id) => vec!["set-default".into(), id.to_string()],
}
}
}
pub fn spawn(
name: String,
settings: AudioModule,
sender: Sender<ModuleEvent>,
stop: Arc<AtomicBool>,
) -> (JoinHandle<()>, mpsc::Sender<AudioAction>) {
let (actions, mut requests) = mpsc::channel(32);
let worker = thread::Builder::new().name(format!("lightbar-audio-{name}")).spawn(move || {
let runtime = match tokio::runtime::Builder::new_current_thread().enable_all().build() {
Ok(runtime) => runtime,
Err(error) => {
tracing::error!(module = %name, %error, "could not start audio runtime");
return;
}
};
runtime.block_on(async {
let mut failures = 0_u32;
while !stop.load(Ordering::Acquire) {
let result = monitor(&name, &settings, &sender, &stop, &mut requests, &mut failures).await;
if stop.load(Ordering::Acquire) { break; }
if let Err(error) = result {
tracing::warn!(module = %name, %error, "audio monitor stopped; reconnecting");
}
let placeholder = match settings.common.unavailable {
UnavailableBehavior::Hide => None,
UnavailableBehavior::Placeholder => Some("Audio unavailable"),
};
if sender.send(ModuleEvent { module: name.clone(), snapshot: ModuleSnapshot::unavailable(placeholder) }).is_err() { break; }
// Commands queued for a disconnected generation must not affect a new device.
while requests.try_recv().is_ok() {
tracing::warn!(module = %name, "discarded audio action after disconnect");
}
failures = failures.saturating_add(1);
let deadline = tokio::time::Instant::now() + super::process::backoff(failures);
while !stop.load(Ordering::Acquire) && tokio::time::Instant::now() < deadline {
tokio::select! {
() = tokio::time::sleep_until(deadline) => break,
action = requests.recv() => {
if action.is_none() { return; }
tracing::warn!(module = %name, "discarded audio action while disconnected");
}
}
}
}
});
}).expect("audio worker thread");
(worker, actions)
}
async fn monitor(
name: &str,
settings: &AudioModule,
sender: &Sender<ModuleEvent>,
stop: &AtomicBool,
requests: &mut mpsc::Receiver<AudioAction>,
failures: &mut u32,
) -> Result<()> {
let mut child = Command::new("pw-dump")
.args(["--monitor", "--no-colors"])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.kill_on_drop(true)
.spawn()
.context("could not start pw-dump")?;
let mut stdout = child
.stdout
.take()
.context("pw-dump stdout was not piped")?;
let result = async {
let mut frames = JsonFrames::default();
let mut state = AudioState::default();
let mut last = None;
let mut buffer = [0_u8; 8192];
while !stop.load(Ordering::Acquire) {
tokio::select! {
read = stdout.read(&mut buffer) => {
let read = read.context("could not read pw-dump events")?;
if read == 0 { bail!("pw-dump event stream closed"); }
for batch in frames.push(&buffer[..read])? {
if state.apply(batch)? {
// pw-dump's metadata deltas omit deleted keys. Re-read the
// default metadata to distinguish a removal from no change.
let defaults = read_defaults().await?;
state.replace_defaults(defaults)?;
}
let snapshot = state.snapshot(settings);
if last.as_ref() != Some(&snapshot) {
sender.send(ModuleEvent { module: name.to_owned(), snapshot: snapshot.clone() })
.map_err(|_| anyhow::anyhow!("UI event channel closed"))?;
last = Some(snapshot);
}
*failures = 0;
}
}
action = requests.recv() => {
let Some(action) = action else { break; };
if let Err(error) = run_action(action).await { tracing::warn!(module = name, %error, "audio control failed"); }
}
}
}
Ok(())
}.await;
// Always reap the monitor, including malformed input and closed UI channels.
if child.try_wait()?.is_none() {
child.kill().await.context("could not stop pw-dump")?;
}
child.wait().await.context("could not reap pw-dump")?;
result
}
async fn run_action(action: AudioAction) -> Result<()> {
let args = action.arguments();
let mut child = Command::new("wpctl")
.args(&args)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::inherit())
.kill_on_drop(true)
.spawn()
.with_context(|| format!("could not run wpctl {}", args[0]))?;
if let Ok(status) = tokio::time::timeout(COMMAND_TIMEOUT, child.wait()).await {
let status = status.context("could not wait for wpctl")?;
if !status.success() {
bail!("wpctl {} exited with {status}", args[0]);
}
Ok(())
} else {
child
.kill()
.await
.context("could not stop timed-out wpctl")?;
child.wait().await.context("could not reap wpctl")?;
bail!("wpctl {} timed out", args[0]);
}
}
async fn read_defaults() -> Result<Vec<Json>> {
let mut child = Command::new("pw-dump")
.args(["--no-colors", "default"])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.kill_on_drop(true)
.spawn()
.context("could not read default audio metadata")?;
let stdout = child
.stdout
.take()
.context("metadata stdout was not piped")?;
let operation = async {
let mut bytes = Vec::new();
stdout
.take(MAX_FRAME_BYTES as u64 + 1)
.read_to_end(&mut bytes)
.await?;
if bytes.len() > MAX_FRAME_BYTES {
bail!("default audio metadata exceeded the size limit");
}
let status = child.wait().await?;
if !status.success() {
bail!("default audio metadata command exited with {status}");
}
serde_json::from_slice(&bytes).context("invalid default audio metadata")
};
let result = tokio::time::timeout(COMMAND_TIMEOUT, operation).await;
if child.try_wait()?.is_none() {
child
.kill()
.await
.context("could not stop audio metadata command")?;
}
child
.wait()
.await
.context("could not reap audio metadata command")?;
result.context("default audio metadata query timed out")?
}
#[derive(Default)]
struct JsonFrames {
bytes: Vec<u8>,
depth: usize,
quoted: bool,
escaped: bool,
}
impl JsonFrames {
fn push(&mut self, bytes: &[u8]) -> Result<Vec<Vec<Json>>> {
let mut batches = Vec::new();
for &byte in bytes {
if self.bytes.is_empty() {
if byte.is_ascii_whitespace() {
continue;
}
if byte != b'[' {
bail!("pw-dump event must be a JSON array");
}
}
if self.bytes.len() >= MAX_FRAME_BYTES {
bail!("pw-dump event exceeded {MAX_FRAME_BYTES} bytes");
}
self.bytes.push(byte);
if self.quoted {
if self.escaped {
self.escaped = false;
} else if byte == b'\\' {
self.escaped = true;
} else if byte == b'"' {
self.quoted = false;
}
} else {
match byte {
b'"' => self.quoted = true,
b'[' | b'{' => self.depth += 1,
b']' | b'}' => {
self.depth = self
.depth
.checked_sub(1)
.context("unbalanced pw-dump JSON")?;
if self.depth == 0 {
batches.push(
serde_json::from_slice(&self.bytes)
.context("invalid pw-dump JSON")?,
);
self.bytes.clear();
}
}
_ => {}
}
}
}
Ok(batches)
}
}
#[derive(Default)]
struct AudioState {
objects: BTreeMap<u32, Json>,
metadata_id: Option<u32>,
default_output: Option<String>,
default_input: Option<String>,
}
impl AudioState {
fn apply(&mut self, batch: Vec<Json>) -> Result<bool> {
let mut metadata_changed = false;
for object in batch {
let id = object["id"]
.as_u64()
.and_then(|id| u32::try_from(id).ok())
.context("pw-dump object has no valid id")?;
if object.get("info").is_some_and(Json::is_null) {
self.objects.remove(&id);
if self.metadata_id == Some(id) {
self.metadata_id = None;
self.default_output = None;
self.default_input = None;
}
continue;
}
if object["props"]["metadata.name"] == "default" || self.metadata_id == Some(id) {
metadata_changed = true;
self.metadata_id = Some(id);
if let Some(entries) = object["metadata"].as_array() {
for entry in entries {
if entry["subject"].as_u64() != Some(0) {
continue;
}
let target = match entry["key"].as_str() {
Some("default.audio.sink") => &mut self.default_output,
Some("default.audio.source") => &mut self.default_input,
_ => continue,
};
*target = entry["value"]["name"].as_str().map(str::to_owned);
}
}
continue;
}
let class = object["info"]["props"]["media.class"].as_str();
if matches!(
class,
Some("Audio/Sink" | "Audio/Source" | "Audio/Source/Virtual" | "Audio/Device")
) || self.objects.contains_key(&id)
{
let stored = self
.objects
.entry(id)
.or_insert_with(|| serde_json::json!({}));
// Node properties are complete; parameter keys are incremental.
if let Some(info) = object["info"].as_object() {
if stored.get("info").is_none() {
stored["info"] = serde_json::json!({});
}
if let Some(props) = info.get("props") {
stored["info"]["props"] = props.clone();
}
if let Some(params) = info.get("params").and_then(Json::as_object) {
if stored["info"].get("params").is_none() {
stored["info"]["params"] = serde_json::json!({});
}
for (key, value) in params {
stored["info"]["params"][key] = value.clone();
}
}
}
}
if self.objects.len() > MAX_OBJECTS {
bail!("too many audio objects");
}
}
Ok(metadata_changed)
}
fn replace_defaults(&mut self, snapshot: Vec<Json>) -> Result<()> {
self.default_output = None;
self.default_input = None;
self.metadata_id = None;
self.apply(snapshot)?;
Ok(())
}
fn devices(&self) -> Vec<AudioDevice> {
let mut devices = Vec::new();
for (&id, object) in &self.objects {
let props = &object["info"]["props"];
let input = match props["media.class"].as_str() {
Some("Audio/Sink") => false,
Some("Audio/Source" | "Audio/Source/Virtual") => true,
_ => continue,
};
let Some(name) = props["node.name"].as_str() else {
continue;
};
let description = props["node.description"]
.as_str()
.or_else(|| props["node.nick"].as_str())
.unwrap_or(name);
let route_volume = number(&props["device.id"])
.and_then(|id| self.objects.get(&id))
.and_then(|device| device["info"]["params"]["Route"].as_array())
.and_then(|routes| {
routes.iter().find(|route| {
number(&route["device"]) == number(&props["card.profile.device"])
&& number(&route["device"]).is_some()
})
})
.and_then(|route| volume(&route["props"]));
let node_volume = || {
object["info"]["params"]["Props"]
.as_array()
.and_then(|params| params.iter().find_map(volume))
};
let level = route_volume.or_else(node_volume);
let default = if input {
&self.default_input
} else {
&self.default_output
};
devices.push(AudioDevice {
id,
name: name.to_owned(),
description: description.to_owned(),
input,
is_default: default.as_deref() == Some(name),
volume: level.map(|(volume, _)| volume),
muted: level.map(|(_, muted)| muted),
});
}
devices.sort_by(|a, b| {
a.input
.cmp(&b.input)
.then_with(|| a.description.cmp(&b.description))
.then_with(|| a.name.cmp(&b.name))
});
devices
}
fn snapshot(&self, settings: &AudioModule) -> ModuleSnapshot {
let model = AudioModel {
devices: self.devices(),
..AudioModel::default()
};
let mut snapshot = ModuleSnapshot {
visible: true,
..ModuleSnapshot::default()
};
let mut descriptions = Vec::new();
for (input, id, label) in [(false, "output", "Output"), (true, "input", "Microphone")] {
let device = model.default_device(input);
let mut segment = Segment::new(format!("{label} N/A"));
segment.id = Some(id.to_owned());
if let Some(device) = device {
let format = match (input, device.muted == Some(true)) {
(false, false) => &settings.format,
(false, true) => &settings.muted_format,
(true, false) => &settings.microphone_format,
(true, true) => &settings.microphone_muted_format,
};
segment.text = expand(
format,
&BTreeMap::from([
(
"volume",
device
.volume
.map_or_else(|| Value::from("N/A"), Value::from),
),
("device", Value::from(device.description.clone())),
("name", Value::from(device.name.clone())),
]),
);
descriptions.push(format!(
"{label}: {} · {}{}",
device.description,
device.volume.map_or_else(
|| "Volume unavailable".to_owned(),
|volume| format!("{volume:.0}%")
),
if device.muted == Some(true) {
" · Muted"
} else {
""
}
));
} else {
descriptions.push(format!("{label}: no default device"));
}
if device.is_none_or(|device| device.volume.is_none() || device.muted == Some(true)) {
"disabled".clone_into(&mut segment.state);
}
snapshot.segments.push(segment);
}
snapshot.tooltip = Some(descriptions.join("\n"));
snapshot.popup = Some(PopupModel {
title: "Audio".to_owned(),
content: PopupContent::Audio(model),
});
snapshot
}
}
fn number(value: &Json) -> Option<u32> {
value
.as_u64()
.and_then(|value| u32::try_from(value).ok())
.or_else(|| value.as_str()?.parse().ok())
}
fn volume(props: &Json) -> Option<(f64, bool)> {
let muted = props["mute"].as_bool()?;
let linear = props["channelVolumes"]
.as_array()
.and_then(|channels| channels.first())
.and_then(Json::as_f64)
.or_else(|| props["volume"].as_f64())?;
// WirePlumber exposes the first channel on a cubic scale, including route volume.
(linear.is_finite() && linear >= 0.0).then(|| ((linear.cbrt() * 100.0).round(), muted))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn fixture() -> AudioState {
let mut state = AudioState::default();
state
.apply(serde_json::from_str(include_str!("../../tests/fixtures/audio.json")).unwrap())
.unwrap();
state
}
#[test]
fn reports_actual_defaults_route_volume_and_independent_mute() {
let state = fixture();
let snapshot = state.snapshot(&AudioModule::default());
let PopupContent::Audio(audio) = snapshot.popup.unwrap().content else {
panic!("audio popup missing");
};
assert_eq!(
audio.devices.len(),
3,
"application streams are not devices"
);
let output = audio.default_device(false).unwrap();
assert_eq!(
(output.id, output.volume, output.muted),
(10, Some(50.0), Some(false))
);
let input = audio.default_device(true).unwrap();
assert_eq!(
(input.id, input.volume, input.muted),
(11, Some(80.0), Some(true))
);
assert_eq!(snapshot.segments[0].state, "normal");
assert_eq!(snapshot.segments[1].state, "disabled");
assert!(
snapshot
.tooltip
.unwrap()
.contains("USB Microphone · 80% · Muted")
);
}
#[test]
fn follows_external_defaults_and_keeps_other_endpoint() {
let mut state = fixture();
state.apply(vec![json!({"id":1,"metadata":[{"subject":0,"key":"default.audio.sink","value":{"name":"headphones"}}]})]).unwrap();
let devices = state.devices();
assert!(
devices
.iter()
.any(|device| device.id == 12 && device.is_default)
);
assert!(
devices
.iter()
.any(|device| device.id == 11 && device.is_default)
);
state.apply(vec![json!({"id":11,"info":null})]).unwrap();
let snapshot = state.snapshot(&AudioModule::default());
assert_eq!(snapshot.segments[1].text, "Microphone N/A");
let PopupContent::Audio(audio) = snapshot.popup.unwrap().content else {
panic!();
};
assert!(audio.default_device(false).is_some());
assert!(audio.default_device(true).is_none());
assert!(
!audio
.controls()
.iter()
.any(|control| control.id.as_deref() == Some("audio-mute:11"))
);
state.apply(vec![json!({"id":1,"info":null})]).unwrap();
assert!(!state.devices().iter().any(|device| device.is_default));
}
#[test]
fn removing_a_default_clears_only_that_endpoint() {
let mut state = fixture();
assert!(state.apply(vec![json!({"id":1,"metadata":[]})]).unwrap());
state
.replace_defaults(vec![
json!({"id":1,"props":{"metadata.name":"default"},"metadata":[
{"subject":0,"key":"default.audio.sink","value":{"name":"speakers"}}
]}),
])
.unwrap();
let devices = state.devices();
assert!(
devices
.iter()
.any(|device| device.id == 10 && device.is_default)
);
assert!(
!devices
.iter()
.any(|device| device.input && device.is_default)
);
}
#[test]
fn preserves_final_event_in_a_burst_and_handles_partial_parameters() {
let mut frames = JsonFrames::default();
let mut state = fixture();
let stream = concat!(
"[{\"id\":11,\"info\":{\"params\":{\"Props\":[{\"mute\":false,\"channelVolumes\":[0.001]}]}}}]\n",
"[{\"id\":11,\"info\":{\"params\":{\"Props\":[{\"mute\":false,\"channelVolumes\":[0.729]}]}}}]\n"
);
for chunk in stream.as_bytes().chunks(7) {
for batch in frames.push(chunk).unwrap() {
state.apply(batch).unwrap();
}
}
let microphone = state
.devices()
.into_iter()
.find(|device| device.id == 11)
.unwrap();
assert_eq!(
(microphone.volume, microphone.muted),
(Some(90.0), Some(false))
);
assert_eq!(microphone.description, "USB Microphone");
}
#[test]
fn stream_parser_handles_escaped_names_and_limits() {
let data = json!([{"id":1,"text":"Device [USB] \"name\" \\ test"}]).to_string();
let mut frames = JsonFrames::default();
let mut batches = Vec::new();
for byte in data.as_bytes() {
batches.extend(frames.push(&[*byte]).unwrap());
}
assert_eq!(batches.len(), 1);
assert_eq!(batches[0][0]["text"], "Device [USB] \"name\" \\ test");
assert!(JsonFrames::default().push(b"not json").is_err());
assert!(JsonFrames::default().push(b"[{broken}]").is_err());
let mut oversized = JsonFrames::default();
oversized.push(b"[\"").unwrap();
assert!(oversized.push(&vec![b'x'; MAX_FRAME_BYTES]).is_err());
}
#[test]
fn volume_keeps_overamplification_and_distinguishes_silence_from_mute() {
assert_eq!(
volume(&json!({"mute":false,"channelVolumes":[0.0]})),
Some((0.0, false))
);
assert_eq!(
volume(&json!({"mute":true,"channelVolumes":[8.0]})),
Some((200.0, true))
);
assert_eq!(
volume(&json!({"mute":false,"volume":0.125})),
Some((50.0, false))
);
assert_eq!(volume(&json!({"mute":false,"channelVolumes":[-1.0]})), None);
assert_eq!(volume(&json!({"channelVolumes":[1.0]})), None);
}
#[test]
fn audio_commands_target_the_selected_endpoint_and_limit_volume() {
assert_eq!(
AudioAction::ToggleMute(AudioTarget::DefaultInput).arguments(),
["set-mute", "@DEFAULT_AUDIO_SOURCE@", "toggle"]
);
assert_eq!(
AudioAction::ChangeVolume(AudioTarget::DefaultOutput, 7).arguments(),
["set-volume", "-l", "1.5", "@DEFAULT_AUDIO_SINK@", "7%+"]
);
assert_eq!(
AudioAction::ChangeVolume(AudioTarget::Device(11), -7).arguments(),
["set-volume", "-l", "1.5", "11", "7%-"]
);
assert_eq!(
AudioAction::SetDefault(12).arguments(),
["set-default", "12"]
);
assert_eq!(
AudioAction::SetVolume(11, 180.0).arguments(),
["set-volume", "11", "1.500"]
);
assert_eq!(
AudioAction::SetVolume(11, -10.0).arguments(),
["set-volume", "11", "0.000"]
);
}
#[test]
fn custom_formats_include_names_and_numeric_volume() {
let settings = AudioModule {
format: "{device}: {volume:.1}%".to_owned(),
microphone_muted_format: "{name}: muted at {volume}%".to_owned(),
..AudioModule::default()
};
let snapshot = fixture().snapshot(&settings);
assert_eq!(snapshot.segments[0].text, "Desktop Speakers: 50.0%");
assert_eq!(snapshot.segments[1].text, "microphone: muted at 80%");
}
#[test]
#[ignore = "requires a running PipeWire/WirePlumber session; read-only"]
fn live_pipewire_matches_wpctl_and_worker_stops_when_idle() {
let output = std::process::Command::new("pw-dump")
.arg("--no-colors")
.output()
.unwrap();
assert!(output.status.success());
let mut state = AudioState::default();
state
.apply(serde_json::from_slice(&output.stdout).unwrap())
.unwrap();
let devices = state.devices();
for target in [AudioTarget::DefaultOutput, AudioTarget::DefaultInput] {
let output = std::process::Command::new("wpctl")
.args(["inspect", &target.argument()])
.output()
.unwrap();
assert!(output.status.success());
let text = String::from_utf8(output.stdout).unwrap();
let id: u32 = text
.split_whitespace()
.nth(1)
.unwrap()
.trim_end_matches(',')
.parse()
.unwrap();
let device = devices.iter().find(|device| device.id == id).unwrap();
assert!(device.is_default);
let output = std::process::Command::new("wpctl")
.args(["get-volume", &target.argument()])
.output()
.unwrap();
assert!(output.status.success());
let text = String::from_utf8(output.stdout).unwrap();
let percent = (text
.split_whitespace()
.nth(1)
.unwrap()
.parse::<f64>()
.unwrap()
* 100.0)
.round();
assert_eq!(device.volume, Some(percent));
assert_eq!(device.muted, Some(text.contains("[MUTED]")));
}
let (sender, receiver) = calloop::channel::channel();
let stop = Arc::new(AtomicBool::new(false));
let (worker, actions) = spawn(
"audio-test".to_owned(),
AudioModule::default(),
sender,
Arc::clone(&stop),
);
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let first_event = loop {
if let Ok(event) = receiver.try_recv() {
break Some(event);
}
if std::time::Instant::now() >= deadline {
break None;
}
thread::sleep(Duration::from_millis(10));
};
stop.store(true, Ordering::Release);
drop(actions);
let start = std::time::Instant::now();
worker.join().unwrap();
assert!(start.elapsed() < Duration::from_secs(1));
assert!(first_event.unwrap().snapshot.popup.is_some());
}
}