clawhdf5: a dropped FileEditor releases its lock at once
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]>
This commit is contained in:
@@ -20,7 +20,7 @@ const LOCK_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(1
|
||||
/// never leaves a stale lock behind; the empty lock file itself is harmless).
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct StoreLock {
|
||||
_file: File,
|
||||
file: File,
|
||||
}
|
||||
|
||||
impl StoreLock {
|
||||
@@ -42,7 +42,7 @@ impl StoreLock {
|
||||
let mut attempts_left = LOCK_RETRIES;
|
||||
loop {
|
||||
match file.try_lock() {
|
||||
Ok(()) => return Ok(Self { _file: file }),
|
||||
Ok(()) => return Ok(Self { file }),
|
||||
Err(TryLockError::WouldBlock) if attempts_left > 0 => {
|
||||
attempts_left -= 1;
|
||||
std::thread::sleep(LOCK_RETRY_DELAY);
|
||||
@@ -60,6 +60,16 @@ impl StoreLock {
|
||||
}
|
||||
}
|
||||
|
||||
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::*;
|
||||
|
||||
@@ -61,6 +61,9 @@ const MSG_FLAG_DONTSHARE: u8 = 0x04;
|
||||
/// Opening takes an exclusive advisory lock on the file (`flock`, the lock
|
||||
/// libhdf5 itself takes when file locking is on), so a second editor, or
|
||||
/// h5py opening the file for writing, fails until the editor is dropped.
|
||||
/// The drop unlocks the file before closing it, so the lock is gone as soon
|
||||
/// as the drop returns, even in a program whose other threads spawn
|
||||
/// processes (a forked child shares the locked descriptor until it execs).
|
||||
/// Readers that do not lock ([`File`]) can still open it, but see a file
|
||||
/// that may be mid-update.
|
||||
///
|
||||
@@ -114,6 +117,19 @@ pub struct FileEditor {
|
||||
free: FreeList,
|
||||
}
|
||||
|
||||
impl Drop for FileEditor {
|
||||
/// Releases the lock explicitly before the file is closed. A `flock`
|
||||
/// belongs to the open file description and lasts until every
|
||||
/// descriptor of it is closed; a process another thread forks (any
|
||||
/// `std::process::Command`) inherits the descriptor and keeps it until
|
||||
/// it execs, so closing alone could leave the file locked for a moment
|
||||
/// after the drop. Unlocking through our descriptor releases the lock
|
||||
/// for all of them.
|
||||
fn drop(&mut self) {
|
||||
let _ = self.file.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
/// Where a layout message keeps the fields an edit may change (offsets in
|
||||
/// the message body).
|
||||
#[derive(Debug, Default, Clone, Copy)]
|
||||
|
||||
@@ -218,3 +218,49 @@ fn edits_go_to_the_file_held_not_the_path() {
|
||||
assert_eq!(f.dataset("big").unwrap().read_f64().unwrap(), [1.5; 5000]);
|
||||
assert!(matches!(f.root().attr("note").unwrap(), Some(AttrValue::F64Array(v)) if v == vals));
|
||||
}
|
||||
|
||||
/// Dropping an editor releases its lock at once, even while other threads
|
||||
/// keep spawning processes: a child inherits the locked descriptor between
|
||||
/// `fork` and `exec` (close-on-exec closes it only at `exec`), and a
|
||||
/// `flock` lasts while any descriptor of the open file does, so without the
|
||||
/// explicit unlock in `Drop` a reopen right after the drop was sometimes
|
||||
/// refused with `Error::Locked`.
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn drop_releases_the_lock_while_other_threads_spawn_processes() {
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = sample(dir.path());
|
||||
let stop = Arc::new(AtomicBool::new(false));
|
||||
let spawners: Vec<_> = (0..4)
|
||||
.map(|_| {
|
||||
let stop = Arc::clone(&stop);
|
||||
std::thread::spawn(move || {
|
||||
let mut n = 0u64;
|
||||
while !stop.load(Ordering::Relaxed) {
|
||||
let _ = std::process::Command::new("true").status();
|
||||
n += 1;
|
||||
}
|
||||
n
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
let rounds = 2000;
|
||||
let mut refused = 0;
|
||||
// Every open but the first follows the drop of the previous editor.
|
||||
for _ in 0..rounds {
|
||||
match FileEditor::open(&path) {
|
||||
Ok(ed) => drop(ed),
|
||||
Err(Error::Locked(_)) => refused += 1,
|
||||
Err(e) => panic!("{e}"),
|
||||
}
|
||||
}
|
||||
stop.store(true, Ordering::Relaxed);
|
||||
let spawned: u64 = spawners.into_iter().map(|t| t.join().unwrap()).sum();
|
||||
assert!(spawned > 0);
|
||||
assert_eq!(
|
||||
refused, 0,
|
||||
"{refused} of {rounds} reopens refused ({spawned} processes spawned)"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user