use std::cell::{Cell, RefCell};
use std::collections::HashMap;
use std::rc::Rc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;
use log::{debug, info, warn};
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 CONNECT_TIMEOUT: Duration = Duration::from_secs(2);
const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10);
const RETRY_DELAY: Duration = Duration::from_secs(1);
const MAX_RETRY_DELAY: Duration = Duration::from_secs(30);
const RESYNC_TICKS: u32 = 15;
const PW_ID_CORE: u32 = 0;
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)
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct PropsUpdate {
pub volume: Option<f32>,
pub muted: Option<bool>,
}
pub fn parse_props_pod(pod_bytes: &[u8]) -> Option<PropsUpdate> {
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 volume = obj
.find_prop(pipewire::spa::utils::Id(SPA_PROP_channelVolumes))
.and_then(|channel_volumes_pod| {
let value: Value =
pipewire::spa::pod::deserialize::PodDeserializer::deserialize_any_from(
channel_volumes_pod.value().as_bytes(),
)
.ok()
.map(|(_, v)| v)?;
match &value {
Value::ValueArray(pipewire::spa::pod::ValueArray::Float(v)) => {
average_to_user_volume(v)
}
Value::Float(v) => average_to_user_volume(&[*v]),
_ => None,
}
});
let muted = obj
.find_prop(pipewire::spa::utils::Id(SPA_PROP_mute))
.and_then(|mute_pod| mute_pod.value().get_bool().ok());
if volume.is_none() && muted.is_none() {
return None;
}
Some(PropsUpdate { volume, muted })
}
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 Emitter {
tx: MsgSender,
closed: Arc<AtomicBool>,
}
impl Emitter {
fn new(tx: MsgSender) -> Self {
Self {
tx,
closed: Arc::new(AtomicBool::new(false)),
}
}
fn send(&self, msg: Msg) {
if self.tx.send(msg).is_err() {
debug!("volume: tx.send returned Closed");
self.closed.store(true, Ordering::Relaxed);
}
}
fn is_closed(&self) -> bool {
self.closed.load(Ordering::Relaxed)
}
}
#[derive(Clone)]
struct SinkEntry {
name: String,
device: DeviceKind,
running: bool,
volume: Option<f32>,
muted: Option<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 emitter = Emitter::new(tx);
thread::Builder::new()
.name("tablero-volume".to_string())
.spawn(move || {
let result = supervise(&emitter, 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()),
}
}
struct SessionOutcome {
handshook: bool,
result: ProducerResult,
}
fn supervise(tx: &Emitter, period: Duration) -> ProducerResult {
pipewire::init();
let mut delay = RETRY_DELAY;
let mut failures: u32 = 0;
loop {
let outcome = run_session(tx, period);
if tx.is_closed() {
return Ok(());
}
if outcome.handshook {
failures = 0;
delay = RETRY_DELAY;
}
failures += 1;
let ended = match &outcome.result {
Ok(()) => "ended".to_string(),
Err(e) => format!("failed: {e}"),
};
if failures == 1 {
warn!("volume: PipeWire session {ended}; reconnecting in {delay:?}");
} else {
debug!("volume: PipeWire session {ended} (attempt {failures}); retrying in {delay:?}");
tx.send(Msg::Volume(None));
}
thread::sleep(delay);
delay = (delay * 2).min(MAX_RETRY_DELAY);
}
}
fn run_session(tx: &Emitter, period: Duration) -> SessionOutcome {
let handshook = Rc::new(Cell::new(false));
let result = run_main_loop(tx, period, &handshook);
SessionOutcome {
handshook: handshook.get(),
result,
}
}
fn run_main_loop(tx: &Emitter, period: Duration, handshook: &Rc<Cell<bool>>) -> ProducerResult {
let mainloop = MainLoopRc::new(None)?;
let context = ContextRc::new(&mainloop, None)?;
let core = context.connect_rc(None)?;
let registry = core.get_registry_rc()?;
debug!("volume: connected to PipeWire; awaiting the core handshake");
let pending_sync = Rc::new(Cell::new(false));
let handshook_for_done = handshook.clone();
let pending_for_done = pending_sync.clone();
let mainloop_for_error = mainloop.downgrade();
let _core_listener = core
.add_listener_local()
.done(move |id, _seq| {
if id != PW_ID_CORE {
return;
}
pending_for_done.set(false);
if !handshook_for_done.replace(true) {
info!("volume: PipeWire core handshake complete");
}
})
.error(move |id, _seq, res, message| {
warn!("volume: PipeWire reported an error on object {id}: {message} ({res})");
if id == PW_ID_CORE
&& let Some(mainloop) = mainloop_for_error.upgrade()
{
mainloop.quit();
}
})
.register();
core.sync(PW_ID_CORE as i32)?;
pending_sync.set(true);
let core_for_watchdog = core.clone();
let mainloop_for_watchdog = mainloop.downgrade();
let pending_for_watchdog = pending_sync.clone();
let watchdog = mainloop.loop_().add_timer(move |_expirations| {
if pending_for_watchdog.get() {
warn!("volume: PipeWire did not answer a sync; dropping the connection");
if let Some(mainloop) = mainloop_for_watchdog.upgrade() {
mainloop.quit();
}
return;
}
match core_for_watchdog.sync(PW_ID_CORE as i32) {
Ok(_) => pending_for_watchdog.set(true),
Err(e) => {
warn!("volume: PipeWire sync failed: {e}; dropping the connection");
if let Some(mainloop) = mainloop_for_watchdog.upgrade() {
mainloop.quit();
}
}
}
});
watchdog.update_timer(Some(CONNECT_TIMEOUT), Some(HEARTBEAT_INTERVAL));
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 mainloop_for_timer = mainloop.downgrade();
let timer = mainloop.loop_().add_timer(move |_expirations| {
on_poll_tick(&state_for_timer, &tx_for_timer);
if tx_for_timer.is_closed()
&& let Some(mainloop) = mainloop_for_timer.upgrade()
{
debug!("volume: render loop went away; ending the PipeWire session");
mainloop.quit();
}
});
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>>,
ticks_since_resync: u32,
}
fn on_poll_tick(state: &Arc<Mutex<State>>, tx: &Emitter) {
{
let mut s = state.lock().expect("volume state poisoned");
s.ticks_since_resync += 1;
if s.ticks_since_resync >= RESYNC_TICKS {
s.ticks_since_resync = 0;
s.last_emitted = None;
}
}
on_tick(state, tx);
}
fn on_tick(state: &Arc<Mutex<State>>, tx: &Emitter) {
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 first = {
let mut state = state.lock().expect("volume state poisoned");
let first = !matches!(state.last_emitted, Some(Some(_)));
state.last_emitted = Some(snapshot);
first
};
if let (true, Some(volume)) = (first, snapshot) {
info!(
"volume: active sink at {}%{}",
volume.level(),
if volume.muted() { " (muted)" } else { "" }
);
}
tx.send(Msg::Volume(snapshot));
}
}
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 = entry.volume?;
Some(Volume::new(
level,
entry.muted.unwrap_or(false),
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: Emitter,
) -> 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();
debug!("volume: found audio sink {id} ({name})");
{
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,
muted: 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 Some(update) = parse_props_pod(param.as_bytes()) else {
debug!("volume: sink {id} sent a Props POD with no volume or mute");
return;
};
{
let mut s = state_for_param.lock().expect("volume state poisoned");
if let Some(entry) = s.sinks.get_mut(&id) {
if let Some(volume) = update.volume {
entry.volume = Some(volume);
}
if let Some(muted) = update.muted {
entry.muted = Some(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: Emitter,
) -> 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(err) => {
debug!("volume: could not bind metadata global {}: {err}", obj.id);
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);
match &name {
Some(name) => debug!("volume: default sink is {name}"),
None => debug!("volume: default-sink metadata had no usable 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);
}
struct TestProps {
channels: Option<Vec<f32>>,
muted: Option<bool>,
}
impl pipewire::spa::pod::serialize::PodSerialize for TestProps {
fn serialize<O: std::io::Write + std::io::Seek>(
&self,
serializer: pipewire::spa::pod::serialize::PodSerializer<O>,
) -> Result<
pipewire::spa::pod::serialize::SerializeSuccess<O>,
pipewire::spa::pod::serialize::GenError,
> {
use pipewire::spa::pod::PropertyFlags;
use pipewire::spa::sys::{
SPA_PARAM_Props, SPA_PROP_channelVolumes, SPA_PROP_mute, SPA_TYPE_OBJECT_Props,
};
let mut obj_serializer =
serializer.serialize_object(SPA_TYPE_OBJECT_Props, SPA_PARAM_Props)?;
if let Some(channels) = &self.channels {
obj_serializer.serialize_property(
SPA_PROP_channelVolumes,
channels.as_slice(),
PropertyFlags::empty(),
)?;
}
if let Some(muted) = &self.muted {
obj_serializer.serialize_property(SPA_PROP_mute, muted, PropertyFlags::empty())?;
}
obj_serializer.end()
}
}
fn serialize_props(props: &TestProps) -> Vec<u8> {
pipewire::spa::pod::serialize::PodSerializer::serialize(
std::io::Cursor::new(Vec::new()),
props,
)
.expect("serialize props")
.0
.into_inner()
}
#[test]
fn spa_prop_keys_match_the_shipped_headers() {
use pipewire::spa::sys::{SPA_PROP_channelVolumes, SPA_PROP_mute};
assert_eq!(
SPA_PROP_channelVolumes, 65544,
"SPA_PROP_channelVolumes drifted"
);
assert_eq!(SPA_PROP_mute, 65540, "SPA_PROP_mute drifted");
}
#[test]
fn parse_props_pod_round_trip_extracts_volume_and_mute() {
let cubic = 0.5_f32.powi(3);
let bytes = serialize_props(&TestProps {
channels: Some(vec![cubic, cubic]),
muted: Some(true),
});
assert_eq!(
parse_props_pod(&bytes),
Some(PropsUpdate {
volume: Some(0.5),
muted: Some(true),
})
);
}
#[test]
fn parse_props_pod_accepts_a_volume_only_pod() {
let cubic = 0.3_f32.powi(3);
let bytes = serialize_props(&TestProps {
channels: Some(vec![cubic, cubic]),
muted: None,
});
let parsed = parse_props_pod(&bytes).expect("volume-only POD parses");
assert!((parsed.volume.expect("volume") - 0.3).abs() < 1e-5);
assert_eq!(parsed.muted, None);
}
#[test]
fn parse_props_pod_accepts_a_mute_only_pod() {
let bytes = serialize_props(&TestProps {
channels: None,
muted: Some(true),
});
assert_eq!(
parse_props_pod(&bytes),
Some(PropsUpdate {
volume: None,
muted: Some(true),
})
);
}
#[test]
fn parse_props_pod_returns_none_when_neither_property_is_present() {
let bytes = serialize_props(&TestProps {
channels: None,
muted: None,
});
assert_eq!(parse_props_pod(&bytes), None);
}
#[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),
muted: Some(false),
},
);
state.sinks.insert(
2,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: false,
volume: Some(0.3),
muted: Some(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),
muted: Some(false),
},
);
state.sinks.insert(
2,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: false,
volume: Some(0.3),
muted: Some(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,
muted: None,
},
);
let state = Arc::new(Mutex::new(state));
assert!(compute_snapshot(&state).is_none());
}
#[test]
fn compute_snapshot_shows_a_sink_that_has_not_reported_mute_yet() {
let mut state = State::default();
state.sinks.insert(
1,
SinkEntry {
name: "speakers".to_string(),
device: DeviceKind::Speakers,
running: true,
volume: Some(0.42),
muted: None,
},
);
let state = Arc::new(Mutex::new(state));
let snap = compute_snapshot(&state).expect("snapshot");
assert_eq!(snap.level(), 42);
assert!(!snap.muted());
}
}