Skip to main content

amont_runtime/
host_slots.rs

1//! Host-wide slots for heavy checks (ADR-0009, `hooks.host-concurrency`).
2//!
3//! Several worktrees, sessions and agents on one machine each run their own
4//! clippy, test suite and type checker, and each of those already uses every
5//! core. Run together they thrash the same cores and the same cargo lock,
6//! and the kill clocks counted that thrash as the check's own time. A check
7//! that compiles or executes the product is therefore Heavy, and takes one
8//! of `amont.hostSlots` slots — `flock` on a file under one fixed per-user
9//! directory — before its first tool runs. While it waits, its clocks have
10//! not started, because its tool has not.
11//!
12//! The slot is taken LAZILY, at the first tool spawn of a heavy check, not
13//! when the check starts: a clippy with no Rust staged returns at once and
14//! must not wait behind somebody else's suite to do nothing.
15//!
16//! Fairness is polling, not first come first served.
17//! holds-until: more than a handful of gates contend on one host, when a
18//! queue file with tickets is needed.
19
20use std::cell::RefCell;
21use std::path::PathBuf;
22use std::time::{Duration, Instant};
23
24/// The environment variable a held slot exports to every tool it runs, and
25/// that a nested amont reads to skip queueing: amont's own
26/// `pre-push-cargo-test` runs fixtures that run amont, and those must not
27/// wait for the slot their own parent holds.
28pub const HELD_ENV: &str = "AMONT_HOST_SLOT";
29
30/// A diagnostic, and the seam the tests use: the slot directory, instead of
31/// the fixed per-user one.
32pub const DIR_ENV: &str = "AMONT_SLOT_DIR";
33
34/// The checks that compile or execute the product. Pinned by a registry
35/// test: every name here is a built-in.
36pub const HEAVY: &[&str] = &[
37    "pre-commit-clippy",
38    "pre-commit-go-vet",
39    "pre-commit-pyright",
40    "pre-push-run-tests-js",
41    "pre-push-cargo-test",
42    "pre-push-go-test",
43    "pre-push-pytest",
44];
45
46/// Whether a check is Heavy.
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum Weight {
49    Light,
50    Heavy,
51}
52
53pub fn weight_of(name: &str) -> Weight {
54    if HEAVY.contains(&name) {
55        Weight::Heavy
56    } else {
57        Weight::Light
58    }
59}
60
61/// Why a heavy check ran without a slot.
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub enum Reason {
64    /// `amont.hostSlots 0`.
65    Disabled,
66    /// A slot is already held by the amont that runs this one.
67    Nested,
68    /// The slot directory is not one this user alone owns.
69    DirUnsafe(String),
70    /// No `flock` here: Windows, the stated divergence.
71    Unsupported,
72    /// No slot freed within `amont.timeout`.
73    TimedOut(Duration),
74    /// A slot file could not be opened or locked for a reason other than
75    /// "another holder has it".
76    Io(String),
77}
78
79/// What [`acquire`] got.
80pub enum Acquired {
81    Slot(Slot),
82    Unqueued(Reason),
83}
84
85/// A held slot: released when dropped, and by the kernel if amont dies.
86pub struct Slot {
87    _file: std::fs::File,
88}
89
90/// `amont.hostSlots`: a host key, default a quarter of the cores and at
91/// least one, `0` disables.
92pub fn configured(settings: &crate::config::Settings) -> u64 {
93    let cores = std::thread::available_parallelism().map_or(1, |n| n.get() as i64);
94    crate::config::host_integer_or(settings, "amont.hostSlots", (cores / 4).max(1), 0..=64) as u64
95}
96
97/// The slot directory: `AMONT_SLOT_DIR`, else `$XDG_RUNTIME_DIR/amont-slots`
98/// on Linux, else `/tmp/amont-slots-<uid>`. A FIXED path, never
99/// `std::env::temp_dir()`, which follows a `$TMPDIR` that macOS and agent
100/// sandboxes set per session — every session would get a queue of its own.
101pub fn dir() -> PathBuf {
102    if let Some(d) = std::env::var_os(DIR_ENV).filter(|d| !d.is_empty()) {
103        return PathBuf::from(d);
104    }
105    #[cfg(target_os = "linux")]
106    if let Some(d) = std::env::var_os("XDG_RUNTIME_DIR").filter(|d| !d.is_empty()) {
107        return PathBuf::from(d).join("amont-slots");
108    }
109    PathBuf::from(format!("/tmp/amont-slots-{}", platform::uid()))
110}
111
112/// Make `dir` if needed, then refuse it unless it is a real directory this
113/// user owns with mode 0700: on a shared `/tmp`, another user could create
114/// it first and hold every slot.
115pub fn safe_dir(dir: &std::path::Path) -> Result<(), String> {
116    platform::make_private(dir)?;
117    platform::check_private(dir)
118}
119
120/// Take one of `n` slots, waiting at most `budget`.
121pub fn acquire(n: u64, budget: Duration) -> Acquired {
122    if n == 0 {
123        return Acquired::Unqueued(Reason::Disabled);
124    }
125    if !platform::SUPPORTED {
126        return Acquired::Unqueued(Reason::Unsupported);
127    }
128    let dir = dir();
129    if let Err(why) = safe_dir(&dir) {
130        return Acquired::Unqueued(Reason::DirUnsafe(why));
131    }
132    let started = Instant::now();
133    loop {
134        for i in 0..n {
135            let path = dir.join(format!("slot-{i}"));
136            match platform::try_lock(&path) {
137                Ok(Some(file)) => return Acquired::Slot(Slot { _file: file }),
138                Ok(None) => {}
139                // Not "busy": an error waiting would only repeat, and a
140                // full queue must not be claimed when nobody was seen.
141                Err(e) => {
142                    return Acquired::Unqueued(Reason::Io(format!("{}: {e}", path.display())))
143                }
144            }
145        }
146        if started.elapsed() >= budget {
147            return Acquired::Unqueued(Reason::TimedOut(started.elapsed()));
148        }
149        std::thread::sleep(Duration::from_millis(200));
150    }
151}
152
153/// What the check running on THIS thread is, as far as slots go.
154#[derive(Default)]
155struct Current {
156    /// The check's name, for the messages its tools' runners print.
157    name: String,
158    heavy: bool,
159    /// The slot question has been answered for this check (held, or
160    /// unqueued for a reason).
161    settled: bool,
162    held: Option<Slot>,
163    queued: Duration,
164}
165
166thread_local! {
167    static CURRENT: RefCell<Current> = RefCell::new(Current::default());
168}
169
170/// Installed by the dispatcher around one check's `run`; clears this
171/// thread's slot state, and releases a held slot, when dropped.
172pub struct CheckGuard;
173
174impl CheckGuard {
175    /// How long this check waited for a slot so far.
176    pub fn queued(&self) -> Duration {
177        CURRENT.with(|c| c.borrow().queued)
178    }
179}
180
181impl Drop for CheckGuard {
182    fn drop(&mut self) {
183        CURRENT.with(|c| *c.borrow_mut() = Current::default());
184    }
185}
186
187/// The dispatcher is about to run `name` on this thread.
188pub fn enter_check(name: &str) -> CheckGuard {
189    CURRENT.with(|c| {
190        *c.borrow_mut() = Current {
191            name: name.to_string(),
192            heavy: weight_of(name) == Weight::Heavy,
193            ..Current::default()
194        }
195    });
196    CheckGuard
197}
198
199/// The name of the check running on this thread, without its stage
200/// prefix (`clippy`), when the dispatcher installed one.
201pub fn current_check() -> Option<String> {
202    CURRENT.with(|c| {
203        let c = c.borrow();
204        (!c.name.is_empty()).then(|| crate::short_name(&c.name).to_string())
205    })
206}
207
208/// Whether this process was started by an amont that holds a slot.
209fn nested() -> bool {
210    std::env::var_os(HELD_ENV).is_some_and(|v| !v.is_empty())
211}
212
213/// Called before a tool is spawned. For the first tool of a heavy check,
214/// waits for a slot, showing the wait in the check's row. Returns whether
215/// this thread holds a slot, so the caller exports [`HELD_ENV`].
216pub fn before_spawn(settings: &crate::config::Settings) -> bool {
217    let pending = CURRENT.with(|c| {
218        let c = c.borrow();
219        c.heavy && !c.settled
220    });
221    if pending {
222        // Nested first: a nested amont must not pay the host-key scan for
223        // a slot count it is about to ignore.
224        let n = if nested() { 0 } else { configured(settings) };
225        let outcome = if nested() {
226            Acquired::Unqueued(Reason::Nested)
227        } else if n == 0 {
228            Acquired::Unqueued(Reason::Disabled)
229        } else {
230            let budget = match crate::hooks::common::check_timeout(settings) {
231                0 => 3600,
232                s => s,
233            };
234            let sink = crate::live::current_sink();
235            if let Some((stage, idx)) = &sink {
236                stage.queue(*idx, Some((Instant::now(), n)));
237            }
238            let started = Instant::now();
239            let got = acquire(n, Duration::from_secs(budget));
240            if let Some((stage, idx)) = &sink {
241                stage.queue(*idx, None);
242            }
243            let waited = started.elapsed();
244            CURRENT.with(|c| c.borrow_mut().queued += waited);
245            // Piped output never shows the region, and a heartbeat only
246            // comes after a minute: say once where the time went.
247            if matches!(got, Acquired::Slot(_)) && waited >= Duration::from_secs(1) {
248                crate::hooks::common::say(&format!(
249                    "  (waited {} for a host slot, {} {n})",
250                    crate::hooks::common::human_secs(waited.as_secs()),
251                    crate::ui::highlight("amont.hostSlots")
252                ));
253            }
254            got
255        };
256        let held = match outcome {
257            Acquired::Slot(slot) => Some(slot),
258            Acquired::Unqueued(reason) => {
259                note(&reason, n);
260                None
261            }
262        };
263        CURRENT.with(|c| {
264            let mut c = c.borrow_mut();
265            c.held = held;
266            c.settled = true;
267        });
268    }
269    CURRENT.with(|c| c.borrow().held.is_some())
270}
271
272/// One line for the reasons a person can act on; nothing for the ones
273/// that are the normal case (off, nested, the platform).
274fn note(reason: &Reason, n: u64) {
275    match reason {
276        Reason::DirUnsafe(why) => crate::hooks::common::warn(&format!(
277            "host slots are off for this check: {why}. Remove it, or set {} to a directory \
278             only you own",
279            crate::ui::highlight(DIR_ENV)
280        )),
281        Reason::TimedOut(waited) => crate::hooks::common::warn(&format!(
282            "no host slot freed in {} ({} heavy checks were running on this machine, \
283             {} {n}); running it anyway",
284            crate::hooks::common::human_secs(waited.as_secs()),
285            n,
286            crate::ui::highlight("amont.hostSlots")
287        )),
288        Reason::Io(why) => {
289            crate::hooks::common::warn(&format!("host slots are off for this check: {why}"))
290        }
291        Reason::Disabled | Reason::Nested | Reason::Unsupported => {}
292    }
293}
294
295#[cfg(unix)]
296mod platform {
297    use std::os::unix::fs::{DirBuilderExt, MetadataExt, PermissionsExt};
298    use std::os::unix::io::AsRawFd;
299    use std::path::Path;
300
301    pub const SUPPORTED: bool = true;
302
303    extern "C" {
304        #[link_name = "flock"]
305        fn libc_flock(fd: i32, op: i32) -> i32;
306        #[link_name = "getuid"]
307        fn libc_getuid() -> u32;
308    }
309    const LOCK_EX: i32 = 2;
310    const LOCK_NB: i32 = 4;
311
312    pub fn uid() -> u32 {
313        // SAFETY: getuid takes nothing and cannot fail.
314        unsafe { libc_getuid() }
315    }
316
317    pub fn make_private(dir: &Path) -> Result<(), String> {
318        if dir.symlink_metadata().is_ok() {
319            return Ok(());
320        }
321        std::fs::DirBuilder::new()
322            .recursive(true)
323            .mode(0o700)
324            .create(dir)
325            .map_err(|e| format!("cannot create {}: {e}", dir.display()))
326    }
327
328    pub fn check_private(dir: &Path) -> Result<(), String> {
329        let meta = dir
330            .symlink_metadata()
331            .map_err(|e| format!("cannot read {}: {e}", dir.display()))?;
332        if meta.file_type().is_symlink() {
333            return Err(format!("{} is a symbolic link", dir.display()));
334        }
335        if !meta.is_dir() {
336            return Err(format!("{} is not a directory", dir.display()));
337        }
338        if meta.uid() != uid() {
339            return Err(format!("{} belongs to another user", dir.display()));
340        }
341        let mode = meta.permissions().mode() & 0o777;
342        if mode != 0o700 {
343            return Err(format!("{} has mode {mode:o}, not 700", dir.display()));
344        }
345        Ok(())
346    }
347
348    /// `Ok(None)` only when another holder has the lock (EWOULDBLOCK);
349    /// every other failure is an error the caller reports.
350    pub fn try_lock(path: &Path) -> std::io::Result<Option<std::fs::File>> {
351        let file = std::fs::OpenOptions::new()
352            .create(true)
353            .truncate(false)
354            .write(true)
355            .open(path)?;
356        // SAFETY: a valid fd we own; flock touches no memory.
357        if unsafe { libc_flock(file.as_raw_fd(), LOCK_EX | LOCK_NB) } == 0 {
358            return Ok(Some(file));
359        }
360        let e = std::io::Error::last_os_error();
361        if e.kind() == std::io::ErrorKind::WouldBlock {
362            Ok(None)
363        } else {
364            Err(e)
365        }
366    }
367}
368
369#[cfg(not(unix))]
370mod platform {
371    use std::path::Path;
372
373    pub const SUPPORTED: bool = false;
374
375    pub fn uid() -> u32 {
376        0
377    }
378    pub fn make_private(_dir: &Path) -> Result<(), String> {
379        Ok(())
380    }
381    pub fn check_private(_dir: &Path) -> Result<(), String> {
382        Ok(())
383    }
384    pub fn try_lock(_path: &Path) -> std::io::Result<Option<std::fs::File>> {
385        Ok(None)
386    }
387}
388
389#[cfg(all(test, unix))]
390mod tests {
391    use super::*;
392
393    fn scratch(tag: &str) -> PathBuf {
394        let d = std::env::temp_dir().join(format!("amont-slots-test-{}-{tag}", std::process::id()));
395        let _ = std::fs::remove_dir_all(&d);
396        d
397    }
398
399    /// A fresh directory is made 0700 and accepted; a 0755 one, a symlink
400    /// and a file are refused.
401    #[test]
402    fn only_a_private_directory_is_accepted() {
403        use std::os::unix::fs::PermissionsExt;
404        let d = scratch("private");
405        assert_eq!(safe_dir(&d), Ok(()));
406        std::fs::set_permissions(&d, std::fs::Permissions::from_mode(0o755)).unwrap();
407        assert!(safe_dir(&d).unwrap_err().contains("mode 755"));
408        let link = scratch("link");
409        std::os::unix::fs::symlink(&d, &link).unwrap();
410        assert!(safe_dir(&link).unwrap_err().contains("symbolic link"));
411        let _ = std::fs::remove_file(&link);
412        let _ = std::fs::remove_dir_all(&d);
413    }
414
415    /// Two slots: two holders get one each, a third finds none within its
416    /// budget, and a dropped slot is free again.
417    #[test]
418    fn slots_are_exclusive_and_released_on_drop() {
419        let d = scratch("exclusive");
420        safe_dir(&d).unwrap();
421        let a = platform::try_lock(&d.join("slot-0"))
422            .expect("no error")
423            .expect("first");
424        assert!(platform::try_lock(&d.join("slot-0"))
425            .expect("busy is not an error")
426            .is_none());
427        drop(a);
428        // Released on the last close. A thread of another test that forks
429        // in the instant before its exec holds a copy of every descriptor
430        // until then (O_CLOEXEC closes it at exec, not at fork), so allow
431        // that instant rather than assert a release that is a race away.
432        let freed = (0..100).any(|_| {
433            matches!(platform::try_lock(&d.join("slot-0")), Ok(Some(_))) || {
434                std::thread::sleep(Duration::from_millis(20));
435                false
436            }
437        });
438        assert!(freed, "a dropped slot was not released within 2 s");
439        let _ = std::fs::remove_dir_all(&d);
440    }
441
442    #[test]
443    fn heavy_is_the_list() {
444        assert_eq!(weight_of("pre-commit-clippy"), Weight::Heavy);
445        assert_eq!(weight_of("pre-commit-cargo-fmt"), Weight::Light);
446    }
447
448    /// Every heavy name is a built-in: a renamed check must not silently
449    /// stop queueing.
450    #[test]
451    fn every_heavy_name_is_a_builtin() {
452        for name in HEAVY {
453            assert!(
454                crate::registry::CHECKS.iter().any(|c| c.name == *name),
455                "{name} is not a built-in check"
456            );
457        }
458    }
459
460    /// `.cargo/config.toml` reaches the test binaries: without it the
461    /// integration fixtures would queue on the real host slots.
462    #[test]
463    fn cargo_pins_the_host_knobs_for_tests() {
464        assert_eq!(std::env::var(HELD_ENV).as_deref(), Ok("held"));
465        assert_eq!(std::env::var("AMONT_IDLE_LOAD_SCALE").as_deref(), Ok("1"));
466    }
467
468    /// A light check never asks. A heavy one run by an amont that holds a
469    /// slot (the test process has `AMONT_HOST_SLOT=held` from
470    /// `.cargo/config.toml`) is nested: it settles without a slot.
471    #[test]
472    fn a_light_check_never_takes_a_slot() {
473        let settings = crate::config::Settings::default();
474        let _g = enter_check("pre-commit-cargo-fmt");
475        assert!(!before_spawn(&settings));
476    }
477
478    #[test]
479    fn a_nested_heavy_check_settles_without_a_slot() {
480        let settings = crate::config::Settings::default();
481        let _g = enter_check("pre-commit-clippy");
482        assert_eq!(current_check().as_deref(), Some("clippy"));
483        assert!(!before_spawn(&settings));
484        assert!(CURRENT.with(|c| c.borrow().settled));
485    }
486
487    /// An unopenable slot file is an error, not a full queue.
488    #[test]
489    fn an_unopenable_slot_is_an_error_not_busy() {
490        let d = scratch("missing").join("absent-subdir").join("slot-0");
491        assert!(platform::try_lock(&d).is_err());
492    }
493}