1use crate::phase_id::PhaseId;
14use std::fs::{self, File};
15use std::io::{self, Write};
16use std::path::{Path, PathBuf};
17
18#[derive(Debug, thiserror::Error)]
20pub enum LockError {
21 #[error("lock already held by pid {pid} at {path}")]
23 Contended { pid: String, path: PathBuf },
24 #[error("lock I/O failed: {0}")]
26 Io(#[from] io::Error),
27}
28
29pub fn acquire(project_root: &Path, phase: PhaseId) -> Result<LockGuard, LockError> {
34 acquire_path(lock_path(project_root, phase))
35}
36
37pub fn acquire_project(project_root: &Path) -> Result<LockGuard, LockError> {
47 acquire_path(project_lock_path(project_root))
48}
49
50pub fn acquire_project_blocking(
56 project_root: &Path,
57 timeout: std::time::Duration,
58) -> Result<LockGuard, LockError> {
59 let start = std::time::Instant::now();
60 let mut backoff = std::time::Duration::from_millis(100);
61 loop {
62 match acquire_project(project_root) {
63 Ok(guard) => return Ok(guard),
64 Err(err @ LockError::Contended { .. }) => {
65 if start.elapsed() >= timeout {
66 return Err(err);
67 }
68 std::thread::sleep(backoff.min(timeout.saturating_sub(start.elapsed())));
69 backoff = (backoff * 2).min(std::time::Duration::from_secs(2));
70 }
71 Err(err) => return Err(err),
72 }
73 }
74}
75
76fn acquire_path(path: PathBuf) -> Result<LockGuard, LockError> {
77 let parent = path.parent().ok_or_else(|| {
78 io::Error::new(
79 io::ErrorKind::InvalidInput,
80 "lock path has no parent directory",
81 )
82 })?;
83 crate::workflow::ensure_devflow_dir(parent)?;
84
85 match File::create_new(&path) {
86 Ok(mut f) => {
87 write!(f, "{}", lock_contents())?;
88 Ok(LockGuard { path })
89 }
90 Err(err) if err.kind() == io::ErrorKind::AlreadyExists => {
91 let pid = read_holder_pid(&path);
92 if !pid_is_alive(&pid) {
101 tracing::warn!(
102 "reclaiming stale devflow lock at {} (holder pid {pid} is not alive)",
103 path.display()
104 );
105 let _ = fs::remove_file(&path);
106 return match File::create_new(&path) {
107 Ok(mut f) => {
108 write!(f, "{}", lock_contents())?;
109 Ok(LockGuard { path })
110 }
111 Err(err) if err.kind() == io::ErrorKind::AlreadyExists => {
112 let pid = read_holder_pid(&path);
113 Err(LockError::Contended { pid, path })
114 }
115 Err(err) => Err(err.into()),
116 };
117 }
118 Err(LockError::Contended { pid, path })
119 }
120 Err(err) => Err(err.into()),
121 }
122}
123
124fn lock_contents() -> String {
139 let pid = std::process::id();
140 match crate::agent::process_start_time(pid) {
141 Some(start) => format!("{pid}\n{start}"),
142 None => format!("{pid}"),
146 }
147}
148
149fn read_holder_pid(path: &Path) -> String {
152 fs::read_to_string(path)
153 .ok()
154 .and_then(|text| text.lines().next().map(|line| line.trim().to_string()))
155 .filter(|pid| !pid.is_empty())
156 .unwrap_or_else(|| "unknown".into())
157}
158
159fn read_holder_start_time(path: &Path) -> Option<u64> {
162 fs::read_to_string(path)
163 .ok()?
164 .lines()
165 .nth(1)?
166 .trim()
167 .parse::<u64>()
168 .ok()
169}
170
171pub fn holder_identity(project_root: &Path, phase: PhaseId) -> Option<(u32, Option<u64>)> {
179 let path = lock_path(project_root, phase);
180 let pid = read_holder_pid(&path).parse::<u32>().ok()?;
181 Some((pid, read_holder_start_time(&path)))
182}
183
184fn pid_is_alive(pid: &str) -> bool {
193 pid.parse::<u32>().is_ok_and(crate::agent::agent_running)
194}
195
196pub fn holder(project_root: &Path, phase: PhaseId) -> Option<(String, PathBuf)> {
199 let path = lock_path(project_root, phase);
200 fs::read_to_string(&path).ok()?;
203 let pid = read_holder_pid(&path);
204 let pid = if pid == "unknown" { String::new() } else { pid };
205 if pid.is_empty() {
206 let _ = fs::remove_file(&path);
208 return None;
209 }
210 Some((pid, path))
211}
212
213fn release(path: &Path) {
216 let _ = fs::remove_file(path);
217}
218
219#[derive(Debug)]
221pub struct LockGuard {
222 path: PathBuf,
223}
224
225impl Drop for LockGuard {
226 fn drop(&mut self) {
227 release(&self.path);
228 }
229}
230
231const LOCK_FILE_PREFIX: &str = "lock-";
236
237pub(crate) fn lock_path(project_root: &Path, phase: PhaseId) -> PathBuf {
238 project_root.join(".devflow").join(format!(
239 "{LOCK_FILE_PREFIX}{padded}",
240 padded = phase.padded()
241 ))
242}
243
244pub(crate) fn project_lock_path(project_root: &Path) -> PathBuf {
245 project_root
246 .join(".devflow")
247 .join(format!("{LOCK_FILE_PREFIX}project"))
248}
249
250pub fn remove_stale_locks(project_root: &Path) -> Vec<String> {
259 let mut warnings = Vec::new();
260 let devflow_dir = project_root.join(".devflow");
261 let Ok(entries) = fs::read_dir(&devflow_dir) else {
262 return warnings;
263 };
264 for entry in entries.flatten() {
265 let name = entry.file_name();
266 let Some(name) = name.to_str() else { continue };
267 if !name.starts_with(LOCK_FILE_PREFIX) {
268 continue;
269 }
270 let path = entry.path();
271 let holder_pid = read_holder_pid(&path);
273 if pid_is_alive(&holder_pid) {
274 warnings.push(format!(
275 "kept {} — holder pid {holder_pid} is still alive",
276 path.display()
277 ));
278 continue;
279 }
280 if let Err(err) = fs::remove_file(&path) {
281 warnings.push(format!("could not remove {}: {err}", path.display()));
282 }
283 }
284 warnings
285}
286
287#[cfg(test)]
288mod tests {
289 use super::*;
290
291 #[test]
292 fn acquire_creates_lock_and_records_pid() {
293 let dir = tempfile::tempdir().unwrap();
294 let guard = acquire(dir.path(), PhaseId::new(1)).expect("acquire");
295
296 let (pid, path) = holder(dir.path(), PhaseId::new(1)).expect("holder present");
297 assert_eq!(pid, std::process::id().to_string());
298 assert!(path.exists());
299 drop(guard);
300 }
301
302 #[test]
303 fn acquire_creates_devflow_directory_when_absent() {
304 let dir = tempfile::tempdir().unwrap();
305 assert!(!dir.path().join(".devflow").exists());
306 let _guard = acquire(dir.path(), PhaseId::new(1)).expect("acquire");
307 assert!(dir.path().join(".devflow").exists());
308 }
309
310 #[test]
311 fn second_acquire_is_contended() {
312 let dir = tempfile::tempdir().unwrap();
313 let _guard = acquire(dir.path(), PhaseId::new(1)).expect("first acquire");
314
315 match acquire(dir.path(), PhaseId::new(1)) {
316 Err(LockError::Contended { pid, .. }) => {
317 assert_eq!(pid, std::process::id().to_string());
318 }
319 Ok(_) => panic!("second acquire must fail"),
320 Err(other) => panic!("expected Contended, got {other:?}"),
321 }
322 }
323
324 #[test]
329 fn different_phases_do_not_contend() {
330 let dir = tempfile::tempdir().unwrap();
331 let _guard_a = acquire(dir.path(), PhaseId::new(1)).expect("acquire phase 1");
332 let _guard_b =
333 acquire(dir.path(), PhaseId::new(2)).expect("acquire phase 2 must not contend");
334 }
335
336 #[test]
337 fn dropping_guard_releases_lock() {
338 let dir = tempfile::tempdir().unwrap();
339 {
340 let _guard = acquire(dir.path(), PhaseId::new(1)).expect("acquire");
341 assert!(holder(dir.path(), PhaseId::new(1)).is_some());
342 }
343 assert!(holder(dir.path(), PhaseId::new(1)).is_none());
345 let _again = acquire(dir.path(), PhaseId::new(1)).expect("re-acquire after release");
346 }
347
348 #[test]
349 fn holder_is_none_without_lock_file() {
350 let dir = tempfile::tempdir().unwrap();
351 assert!(holder(dir.path(), PhaseId::new(1)).is_none());
352 }
353
354 #[test]
355 fn holder_cleans_up_empty_lock_file() {
356 let dir = tempfile::tempdir().unwrap();
357 let path = lock_path(dir.path(), PhaseId::new(1));
358 fs::create_dir_all(path.parent().unwrap()).unwrap();
359 fs::write(&path, " \n").unwrap();
360
361 assert!(holder(dir.path(), PhaseId::new(1)).is_none());
362 assert!(!path.exists());
364 let _guard = acquire(dir.path(), PhaseId::new(1)).expect("acquire after stale cleanup");
365 }
366
367 #[test]
371 fn acquire_reclaims_lock_from_dead_holder() {
372 let dir = tempfile::tempdir().unwrap();
373 let path = lock_path(dir.path(), PhaseId::new(1));
374 fs::create_dir_all(path.parent().unwrap()).unwrap();
375 fs::write(&path, "9999999").unwrap();
377
378 let guard = acquire(dir.path(), PhaseId::new(1)).expect("stale lock must be reclaimed");
379 let (pid, _) = holder(dir.path(), PhaseId::new(1)).expect("holder present");
380 assert_eq!(pid, std::process::id().to_string());
381 drop(guard);
382 }
383
384 #[test]
385 fn acquire_reclaims_lock_with_corrupt_pid() {
386 let dir = tempfile::tempdir().unwrap();
387 let path = lock_path(dir.path(), PhaseId::new(1));
388 fs::create_dir_all(path.parent().unwrap()).unwrap();
389 fs::write(&path, "not-a-pid").unwrap();
390
391 acquire(dir.path(), PhaseId::new(1)).expect("corrupt lock must be reclaimed");
392 }
393
394 #[test]
398 fn remove_stale_locks_keeps_live_holder_and_sweeps_dead() {
399 let dir = tempfile::tempdir().unwrap();
400 let live = lock_path(dir.path(), PhaseId::new(1));
401 let dead = lock_path(dir.path(), PhaseId::new(2));
402 fs::create_dir_all(live.parent().unwrap()).unwrap();
403 fs::write(&live, std::process::id().to_string()).unwrap();
404 fs::write(&dead, "9999999").unwrap();
405
406 let warnings = remove_stale_locks(dir.path());
407
408 assert!(live.exists(), "live holder's lock must be kept");
409 assert!(!dead.exists(), "dead holder's lock must be swept");
410 assert_eq!(warnings.len(), 1, "keeping a live lock must be reported");
411 assert!(warnings[0].contains("still alive"));
412 }
413
414 #[test]
419 fn project_lock_is_independent_of_phase_locks() {
420 let dir = tempfile::tempdir().unwrap();
421 let _phase = acquire(dir.path(), PhaseId::new(1)).expect("phase lock");
422 let _project = acquire_project(dir.path()).expect("project lock must not contend");
423 }
424
425 #[test]
426 fn project_lock_contends_with_itself() {
427 let dir = tempfile::tempdir().unwrap();
428 let _held = acquire_project(dir.path()).expect("first acquire");
429 assert!(matches!(
430 acquire_project(dir.path()),
431 Err(LockError::Contended { .. })
432 ));
433 }
434
435 #[test]
438 fn project_lock_blocking_waits_for_release() {
439 let dir = tempfile::tempdir().unwrap();
440 let held = acquire_project(dir.path()).expect("first acquire");
441 let root = dir.path().to_path_buf();
442
443 std::thread::scope(|scope| {
444 let waiter = scope
445 .spawn(move || acquire_project_blocking(&root, std::time::Duration::from_secs(10)));
446 std::thread::sleep(std::time::Duration::from_millis(300));
449 drop(held);
450 waiter
451 .join()
452 .expect("waiter thread")
453 .expect("blocking acquire must succeed once the holder releases");
454 });
455 }
456
457 #[test]
458 fn project_lock_blocking_times_out_against_live_holder() {
459 let dir = tempfile::tempdir().unwrap();
460 let _held = acquire_project(dir.path()).expect("first acquire");
461 let err = acquire_project_blocking(dir.path(), std::time::Duration::from_millis(300))
462 .expect_err("must time out while the live holder keeps the lock");
463 assert!(matches!(err, LockError::Contended { .. }));
464 }
465
466 #[test]
471 fn acquire_reclaims_lock_with_pid_zero() {
472 let dir = tempfile::tempdir().unwrap();
473 let path = lock_path(dir.path(), PhaseId::new(1));
474 fs::create_dir_all(path.parent().unwrap()).unwrap();
475 fs::write(&path, "0").unwrap();
476
477 acquire(dir.path(), PhaseId::new(1)).expect("pid-0 lock must be reclaimed");
478 }
479}