1use 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 #[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 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
53pub 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 ¤t {
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
139fn 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 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 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}