use std::io::Write;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use crate::errors::CliError;
#[derive(Serialize, Deserialize)]
struct LockFile {
pid: u32,
started_at: String,
operation: String,
#[serde(default)]
nonce: String,
}
const STALE_THRESHOLD_SECS: i64 = 3600;
fn is_stale(started_at: &str) -> bool {
chrono::DateTime::parse_from_rfc3339(started_at)
.map(|t| chrono::Utc::now().signed_duration_since(t).num_seconds() > STALE_THRESHOLD_SECS)
.unwrap_or(true)
}
#[cfg(unix)]
fn process_alive(pid: u32) -> bool {
if pid == 0 || pid > i32::MAX as u32 {
return false;
}
unsafe { libc::kill(pid as i32, 0) == 0 }
}
#[cfg(windows)]
fn process_alive(pid: u32) -> bool {
use windows_sys::Win32::Foundation::CloseHandle;
use windows_sys::Win32::System::Threading::{OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION};
unsafe {
let handle = OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid);
if handle.is_null() {
false
} else {
CloseHandle(handle);
true
}
}
}
#[cfg(not(any(unix, windows)))]
fn process_alive(_pid: u32) -> bool {
true
}
pub struct DuplicateGuard {
lock_path: PathBuf,
operation: String,
nonce: String,
acquired: bool,
}
impl DuplicateGuard {
pub fn new(data_dir: &Path, operation: &str) -> Self {
let lock_dir = data_dir.join("locks");
let _ = std::fs::create_dir_all(&lock_dir);
Self {
lock_path: lock_dir.join(format!("{operation}.lock")),
operation: operation.to_string(),
nonce: uuid::Uuid::new_v4().to_string(),
acquired: false,
}
}
pub fn acquire(&mut self, force: bool) -> Result<(), CliError> {
for _ in 0..16 {
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&self.lock_path)
{
Ok(mut file) => {
let lock = LockFile {
pid: std::process::id(),
started_at: chrono::Utc::now().to_rfc3339(),
operation: self.operation.clone(),
nonce: self.nonce.clone(),
};
file.write_all(serde_json::to_string(&lock)?.as_bytes())?;
self.acquired = true;
return Ok(());
}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
let held = std::fs::read_to_string(&self.lock_path)
.ok()
.and_then(|c| serde_json::from_str::<LockFile>(&c).ok());
let reclaimable = match &held {
Some(lock) => {
force || !process_alive(lock.pid) || is_stale(&lock.started_at)
}
None => true,
};
if !reclaimable {
let lock = held.expect("Some when not reclaimable");
return Err(CliError::InvalidInput(format!(
"Operation '{}' already running (pid {}). Use --force to override.",
lock.operation, lock.pid
)));
}
match std::fs::remove_file(&self.lock_path) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(e.into()),
}
}
Err(e) => return Err(e.into()),
}
}
Err(CliError::InvalidInput(format!(
"could not acquire '{}' lock — contended lock file at {}",
self.operation,
self.lock_path.display()
)))
}
pub fn release(&mut self) {
if !self.acquired {
return;
}
self.acquired = false;
let still_ours = std::fs::read_to_string(&self.lock_path)
.ok()
.and_then(|c| serde_json::from_str::<LockFile>(&c).ok())
.is_some_and(|lock| lock.nonce == self.nonce);
if still_ours {
let _ = std::fs::remove_file(&self.lock_path);
}
}
}
impl Drop for DuplicateGuard {
fn drop(&mut self) {
self.release();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn second_acquire_conflicts_and_force_overrides() {
let tmp = tempfile::tempdir().unwrap();
let mut first = DuplicateGuard::new(tmp.path(), "op");
first.acquire(false).unwrap();
let mut second = DuplicateGuard::new(tmp.path(), "op");
let err = second.acquire(false).unwrap_err();
assert_eq!(err.exit_code(), 3);
second.acquire(true).unwrap();
}
#[test]
fn lock_released_on_drop() {
let tmp = tempfile::tempdir().unwrap();
let lock_path = {
let mut guard = DuplicateGuard::new(tmp.path(), "op");
guard.acquire(false).unwrap();
guard.lock_path.clone()
};
assert!(!lock_path.exists());
}
#[test]
fn failed_acquire_does_not_release_holders_lock() {
let tmp = tempfile::tempdir().unwrap();
let mut holder = DuplicateGuard::new(tmp.path(), "op");
holder.acquire(false).unwrap();
{
let mut retry = DuplicateGuard::new(tmp.path(), "op");
retry.acquire(false).unwrap_err();
}
assert!(holder.lock_path.exists());
let mut retry2 = DuplicateGuard::new(tmp.path(), "op");
retry2.acquire(false).unwrap_err();
}
#[test]
fn stale_lock_is_overwritten() {
let tmp = tempfile::tempdir().unwrap();
let lock_dir = tmp.path().join("locks");
std::fs::create_dir_all(&lock_dir).unwrap();
let two_hours_ago = chrono::Utc::now() - chrono::Duration::hours(2);
let stale = serde_json::json!({
"pid": std::process::id(),
"started_at": two_hours_ago.to_rfc3339(),
"operation": "op",
});
std::fs::write(lock_dir.join("op.lock"), stale.to_string()).unwrap();
let mut guard = DuplicateGuard::new(tmp.path(), "op");
guard.acquire(false).unwrap();
}
#[test]
fn concurrent_acquire_only_one_wins() {
let tmp = tempfile::tempdir().unwrap();
let mut a = DuplicateGuard::new(tmp.path(), "gen");
a.acquire(false).unwrap();
let mut b = DuplicateGuard::new(tmp.path(), "gen");
let err = b.acquire(false).unwrap_err();
assert_eq!(
err.exit_code(),
3,
"second concurrent acquire must be exit 3"
);
assert!(!b.acquired, "the loser must not believe it holds the lock");
}
#[test]
fn release_does_not_delete_a_reclaimed_lock() {
let tmp = tempfile::tempdir().unwrap();
let mut b = DuplicateGuard::new(tmp.path(), "op");
b.acquire(false).unwrap();
let mut a = DuplicateGuard::new(tmp.path(), "op");
a.acquire(true).unwrap();
b.release();
assert!(
a.lock_path.exists(),
"stale releaser deleted the successor's lock"
);
a.release();
assert!(!a.lock_path.exists());
}
#[cfg(unix)]
#[test]
fn dead_pid_lock_reclaimed_regardless_of_age() {
let tmp = tempfile::tempdir().unwrap();
let lock_dir = tmp.path().join("locks");
std::fs::create_dir_all(&lock_dir).unwrap();
let mut child = std::process::Command::new("true").spawn().unwrap();
let dead_pid = child.id();
child.wait().unwrap();
let fresh = serde_json::json!({
"pid": dead_pid,
"started_at": chrono::Utc::now().to_rfc3339(),
"operation": "op",
"nonce": "someone-else",
});
std::fs::write(lock_dir.join("op.lock"), fresh.to_string()).unwrap();
let mut guard = DuplicateGuard::new(tmp.path(), "op");
guard.acquire(false).unwrap();
}
}