use std::fs;
use std::io::{Read, Seek, Write};
use std::os::unix::fs::{DirBuilderExt, FileTypeExt, PermissionsExt};
use std::os::unix::io::AsRawFd;
use std::os::unix::net::UnixListener;
use std::path::Path;
use std::sync::{Arc, atomic::Ordering, mpsc};
use std::thread;
use std::time::Duration;
use pipewire::channel;
use tracing::{debug, error, info};
use crate::paths::{lock_path, socket_path, validate_runtime_dir};
use crate::pipeline::{Pipeline, SAMPLE_RATE};
use crate::protocol::PushEvent;
use crate::state::{PwCommand, PwEvent};
use crate::{AppResult, pw};
use super::auth::peer_is_self;
use super::state::{ClientHandle, DaemonState};
use super::transport::handle_client;
fn acquire_lock(path: &Path) -> AppResult<fs::File> {
let mut lock_file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(path)?;
if unsafe { libc::flock(lock_file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } == -1 {
let mut pid = String::new();
let _ = lock_file.read_to_string(&mut pid);
eprintln!(
"Daemon already running (pid {}). Use `eqtui stop` to stop it first.",
pid.trim()
);
std::process::exit(1);
}
lock_file.set_len(0)?; lock_file.seek(std::io::SeekFrom::Start(0))?;
writeln!(&lock_file, "{}", std::process::id())?;
lock_file.sync_all()?;
Ok(lock_file)
}
pub fn run_daemon() -> AppResult<()> {
tracing::info!("Daemon starting up");
let run_dir = validate_runtime_dir()?;
let eqtui_dir = run_dir.join("eqtui");
fs::DirBuilder::new()
.recursive(true)
.mode(0o700)
.create(&eqtui_dir)?;
fs::set_permissions(&eqtui_dir, fs::Permissions::from_mode(0o700))?;
let lock_path = lock_path()?;
let _lock_file = acquire_lock(&lock_path)?;
let socket_path = socket_path()?;
let pipeline = Arc::new(Pipeline::new(SAMPLE_RATE));
let state = Arc::new(DaemonState::new(pipeline.clone()));
let (pw_tx, pw_rx) = mpsc::channel::<PwEvent>();
let (cmd_tx, cmd_rx) = channel::channel::<PwCommand>();
let pw_pipeline = pipeline.clone();
let pw_shutdown = state.shutting_down.clone(); let pw_thread = thread::Builder::new().name("pw".into()).spawn(move || {
pw::run(pw_tx, cmd_rx, pw_pipeline, pw_shutdown);
})?;
let bridge_state = state.clone();
let bridge_socket = socket_path.clone();
let bridge = thread::Builder::new()
.name("pw-bridge".into())
.spawn(move || {
while let Ok(event) = pw_rx.recv() {
bridge_state.handle_pw_event(event);
}
if !bridge_state.shutting_down.load(Ordering::Acquire) {
error!("PW event channel closed unexpectedly — shutting down daemon");
bridge_state.shutting_down.store(true, Ordering::Release);
if let Err(e) = std::os::unix::net::UnixStream::connect(&bridge_socket) {
debug!(%e, "Failed to connect to socket to unblock accept loop");
}
}
})?;
let peak_state = state.clone();
let peak = thread::Builder::new()
.name("peak-broadcast".into())
.spawn(move || {
loop {
thread::sleep(Duration::from_millis(66));
if peak_state.shutting_down.load(Ordering::Acquire) {
break;
}
let (l, r) = peak_state.pipeline.peaks();
peak_state.push_event(PushEvent::PeakUpdate { l, r });
}
})?;
match fs::symlink_metadata(&socket_path) {
Ok(md) if md.file_type().is_socket() => fs::remove_file(&socket_path)?,
Ok(_) => {
return Err(std::io::Error::new(
std::io::ErrorKind::AlreadyExists,
format!(
"{} exists and is not a socket; refusing to remove",
socket_path.display()
),
)
.into());
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(e.into()),
}
let listener = UnixListener::bind(&socket_path)?;
fs::set_permissions(&socket_path, fs::Permissions::from_mode(0o600))?;
info!("Daemon listening on {}", socket_path.display());
let mut client_id_counter: u64 = 0;
let mut client_handles: Vec<thread::JoinHandle<()>> = Vec::new();
for stream in listener.incoming() {
if state.shutting_down.load(Ordering::Acquire) {
break;
}
let stream = match stream {
Ok(s) => s,
Err(e) => {
error!(%e, "Accept error");
continue;
}
};
let euid = unsafe { libc::geteuid() };
if !peer_is_self(&stream, euid) {
continue; }
client_handles.retain(|h| !h.is_finished());
let handler_state = state.clone();
let handler_cmd_tx = cmd_tx.clone();
let h = thread::Builder::new()
.name(format!("client-{client_id_counter}"))
.spawn(move || {
handle_client(stream, handler_state, handler_cmd_tx, client_id_counter);
})?;
client_handles.push(h);
client_id_counter += 1;
}
info!("Daemon shutting down");
let _ = cmd_tx.send(PwCommand::Terminate);
if let Err(e) = pw_thread.join() {
error!("pw thread panicked: {e:?}");
}
if let Err(e) = bridge.join() {
error!("pw-bridge panicked: {e:?}");
}
if let Err(e) = peak.join() {
error!("peak-broadcast panicked: {e:?}");
}
let drained: Vec<ClientHandle> = state.clients.lock().unwrap().drain(..).collect();
for c in drained {
let _ = c.stream.shutdown(std::net::Shutdown::Both);
}
for h in client_handles {
let _ = h.join();
}
let _ = fs::remove_file(&socket_path);
info!("Daemon exited cleanly");
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn second_instance_cannot_corrupt_lock() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("eqtui.lock");
let _first = acquire_lock(&path).unwrap();
let first_pid = std::fs::read_to_string(&path).unwrap();
assert_eq!(first_pid.trim(), std::process::id().to_string());
let second = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(&path)
.unwrap();
let rc = unsafe { libc::flock(second.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
assert_eq!(rc, -1);
assert_eq!(std::fs::read_to_string(&path).unwrap(), first_pid);
}
}