use std::fs::{File, OpenOptions};
use std::io;
use std::path::Path;
use std::time::{Duration, Instant};
use fs4::fs_std::FileExt;
pub const LOCK_DEADLINE: Duration = Duration::from_secs(5);
const POLL_INTERVAL: Duration = Duration::from_millis(50);
pub struct WorkloadLock {
file: File,
}
impl Drop for WorkloadLock {
fn drop(&mut self) {
let _ = FileExt::unlock(&self.file);
}
}
#[cfg(windows)]
fn lock_target(path: &Path) -> std::path::PathBuf {
let mut os = path.as_os_str().to_os_string();
os.push(".lock");
std::path::PathBuf::from(os)
}
pub fn acquire(path: &Path) -> io::Result<WorkloadLock> {
#[cfg(windows)]
let file = {
if !path.exists() {
return Err(io::Error::new(
io::ErrorKind::NotFound,
format!(
"workload lock: cannot open '{}' for locking: not found",
path.display()
),
));
}
let lock_path = lock_target(path);
OpenOptions::new()
.read(true)
.write(true)
.create(true)
.open(&lock_path)
.map_err(|e| {
io::Error::new(
e.kind(),
format!(
"workload lock: cannot open '{}' for locking: {e}",
lock_path.display()
),
)
})?
};
#[cfg(not(windows))]
let file = OpenOptions::new().read(true).open(path).map_err(|e| {
io::Error::new(
e.kind(),
format!(
"workload lock: cannot open '{}' for locking: {e}",
path.display()
),
)
})?;
let deadline = Instant::now() + LOCK_DEADLINE;
loop {
match FileExt::try_lock_exclusive(&file) {
Ok(true) => return Ok(WorkloadLock { file }),
Ok(false) => {
if Instant::now() >= deadline {
return Err(io::Error::new(
io::ErrorKind::WouldBlock,
format!(
"workload lock on '{}': another nmbrs process holds the \
exclusive lock; waited {LOCK_DEADLINE:?} and gave up. \
Re-run after that process completes, or kill it if it's \
stuck.",
path.display(),
),
));
}
std::thread::sleep(POLL_INTERVAL);
}
Err(e) => {
return Err(io::Error::new(
e.kind(),
format!("workload lock on '{}': {e}", path.display()),
));
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::Duration;
fn touch(path: &Path) {
std::fs::write(path, b"# workload\n").unwrap();
}
#[test]
fn acquire_then_drop_round_trips() {
let dir = tempdir_relative("lock_round_trips");
let path = dir.join("w.yaml");
touch(&path);
let l1 = acquire(&path).expect("first lock");
drop(l1);
let _l2 = acquire(&path).expect("second lock");
}
#[test]
fn second_acquire_blocks_then_succeeds_when_first_releases() {
let dir = tempdir_relative("lock_block_release");
let path = dir.join("w.yaml");
touch(&path);
let l1 = acquire(&path).expect("first lock");
let path_clone = path.clone();
let waiter_done = Arc::new(AtomicBool::new(false));
let flag = waiter_done.clone();
let handle = thread::spawn(move || {
let _l2 = acquire(&path_clone).expect("second lock");
flag.store(true, Ordering::SeqCst);
});
thread::sleep(Duration::from_millis(100));
assert!(
!waiter_done.load(Ordering::SeqCst),
"second acquirer must still be waiting"
);
drop(l1);
handle.join().expect("waiter joined");
assert!(waiter_done.load(Ordering::SeqCst));
}
#[test]
fn deadline_exceeded_returns_wouldblock() {
let dir = tempdir_relative("lock_deadline");
let path = dir.join("w.yaml");
touch(&path);
let _l1 = acquire(&path).expect("first lock");
#[cfg(windows)]
let probe = super::lock_target(&path);
#[cfg(not(windows))]
let probe = path.clone();
let f = std::fs::File::open(&probe).unwrap();
match FileExt::try_lock_exclusive(&f) {
Ok(false) => {} other => panic!("expected try_lock_exclusive=Ok(false), got {other:?}"),
}
}
fn tempdir_relative(label: &str) -> std::path::PathBuf {
let p =
std::env::temp_dir().join(format!("nmbrs-edit-lock-{label}-{}", std::process::id()));
std::fs::create_dir_all(&p).unwrap();
p
}
}