use std::collections::HashSet;
use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
use crate::setup_core::error::{Error, ReasonCode, Result};
pub const LOCK_FILE_NAME: &str = "target.lock";
static HELD_TARGETS: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
fn held_targets() -> &'static Mutex<HashSet<PathBuf>> {
HELD_TARGETS.get_or_init(|| Mutex::new(HashSet::new()))
}
fn claim_in_process(path: &Path) -> bool {
let mut held = held_targets()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
held.insert(path.to_path_buf())
}
fn release_in_process(path: &Path) {
let mut held = held_targets()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
held.remove(path);
}
#[derive(Debug)]
pub struct TargetLock {
file: File,
path: PathBuf,
}
impl TargetLock {
pub fn acquire(control_directory: &Path) -> Result<Self> {
let path = control_directory.join(LOCK_FILE_NAME);
if !claim_in_process(&path) {
return Err(Error::new(
ReasonCode::LockUnavailable,
format!("this process already holds {}", path.display()),
));
}
let acquire_os_lock = || -> Result<File> {
let file = OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(&path)
.map_err(|source| {
Error::new(
ReasonCode::StateUnavailable,
format!("cannot open lock file {}", path.display()),
)
.with_source(source)
})?;
file.try_lock().map_err(|source| {
Error::new(
ReasonCode::LockUnavailable,
format!("another process holds {}", path.display()),
)
.with_source(source)
})?;
Ok(file)
};
match acquire_os_lock() {
Ok(file) => Ok(Self { file, path }),
Err(error) => {
release_in_process(&path);
Err(error)
}
}
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
pub fn annotate(&mut self, note: &str) -> Result<()> {
self.file
.set_len(0)
.and_then(|()| self.file.write_all(note.as_bytes()))
.map_err(|source| {
Error::new(
ReasonCode::StateUnavailable,
format!("cannot annotate {}", self.path.display()),
)
.with_source(source)
})
}
}
impl Drop for TargetLock {
fn drop(&mut self) {
let _ = self.file.unlock();
release_in_process(&self.path);
}
}
pub fn atomic_write(path: &Path, bytes: &[u8]) -> Result<()> {
let Some(parent) = path.parent() else {
return Err(Error::new(
ReasonCode::StateUnavailable,
format!("{} has no parent directory", path.display()),
));
};
let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else {
return Err(Error::new(
ReasonCode::StateUnavailable,
format!("{} has no usable file name", path.display()),
));
};
fs::create_dir_all(parent).map_err(|source| {
Error::new(
ReasonCode::StateUnavailable,
format!("cannot create {}", parent.display()),
)
.with_source(source)
})?;
let temporary = parent.join(format!(".{file_name}.staging"));
let write = || -> std::io::Result<()> {
let mut file = File::create(&temporary)?;
file.write_all(bytes)?;
file.sync_all()
};
write().map_err(|source| {
let _ = fs::remove_file(&temporary);
Error::new(
ReasonCode::StateUnavailable,
format!("cannot stage {}", temporary.display()),
)
.with_source(source)
})?;
fs::rename(&temporary, path).map_err(|source| {
let _ = fs::remove_file(&temporary);
Error::new(
ReasonCode::StateUnavailable,
format!(
"cannot promote {} to {}",
temporary.display(),
path.display()
),
)
.with_source(source)
})?;
sync_directory(parent);
Ok(())
}
#[cfg(unix)]
fn sync_directory(path: &Path) {
if let Ok(directory) = File::open(path) {
let _ = directory.sync_all();
}
}
#[cfg(not(unix))]
fn sync_directory(_path: &Path) {
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::panic)]
use super::*;
fn scratch(name: &str) -> PathBuf {
let base =
std::env::temp_dir().join(format!("setup-core-lock-{name}-{}", std::process::id()));
let _ = fs::remove_dir_all(&base);
fs::create_dir_all(&base).unwrap();
base
}
#[test]
fn a_second_holder_in_this_process_is_refused_rather_than_silently_merged() {
let control = scratch("contended");
let first = TargetLock::acquire(&control).unwrap();
let error = TargetLock::acquire(&control).unwrap_err();
assert_eq!(error.reason(), ReasonCode::LockUnavailable);
drop(first);
}
#[test]
fn a_refused_acquisition_does_not_strand_the_claim_it_failed_to_complete() {
let control = scratch("no-leak");
let first = TargetLock::acquire(&control).unwrap();
assert!(TargetLock::acquire(&control).is_err());
assert!(TargetLock::acquire(&control).is_err());
drop(first);
assert!(TargetLock::acquire(&control).is_ok());
}
#[test]
fn the_lock_is_released_when_the_holder_is_dropped() {
let control = scratch("released");
drop(TargetLock::acquire(&control).unwrap());
let second = TargetLock::acquire(&control);
assert!(second.is_ok());
}
#[test]
fn two_targets_do_not_contend() {
let one = scratch("target-one");
let two = scratch("target-two");
let first = TargetLock::acquire(&one).unwrap();
let second = TargetLock::acquire(&two);
assert!(second.is_ok());
drop(first);
}
#[test]
fn an_atomic_write_creates_the_directories_its_path_names() {
let base = scratch("atomic-nested");
let file = base.join("antigravity-cli").join("settings.json");
atomic_write(&file, b"{}").unwrap();
assert_eq!(fs::read(&file).unwrap(), b"{}");
let deeper = base.join("a").join("b").join("c").join("leaf");
atomic_write(&deeper, b"deep").unwrap();
assert_eq!(fs::read(&deeper).unwrap(), b"deep");
}
#[test]
fn an_atomic_write_replaces_contents_and_leaves_no_staging_file() {
let base = scratch("atomic");
let file = base.join("state.json");
atomic_write(&file, b"first").unwrap();
atomic_write(&file, b"second").unwrap();
assert_eq!(fs::read(&file).unwrap(), b"second");
assert!(!base.join(".state.json.staging").exists());
}
}