Skip to main content

mkit_core/
repo_lock.rs

1//! Repo-level lockfile helper (named `repo_lock` to avoid collision
2//! with `std::sync::*Lock`).
3//!
4//! Pattern (mirrors `mkit-transport-file`'s `RefLock`): open-or-create a
5//! **never-unlinked** sentinel file at `<dir>/<name>` and take an
6//! OS-level exclusive advisory lock on it via `std::fs::File::lock`.
7//! Mutual exclusion comes entirely from the kernel lock, never from the
8//! sentinel's presence/absence on disk — so the file is *not* removed on
9//! release. That structurally removes the stale-vs-live ambiguity a
10//! delete-on-release design has: a sentinel orphaned by a SIGKILL'd
11//! `mkit` looks identical on disk to one backing a live holder, but the
12//! kernel already knows the difference, because it releases the
13//! process's `flock` the moment its file descriptors close (including on
14//! sudden process death). A waiter that actually blocks on that kernel
15//! lock therefore never confuses the two: `ls .mkit/*.lock` will show
16//! the file whether or not anyone holds it — use `lsof`/`fuser` (or a
17//! bounded acquire attempt) to check liveness, not existence.
18//!
19//! Acquisition first tries a non-blocking [`File::try_lock`] (the
20//! uncontended fast path costs no thread spawn). On contention, a helper
21//! thread performs the real blocking `lock()` call and reports back over
22//! a channel; the caller bounds the wait with `recv_timeout(timeout)`.
23//! If `timeout` elapses first, the caller gives up and returns
24//! [`LockError::Busy`] — but the helper thread is not abandoned holding
25//! anything: when it eventually wins the kernel lock (because the
26//! original, live holder finally released), its `send` back to the
27//! (already-dropped) receiver fails, handing the locked `File` back to
28//! the helper thread, which drops it immediately, releasing the kernel
29//! lock rather than leaking it into an unreachable holder.
30//!
31//! POSIX-only intent (macOS + Linux). `std::fs::File::lock` is also
32//! supported on Windows since Rust 1.89, so this works there too — the
33//! lock semantics are equivalent (mandatory `LockFileEx` rather than
34//! advisory `flock`).
35//!
36//! # Network-filesystem caveat
37//!
38//! Every guarantee this module makes — including the fail-closed
39//! GC-vs-writer exclusion `SPEC-GC.md` relies on — assumes `<dir>` is a
40//! genuinely local filesystem (or a well-behaved local-only virtual one
41//! such as tmpfs). `flock`/`fcntl` advisory locking is **not reliably
42//! coherent across NFS clients**: `NFSv3` offloads locking to a separate,
43//! frequently-absent `rpc.statd`/NLM side channel that many exports run
44//! without, and even where present, servers vary in whether a lock
45//! actually excludes a *different host's* holder rather than only
46//! callers on the same NFS client. `NFSv4`'s advertised in-protocol
47//! locking is closer to POSIX semantics but still depends on
48//! server/client support and is not something this module verifies at
49//! runtime. The same caution applies to SMB/CIFS mounts and to most
50//! FUSE-backed network filesystems. A `.mkit/` directory served from
51//! such a mount can silently lose the mutual-exclusion property this
52//! module documents above: two processes on different hosts (or even
53//! the same host, depending on client/server behavior) may both believe
54//! they hold the lock. This module performs no detection of the
55//! underlying filesystem type and has no fallback locking strategy for
56//! this case — repositories that must be shared across hosts should use
57//! one of mkit's network transports (see `docs/specs/SPEC-TRANSPORT.md`)
58//! rather than a shared network-mounted `.mkit/` directory. See
59//! `docs/THREAT-MODEL.md` §7 for how this affects mkit's fail-closed
60//! locking claims.
61
62use std::fs::{File, OpenOptions};
63use std::io;
64use std::path::{Path, PathBuf};
65use std::sync::mpsc;
66use std::time::{Duration, Instant};
67
68/// Default total wall-clock timeout (≈5s). Long enough that a slow
69/// commit in another process finishes; short enough that a stale lock
70/// from a SIGKILL'd `mkit` does not wedge the user for more than a moment.
71pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(5);
72
73/// Maximum filename length for a lock name.
74const MAX_NAME_LEN: usize = 255;
75
76/// Errors returned by [`acquire`].
77#[derive(Debug, thiserror::Error)]
78pub enum LockError {
79    /// Timeout exhausted; another holder still owns the lock.
80    #[error("lock '{0}' busy after timeout")]
81    Busy(String),
82    /// `name` is empty or longer than the platform-safe filename cap.
83    #[error("lock name length {0} is invalid (must be 1..={MAX_NAME_LEN})")]
84    NameLength(usize),
85    /// `name` contains a path separator (`/`, `\`) or a NUL byte. A
86    /// length-based classification would be misleading here — the
87    /// value's length is fine; it's the contents that are wrong.
88    #[error("lock name contains an invalid character (`/`, `\\`, or NUL): {0:?}")]
89    InvalidName(String),
90    /// Underlying filesystem failure (disk full, permission denied, …).
91    #[error(transparent)]
92    Io(#[from] io::Error),
93}
94
95/// Result alias used throughout this module.
96pub type LockResult<T> = Result<T, LockError>;
97
98/// Holder for an acquired repo lock. Releases the kernel lock on `Drop`.
99///
100/// `release()` is the explicit form; calling it is optional because
101/// `Drop` does the same work. After `release()` is called, `Drop` is a
102/// cheap no-op.
103///
104/// The sentinel file at [`Self::path`] is **never** removed by
105/// `release`/`Drop` — see the module doc for why a never-unlinked
106/// sentinel is what makes stale-vs-live holder confusion structurally
107/// impossible.
108#[must_use = "RepoLock releases on drop; bind it to a name to keep the lock"]
109#[derive(Debug)]
110pub struct RepoLock {
111    /// Held file with the OS-level exclusive lock applied.
112    /// `None` after `release()`.
113    file: Option<File>,
114    /// Absolute path to the (never-unlinked) lockfile, for diagnostics
115    /// via [`Self::path`].
116    path: PathBuf,
117}
118
119impl RepoLock {
120    /// Returns the absolute path of the held lock file, for diagnostics.
121    #[must_use]
122    pub fn path(&self) -> &Path {
123        &self.path
124    }
125
126    /// Release the lock: drop the OS lock. The sentinel file itself is
127    /// left on disk (see module doc). Safe to call multiple times —
128    /// subsequent calls are no-ops.
129    pub fn release(&mut self) {
130        if let Some(file) = self.file.take() {
131            // `unlock()` is best-effort; `Drop` of the file would also
132            // release the kernel lock. We still call it explicitly so a
133            // mid-test reader can re-acquire on the same handle if it
134            // wants to.
135            let _ = file.unlock();
136            #[cfg(not(target_arch = "wasm32"))] // wasm File has no destructor
137            drop(file);
138        }
139    }
140}
141
142impl Drop for RepoLock {
143    fn drop(&mut self) {
144        self.release();
145    }
146}
147
148/// Acquire a repo-level lock at `<dir>/<name>`. Waits up to `timeout`
149/// for an existing holder to release, blocking on the kernel lock
150/// (no polling) rather than spinning. Returns a guard that `Drop`s into
151/// a release.
152///
153/// `dir` is usually the `.mkit/` directory (not the worktree root).
154/// `name` is the lockfile basename, e.g. `"index.lock"`.
155///
156/// # Errors
157/// - [`LockError::Busy`] if `timeout` elapses without the lock becoming
158///   available.
159/// - [`LockError::NameLength`] if `name` is empty or longer than 255.
160/// - [`LockError::InvalidName`] if `name` contains a path separator
161///   (`/`, `\`) or a NUL byte.
162/// - [`LockError::Io`] for underlying filesystem failures.
163pub fn acquire(dir: &Path, name: &str, timeout: Duration) -> LockResult<RepoLock> {
164    acquire_with(dir, name, timeout, File::try_lock, File::lock)
165}
166
167/// Acquire a **shared** repo-level lock at `<dir>/<name>`. Any number of
168/// shared holders may coexist; a shared lock only conflicts with an
169/// [`acquire`] (exclusive) holder or an in-progress [`probe_exclusive`].
170///
171/// Used by `mkit serve` (#655/MKIT-11): every live `serve` process holds
172/// a shared lock on `<common_dir>/serve.lock` for its lifetime, so
173/// [`probe_exclusive`] on that same name reports "busy" for as long as
174/// at least one `serve` is alive, without serve instances excluding each
175/// other (SPEC-CONCURRENCY §3.1 documents multiple concurrent `serve`
176/// processes against one root as a supported deployment).
177///
178/// # Errors
179/// See [`acquire`].
180pub fn acquire_shared(dir: &Path, name: &str, timeout: Duration) -> LockResult<RepoLock> {
181    acquire_with(dir, name, timeout, File::try_lock_shared, File::lock_shared)
182}
183
184/// Convenience wrapper: acquire with the default timeout.
185///
186/// # Errors
187/// See [`acquire`].
188pub fn acquire_default(dir: &Path, name: &str) -> LockResult<RepoLock> {
189    acquire(dir, name, DEFAULT_TIMEOUT)
190}
191
192/// Non-blocking probe: is `<dir>/<name>` currently held by anyone (shared
193/// or exclusive)? Attempts a non-blocking **exclusive** `try_lock`
194/// (which only ever succeeds when no holder — shared or exclusive — is
195/// present) and immediately releases it on success, so the probe leaves
196/// no lock behind either way.
197///
198/// Returns `Ok(true)` when the lock was free (the probe's own momentary
199/// exclusive hold has already been dropped by the time this returns), or
200/// `Ok(false)` when another process currently holds it.
201///
202/// Used by `mkit serve`'s local-command guard (#655/MKIT-11) to detect
203/// "is at least one `serve` alive against this root?" without blocking —
204/// a local command that finds the lock busy proceeds anyway and only
205/// emits a warning (see `mkit-cli`'s `commands::warn_if_served`).
206///
207/// # Errors
208/// - [`LockError::NameLength`]/[`LockError::InvalidName`] as in
209///   [`acquire`].
210/// - [`LockError::Io`] for underlying filesystem failures.
211pub fn probe_exclusive(dir: &Path, name: &str) -> LockResult<bool> {
212    let (_path, file) = open_sentinel(dir, name)?;
213    match file.try_lock() {
214        Ok(()) => {
215            // Drop releases the kernel lock (and the fd); the sentinel
216            // file at `_path` is intentionally left in place.
217            #[cfg(not(target_arch = "wasm32"))] // wasm File has no destructor
218            drop(file);
219            Ok(true)
220        }
221        Err(std::fs::TryLockError::Error(e)) => Err(LockError::Io(e)),
222        Err(std::fs::TryLockError::WouldBlock) => Ok(false),
223    }
224}
225
226/// Validate `name` and open-or-create the never-unlinked sentinel file at
227/// `<dir>/<name>`. Shared by every acquire/probe entry point so the
228/// validation rules and sentinel-opening semantics never drift between
229/// them (see the module doc for why the sentinel is never unlinked).
230fn open_sentinel(dir: &Path, name: &str) -> LockResult<(PathBuf, File)> {
231    if name.is_empty() || name.len() > MAX_NAME_LEN {
232        return Err(LockError::NameLength(name.len()));
233    }
234    // Reject path separators and NUL so callers cannot escape `dir`
235    // nor embed bytes the platform filesystem treats specially. These
236    // are CONTENT violations, not LENGTH violations — hence a
237    // dedicated variant.
238    if name.contains('/') || name.contains('\\') || name.contains('\0') {
239        return Err(LockError::InvalidName(name.to_string()));
240    }
241    let path = dir.join(name);
242    let file = OpenOptions::new()
243        .read(true)
244        .write(true)
245        .create(true)
246        .truncate(false)
247        .open(&path)?;
248    Ok((path, file))
249}
250
251/// Shared acquisition core for [`acquire`] and [`acquire_shared`],
252/// parameterized over the lock-mode's non-blocking (`try_lock`) and
253/// blocking (`lock`) primitives so the two modes cannot drift on
254/// validation, sentinel handling, the fast/slow-path split, or the
255/// bounded-wait/no-leak behavior documented in the module doc.
256fn acquire_with(
257    dir: &Path,
258    name: &str,
259    timeout: Duration,
260    try_lock: fn(&File) -> Result<(), std::fs::TryLockError>,
261    lock: fn(&File) -> io::Result<()>,
262) -> LockResult<RepoLock> {
263    let start = Instant::now();
264    let (path, file) = open_sentinel(dir, name)?;
265
266    // Fast path: try the non-blocking lock first so the common
267    // uncontended case never pays for a thread spawn.
268    match try_lock(&file) {
269        Ok(()) => {
270            return Ok(RepoLock {
271                file: Some(file),
272                path,
273            });
274        }
275        Err(std::fs::TryLockError::Error(e)) => return Err(LockError::Io(e)),
276        Err(std::fs::TryLockError::WouldBlock) => {}
277    }
278
279    let remaining = timeout.saturating_sub(start.elapsed());
280    if remaining.is_zero() {
281        return Err(LockError::Busy(name.to_string()));
282    }
283
284    // Slow path: block on the kernel lock for real. A helper thread
285    // performs the blocking `lock()` call (which cannot itself be
286    // bounded by a timeout) and reports back over a channel; this
287    // thread bounds the *wait* with `recv_timeout`. See the module doc
288    // for how an abandoned wait (timeout wins the race) avoids leaking
289    // the lock into an unreachable holder.
290    let (tx, rx) = mpsc::channel();
291    std::thread::spawn(move || {
292        let result = lock(&file).map(|()| file);
293        let _ = tx.send(result);
294    });
295
296    match rx.recv_timeout(remaining) {
297        Ok(Ok(file)) => Ok(RepoLock {
298            file: Some(file),
299            path,
300        }),
301        Ok(Err(e)) => Err(LockError::Io(e)),
302        Err(mpsc::RecvTimeoutError::Timeout) => Err(LockError::Busy(name.to_string())),
303        Err(mpsc::RecvTimeoutError::Disconnected) => Err(LockError::Io(io::Error::other(
304            "repo lock wait thread exited without reporting a result",
305        ))),
306    }
307}
308
309#[cfg(test)]
310mod tests {
311    use super::*;
312    use tempfile::TempDir;
313
314    #[test]
315    fn acquire_release_round_trip() {
316        let dir = TempDir::new().unwrap();
317        {
318            let lock = acquire_default(dir.path(), "index.lock").unwrap();
319            assert!(lock.path().is_file());
320            assert_eq!(lock.path().file_name().unwrap(), "index.lock");
321        }
322        // Never-unlinked sentinel: the file persists after Drop/release
323        // (see module doc) — only the kernel lock is released. Re-acquire
324        // promptly to prove the lock itself is actually free.
325        assert!(
326            dir.path().join("index.lock").exists(),
327            "sentinel file must persist after release"
328        );
329        let l2 = acquire(dir.path(), "index.lock", Duration::from_millis(200)).unwrap();
330        drop(l2);
331    }
332
333    #[test]
334    fn second_acquire_after_release_succeeds() {
335        let dir = TempDir::new().unwrap();
336        let l1 = acquire_default(dir.path(), "index.lock").unwrap();
337        drop(l1);
338        let l2 = acquire_default(dir.path(), "index.lock").unwrap();
339        assert!(l2.path().is_file());
340    }
341
342    #[test]
343    fn acquire_while_held_returns_busy_after_short_timeout() {
344        let dir = TempDir::new().unwrap();
345        let _l1 = acquire_default(dir.path(), "index.lock").unwrap();
346        let err = acquire(dir.path(), "index.lock", Duration::from_millis(150)).unwrap_err();
347        assert!(matches!(err, LockError::Busy(_)));
348    }
349
350    #[test]
351    fn release_is_idempotent() {
352        let dir = TempDir::new().unwrap();
353        let mut lock = acquire_default(dir.path(), "index.lock").unwrap();
354        lock.release();
355        lock.release(); // No-op, no panic.
356        // Sentinel persists (never unlinked); the kernel lock is free.
357        assert!(dir.path().join("index.lock").exists());
358        let l2 = acquire(dir.path(), "index.lock", Duration::from_millis(200)).unwrap();
359        drop(l2);
360    }
361
362    #[test]
363    fn acquire_rejects_empty_name() {
364        let dir = TempDir::new().unwrap();
365        let err = acquire(dir.path(), "", DEFAULT_TIMEOUT).unwrap_err();
366        assert!(matches!(err, LockError::NameLength(0)));
367    }
368
369    #[test]
370    fn acquire_rejects_oversize_name() {
371        let dir = TempDir::new().unwrap();
372        let huge = "a".repeat(300);
373        let err = acquire(dir.path(), &huge, DEFAULT_TIMEOUT).unwrap_err();
374        assert!(matches!(err, LockError::NameLength(300)));
375    }
376
377    #[test]
378    fn acquire_rejects_separators() {
379        let dir = TempDir::new().unwrap();
380        assert!(matches!(
381            acquire(dir.path(), "../escape", DEFAULT_TIMEOUT).unwrap_err(),
382            LockError::InvalidName(_)
383        ));
384        assert!(matches!(
385            acquire(dir.path(), "sub/lock", DEFAULT_TIMEOUT).unwrap_err(),
386            LockError::InvalidName(_)
387        ));
388    }
389
390    #[test]
391    fn acquire_rejects_backslash_and_nul() {
392        let dir = TempDir::new().unwrap();
393        assert!(matches!(
394            acquire(dir.path(), "has\\backslash", DEFAULT_TIMEOUT).unwrap_err(),
395            LockError::InvalidName(_)
396        ));
397        assert!(matches!(
398            acquire(dir.path(), "has\0nul", DEFAULT_TIMEOUT).unwrap_err(),
399            LockError::InvalidName(_)
400        ));
401    }
402
403    #[test]
404    fn two_distinct_lock_names_coexist() {
405        let dir = TempDir::new().unwrap();
406        let _a = acquire_default(dir.path(), "a.lock").unwrap();
407        let _b = acquire_default(dir.path(), "b.lock").unwrap();
408        assert!(dir.path().join("a.lock").is_file());
409        assert!(dir.path().join("b.lock").is_file());
410    }
411
412    // -----------------------------------------------------------------
413    // #635 / INV-16 / INV-17 — blocking-wait regression tests.
414    // -----------------------------------------------------------------
415
416    /// INV-16: a lockfile orphaned by a process that died mid-hold must
417    /// not permanently wedge future acquires.
418    ///
419    /// We simulate "a process acquired the lock and was SIGKILL'd before
420    /// it could run its release/unlink cleanup" without a real
421    /// subprocess: create + `flock` the sentinel directly (bypassing
422    /// `acquire`, which would also register the normal release path),
423    /// then `drop` the handle without unlinking the file. Dropping the
424    /// `File` closes its descriptor, which releases the kernel `flock`
425    /// exactly as a killed process's fd closing would — but the sentinel
426    /// file itself is left behind on disk, exactly as `O_EXCL`-created
427    /// lockfiles are on a real SIGKILL.
428    ///
429    /// Against the old poll-`O_EXCL` wait loop this wedges forever: every
430    /// `acquire` sees `AlreadyExists` on `create_new`, never checks
431    /// whether the kernel lock is actually free, and burns the full
432    /// timeout before failing `Busy` — even though nobody holds it.
433    #[test]
434    fn orphaned_lock_does_not_wedge_future_acquire() {
435        let dir = TempDir::new().unwrap();
436        let path = dir.path().join("index.lock");
437
438        {
439            let orphan = OpenOptions::new()
440                .write(true)
441                .create_new(true)
442                .open(&path)
443                .unwrap();
444            orphan.lock().unwrap();
445            // Simulates the fd closing on process death: releases the
446            // kernel lock but leaves the sentinel file on disk.
447            drop(orphan);
448        }
449        assert!(
450            path.exists(),
451            "orphaned sentinel file must remain on disk after the simulated death"
452        );
453
454        // Nobody holds the kernel lock anymore, so a fresh acquire must
455        // succeed well within a bounded timeout rather than burning it.
456        let start = Instant::now();
457        let lock = acquire(dir.path(), "index.lock", Duration::from_secs(2))
458            .expect("acquire must succeed once the orphaned holder's kernel lock is free");
459        let elapsed = start.elapsed();
460        assert!(
461            elapsed < Duration::from_millis(500),
462            "acquire should not pay anywhere near the timeout on an orphaned lock, took {elapsed:?}"
463        );
464        drop(lock);
465    }
466
467    // -----------------------------------------------------------------
468    // MKIT-11 — shared lock + non-blocking exclusive probe.
469    // -----------------------------------------------------------------
470
471    /// Two shared holders may coexist on the same lock name: this is
472    /// exactly the semantics `mkit serve` needs, since SPEC-TRANSPORT
473    /// supports multiple concurrent `serve` processes against one root.
474    #[test]
475    fn two_shared_acquires_coexist() {
476        let dir = TempDir::new().unwrap();
477        let a = acquire_shared(dir.path(), "serve.lock", Duration::from_millis(200)).unwrap();
478        let b = acquire_shared(dir.path(), "serve.lock", Duration::from_millis(200)).unwrap();
479        drop(a);
480        drop(b);
481    }
482
483    /// `probe_exclusive` reports the lock as unavailable (`Ok(false)`)
484    /// while a shared guard is alive, and available (`Ok(true)`) again
485    /// once it is dropped.
486    #[test]
487    fn probe_exclusive_reflects_a_live_shared_holder() {
488        let dir = TempDir::new().unwrap();
489        let shared = acquire_shared(dir.path(), "serve.lock", DEFAULT_TIMEOUT).unwrap();
490        assert!(
491            !probe_exclusive(dir.path(), "serve.lock").unwrap(),
492            "probe must report busy while a shared holder is alive"
493        );
494        drop(shared);
495        assert!(
496            probe_exclusive(dir.path(), "serve.lock").unwrap(),
497            "probe must report free once the shared holder releases"
498        );
499    }
500
501    /// With nobody holding the lock, probing it must not itself leave a
502    /// lingering hold behind (i.e. it must be a true probe, not a leak).
503    #[test]
504    fn probe_exclusive_does_not_leak_the_lock_it_takes() {
505        let dir = TempDir::new().unwrap();
506        assert!(probe_exclusive(dir.path(), "serve.lock").unwrap());
507        // If the previous probe leaked its lock, this exclusive acquire
508        // would time out.
509        let l = acquire(dir.path(), "serve.lock", Duration::from_millis(200)).unwrap();
510        drop(l);
511    }
512
513    /// INV-8/INV-16: a genuinely blocking waiter must observe the actual
514    /// release event promptly, not merely "eventually, after polling
515    /// catches up." A holder thread releases well before the acquirer's
516    /// timeout; the acquirer must return success shortly after that
517    /// release rather than needing to reach the timeout.
518    #[test]
519    fn waiter_observes_release_from_a_live_holder_without_timing_out() {
520        let dir = TempDir::new().unwrap();
521        let dir_path = dir.path().to_path_buf();
522
523        let holder = acquire_default(&dir_path, "index.lock").unwrap();
524        let (release_tx, release_rx) = mpsc::channel::<()>();
525        let holder_thread = std::thread::spawn(move || {
526            // Hold the lock until told to let go.
527            let _ = release_rx.recv();
528            drop(holder);
529        });
530
531        let waiter_dir = dir_path.clone();
532        let waiter = std::thread::spawn(move || {
533            let start = Instant::now();
534            let lock = acquire(&waiter_dir, "index.lock", Duration::from_secs(5))
535                .expect("waiter must observe the release rather than timing out");
536            (lock, start.elapsed())
537        });
538
539        // Let the holder sit on the lock briefly, then release it. The
540        // waiter's bounded 5s timeout is generous; what matters is that
541        // it does not need anywhere near that long once release happens.
542        std::thread::sleep(Duration::from_millis(150));
543        release_tx.send(()).unwrap();
544        holder_thread.join().unwrap();
545
546        let (lock, elapsed) = waiter.join().unwrap();
547        assert!(
548            elapsed < Duration::from_secs(2),
549            "waiter should wake up promptly on release, not approach the 5s timeout, took {elapsed:?}"
550        );
551        drop(lock);
552    }
553}