Skip to main content

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}