FileEditor's flock belongs to the open file description. When another thread forks to spawn a process, the child shares the locked descriptor until it execs, so a reopen right after the drop could be refused with Error::Locked (a one-off failure of edit_interop::editor_locks_the_file in a parallel test run). Drop now unlocks before closing, which releases the lock for every descriptor sharing it. Reproducer edit_tests::drop_releases_the_lock_while_other_threads_spawn_processes (4 threads running `true`, 2000 open/drop rounds): 1483 of 2000 reopens refused before, 0 in 30 runs after (tank). An OFD lock would not help: it is inherited across fork the same way and does not conflict with libhdf5's flock. The agent store's lock file unlocks on drop too. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
90 lines
3.2 KiB
Rust
90 lines
3.2 KiB
Rust
//! Single-writer guard for a memory store.
|
|
//!
|
|
//! `HDF5Memory` keeps the whole store in memory and rewrites the `.h5` file at
|
|
//! every checkpoint, so two handles on one store (two processes, or two opens
|
|
//! in one process) silently destroy each other's data: whoever checkpoints
|
|
//! last wins, and both append to the same WAL with independent CRC chains.
|
|
//! The lock turns that into an immediate, explicit error.
|
|
|
|
use std::fs::{File, OpenOptions, TryLockError};
|
|
use std::path::{Path, PathBuf};
|
|
|
|
use crate::MemoryError;
|
|
|
|
const LOCK_RETRIES: u32 = 25;
|
|
const LOCK_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(10);
|
|
|
|
/// An exclusive advisory lock on `<store>.h5.lock`, held for the lifetime of
|
|
/// the owning `HDF5Memory` and released when it is dropped (or when the
|
|
/// process dies — the OS drops the lock with the file descriptor, so a crash
|
|
/// never leaves a stale lock behind; the empty lock file itself is harmless).
|
|
#[derive(Debug)]
|
|
pub(crate) struct StoreLock {
|
|
file: File,
|
|
}
|
|
|
|
impl StoreLock {
|
|
pub(crate) fn lock_path(store: &Path) -> PathBuf {
|
|
store.with_extension("h5.lock")
|
|
}
|
|
|
|
pub(crate) fn acquire(store: &Path) -> Result<Self, MemoryError> {
|
|
let path = Self::lock_path(store);
|
|
let file = OpenOptions::new()
|
|
.create(true)
|
|
.truncate(false)
|
|
.write(true)
|
|
.open(&path)?;
|
|
// A previous owner may be mid-teardown (e.g. an `AsyncHDF5Memory`
|
|
// dropped without `shutdown()`: its background task releases the
|
|
// store a moment later), so give the lock a short, bounded grace
|
|
// period before reporting a genuine second writer.
|
|
let mut attempts_left = LOCK_RETRIES;
|
|
loop {
|
|
match file.try_lock() {
|
|
Ok(()) => return Ok(Self { file }),
|
|
Err(TryLockError::WouldBlock) if attempts_left > 0 => {
|
|
attempts_left -= 1;
|
|
std::thread::sleep(LOCK_RETRY_DELAY);
|
|
}
|
|
Err(TryLockError::WouldBlock) => {
|
|
return Err(MemoryError::Locked(format!(
|
|
"{} is already open in this or another process (lock file {})",
|
|
store.display(),
|
|
path.display()
|
|
)));
|
|
}
|
|
Err(TryLockError::Error(e)) => return Err(MemoryError::Io(e)),
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Drop for StoreLock {
|
|
/// Unlocks before the file is closed: a process another thread forks
|
|
/// inherits the descriptor until it execs, and a `flock` lasts while any
|
|
/// descriptor of the open file does, so closing alone could keep the
|
|
/// store locked for a moment after the drop (see `FileEditor`'s `Drop`).
|
|
fn drop(&mut self) {
|
|
let _ = self.file.unlock();
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn second_acquire_fails_until_first_is_dropped() {
|
|
let dir = tempfile::TempDir::new().unwrap();
|
|
let store = dir.path().join("s.h5");
|
|
let first = StoreLock::acquire(&store).unwrap();
|
|
assert!(matches!(
|
|
StoreLock::acquire(&store),
|
|
Err(MemoryError::Locked(_))
|
|
));
|
|
drop(first);
|
|
StoreLock::acquire(&store).unwrap();
|
|
}
|
|
}
|