Skip to main content

rightkit_qa/
lock.rs

1//! Run lock: one QA run per lock directory. A lock is reclaimed only when its owner is
2//! provably gone (pid dead, or the pid now belongs to a different process); a live owner is
3//! never displaced however long its run takes. Age applies only to a lock directory that
4//! never received an owner record (a crash between `mkdir` and the owner write).
5use crate::process::is_alive;
6use crate::util::{err, new_id, now_iso, now_ms, Result};
7use serde::{Deserialize, Serialize};
8use std::fs;
9use std::path::{Path, PathBuf};
10use std::time::Duration;
11
12#[derive(Debug, Clone, Serialize, Deserialize)]
13pub struct RunOwner {
14    pub label: String,
15    pub pid: u32,
16    pub token: String,
17    pub started_at: String,
18    pub started_at_ms: u128,
19    /// Command line of the owner at acquisition; guards against pid reuse.
20    #[serde(default)]
21    pub fingerprint: Option<String>,
22}
23
24pub struct RunOwnership {
25    pub owner: RunOwner,
26    dir: PathBuf,
27    released: bool,
28}
29
30#[derive(Debug, Clone)]
31pub struct LockOptions {
32    pub timeout: Duration,
33    /// Only for a lock directory with no readable owner record: how long before it counts as abandoned.
34    /// Never applied to a lock whose owner is alive.
35    pub stale: Duration,
36    pub poll: Duration,
37}
38
39impl Default for LockOptions {
40    fn default() -> Self {
41        Self {
42            timeout: Duration::from_secs(120),
43            stale: Duration::from_secs(60),
44            poll: Duration::from_millis(250),
45        }
46    }
47}
48
49pub fn acquire(lock_dir: &Path, label: &str, opts: &LockOptions) -> Result<RunOwnership> {
50    acquire_with(lock_dir, label, opts, &owner_alive)
51}
52
53/// The owner is alive when its pid runs and still has the command line recorded at acquisition.
54pub fn owner_alive(o: &RunOwner) -> bool {
55    if !is_alive(o.pid) {
56        return false;
57    }
58    match (&o.fingerprint, crate::process::fingerprint(o.pid)) {
59        (Some(want), Some(now)) => *want == now,
60        _ => true,
61    }
62}
63
64pub fn acquire_with(
65    lock_dir: &Path,
66    label: &str,
67    opts: &LockOptions,
68    alive: &dyn Fn(&RunOwner) -> bool,
69) -> Result<RunOwnership> {
70    let owner_path = lock_dir.join("owner.json");
71    let waiting_since = std::time::Instant::now();
72    let owner = RunOwner {
73        label: if label.trim().is_empty() {
74            "qa-run".into()
75        } else {
76            label.trim().into()
77        },
78        pid: std::process::id(),
79        token: new_id(),
80        started_at: now_iso(),
81        started_at_ms: now_ms(),
82        fingerprint: crate::process::fingerprint(std::process::id()),
83    };
84    if let Some(parent) = lock_dir.parent() {
85        fs::create_dir_all(parent)?;
86    }
87    loop {
88        match fs::create_dir(lock_dir) {
89            Ok(()) => {
90                let tmp = lock_dir.join(format!("owner.json.tmp-{}", owner.token));
91                fs::write(&tmp, serde_json::to_vec_pretty(&owner)?)?;
92                fs::rename(&tmp, &owner_path)?;
93                return Ok(RunOwnership {
94                    owner,
95                    dir: lock_dir.to_path_buf(),
96                    released: false,
97                });
98            }
99            Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
100                let current = read_owner(&owner_path);
101                let reclaim = match &current {
102                    Some(o) => !alive(o),
103                    None => {
104                        let age = fs::metadata(lock_dir)
105                            .and_then(|m| m.modified())
106                            .ok()
107                            .and_then(|t| t.elapsed().ok())
108                            .unwrap_or(Duration::ZERO);
109                        age > opts.stale
110                    }
111                };
112                if reclaim {
113                    reclaim_dir(lock_dir, current.as_ref().map(|o| o.token.as_str()));
114                    continue;
115                }
116                if waiting_since.elapsed() >= opts.timeout {
117                    return err(format!(
118                        "timed out waiting for QA run lock held by {} pid {}",
119                        current
120                            .as_ref()
121                            .map(|o| o.label.as_str())
122                            .unwrap_or("unknown"),
123                        current
124                            .as_ref()
125                            .map(|o| o.pid.to_string())
126                            .unwrap_or_else(|| "unknown".into())
127                    ));
128                }
129                std::thread::sleep(opts.poll);
130            }
131            Err(e) => {
132                let _ = fs::remove_dir_all(lock_dir);
133                return Err(e.into());
134            }
135        }
136    }
137}
138
139/// Atomically move the observed-dead lock aside, then delete it. The rename has one
140/// winner; if what moved turns out to be a different owner's fresh lock it is put back.
141fn reclaim_dir(lock_dir: &Path, observed: Option<&str>) {
142    let tomb = lock_dir.with_extension(format!("reclaimed-{}", new_id()));
143    if fs::rename(lock_dir, &tomb).is_err() {
144        return;
145    }
146    let moved = read_owner(&tomb.join("owner.json")).map(|o| o.token);
147    if moved.is_some() && moved.as_deref() != observed {
148        let _ = fs::rename(&tomb, lock_dir);
149        return;
150    }
151    let _ = fs::remove_dir_all(&tomb);
152}
153
154fn read_owner(path: &Path) -> Option<RunOwner> {
155    serde_json::from_slice(&fs::read(path).ok()?).ok()
156}
157
158impl RunOwnership {
159    /// Remove the lock only if this run still owns it.
160    pub fn release(&mut self) {
161        if self.released {
162            return;
163        }
164        self.released = true;
165        if read_owner(&self.dir.join("owner.json"))
166            .map(|o| o.token == self.owner.token)
167            .unwrap_or(false)
168        {
169            let _ = fs::remove_dir_all(&self.dir);
170        }
171    }
172}
173
174impl Drop for RunOwnership {
175    fn drop(&mut self) {
176        self.release();
177    }
178}
179
180#[cfg(test)]
181mod tests {
182    use super::*;
183
184    fn tmp(name: &str) -> PathBuf {
185        let d = std::env::temp_dir().join(format!("rkqa-lock-{name}-{}", new_id()));
186        fs::create_dir_all(&d).unwrap();
187        d
188    }
189
190    #[test]
191    fn second_run_times_out_while_first_holds() {
192        let d = tmp("hold").join("native.lock");
193        let first = acquire(&d, "a", &LockOptions::default()).unwrap();
194        let quick = LockOptions {
195            timeout: Duration::from_millis(60),
196            poll: Duration::from_millis(10),
197            ..Default::default()
198        };
199        let e = acquire(&d, "b", &quick).err().expect("must time out");
200        assert!(e.0.contains("held by a"), "{e}");
201        drop(first);
202        assert!(acquire(&d, "c", &quick).is_ok());
203    }
204
205    #[test]
206    fn live_owner_is_never_reclaimed_by_age() {
207        let d = tmp("live").join("native.lock");
208        let _first = acquire(&d, "a", &LockOptions::default()).unwrap();
209        // Age the owner record far past any limit; the owner (this process) is alive.
210        let mut o: RunOwner =
211            serde_json::from_slice(&fs::read(d.join("owner.json")).unwrap()).unwrap();
212        o.started_at_ms = 1;
213        fs::write(d.join("owner.json"), serde_json::to_vec(&o).unwrap()).unwrap();
214        let tight = LockOptions {
215            timeout: Duration::from_millis(80),
216            poll: Duration::from_millis(10),
217            stale: Duration::from_millis(1),
218        };
219        assert!(
220            acquire(&d, "b", &tight).is_err(),
221            "a live owner must not be displaced by age"
222        );
223        assert_eq!(read_owner(&d.join("owner.json")).unwrap().label, "a");
224    }
225
226    #[test]
227    fn dead_owner_is_reclaimed() {
228        let d = tmp("dead").join("native.lock");
229        let _first = acquire_with(&d, "a", &LockOptions::default(), &|_: &RunOwner| true).unwrap();
230        let second = acquire_with(&d, "b", &LockOptions::default(), &|_: &RunOwner| false).unwrap();
231        assert_eq!(second.owner.label, "b");
232    }
233}