use std::cell::RefCell;
use std::collections::HashMap;
use std::rc::Rc;
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;
use log::debug;
use pipewire::context::ContextRc;
use pipewire::main_loop::MainLoopRc;
use pipewire::node::{Node, NodeInfoRef, NodeState};
use pipewire::proxy::Listener;
use pipewire::registry::{GlobalObject, RegistryRc};
use pipewire::spa::pod::Pod;
use pipewire::spa::utils::dict::DictRef;
use pipewire::types::ObjectType;
use tokio::sync::oneshot;
use crate::widget::{DeviceKind, Msg, Volume};
use crate::producer::{MsgSender, Producer, ProducerFuture, ProducerResult};
const DEFAULT_INTERVAL: Duration = Duration::from_secs(2);
const DEFAULT_AUDIO_SINK_KEY: &str = "default.audio.sink";
pub fn device_kind_from_icon_name(icon: Option<&str>) -> DeviceKind {
let Some(icon) = icon else {
return DeviceKind::Other;
};
if icon.starts_with("audio-headset") {
DeviceKind::Headset
} else if icon.starts_with("audio-headphone") {
DeviceKind::Headphones
} else if icon.starts_with("audio-speaker") {
DeviceKind::Speakers
} else if icon.starts_with("audio-card-hdmi")
|| icon.starts_with("audio-card-dp")
|| icon.starts_with("audio-card-displayport")
{
DeviceKind::Monitor
} else if icon.starts_with("audio-card") {
DeviceKind::Speakers
} else if icon.starts_with("video-display") || icon.starts_with("video-monitor") {
DeviceKind::Monitor
} else if icon.starts_with("audio-handsfree")
|| icon.starts_with("audio-cellphone")
|| icon == "phone"
{
DeviceKind::Phone
} else if icon.starts_with("tv") || icon.starts_with("television") {
DeviceKind::Tv
} else {
DeviceKind::Other
}
}
pub fn device_kind_from_form_factor(form: Option<&str>) -> DeviceKind {
let Some(form) = form else {
return DeviceKind::Other;
};
match form {
"headphone" | "earphone" => DeviceKind::Headphones,
"headset" => DeviceKind::Headset,
"speaker" | "desk" | "hifi" | "computer" | "portable" => DeviceKind::Speakers,
"monitor" => DeviceKind::Monitor,
"tv" => DeviceKind::Tv,
"handset" | "phone" | "car" => DeviceKind::Phone,
_ => DeviceKind::Other,
}
}
pub fn device_kind_from_props(props: &DictRef) -> DeviceKind {
let icon = props.get("device.icon-name");
let kind = device_kind_from_icon_name(icon);
if kind != DeviceKind::Other {
return kind;
}
device_kind_from_form_factor(props.get("device.form-factor"))
}
pub fn is_audio_sink(props: &DictRef) -> bool {
props.get("media.class") == Some("Audio/Sink")
}
pub fn pick_active_sink_id(candidates: &[(u32, &str, bool)]) -> Option<u32> {
candidates
.iter()
.filter(|(_, _, running)| *running)
.min_by(|a, b| a.1.cmp(b.1))
.map(|(id, _, _)| *id)
}
fn is_node_running(state: &NodeState<'_>) -> bool {
matches!(state, NodeState::Running)
}
pub fn parse_props_pod(pod_bytes: &[u8]) -> Option<(f32, bool)> {
use pipewire::spa::pod::Value;
use pipewire::spa::sys::{SPA_PROP_channelVolumes, SPA_PROP_mute};
let pod = Pod::from_bytes(pod_bytes)?;
let obj = pod.as_object().ok()?;
let channel_volumes_pod = obj.find_prop(pipewire::spa::utils::Id(SPA_PROP_channelVolumes))?;
let mute_pod = obj.find_prop(pipewire::spa::utils::Id(SPA_PROP_mute))?;
let value: Value = pipewire::spa::pod::deserialize::PodDeserializer::deserialize_any_from(
channel_volumes_pod.value().as_bytes(),
)
.ok()
.map(|(_, v)| v)?;
let user_volume = match &value {
Value::ValueArray(pipewire::spa::pod::ValueArray::Float(v)) => average_to_user_volume(v)?,
Value::Float(v) => average_to_user_volume(&[*v])?,
_ => return None,
};
let mute = mute_pod.value().get_bool().ok()?;
Some((user_volume, mute))
}
pub fn average_to_user_volume(channels: &[f32]) -> Option<f32> {
if channels.is_empty() {
return None;
}
let sum: f32 = channels.iter().filter(|v| v.is_finite() && **v > 0.0).sum();
let n = channels
.iter()
.filter(|v| v.is_finite() && **v > 0.0)
.count() as f32;
if n == 0.0 {
return Some(0.0);
}
let avg = sum / n;
let cbrt = avg.max(0.0).cbrt();
Some(cbrt.clamp(0.0, 1.0))
}
#[derive(Clone)]
struct SinkEntry {
name: String,
device: DeviceKind,
running: bool,
volume: Option<(f32, bool)>,
}
pub struct VolumeProducer {
interval: Duration,
}
impl VolumeProducer {
pub fn new() -> Self {
Self {
interval: DEFAULT_INTERVAL,
}
}
pub fn with_interval(interval: Duration) -> Self {
Self { interval }
}
}
impl Default for VolumeProducer {
fn default() -> Self {
Self::new()
}
}
impl Producer for VolumeProducer {
fn name(&self) -> String {
"volume".to_string()
}
fn run(self: Box<Self>, tx: MsgSender) -> ProducerFuture {
Box::pin(run(tx, self.interval))
}
}
async fn run(tx: MsgSender, period: Duration) -> ProducerResult {
let (done_tx, done_rx) = oneshot::channel();
let tx_for_thread = tx.clone();
thread::Builder::new()
.name("tablero-volume".to_string())
.spawn(move || {
let result = run_main_loop(tx_for_thread, period);
let _ = done_tx.send(result);
})
.map_err(|e| -> Box<dyn std::error::Error + Send + Sync> { Box::new(e) })?;
match done_rx.await {
Ok(result) => result,
Err(_) => Err("volume producer thread dropped its result".into()),
}
}
fn run_main_loop(tx: MsgSender, period: Duration) -> ProducerResult {
pipewire::init();
let mainloop = MainLoopRc::new(None)?;
let context = ContextRc::new(&mainloop, None)?;
let core = context.connect_rc(None)?;
let registry = core.get_registry_rc()?;
let state: Arc<Mutex<State>> = Arc::new(Mutex::new(State::default()));
let state_for_timer = state.clone();
let tx_for_timer = tx.clone();
let node_proxies: Rc<RefCell<Vec<Node>>> = Rc::new(RefCell::new(Vec::new()));
let node_listeners: Rc<RefCell<Vec<Box<dyn Listener>>>> = Rc::new(RefCell::new(Vec::new()));
let node_proxies_for_registry = node_proxies.clone();
let node_listeners_for_registry = node_listeners.clone();
let node_listeners_for_metadata = node_listeners.clone();
let timer = mainloop.loop_().add_timer(move |_expirations| {
on_tick(&state_for_timer, &tx_for_timer);
});
timer.update_timer(Some(period), Some(period));
let _registry_listener: pipewire::registry::Listener = attach_registry_listener(
registry.clone(),
state.clone(),
node_proxies_for_registry,
node_listeners_for_registry,
tx.clone(),
);
let metadata_proxies: Rc<RefCell<Vec<pipewire::metadata::Metadata>>> =
Rc::new(RefCell::new(Vec::new()));
let _metadata_listener: pipewire::registry::Listener = attach_metadata_listener(
®istry,
state.clone(),
metadata_proxies,
node_listeners_for_metadata,
tx.clone(),
);
mainloop.run();
Ok(())
}
#[derive(Default)]
struct State {
sinks: HashMap<u32, SinkEntry>,
default_sink_name: Option<String>,
last_emitted: Option<Option<Volume>>,
}
fn on_tick(state: &Arc<Mutex<State>>, tx: &MsgSender) {
let snapshot = compute_snapshot(state);
let changed = {
let state = state.lock().expect("volume state poisoned");
state.last_emitted.as_ref() != Some(&snapshot)
};
if changed {
{
let mut state = state.lock().expect("volume state poisoned");
state.last_emitted = Some(snapshot);
}
if tx.send(Msg::Volume(snapshot)).is_err() {
debug!("volume: tx.send returned Closed");
}
}
}
fn compute_snapshot(state: &Arc<Mutex<State>>) -> Option<Volume> {
let state = state.lock().expect("volume state poisoned");
let sinks = &state.sinks;
if sinks.is_empty() {
return None;
}
let candidates: Vec<(u32, &str, bool)> = sinks
.iter()
.map(|(id, entry)| (*id, entry.name.as_str(), entry.running))
.collect();
let active = pick_active_sink_id(&candidates)
.or_else(|| {
let default_name = state.default_sink_name.as_deref()?;
let mut sorted = candidates.clone();
sorted.sort_by(|a, b| a.1.cmp(b.1));
sorted
.into_iter()
.find(|(_, name, _)| *name == default_name)
.map(|(id, _, _)| id)
})
.or_else(|| {
let mut sorted = candidates.clone();
sorted.sort_by(|a, b| a.1.cmp(b.1));
sorted.into_iter().next().map(|(id, _, _)| id)
});
let id = active?;
let entry = sinks.get(&id)?;
let (level, muted) = entry.volume?;
Some(Volume::new(level, muted, entry.device))
}
fn attach_registry_listener(
registry: RegistryRc,
state: Arc<Mutex<State>>,
node_proxies: Rc<RefCell<Vec<Node>>>,
node_listeners: Rc<RefCell<Vec<Box<dyn Listener>>>>,
tx: MsgSender,
) -> pipewire::registry::Listener {
let registry_weak = registry.downgrade();
let state_for_global = state.clone();
let state_for_remove = state;
let node_proxies_for_cb = node_proxies.clone();
let node_listeners_for_cb = node_listeners.clone();
let tx_for_global = tx.clone();
let tx_for_remove = tx;
registry
.add_listener_local()
.global(move |obj| {
if obj.type_ != ObjectType::Node {
return;
}
let Some(props) = obj.props else {
return;
};
if !is_audio_sink(props) {
return;
}
let Some(registry) = registry_weak.upgrade() else {
return;
};
let node: Node = match registry.bind(obj) {
Ok(n) => n,
Err(e) => {
debug!("volume: bind sink node failed: {e}");
return;
}
};
let id = obj.id;
let name = props.get("node.name").unwrap_or("?").to_string();
{
let mut s = state_for_global.lock().expect("volume state poisoned");
s.sinks.insert(
id,
SinkEntry {
name: name.clone(),
device: DeviceKind::Other,
running: false,
volume: None,
},
);
}
let state_for_info = state_for_global.clone();
let state_for_param = state_for_global.clone();
let tx_for_info = tx_for_global.clone();
let tx_for_param = tx_for_global.clone();
let state_for_tick_after_info = state_for_global.clone();
let state_for_tick_after_param = state_for_global.clone();
let node_listener = node
.add_listener_local()
.info(move |info: &NodeInfoRef| {
let running = is_node_running(&info.state());
let props = info.props();
let device = props
.map(device_kind_from_props)
.unwrap_or(DeviceKind::Other);
{
let mut s = state_for_info.lock().expect("volume state poisoned");
if let Some(entry) = s.sinks.get_mut(&id) {
entry.running = running;
if device != DeviceKind::Other || entry.device == DeviceKind::Other {
entry.device = device;
}
}
}
on_tick(&state_for_tick_after_info, &tx_for_info);
})
.param(move |_seq, _id, _index, _next, param: Option<&Pod>| {
let Some(param) = param else {
return;
};
let parsed = parse_props_pod(param.as_bytes());
let Some((level, muted)) = parsed else {
return;
};
{
let mut s = state_for_param.lock().expect("volume state poisoned");
if let Some(entry) = s.sinks.get_mut(&id) {
entry.volume = Some((level, muted));
}
}
on_tick(&state_for_tick_after_param, &tx_for_param);
})
.register();
node.subscribe_params(&[pipewire::spa::param::ParamType::Props]);
node.enum_params(0, Some(pipewire::spa::param::ParamType::Props), 0, 1);
node_listeners_for_cb
.borrow_mut()
.push(Box::new(node_listener));
node_proxies_for_cb.borrow_mut().push(node);
on_tick(&state_for_global, &tx_for_global);
})
.global_remove(move |id| {
{
let mut s = state_for_remove.lock().expect("volume state poisoned");
s.sinks.remove(&id);
}
on_tick(&state_for_remove, &tx_for_remove);
})
.register()
}
fn attach_metadata_listener(
registry: &RegistryRc,
state: Arc<Mutex<State>>,
node_proxies: Rc<RefCell<Vec<pipewire::metadata::Metadata>>>,
node_listeners: Rc<RefCell<Vec<Box<dyn Listener>>>>,
tx: MsgSender,
) -> pipewire::registry::Listener {
use pipewire::metadata::Metadata;
let registry_weak = registry.downgrade();
let node_listeners_for_cb = node_listeners.clone();
let tx_for_metadata = tx;
registry
.add_listener_local()
.global(move |obj: &GlobalObject<&DictRef>| {
if obj.type_ != ObjectType::Metadata {
return;
}
let Some(registry) = registry_weak.upgrade() else {
return;
};
let metadata: Metadata = match registry.bind(obj) {
Ok(m) => m,
Err(_) => return,
};
let state_for_property = state.clone();
let state_for_tick = state.clone();
let tx_for_tick = tx_for_metadata.clone();
let md_listener = metadata
.add_listener_local()
.property(move |_subject, key, _type_, value| {
if key != Some(DEFAULT_AUDIO_SINK_KEY) {
return 0;
}
let Some(value) = value else {
return 0;
};
let name = parse_metadata_default_name(value);
{
let mut s = state_for_property.lock().expect("volume state poisoned");
s.default_sink_name = name;
}
on_tick(&state_for_tick, &tx_for_tick);
0
})
.register();
node_listeners_for_cb
.borrow_mut()
.push(Box::new(md_listener));
node_proxies.borrow_mut().push(metadata);
})
.register()
}
fn parse_metadata_default_name(value: &str) -> Option<String> {
let key = "\"name\"";
let name_start = value.find(key)?;
let after_key = &value[name_start + key.len()..];
let colon = after_key.find(':')?;
let after_colon = after_key[colon + 1..].trim_start();
let open_quote = after_colon.find('"')?;
let after_quote = &after_colon[open_quote + 1..];
let close_quote = after_quote.find('"')?;
Some(after_quote[..close_quote].to_string())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn icon_name_headphones_prefix_maps_to_headphones() {
assert_eq!(
device_kind_from_icon_name(Some("audio-headphones")),
DeviceKind::Headphones
);
assert_eq!(
device_kind_from_icon_name(Some("audio-headphones-analog")),
DeviceKind::Headphones
);
}
#[test]
fn icon_name_headset_prefix_maps_to_headset() {
assert_eq!(
device_kind_from_icon_name(Some("audio-headset")),
DeviceKind::Headset
);
}
#[test]
fn icon_name_speakers_prefix_maps_to_speakers() {
assert_eq!(
device_kind_from_icon_name(Some("audio-speakers")),
DeviceKind::Speakers
);
}
#[test]
fn icon_name_video_display_maps_to_monitor() {
assert_eq!(
device_kind_from_icon_name(Some("video-display")),
DeviceKind::Monitor
);
assert_eq!(
device_kind_from_icon_name(Some("video-monitor")),
DeviceKind::Monitor
);
}
#[test]
fn icon_name_phone_maps_to_phone() {
assert_eq!(
device_kind_from_icon_name(Some("audio-handsfree")),
DeviceKind::Phone
);
assert_eq!(device_kind_from_icon_name(Some("phone")), DeviceKind::Phone);
}
#[test]
fn icon_name_tv_maps_to_tv() {
assert_eq!(device_kind_from_icon_name(Some("tv")), DeviceKind::Tv);
assert_eq!(
device_kind_from_icon_name(Some("television")),
DeviceKind::Tv
);
}
#[test]
fn icon_name_missing_or_unknown_falls_back_to_other() {
assert_eq!(device_kind_from_icon_name(None), DeviceKind::Other);
assert_eq!(
device_kind_from_icon_name(Some("totally-unknown-icon")),
DeviceKind::Other
);
}
#[test]
fn icon_name_audio_card_analog_maps_to_speakers() {
assert_eq!(
device_kind_from_icon_name(Some("audio-card-analog")),
DeviceKind::Speakers
);
assert_eq!(
device_kind_from_icon_name(Some("audio-card")),
DeviceKind::Speakers
);
}
#[test]
fn icon_name_audio_card_hdmi_maps_to_monitor() {
assert_eq!(
device_kind_from_icon_name(Some("audio-card-hdmi")),
DeviceKind::Monitor
);
assert_eq!(
device_kind_from_icon_name(Some("audio-card-dp")),
DeviceKind::Monitor
);
}
#[test]
fn form_factor_known_values_map_to_their_kinds() {
assert_eq!(
device_kind_from_form_factor(Some("headphone")),
DeviceKind::Headphones
);
assert_eq!(
device_kind_from_form_factor(Some("headset")),
DeviceKind::Headset
);
assert_eq!(
device_kind_from_form_factor(Some("speaker")),
DeviceKind::Speakers
);
assert_eq!(
device_kind_from_form_factor(Some("monitor")),
DeviceKind::Monitor
);
assert_eq!(device_kind_from_form_factor(Some("tv")), DeviceKind::Tv);
assert_eq!(
device_kind_from_form_factor(Some("handset")),
DeviceKind::Phone
);
}
#[test]
fn form_factor_missing_or_unknown_falls_back_to_other() {
assert_eq!(device_kind_from_form_factor(None), DeviceKind::Other);
assert_eq!(
device_kind_from_form_factor(Some("wonderful")),
DeviceKind::Other
);
}
#[test]
fn pick_active_sink_returns_none_for_no_candidates() {
assert_eq!(pick_active_sink_id(&[]), None);
}
#[test]
fn pick_active_sink_returns_none_when_nothing_is_running() {
let candidates = &[(1u32, "alpha", false), (2, "beta", false)];
assert_eq!(pick_active_sink_id(candidates), None);
}
#[test]
fn pick_active_sink_picks_the_running_sink() {
let candidates = &[
(1u32, "alpha", false),
(2, "beta", true),
(3, "gamma", false),
];
assert_eq!(pick_active_sink_id(candidates), Some(2));
}
#[test]
fn pick_active_sink_tiebreaks_by_alphabetical_name() {
let a = &[
(1u32, "charlie", true),
(2, "alpha", true),
(3, "bravo", true),
];
assert_eq!(pick_active_sink_id(a), Some(2));
let b = &[
(2u32, "alpha", true),
(3, "bravo", true),
(1, "charlie", true),
];
assert_eq!(pick_active_sink_id(b), Some(2));
}
#[test]
fn pick_active_sink_ignores_non_running_even_with_earlier_name() {
let candidates = &[(1u32, "alpha", false), (2, "beta", true)];
assert_eq!(pick_active_sink_id(candidates), Some(2));
}
#[test]
fn parse_metadata_default_name_extracts_name_from_json() {
let value = r#"{"name":"alsa_output.pci-0000_00_1f.3.analog-stereo"}"#;
assert_eq!(
parse_metadata_default_name(value).as_deref(),
Some("alsa_output.pci-0000_00_1f.3.analog-stereo")
);
}
#[test]
fn parse_metadata_default_name_handles_extra_whitespace() {
let value = r#"{ "name" : "my-sink" }"#;
assert_eq!(
parse_metadata_default_name(value).as_deref(),
Some("my-sink")
);
}
#[test]
fn parse_metadata_default_name_returns_none_for_missing_name() {
assert_eq!(parse_metadata_default_name(r#"{"foo":"bar"}"#), None);
assert_eq!(parse_metadata_default_name(""), None);
}
#[test]
fn parse_props_pod_round_trip_extracts_volume_and_mute() {
use pipewire::spa::pod::PropertyFlags;
use pipewire::spa::pod::serialize::{GenError, PodSerialize, PodSerializer};
use pipewire::spa::sys::{
SPA_PARAM_Props, SPA_PROP_channelVolumes, SPA_PROP_mute, SPA_TYPE_OBJECT_Props,
};
use std::io::Cursor;
struct MyProps {
channels: Vec<f32>,
muted: bool,
}
impl PodSerialize for MyProps {
fn serialize<O: std::io::Write + std::io::Seek>(
&self,
serializer: PodSerializer<O>,
) -> Result<pipewire::spa::pod::serialize::SerializeSuccess<O>, GenError> {
let mut obj_serializer =
serializer.serialize_object(SPA_TYPE_OBJECT_Props, SPA_PARAM_Props)?;
obj_serializer.serialize_property(
SPA_PROP_channelVolumes,
self.channels.as_slice(),
PropertyFlags::empty(),
)?;
obj_serializer.serialize_property(
SPA_PROP_mute,
&self.muted,
PropertyFlags::empty(),
)?;
obj_serializer.end()
}
}
assert_eq!(
SPA_PROP_channelVolumes, 65544,
"SPA_PROP_channelVolumes drifted"
);
assert_eq!(SPA_PROP_mute, 65540, "SPA_PROP_mute drifted");
let cubic = 0.5_f32.powi(3);
let props = MyProps {
channels: vec![cubic, cubic],
muted: true,
};
let bytes = PodSerializer::serialize(Cursor::new(Vec::new()), &props)
.unwrap()
.0
.into_inner();
let parsed = parse_props_pod(&bytes);
assert_eq!(parsed, Some((0.5, true)));
}
#[test]
fn average_to_user_volume_inverts_wireplumbers_cubic_ramp() {
let channels = vec![0.5_f32.powi(3), 0.5_f32.powi(3)];
let user_volume = average_to_user_volume(&channels).unwrap();
assert!((user_volume - 0.5).abs() < 1e-5);
let cubic = 0.3_f32.powi(3);
let user_volume = average_to_user_volume(&[cubic, cubic]).unwrap();
assert!((user_volume - 0.3).abs() < 1e-5);
assert_eq!(average_to_user_volume(&[]), None);
let user_volume = average_to_user_volume(&[f32::NAN, 0.5_f32.powi(3)]).unwrap();
assert!((user_volume - 0.5).abs() < 1e-5);
let user_volume = average_to_user_volume(&[-0.1, 0.5_f32.powi(3)]).unwrap();
assert!((user_volume - 0.5).abs() < 1e-5);
}
#[test]
fn parse_props_pod_returns_none_for_garbage() {
assert_eq!(parse_props_pod(&[]), None);
assert_eq!(parse_props_pod(&[0xff; 16]), None);
}
#[test]
fn compute_snapshot_returns_none_for_no_sinks() {
let state = Arc::new(Mutex::new(State::default()));
assert!(compute_snapshot(&state).is_none());
}
#[test]
fn compute_snapshot_uses_the_running_sink_when_available() {
let mut state = State::default();
state.sinks.insert(
1,
SinkEntry {
name: "headphones".to_string(),
device: DeviceKind::Headphones,
running: true,
volume: Some((0.5, false)),
},
);
state.sinks.insert(
2,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: false,
volume: Some((0.3, true)),
},
);
let state = Arc::new(Mutex::new(state));
let snap = compute_snapshot(&state).expect("snapshot");
assert_eq!(snap.level(), 50);
assert!(!snap.muted());
assert_eq!(snap.device(), DeviceKind::Headphones);
}
#[test]
fn compute_snapshot_falls_back_to_default_sink_when_none_running() {
let mut state = State::default();
state.sinks.insert(
1,
SinkEntry {
name: "headphones".to_string(),
device: DeviceKind::Headphones,
running: false,
volume: Some((0.5, false)),
},
);
state.sinks.insert(
2,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: false,
volume: Some((0.3, true)),
},
);
state.default_sink_name = Some("speakers".to_string());
let state = Arc::new(Mutex::new(state));
let snap = compute_snapshot(&state).expect("snapshot");
assert_eq!(snap.device(), DeviceKind::Speakers);
assert!(snap.muted());
}
#[test]
fn compute_snapshot_returns_none_when_sink_has_no_volume_yet() {
let mut state = State::default();
state.sinks.insert(
1,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: true,
volume: None,
},
);
let state = Arc::new(Mutex::new(state));
assert!(compute_snapshot(&state).is_none());
}
}