taskfleet_core/lock.rs
1//! Per-run advisory `flock` primitive (design.md §4).
2//!
3//! The lock is also the source of a **compile-time witness**, [`LockedRun`],
4//! that the run's exclusive lock is held. The unlocked event-append entry points
5//! ([`crate::append_and_apply_unlocked`] and friends) take a `&LockedRun`
6//! parameter, so the type system — not a "the caller must hold the lock" doc
7//! comment — proves a writer holds the `flock` before it appends. Only the
8//! **exclusive** guard mints a witness ([`RunLock::witness`]); a shared
9//! (`LOCK_SH`) reader has no write capability and cannot produce one, because
10//! [`RunLock`] is a typestate generic ([`Exclusive`] vs [`Shared`]) and
11//! `witness` exists only on `RunLock<Exclusive>`.
12
13use std::fs::{File, OpenOptions};
14use std::io::ErrorKind;
15use std::marker::PhantomData;
16use std::path::Path;
17
18use fs4::FileExt;
19
20use crate::error::{Error, Result};
21use crate::paths::{reject_symlink, RunPaths};
22
23/// Typestate marker: the **exclusive** (`LOCK_EX`) lock. A `RunLock<Exclusive>`
24/// grants write access — it alone can mint a [`LockedRun`] witness via
25/// [`RunLock::witness`]. Uninhabited: it exists only as a type tag.
26pub enum Exclusive {}
27
28/// Typestate marker: the **shared** (`LOCK_SH`) lock. A `RunLock<Shared>` is
29/// read-only and cannot produce a [`LockedRun`] witness, so it can never be
30/// used to reach a write-side entry point. Uninhabited: only a type tag.
31pub enum Shared {}
32
33/// Compile-time proof that the holder is inside a critical section guarded by
34/// the run's **exclusive** `flock`.
35///
36/// The unlocked append primitives ([`append_and_apply_unlocked`],
37/// [`append_event_with_seq`], [`quarantine_corrupt_lines_unlocked`]) take a
38/// `&LockedRun` so they cannot be called without proof the lock is held —
39/// replacing the old "caller must already hold the `RunLock`" contract that was
40/// enforced only by a doc comment. Obtain one from
41/// [`RunLock::with_lock`] (which threads it into the closure) or
42/// [`RunLock::witness`] (for a manually-held [`RunLock::acquire`] guard).
43///
44/// The witness is a zero-sized borrow of the guard: its lifetime `'a` ties it to
45/// the [`RunLock`] it was minted from, so it cannot outlive the lock. It is
46/// deliberately **non-`Send` and non-`Sync`** (the `PhantomData<*const …>`) — a
47/// proof that *this thread* holds the `flock` must not cross a task or thread
48/// boundary, where the lock would no longer apply.
49///
50/// [`append_and_apply_unlocked`]: crate::append_and_apply_unlocked
51/// [`append_event_with_seq`]: crate::events
52/// [`quarantine_corrupt_lines_unlocked`]: crate::quarantine_corrupt_lines_unlocked
53pub struct LockedRun<'a> {
54 // `*const` makes the witness `!Send + !Sync`; the `&'a ()` ties it to the
55 // borrowed guard's lifetime. Zero-sized — it carries no data, only proof.
56 _lock: PhantomData<*const &'a ()>,
57}
58
59/// RAII guard holding the run's `flock`.
60///
61/// Released on drop. Acquired exclusively ([`RunLock::acquire`]) and held across
62/// all writes for a single logical mutation, or shared ([`RunLock::acquire_shared`])
63/// for the duration of a reader's multi-file scan so a concurrent reducer cannot
64/// leave the reader with a half-updated projection set. A shared guard may hold
65/// no lock at all when the run has no `.lock` file yet — see
66/// [`RunLock::acquire_shared`].
67///
68/// The `Mode` typestate ([`Exclusive`] / [`Shared`]) records which kind of lock
69/// is held: only `RunLock<Exclusive>` exposes [`RunLock::witness`], so a shared
70/// reader can never forge the write capability a [`LockedRun`] represents.
71pub struct RunLock<Mode = Exclusive> {
72 file: Option<File>,
73 _mode: PhantomData<Mode>,
74}
75
76impl RunLock<Exclusive> {
77 /// Acquire the exclusive lock on `<run-dir>/.lock`, creating the file if
78 /// needed. Blocks until the lock is available.
79 ///
80 /// Best-effort symlink containment: a `.lock` that is a symlink is refused
81 /// ([`Error::SymlinkStateFile`]) so `flock` cannot be taken on a file
82 /// outside the run tree, which would silently break mutual exclusion. This
83 /// guards the lock file's own final component; a symlinked *run root* is
84 /// caught downstream when the held critical section opens `events.jsonl` /
85 /// the projections (both re-guard the root before writing). See
86 /// [`reject_symlink`](crate::paths) for the check-then-open TOCTOU caveat.
87 pub fn acquire(lock_path: &Path) -> Result<Self> {
88 // Test-only spy: count this acquisition so a test can assert a
89 // multi-write transaction (e.g. `cancel_run`) takes the lock exactly
90 // once, not once per appended event.
91 #[cfg(test)]
92 ACQUIRE_COUNT.with(|c| c.set(c.get() + 1));
93 if let Some(p) = lock_path.parent() {
94 std::fs::create_dir_all(p).map_err(|e| Error::io(p, e))?;
95 }
96 reject_symlink(lock_path, || Error::SymlinkStateFile {
97 name: "lock",
98 path: lock_path.to_path_buf(),
99 })?;
100 let mut opts = OpenOptions::new();
101 opts.create(true).read(true).write(true).truncate(false);
102 // `O_NOFOLLOW`: refuse to take `flock` through a symlinked `.lock`, the
103 // file-level backstop to the `reject_symlink` check above.
104 crate::paths::nofollow(&mut opts);
105 let file = opts.open(lock_path).map_err(|e| Error::io(lock_path, e))?;
106 // Fully-qualified to call fs4's trait method, not `std::fs::File::lock`
107 // (an inherent method stable since 1.89 that would otherwise shadow it
108 // on newer toolchains). fs4 renamed `fs2`'s `lock_exclusive` to `lock`
109 // to mirror std.
110 <File as FileExt>::lock(&file).map_err(|e| Error::io(lock_path, e))?;
111 Ok(Self {
112 file: Some(file),
113 _mode: PhantomData,
114 })
115 }
116
117 /// Acquire the exclusive lock on an existing lock file without creating
118 /// either the lock file or its parent run directory.
119 ///
120 /// This is for mutation endpoints that must reject a missing/deleted run
121 /// without resurrecting it as an empty directory. A canonical run has a
122 /// `.lock` because its event creation path acquired the ordinary exclusive
123 /// lock before publishing projections. Like [`RunLock::acquire`], this
124 /// rejects symlinks and blocks until the exclusive lock is available.
125 pub fn acquire_existing(lock_path: &Path) -> Result<Self> {
126 let file = open_existing_lock(lock_path)?;
127 <File as FileExt>::lock(&file).map_err(|e| Error::io(lock_path, e))?;
128 Ok(Self {
129 file: Some(file),
130 _mode: PhantomData,
131 })
132 }
133
134 /// Try to acquire the exclusive lock on an existing lock file without
135 /// creating state or waiting. Returns `Ok(None)` when another process holds
136 /// the lock. Migration preflight uses this bounded probe so a busy run is a
137 /// prompt refusal rather than an unbounded wait.
138 pub fn try_acquire_existing(lock_path: &Path) -> Result<Option<Self>> {
139 let file = open_existing_lock(lock_path)?;
140 match <File as FileExt>::try_lock(&file) {
141 Ok(()) => Ok(Some(Self {
142 file: Some(file),
143 _mode: PhantomData,
144 })),
145 Err(fs4::TryLockError::WouldBlock) => Ok(None),
146 Err(fs4::TryLockError::Error(error)) => Err(Error::io(lock_path, error)),
147 }
148 }
149
150 /// Mint a [`LockedRun`] witness proving this exclusive guard holds the
151 /// run's `flock`. The witness borrows `self`, so the borrow checker forbids
152 /// it from outliving the guard (and thus the lock). Use this for the
153 /// manually-held [`RunLock::acquire`] pattern — when the locked body needs
154 /// control flow ([`with_lock`](RunLock::with_lock)'s closure cannot express)
155 /// — then pass `&witness` to the unlocked append entry points.
156 // `&self` is load-bearing despite the body not reading it: it borrows the
157 // guard so the returned `LockedRun<'_>` is lifetime-tied to the held lock
158 // and cannot outlive it. That is the whole point — not an accidental unused
159 // receiver, so this method must stay a method, never an associated fn.
160 #[allow(clippy::unused_self)]
161 pub fn witness(&self) -> LockedRun<'_> {
162 LockedRun { _lock: PhantomData }
163 }
164
165 /// Convenience: run `f` with the exclusive lock held, passing it a
166 /// [`LockedRun`] witness it can thread into the unlocked append entry
167 /// points, releasing the lock afterwards.
168 pub fn with_lock<R>(paths: &RunPaths, f: impl FnOnce(&LockedRun) -> Result<R>) -> Result<R> {
169 let guard = Self::acquire(&paths.lock())?;
170 let r = f(&guard.witness());
171 drop(guard);
172 r
173 }
174}
175
176fn open_existing_lock(lock_path: &Path) -> Result<File> {
177 reject_symlink(lock_path, || Error::SymlinkStateFile {
178 name: "lock",
179 path: lock_path.to_path_buf(),
180 })?;
181 let mut opts = OpenOptions::new();
182 opts.read(true).write(true).truncate(false);
183 crate::paths::nofollow(&mut opts);
184 opts.open(lock_path).map_err(|e| Error::io(lock_path, e))
185}
186
187impl RunLock<Shared> {
188 /// Non-blocking read probe on an existing lock. Unlike the ordinary read
189 /// path, absence is an error: upgrade preflight cannot treat a missing
190 /// writer witness in a published run as proof of quiescence.
191 pub fn try_acquire_shared_existing(lock_path: &Path) -> Result<Option<Self>> {
192 let file = open_existing_lock(lock_path)?;
193 match <File as FileExt>::try_lock_shared(&file) {
194 Ok(()) => Ok(Some(Self {
195 file: Some(file),
196 _mode: PhantomData,
197 })),
198 Err(fs4::TryLockError::WouldBlock) => Ok(None),
199 Err(fs4::TryLockError::Error(error)) => Err(Error::io(lock_path, error)),
200 }
201 }
202
203 /// Acquire a **shared** (`LOCK_SH`) lock on an existing `<run-dir>/.lock`,
204 /// blocking until no writer holds the exclusive lock. The mirror of
205 /// [`RunLock::acquire`] for the read side: many readers may hold the shared
206 /// lock at once, but the exclusive lock a reducer takes excludes them all,
207 /// so a multi-file read taken under this lock never observes the torn state
208 /// a mid-flight reducer would otherwise expose (design.md §4).
209 ///
210 /// Unlike [`RunLock::acquire`], this **never creates** the run directory or
211 /// the lock file — a reader must not bring run-tree state into existence. A
212 /// missing `.lock` means no writer has ever locked this run, so a lock-free
213 /// read is already coherent: the returned guard then holds nothing and drops
214 /// to a no-op. The same symlink containment as [`RunLock::acquire`] applies;
215 /// the file is opened read-only with `O_NOFOLLOW`.
216 ///
217 /// Nesting is safe: a shared lock is compatible with other shared locks, so
218 /// a read path that calls another read helper (each on its own descriptor)
219 /// cannot deadlock against itself.
220 pub fn acquire_shared(lock_path: &Path) -> Result<Self> {
221 reject_symlink(lock_path, || Error::SymlinkStateFile {
222 name: "lock",
223 path: lock_path.to_path_buf(),
224 })?;
225 let mut opts = OpenOptions::new();
226 // Read-only, no `create`: a reader never authors the lock file or its
227 // parent run directory.
228 opts.read(true);
229 // `O_NOFOLLOW`: refuse to take `flock` through a symlinked `.lock`, the
230 // file-level backstop to the `reject_symlink` check above.
231 crate::paths::nofollow(&mut opts);
232 let file = match opts.open(lock_path) {
233 Ok(f) => f,
234 Err(e) if e.kind() == ErrorKind::NotFound => {
235 // No `.lock` (or no run dir): nothing to serialize against.
236 return Ok(Self {
237 file: None,
238 _mode: PhantomData,
239 });
240 }
241 Err(e) => return Err(Error::io(lock_path, e)),
242 };
243 // Fully-qualified to call fs4's trait method rather than the inherent
244 // `std::fs::File::lock_shared` (stable 1.89) that would shadow it on
245 // newer toolchains — same reasoning as the exclusive `lock` above.
246 <File as FileExt>::lock_shared(&file).map_err(|e| Error::io(lock_path, e))?;
247 Ok(Self {
248 file: Some(file),
249 _mode: PhantomData,
250 })
251 }
252
253 /// Convenience: run `f` with the shared lock held, releasing afterwards.
254 /// The read-side counterpart to [`RunLock::with_lock`] — but the closure
255 /// gets **no** [`LockedRun`] witness: a shared reader has no write
256 /// capability, so it can never reach a write-side entry point. (It also
257 /// takes the lock path directly rather than a [`RunPaths`], since a reader
258 /// may run before the run dir is fully materialized.)
259 pub fn with_shared_lock<T>(lock_path: &Path, f: impl FnOnce() -> Result<T>) -> Result<T> {
260 let guard = Self::acquire_shared(lock_path)?;
261 let r = f();
262 drop(guard);
263 r
264 }
265}
266
267impl<Mode> Drop for RunLock<Mode> {
268 fn drop(&mut self) {
269 if let Some(f) = self.file.take() {
270 // Best-effort unlock — kernel releases on file close anyway.
271 // Use the fs4 trait method explicitly to avoid clashing with
272 // `std::fs::File::unlock` (available on supported stable toolchains).
273 let _ = <File as FileExt>::unlock(&f);
274 }
275 }
276}
277
278#[cfg(test)]
279thread_local! {
280 /// Test-only spy counter for [`RunLock::acquire`] calls on the current
281 /// thread. `cargo test` runs each test on its own thread and `cancel_run`
282 /// does all its work synchronously on the calling thread, so a test can
283 /// reset this and assert the exact number of lock acquisitions a
284 /// transaction performed without cross-test interference.
285 pub(crate) static ACQUIRE_COUNT: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
286}
287
288#[cfg(test)]
289mod tests {
290 use super::*;
291 use tempfile::TempDir;
292
293 #[test]
294 fn acquire_succeeds_on_a_regular_lock_file() {
295 let tmp = TempDir::new().unwrap();
296 let lock = tmp.path().join(".lock");
297 // First acquire creates the file; a second acquire after drop succeeds.
298 drop(RunLock::acquire(&lock).unwrap());
299 assert!(RunLock::acquire(&lock).is_ok());
300 }
301
302 /// Build a valid `RunPaths` rooted at a fresh temp dir for the witness tests.
303 fn fresh_paths(tmp: &TempDir) -> RunPaths {
304 let run_id = "01jxsnap000000000000000000";
305 let dir = tmp.path().join(run_id);
306 std::fs::create_dir_all(&dir).unwrap();
307 RunPaths::new(dir, run_id).unwrap()
308 }
309
310 /// `with_lock` runs the closure under the exclusive lock, hands it a
311 /// [`LockedRun`] witness, and returns the closure's value.
312 #[test]
313 fn with_lock_passes_a_witness_and_returns_the_closure_value() {
314 let tmp = TempDir::new().unwrap();
315 let paths = fresh_paths(&tmp);
316 let got = RunLock::with_lock(&paths, |_witness: &LockedRun| Ok(7u8)).unwrap();
317 assert_eq!(got, 7);
318 }
319
320 #[test]
321 fn acquire_existing_never_creates_missing_state() {
322 let tmp = TempDir::new().unwrap();
323 let run = tmp.path().join("missing-run");
324 let lock = run.join(".lock");
325 assert!(RunLock::acquire_existing(&lock).is_err());
326 assert!(!run.exists(), "no-create acquire must not author a run dir");
327
328 std::fs::create_dir_all(&run).unwrap();
329 std::fs::write(&lock, []).unwrap();
330 assert!(RunLock::acquire_existing(&lock).is_ok());
331 }
332
333 #[test]
334 fn nonblocking_shared_probe_accepts_readers_but_refuses_writer_or_missing_file() {
335 let tmp = TempDir::new().unwrap();
336 let lock = tmp.path().join(".lock");
337 assert!(RunLock::<Shared>::try_acquire_shared_existing(&lock).is_err());
338 let first = RunLock::<Shared>::acquire_shared(&lock).unwrap();
339 assert!(first.file.is_none());
340 drop(RunLock::acquire(&lock).unwrap());
341 let first = RunLock::<Shared>::try_acquire_shared_existing(&lock)
342 .unwrap()
343 .unwrap();
344 let second = RunLock::<Shared>::try_acquire_shared_existing(&lock).unwrap();
345 assert!(second.is_some());
346 drop(first);
347 drop(second);
348 let writer = RunLock::acquire(&lock).unwrap();
349 assert!(RunLock::<Shared>::try_acquire_shared_existing(&lock)
350 .unwrap()
351 .is_none());
352 drop(writer);
353 }
354
355 /// A manually-held exclusive guard ([`RunLock::acquire`]) can mint a witness
356 /// — the escape hatch for lock-held bodies that `with_lock`'s closure shape
357 /// cannot express.
358 #[test]
359 fn manually_acquired_exclusive_guard_mints_a_witness() {
360 let tmp = TempDir::new().unwrap();
361 let lock = tmp.path().join(".lock");
362 let guard = RunLock::acquire(&lock).unwrap();
363 // The witness borrows the guard, so it cannot outlive the held lock.
364 let _witness: LockedRun<'_> = guard.witness();
365 }
366
367 /// A shared lock on a run with no `.lock` file is a no-op guard, not an
368 /// error: a reader must not create run-tree state, and a missing lock file
369 /// means no writer can race the read anyway.
370 #[test]
371 fn acquire_shared_on_missing_lock_file_is_a_noop_guard() {
372 let tmp = TempDir::new().unwrap();
373 let lock = tmp.path().join(".lock");
374 let guard = RunLock::acquire_shared(&lock).expect("missing lock file is fine");
375 assert!(guard.file.is_none(), "no lock file ⇒ guard holds nothing");
376 // The read must not have created the lock file (or its parent run dir).
377 assert!(!lock.exists(), "a reader must never author the lock file");
378 }
379
380 /// A writer holding the exclusive lock blocks a shared reader until release.
381 #[test]
382 fn exclusive_writer_blocks_shared_reader_until_release() {
383 use std::sync::mpsc;
384 use std::thread;
385 use std::time::Duration;
386
387 let tmp = TempDir::new().unwrap();
388 let lock = tmp.path().join(".lock");
389 // The writer's exclusive acquire creates the lock file the reader opens.
390 let writer = RunLock::acquire(&lock).unwrap();
391
392 let (tx, rx) = mpsc::channel();
393 let lock2 = lock.clone();
394 let reader = thread::spawn(move || {
395 // Blocks until the exclusive lock is released.
396 let _g = RunLock::acquire_shared(&lock2).unwrap();
397 tx.send(()).unwrap();
398 });
399
400 // While the exclusive lock is held, the reader cannot proceed.
401 assert!(
402 rx.recv_timeout(Duration::from_millis(250)).is_err(),
403 "shared reader must block while the exclusive lock is held"
404 );
405 drop(writer);
406 // Once released, the reader acquires the shared lock and reports.
407 assert!(
408 rx.recv_timeout(Duration::from_secs(5)).is_ok(),
409 "shared reader must proceed after the exclusive lock is released"
410 );
411 reader.join().unwrap();
412 }
413
414 /// Two shared readers hold the lock concurrently — neither blocks the other.
415 #[test]
416 fn two_shared_readers_proceed_concurrently() {
417 use std::sync::mpsc;
418 use std::thread;
419 use std::time::Duration;
420
421 let tmp = TempDir::new().unwrap();
422 let lock = tmp.path().join(".lock");
423 // Create the lock file (writer makes it, then releases).
424 drop(RunLock::acquire(&lock).unwrap());
425
426 // First reader takes and holds the shared lock.
427 let r1 = RunLock::acquire_shared(&lock).unwrap();
428 assert!(
429 r1.file.is_some(),
430 "lock file exists ⇒ real shared lock held"
431 );
432
433 // A second reader must acquire it without blocking on the first.
434 let (tx, rx) = mpsc::channel();
435 let lock2 = lock.clone();
436 let r2 = thread::spawn(move || {
437 let _g = RunLock::acquire_shared(&lock2).unwrap();
438 tx.send(()).unwrap();
439 });
440 assert!(
441 rx.recv_timeout(Duration::from_secs(5)).is_ok(),
442 "a second shared reader must not block on the first"
443 );
444 r2.join().unwrap();
445 }
446
447 /// Nested `with_shared_lock` calls (each on its own descriptor) must not
448 /// deadlock — shared locks are compatible with one another.
449 #[test]
450 fn nested_shared_locks_do_not_deadlock() {
451 let tmp = TempDir::new().unwrap();
452 let lock = tmp.path().join(".lock");
453 drop(RunLock::acquire(&lock).unwrap());
454 let r = RunLock::with_shared_lock(&lock, || RunLock::with_shared_lock(&lock, || Ok(42)))
455 .unwrap();
456 assert_eq!(r, 42);
457 }
458
459 #[test]
460 fn nonblocking_existing_probe_reports_contention_without_waiting() {
461 let tmp = TempDir::new().unwrap();
462 let lock = tmp.path().join(".lock");
463 let held = RunLock::acquire(&lock).unwrap();
464 assert!(RunLock::try_acquire_existing(&lock).unwrap().is_none());
465 drop(held);
466 assert!(RunLock::try_acquire_existing(&lock).unwrap().is_some());
467 }
468
469 #[cfg(unix)]
470 #[test]
471 fn acquire_rejects_a_symlinked_lock_file() {
472 // A symlinked `.lock` would take `flock` on a file outside the run,
473 // silently breaking mutual exclusion — refuse it.
474 use std::os::unix::fs::symlink;
475 let tmp = TempDir::new().unwrap();
476 let target = tmp.path().join("outside.lock");
477 let lock = tmp.path().join(".lock");
478 symlink(&target, &lock).unwrap();
479 assert!(matches!(
480 RunLock::acquire(&lock),
481 Err(Error::SymlinkStateFile { name: "lock", .. })
482 ));
483 // The forged lock never touched the symlink target.
484 assert!(!target.exists());
485 }
486}