use pipewire::channel;
use tracing::info;
use crate::paths::socket_path;
use crate::protocol::{Request, Response};
use crate::state::{EqBand, MAX_ABS_PREAMP_DB, MAX_BANDS, PwCommand};
use super::state::DaemonState;
pub(super) fn dispatch_request(
req: Request,
state: &DaemonState,
cmd_tx: &channel::Sender<PwCommand>,
) -> Response {
match req {
Request::GetStatus => Response {
ok: true,
error: None,
status: Some(state.get_status()),
},
Request::SetBands { bands } => {
if bands.len() > MAX_BANDS {
return Response {
ok: false,
error: Some(format!("too many bands: {} (max {MAX_BANDS})", bands.len())),
status: None,
};
}
if let Err(reason) = bands.iter().try_for_each(EqBand::validate) {
return Response {
ok: false,
error: Some(reason),
status: None,
};
}
let count = bands.len();
state.status.lock().unwrap().bands.clone_from(&bands);
let _ = cmd_tx.send(PwCommand::UpdateEq { bands });
info!(count, "Bands queued for EQ update");
Response {
ok: true,
error: None,
status: None,
}
}
Request::SetPreamp { gain } => {
if !gain.is_finite() || gain.abs() > MAX_ABS_PREAMP_DB {
return Response {
ok: false,
error: Some(format!(
"preamp {gain} dB out of range ±{MAX_ABS_PREAMP_DB}"
)),
status: None,
};
}
state.status.lock().unwrap().preamp = gain; state.pipeline.set_preamp(gain); info!(gain, "Preamp updated");
Response {
ok: true,
error: None,
status: None,
}
}
Request::SetBypass { bypass } => {
state.status.lock().unwrap().bypass = bypass;
state.pipeline.set_bypass(bypass);
info!(bypass, "Bypass toggled");
Response {
ok: true,
error: None,
status: None,
}
}
Request::ConnectDevice { node_id } => {
let mut s = state.status.lock().unwrap();
let Some(filter_id) = s.filter_node_id else {
return Response {
ok: false,
error: Some("Filter not ready yet".into()),
status: None,
};
};
if node_id == filter_id {
return Response {
ok: false,
error: Some("Cannot connect filter to itself".into()),
status: None,
};
}
if s.null_sink.module_id() == Some(node_id) {
return Response {
ok: false,
error: Some(
"Cannot connect to the null sink (would create a feedback loop)".into(),
),
status: None,
};
}
if s.connected_devices.contains(&node_id) || s.pending_devices.contains(&node_id) {
return Response {
ok: true,
error: None,
status: None,
};
}
s.pending_devices.push(node_id);
drop(s);
if cmd_tx
.send(PwCommand::ConnectDevice { filter_id, node_id })
.is_err()
{
state
.status
.lock()
.unwrap()
.pending_devices
.retain(|id| *id != node_id);
return Response {
ok: false,
error: Some("PipeWire thread unavailable".into()),
status: None,
};
}
info!(node_id, "Device connect queued");
Response {
ok: true,
error: None,
status: None,
} }
Request::DisconnectDevice { node_id } => {
let mut s = state.status.lock().unwrap();
let Some(filter_id) = s.filter_node_id else {
return Response {
ok: false,
error: Some("Filter not ready yet".into()),
status: None,
};
};
let is_connected = s.connected_devices.contains(&node_id);
let is_pending = s.pending_devices.contains(&node_id);
if !is_connected && !is_pending {
return Response {
ok: false,
error: Some(format!("Device {node_id} is not connected")),
status: None,
};
}
if !is_pending {
s.pending_devices.push(node_id);
}
drop(s);
if cmd_tx
.send(PwCommand::DisconnectDevice { filter_id, node_id })
.is_err()
{
state
.status
.lock()
.unwrap()
.pending_devices
.retain(|id| *id != node_id);
return Response {
ok: false,
error: Some("PipeWire thread unavailable".into()),
status: None,
};
}
info!(node_id, "Device disconnect queued");
Response {
ok: true,
error: None,
status: None,
}
}
Request::Shutdown => {
info!("Shutdown requested by client");
state
.shutting_down
.store(true, std::sync::atomic::Ordering::Release);
let _ = cmd_tx.send(PwCommand::Terminate);
if let Ok(path) = socket_path() {
let _ = std::os::unix::net::UnixStream::connect(path);
}
Response {
ok: true,
error: None,
status: None,
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use crate::pipeline::{Pipeline, SAMPLE_RATE};
#[test]
fn connect_is_pending_until_confirmed() {
let (cmd_tx, cmd_rx) = pipewire::channel::channel::<PwCommand>();
let state = DaemonState::new(Arc::new(Pipeline::new(SAMPLE_RATE)));
state.status.lock().unwrap().filter_node_id = Some(42);
let resp = dispatch_request(Request::ConnectDevice { node_id: 7 }, &state, &cmd_tx);
assert!(resp.ok);
let s = state.get_status();
assert!(s.connected_devices.is_empty()); assert!(s.pending_devices.contains(&7));
state.handle_pw_event(crate::state::PwEvent::LinkResult {
device_id: 7,
connect: true,
ok: true,
});
let s = state.get_status();
assert!(s.connected_devices.contains(&7)); assert!(!s.pending_devices.contains(&7)); drop(cmd_rx);
}
#[test]
fn failed_link_never_shows_connected() {
let (cmd_tx, cmd_rx) = pipewire::channel::channel::<PwCommand>();
let state = DaemonState::new(Arc::new(Pipeline::new(SAMPLE_RATE)));
state.status.lock().unwrap().filter_node_id = Some(42);
let resp = dispatch_request(Request::ConnectDevice { node_id: 7 }, &state, &cmd_tx);
assert!(resp.ok);
state.handle_pw_event(crate::state::PwEvent::LinkResult {
device_id: 7,
connect: true,
ok: false,
});
let s = state.get_status();
assert!(s.connected_devices.is_empty());
assert!(!s.pending_devices.contains(&7)); drop(cmd_rx);
}
#[test]
fn disconnect_requires_connection() {
let (cmd_tx, cmd_rx) = pipewire::channel::channel::<PwCommand>();
let state = DaemonState::new(Arc::new(Pipeline::new(SAMPLE_RATE)));
state.status.lock().unwrap().filter_node_id = Some(42);
let resp = dispatch_request(Request::DisconnectDevice { node_id: 7 }, &state, &cmd_tx);
assert!(!resp.ok);
let _ = dispatch_request(Request::ConnectDevice { node_id: 7 }, &state, &cmd_tx);
state.handle_pw_event(crate::state::PwEvent::LinkResult {
device_id: 7,
connect: true,
ok: true,
});
let resp = dispatch_request(Request::DisconnectDevice { node_id: 7 }, &state, &cmd_tx);
assert!(resp.ok);
drop(cmd_rx);
}
#[test]
fn node_list_prunes_vanished_devices() {
let (cmd_tx, cmd_rx) = pipewire::channel::channel::<PwCommand>();
let state = DaemonState::new(Arc::new(Pipeline::new(SAMPLE_RATE)));
state.status.lock().unwrap().filter_node_id = Some(42);
let _ = dispatch_request(Request::ConnectDevice { node_id: 7 }, &state, &cmd_tx);
assert!(state.get_status().pending_devices.contains(&7));
state.handle_pw_event(crate::state::PwEvent::NodeList(vec![]));
assert!(state.get_status().pending_devices.is_empty());
drop(cmd_rx);
}
}