use std::fs::OpenOptions;
use std::path::Path;
use anyhow::{Context, Result, anyhow};
use super::daemon_bridge::{
DAEMON_POLL_INTERVAL, DAEMON_START_TIMEOUT, DaemonBridgeConfig, poll_until_ready,
probe_health_once, spawn_daemon_detached,
};
#[derive(Debug)]
pub struct StartLock {
file: std::fs::File,
}
impl StartLock {
pub fn acquire_blocking(path: &Path) -> Result<Self> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).with_context(|| {
format!("could not create start-lock directory {}", parent.display())
})?;
}
let file = OpenOptions::new()
.create(true)
.read(true)
.truncate(false)
.write(true)
.open(path)
.with_context(|| format!("could not open start-lock file {}", path.display()))?;
#[cfg(unix)]
{
use std::os::unix::io::AsRawFd;
let fd = file.as_raw_fd();
loop {
let rc = unsafe { libc::flock(fd, libc::LOCK_EX) };
if rc == 0 {
break;
}
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue; }
return Err(anyhow!(
"could not acquire start lock {}: {err}",
path.display()
));
}
}
#[cfg(not(unix))]
{
return Err(anyhow!(
"single-flight daemon start requires flock(2) and is not supported on this platform"
));
}
Ok(Self { file })
}
}
impl Drop for StartLock {
fn drop(&mut self) {
#[cfg(unix)]
{
use std::os::unix::io::AsRawFd;
unsafe { libc::flock(self.file.as_raw_fd(), libc::LOCK_UN) };
}
}
}
pub async fn ensure_daemon_up_single_flight(
config: &DaemonBridgeConfig,
lock_path: &Path,
) -> Result<String> {
let startup_timeout = config.startup_timeout.unwrap_or(DAEMON_START_TIMEOUT);
let poll_interval = config.poll_interval.unwrap_or(DAEMON_POLL_INTERVAL);
if probe_health_once(&config.health_url()).await {
return Ok((config.base_url_fn)());
}
let lock_path_owned = lock_path.to_path_buf();
let _lock = tokio::task::spawn_blocking(move || StartLock::acquire_blocking(&lock_path_owned))
.await
.context("start-lock acquisition task panicked")??;
if probe_health_once(&config.health_url()).await {
return Ok((config.base_url_fn)());
}
spawn_daemon_detached(config)?;
poll_until_ready(config, startup_timeout, poll_interval).await
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::{Duration, Instant};
#[test]
fn lock_is_exclusive_across_handles() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("start.lock");
let first = StartLock::acquire_blocking(&path).expect("first acquire");
let p2 = path.clone();
let (tx, rx) = std::sync::mpsc::channel();
let handle = std::thread::spawn(move || {
let _second = StartLock::acquire_blocking(&p2).expect("second acquire");
tx.send(Instant::now()).expect("send");
});
std::thread::sleep(Duration::from_millis(300));
assert!(
rx.try_recv().is_err(),
"second acquisition must block while the first lock is held"
);
drop(first);
let acquired_at = rx
.recv_timeout(Duration::from_secs(5))
.expect("second acquisition must succeed once the first is released");
handle.join().expect("thread join");
let _ = acquired_at;
}
#[test]
fn lock_released_on_error_path() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("start.lock");
let failed: Result<()> = (|| {
let _lock = StartLock::acquire_blocking(&path)?;
Err(anyhow!("simulated start failure"))
})();
assert!(failed.is_err());
let start = Instant::now();
let _again = StartLock::acquire_blocking(&path).expect("reacquire after error");
assert!(
start.elapsed() < Duration::from_secs(1),
"lock must be free after the error path"
);
}
#[test]
fn lock_file_survives_release() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("start.lock");
drop(StartLock::acquire_blocking(&path).expect("acquire"));
assert!(path.exists(), "lock file must not be unlinked on release");
}
#[test]
fn lock_creates_missing_parent_directory() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("nested/deeper/start.lock");
let _lock = StartLock::acquire_blocking(&path).expect("acquire with missing parent");
assert!(path.exists());
}
}