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    /// Acquire a **shared** (`LOCK_SH`) lock on an existing `<run-dir>/.lock`,
189    /// blocking until no writer holds the exclusive lock. The mirror of
190    /// [`RunLock::acquire`] for the read side: many readers may hold the shared
191    /// lock at once, but the exclusive lock a reducer takes excludes them all,
192    /// so a multi-file read taken under this lock never observes the torn state
193    /// a mid-flight reducer would otherwise expose (design.md §4).
194    ///
195    /// Unlike [`RunLock::acquire`], this **never creates** the run directory or
196    /// the lock file — a reader must not bring run-tree state into existence. A
197    /// missing `.lock` means no writer has ever locked this run, so a lock-free
198    /// read is already coherent: the returned guard then holds nothing and drops
199    /// to a no-op. The same symlink containment as [`RunLock::acquire`] applies;
200    /// the file is opened read-only with `O_NOFOLLOW`.
201    ///
202    /// Nesting is safe: a shared lock is compatible with other shared locks, so
203    /// a read path that calls another read helper (each on its own descriptor)
204    /// cannot deadlock against itself.
205    pub fn acquire_shared(lock_path: &Path) -> Result<Self> {
206        reject_symlink(lock_path, || Error::SymlinkStateFile {
207            name: "lock",
208            path: lock_path.to_path_buf(),
209        })?;
210        let mut opts = OpenOptions::new();
211        // Read-only, no `create`: a reader never authors the lock file or its
212        // parent run directory.
213        opts.read(true);
214        // `O_NOFOLLOW`: refuse to take `flock` through a symlinked `.lock`, the
215        // file-level backstop to the `reject_symlink` check above.
216        crate::paths::nofollow(&mut opts);
217        let file = match opts.open(lock_path) {
218            Ok(f) => f,
219            Err(e) if e.kind() == ErrorKind::NotFound => {
220                // No `.lock` (or no run dir): nothing to serialize against.
221                return Ok(Self {
222                    file: None,
223                    _mode: PhantomData,
224                });
225            }
226            Err(e) => return Err(Error::io(lock_path, e)),
227        };
228        // Fully-qualified to call fs4's trait method rather than the inherent
229        // `std::fs::File::lock_shared` (stable 1.89) that would shadow it on
230        // newer toolchains — same reasoning as the exclusive `lock` above.
231        <File as FileExt>::lock_shared(&file).map_err(|e| Error::io(lock_path, e))?;
232        Ok(Self {
233            file: Some(file),
234            _mode: PhantomData,
235        })
236    }
237
238    /// Convenience: run `f` with the shared lock held, releasing afterwards.
239    /// The read-side counterpart to [`RunLock::with_lock`] — but the closure
240    /// gets **no** [`LockedRun`] witness: a shared reader has no write
241    /// capability, so it can never reach a write-side entry point. (It also
242    /// takes the lock path directly rather than a [`RunPaths`], since a reader
243    /// may run before the run dir is fully materialized.)
244    pub fn with_shared_lock<T>(lock_path: &Path, f: impl FnOnce() -> Result<T>) -> Result<T> {
245        let guard = Self::acquire_shared(lock_path)?;
246        let r = f();
247        drop(guard);
248        r
249    }
250}
251
252impl<Mode> Drop for RunLock<Mode> {
253    fn drop(&mut self) {
254        if let Some(f) = self.file.take() {
255            // Best-effort unlock — kernel releases on file close anyway.
256            // Use the fs4 trait method explicitly to avoid clashing with
257            // `std::fs::File::unlock` (stable since 1.89, above our MSRV).
258            let _ = <File as FileExt>::unlock(&f);
259        }
260    }
261}
262
263#[cfg(test)]
264thread_local! {
265    /// Test-only spy counter for [`RunLock::acquire`] calls on the current
266    /// thread. `cargo test` runs each test on its own thread and `cancel_run`
267    /// does all its work synchronously on the calling thread, so a test can
268    /// reset this and assert the exact number of lock acquisitions a
269    /// transaction performed without cross-test interference.
270    pub(crate) static ACQUIRE_COUNT: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
271}
272
273#[cfg(test)]
274mod tests {
275    use super::*;
276    use tempfile::TempDir;
277
278    #[test]
279    fn acquire_succeeds_on_a_regular_lock_file() {
280        let tmp = TempDir::new().unwrap();
281        let lock = tmp.path().join(".lock");
282        // First acquire creates the file; a second acquire after drop succeeds.
283        drop(RunLock::acquire(&lock).unwrap());
284        assert!(RunLock::acquire(&lock).is_ok());
285    }
286
287    /// Build a valid `RunPaths` rooted at a fresh temp dir for the witness tests.
288    fn fresh_paths(tmp: &TempDir) -> RunPaths {
289        let run_id = "01jxsnap000000000000000000";
290        let dir = tmp.path().join(run_id);
291        std::fs::create_dir_all(&dir).unwrap();
292        RunPaths::new(dir, run_id).unwrap()
293    }
294
295    /// `with_lock` runs the closure under the exclusive lock, hands it a
296    /// [`LockedRun`] witness, and returns the closure's value.
297    #[test]
298    fn with_lock_passes_a_witness_and_returns_the_closure_value() {
299        let tmp = TempDir::new().unwrap();
300        let paths = fresh_paths(&tmp);
301        let got = RunLock::with_lock(&paths, |_witness: &LockedRun| Ok(7u8)).unwrap();
302        assert_eq!(got, 7);
303    }
304
305    #[test]
306    fn acquire_existing_never_creates_missing_state() {
307        let tmp = TempDir::new().unwrap();
308        let run = tmp.path().join("missing-run");
309        let lock = run.join(".lock");
310        assert!(RunLock::acquire_existing(&lock).is_err());
311        assert!(!run.exists(), "no-create acquire must not author a run dir");
312
313        std::fs::create_dir_all(&run).unwrap();
314        std::fs::write(&lock, []).unwrap();
315        assert!(RunLock::acquire_existing(&lock).is_ok());
316    }
317
318    /// A manually-held exclusive guard ([`RunLock::acquire`]) can mint a witness
319    /// — the escape hatch for lock-held bodies that `with_lock`'s closure shape
320    /// cannot express.
321    #[test]
322    fn manually_acquired_exclusive_guard_mints_a_witness() {
323        let tmp = TempDir::new().unwrap();
324        let lock = tmp.path().join(".lock");
325        let guard = RunLock::acquire(&lock).unwrap();
326        // The witness borrows the guard, so it cannot outlive the held lock.
327        let _witness: LockedRun<'_> = guard.witness();
328    }
329
330    /// A shared lock on a run with no `.lock` file is a no-op guard, not an
331    /// error: a reader must not create run-tree state, and a missing lock file
332    /// means no writer can race the read anyway.
333    #[test]
334    fn acquire_shared_on_missing_lock_file_is_a_noop_guard() {
335        let tmp = TempDir::new().unwrap();
336        let lock = tmp.path().join(".lock");
337        let guard = RunLock::acquire_shared(&lock).expect("missing lock file is fine");
338        assert!(guard.file.is_none(), "no lock file ⇒ guard holds nothing");
339        // The read must not have created the lock file (or its parent run dir).
340        assert!(!lock.exists(), "a reader must never author the lock file");
341    }
342
343    /// A writer holding the exclusive lock blocks a shared reader until release.
344    #[test]
345    fn exclusive_writer_blocks_shared_reader_until_release() {
346        use std::sync::mpsc;
347        use std::thread;
348        use std::time::Duration;
349
350        let tmp = TempDir::new().unwrap();
351        let lock = tmp.path().join(".lock");
352        // The writer's exclusive acquire creates the lock file the reader opens.
353        let writer = RunLock::acquire(&lock).unwrap();
354
355        let (tx, rx) = mpsc::channel();
356        let lock2 = lock.clone();
357        let reader = thread::spawn(move || {
358            // Blocks until the exclusive lock is released.
359            let _g = RunLock::acquire_shared(&lock2).unwrap();
360            tx.send(()).unwrap();
361        });
362
363        // While the exclusive lock is held, the reader cannot proceed.
364        assert!(
365            rx.recv_timeout(Duration::from_millis(250)).is_err(),
366            "shared reader must block while the exclusive lock is held"
367        );
368        drop(writer);
369        // Once released, the reader acquires the shared lock and reports.
370        assert!(
371            rx.recv_timeout(Duration::from_secs(5)).is_ok(),
372            "shared reader must proceed after the exclusive lock is released"
373        );
374        reader.join().unwrap();
375    }
376
377    /// Two shared readers hold the lock concurrently — neither blocks the other.
378    #[test]
379    fn two_shared_readers_proceed_concurrently() {
380        use std::sync::mpsc;
381        use std::thread;
382        use std::time::Duration;
383
384        let tmp = TempDir::new().unwrap();
385        let lock = tmp.path().join(".lock");
386        // Create the lock file (writer makes it, then releases).
387        drop(RunLock::acquire(&lock).unwrap());
388
389        // First reader takes and holds the shared lock.
390        let r1 = RunLock::acquire_shared(&lock).unwrap();
391        assert!(
392            r1.file.is_some(),
393            "lock file exists ⇒ real shared lock held"
394        );
395
396        // A second reader must acquire it without blocking on the first.
397        let (tx, rx) = mpsc::channel();
398        let lock2 = lock.clone();
399        let r2 = thread::spawn(move || {
400            let _g = RunLock::acquire_shared(&lock2).unwrap();
401            tx.send(()).unwrap();
402        });
403        assert!(
404            rx.recv_timeout(Duration::from_secs(5)).is_ok(),
405            "a second shared reader must not block on the first"
406        );
407        r2.join().unwrap();
408    }
409
410    /// Nested `with_shared_lock` calls (each on its own descriptor) must not
411    /// deadlock — shared locks are compatible with one another.
412    #[test]
413    fn nested_shared_locks_do_not_deadlock() {
414        let tmp = TempDir::new().unwrap();
415        let lock = tmp.path().join(".lock");
416        drop(RunLock::acquire(&lock).unwrap());
417        let r = RunLock::with_shared_lock(&lock, || RunLock::with_shared_lock(&lock, || Ok(42)))
418            .unwrap();
419        assert_eq!(r, 42);
420    }
421
422    #[test]
423    fn nonblocking_existing_probe_reports_contention_without_waiting() {
424        let tmp = TempDir::new().unwrap();
425        let lock = tmp.path().join(".lock");
426        let held = RunLock::acquire(&lock).unwrap();
427        assert!(RunLock::try_acquire_existing(&lock).unwrap().is_none());
428        drop(held);
429        assert!(RunLock::try_acquire_existing(&lock).unwrap().is_some());
430    }
431
432    #[cfg(unix)]
433    #[test]
434    fn acquire_rejects_a_symlinked_lock_file() {
435        // A symlinked `.lock` would take `flock` on a file outside the run,
436        // silently breaking mutual exclusion — refuse it.
437        use std::os::unix::fs::symlink;
438        let tmp = TempDir::new().unwrap();
439        let target = tmp.path().join("outside.lock");
440        let lock = tmp.path().join(".lock");
441        symlink(&target, &lock).unwrap();
442        assert!(matches!(
443            RunLock::acquire(&lock),
444            Err(Error::SymlinkStateFile { name: "lock", .. })
445        ));
446        // The forged lock never touched the symlink target.
447        assert!(!target.exists());
448    }
449}