1use std::path::{Path, PathBuf};
48use std::time::Duration;
49
50use anyhow::{Context as _, Result, bail};
51use jiff::Timestamp;
52use serde::{Deserialize, Serialize};
53
54use crate::proc::{self, Quiet as _};
55
56#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
58pub struct Owner {
59 pub run: String,
61 pub node: String,
63 pub seat: String,
66 pub pid: u32,
70 pub worktree: String,
73 pub head: String,
75}
76
77impl Owner {
78 #[must_use]
80 pub fn here(run: &str, node: &str, seat: &str, worktree: &Path, head: &str) -> Owner {
81 Owner {
82 run: run.to_owned(),
83 node: node.to_owned(),
84 seat: seat.to_owned(),
85 pid: std::process::id(),
86 worktree: worktree.display().to_string(),
87 head: head.to_owned(),
88 }
89 }
90}
91
92#[derive(Debug, Clone, Serialize, Deserialize)]
94struct LeaseFile {
95 cache_dir: String,
99 owner: Owner,
100 acquired_at: Timestamp,
101}
102
103#[derive(Debug, Clone, PartialEq, Eq)]
105pub enum Status {
106 Free,
108 Active(Owner),
110 Stale(Owner),
112 Unknown,
115}
116
117#[derive(Debug, Clone, PartialEq, Eq)]
119pub enum Busy {
120 Active(Owner),
122 Unknown,
124 Contended,
126}
127
128impl Busy {
129 #[must_use]
131 pub fn describe(&self) -> String {
132 match self {
133 Busy::Active(o) => format!(
134 "held by run {} node {} seat {} (pid {})",
135 o.run, o.node, o.seat, o.pid
136 ),
137 Busy::Unknown => {
138 "an unreadable lease is present; refusing to guess who holds it".to_owned()
139 }
140 Busy::Contended => "lost a race for the lease; retrying".to_owned(),
141 }
142 }
143}
144
145#[derive(Debug)]
149pub struct Guard {
150 path: PathBuf,
151 released: bool,
152}
153
154impl Guard {
155 fn new(path: PathBuf) -> Guard {
156 Guard {
157 path,
158 released: false,
159 }
160 }
161
162 pub fn release(mut self) {
166 self.do_release();
167 }
168
169 fn do_release(&mut self) {
170 if !self.released {
171 let _ = std::fs::remove_file(&self.path);
172 self.released = true;
173 }
174 }
175}
176
177impl Drop for Guard {
178 fn drop(&mut self) {
179 self.do_release();
180 }
181}
182
183fn leases_dir(home: &Path) -> PathBuf {
187 home.join("cache-leases")
188}
189
190fn slug(cache_dir: &Path) -> String {
195 let norm = normalize(cache_dir);
196 let mut readable: String = norm
197 .chars()
198 .map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
199 .collect();
200 readable.truncate(80);
201 use std::hash::{Hash, Hasher};
202 let mut hasher = std::collections::hash_map::DefaultHasher::new();
203 norm.hash(&mut hasher);
204 format!("{readable}-{:08x}", hasher.finish() as u32)
205}
206
207fn normalize(p: &Path) -> String {
213 std::fs::canonicalize(p)
214 .map(|p| p.display().to_string())
215 .unwrap_or_else(|_| p.display().to_string())
216 .replace('\\', "/")
217 .to_ascii_lowercase()
218}
219
220fn lease_path(home: &Path, cache_dir: &Path) -> PathBuf {
221 leases_dir(home).join(format!("{}.json", slug(cache_dir)))
222}
223
224fn identity_path(home: &Path, cache_dir: &Path) -> PathBuf {
225 leases_dir(home).join(format!("{}.identity.json", slug(cache_dir)))
226}
227
228fn catalog_path(home: &Path, cache_dir: &Path) -> PathBuf {
229 leases_dir(home).join(format!("{}.catalog.json", slug(cache_dir)))
230}
231
232#[derive(Debug, Clone, Serialize, Deserialize)]
249struct CatalogRecord {
250 cache_dir: String,
251 last_owner: Owner,
252 last_used_at: Timestamp,
253}
254
255fn record_catalog(home: &Path, cache_dir: &Path, owner: &Owner) {
261 let path = catalog_path(home, cache_dir);
262 let record = CatalogRecord {
263 cache_dir: cache_dir.display().to_string(),
264 last_owner: owner.clone(),
265 last_used_at: Timestamp::now(),
266 };
267 let Ok(body) = serde_json::to_string_pretty(&record) else {
268 return;
269 };
270 let tmp = path.with_extension("json.tmp");
271 if std::fs::write(&tmp, &body).is_ok() {
272 let _ = std::fs::rename(&tmp, &path);
273 }
274}
275
276fn read_catalog(path: &Path) -> Option<CatalogRecord> {
281 let body = std::fs::read_to_string(path).ok()?;
282 serde_json::from_str(&body).ok()
283}
284
285fn read_lease(path: &Path) -> Option<LeaseFile> {
288 let body = std::fs::read_to_string(path).ok()?;
289 serde_json::from_str(&body).ok()
290}
291
292fn peek_cache_dir(path: &Path) -> Option<String> {
297 let body = std::fs::read_to_string(path).ok()?;
298 let value: serde_json::Value = serde_json::from_str(&body).ok()?;
299 value
300 .get("cache_dir")
301 .and_then(|v| v.as_str())
302 .map(str::to_owned)
303}
304
305fn classify(path: &Path) -> Status {
307 classify_with(path, proc::pid_alive)
308}
309
310fn classify_with<F: Fn(u32) -> bool>(path: &Path, alive: F) -> Status {
319 if !path.exists() {
320 return Status::Free;
321 }
322 let Some(lease) = read_lease(path) else {
323 return Status::Unknown;
324 };
325 let this_process = std::process::id();
326 if lease.owner.pid == this_process || alive(lease.owner.pid) {
327 Status::Active(lease.owner)
328 } else {
329 Status::Stale(lease.owner)
330 }
331}
332
333fn write_new(path: &Path, cache_dir: &Path, owner: &Owner) -> std::io::Result<()> {
338 use std::io::Write as _;
339 let mut f = std::fs::OpenOptions::new()
340 .write(true)
341 .create_new(true)
342 .open(path)?;
343 let lease = LeaseFile {
344 cache_dir: cache_dir.display().to_string(),
345 owner: owner.clone(),
346 acquired_at: Timestamp::now(),
347 };
348 let body = serde_json::to_string_pretty(&lease).unwrap_or_default();
349 f.write_all(body.as_bytes())?;
350 Ok(())
351}
352
353pub enum AcquireOutcome {
355 Acquired(Guard),
357 Busy(Busy),
359}
360
361pub fn try_acquire(home: &Path, cache_dir: &Path, owner: &Owner) -> Result<AcquireOutcome> {
365 try_acquire_with(home, cache_dir, owner, proc::pid_alive)
366}
367
368fn try_acquire_with<F: Fn(u32) -> bool + Copy>(
371 home: &Path,
372 cache_dir: &Path,
373 owner: &Owner,
374 alive: F,
375) -> Result<AcquireOutcome> {
376 let dir = leases_dir(home);
377 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
378 let path = dir.join(format!("{}.json", slug(cache_dir)));
379
380 for _ in 0..2 {
386 match write_new(&path, cache_dir, owner) {
387 Ok(()) => {
388 record_catalog(home, cache_dir, owner);
389 return Ok(AcquireOutcome::Acquired(Guard::new(path)));
390 }
391 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {}
392 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
393 }
394 match classify_with(&path, alive) {
395 Status::Free => {} Status::Stale(_) => {
397 let _ = std::fs::remove_file(&path);
398 }
399 Status::Active(o) => return Ok(AcquireOutcome::Busy(Busy::Active(o))),
400 Status::Unknown => return Ok(AcquireOutcome::Busy(Busy::Unknown)),
401 }
402 }
403 Ok(AcquireOutcome::Busy(Busy::Contended))
404}
405
406#[must_use]
411pub fn in_use(home: &Path, cache_dir: &Path) -> bool {
412 let path = lease_path(home, cache_dir);
413 matches!(classify(&path), Status::Active(_) | Status::Unknown)
414}
415
416pub async fn wait_for(
425 home: &Path,
426 cache_dir: &Path,
427 owner: &Owner,
428 budget: Duration,
429 poll: Duration,
430) -> Result<Guard> {
431 let start = std::time::Instant::now();
432 loop {
433 match try_acquire(home, cache_dir, owner)? {
434 AcquireOutcome::Acquired(g) => return Ok(g),
435 AcquireOutcome::Busy(busy) => {
436 let elapsed = start.elapsed();
437 if elapsed >= budget {
438 bail!(
439 "timed out after {}s waiting for the build cache at {} ({})",
440 budget.as_secs(),
441 cache_dir.display(),
442 busy.describe()
443 );
444 }
445 tokio::time::sleep(poll.min(budget - elapsed)).await;
446 }
447 }
448 }
449}
450
451#[derive(Debug, Clone)]
457pub struct Entry {
458 pub cache_dir: String,
460 pub status: EntryStatus,
462}
463
464#[derive(Debug, Clone)]
466pub enum EntryStatus {
467 Active(Owner),
469 Stale(Owner),
471 Unknown,
473 Idle(Owner),
479}
480
481#[must_use]
489pub fn inventory(home: &Path) -> Vec<Entry> {
490 inventory_with(home, proc::pid_alive)
491}
492
493fn inventory_with<F: Fn(u32) -> bool + Copy>(home: &Path, alive: F) -> Vec<Entry> {
496 let dir = leases_dir(home);
497 let Ok(rd) = std::fs::read_dir(&dir) else {
498 return Vec::new();
499 };
500 let mut out = Vec::new();
501 for entry in rd.flatten() {
502 let path = entry.path();
503 let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
504 continue;
505 };
506 if !name.ends_with(".json")
507 || name.ends_with(".identity.json")
508 || name.ends_with(".catalog.json")
509 {
510 continue;
511 }
512 let status = match classify_with(&path, alive) {
513 Status::Free => continue,
514 Status::Active(o) => EntryStatus::Active(o),
515 Status::Stale(o) => EntryStatus::Stale(o),
516 Status::Unknown => EntryStatus::Unknown,
517 };
518 let cache_dir = peek_cache_dir(&path)
524 .unwrap_or_else(|| format!("(unreadable lease file: {})", path.display()));
525 out.push(Entry { cache_dir, status });
526 }
527
528 if let Ok(rd) = std::fs::read_dir(&dir) {
534 for entry in rd.flatten() {
535 let path = entry.path();
536 let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
537 continue;
538 };
539 let Some(stem) = name.strip_suffix(".catalog.json") else {
540 continue;
541 };
542 if dir.join(format!("{stem}.json")).exists() {
543 continue;
544 }
545 if let Some(record) = read_catalog(&path) {
546 out.push(Entry {
547 cache_dir: record.cache_dir,
548 status: EntryStatus::Idle(record.last_owner),
549 });
550 }
551 }
552 }
553
554 out.sort_by(|a, b| a.cache_dir.cmp(&b.cache_dir));
555 out
556}
557
558pub fn maintenance_prune(
567 home: &Path,
568 cache_dir: &Path,
569 limit: u64,
570) -> Result<Option<crate::disk::Prune>> {
571 let owner = Owner {
572 run: "maintenance".to_owned(),
573 node: "prune".to_owned(),
574 seat: "janitor".to_owned(),
575 pid: std::process::id(),
576 worktree: String::new(),
577 head: String::new(),
578 };
579 match try_acquire(home, cache_dir, &owner)? {
580 AcquireOutcome::Busy(_) => Ok(None),
581 AcquireOutcome::Acquired(guard) => {
582 let result = crate::disk::prune_dir(cache_dir, limit)?;
583 guard.release();
584 Ok(Some(result))
585 }
586 }
587}
588
589#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
593pub struct Identity {
594 pub worktree: String,
596 pub head: String,
598}
599
600impl Identity {
601 #[must_use]
603 pub fn new(worktree: &Path, head: &str) -> Identity {
604 Identity {
605 worktree: worktree.display().to_string(),
606 head: head.to_owned(),
607 }
608 }
609}
610
611#[must_use]
616pub fn needs_refresh(home: &Path, cache_dir: &Path, current: &Identity) -> bool {
617 let path = identity_path(home, cache_dir);
618 let Ok(body) = std::fs::read_to_string(path) else {
619 return true;
620 };
621 match serde_json::from_str::<Identity>(&body) {
622 Ok(recorded) => &recorded != current,
623 Err(_) => true,
624 }
625}
626
627pub fn record_identity(home: &Path, cache_dir: &Path, identity: &Identity) -> Result<()> {
629 let path = identity_path(home, cache_dir);
630 if let Some(parent) = path.parent() {
631 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
632 }
633 let body = serde_json::to_string_pretty(identity).context("serialize cache identity")?;
634 let tmp = path.with_extension("json.tmp");
635 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
636 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
637 Ok(())
638}
639
640pub fn invalidate_identity(home: &Path, cache_dir: &Path) {
654 let _ = std::fs::remove_file(identity_path(home, cache_dir));
655}
656
657#[must_use]
663pub fn parse_workspace_package_names(metadata_json: &str) -> Vec<String> {
664 let Ok(value) = serde_json::from_str::<serde_json::Value>(metadata_json) else {
665 return Vec::new();
666 };
667 value
668 .get("packages")
669 .and_then(|p| p.as_array())
670 .map(|packages| {
671 packages
672 .iter()
673 .filter_map(|p| p.get("name").and_then(|n| n.as_str()))
674 .map(str::to_owned)
675 .collect()
676 })
677 .unwrap_or_default()
678}
679
680#[must_use]
684pub fn extract_manifest_paths(commands: &[String]) -> Vec<PathBuf> {
685 const KEY: &str = "--manifest-path";
686 let mut out: Vec<PathBuf> = Vec::new();
687 for command in commands {
688 let mut rest = command.as_str();
689 while let Some(at) = rest.find(KEY) {
690 rest = &rest[at + KEY.len()..];
691 let value = if let Some(r) = rest.strip_prefix('=') {
692 r
693 } else if rest.starts_with(char::is_whitespace) {
694 rest.trim_start()
695 } else {
696 continue;
697 };
698 let value = if let Some(s) = value.strip_prefix('\'') {
699 s.split('\'').next().unwrap_or("")
700 } else if let Some(s) = value.strip_prefix('"') {
701 s.split('"').next().unwrap_or("")
702 } else {
703 let end = value.find(char::is_whitespace).unwrap_or(value.len());
704 &value[..end]
705 };
706 if !value.is_empty() {
707 let path = PathBuf::from(value);
708 if !out.contains(&path) {
709 out.push(path);
710 }
711 }
712 }
713 }
714 out
715}
716
717#[derive(Debug, Clone, PartialEq, Eq)]
719pub enum Target {
720 Root,
722 Manifest(PathBuf),
724}
725
726fn plan_targets(
731 worktree: &Path,
732 manifests: &[PathBuf],
733 exists: impl Fn(&Path) -> Result<bool>,
734) -> Result<Vec<Target>> {
735 let mut targets = Vec::new();
736 if exists(&worktree.join("Cargo.toml"))? {
737 targets.push(Target::Root);
738 }
739 for m in manifests {
740 let full = if m.is_absolute() {
741 m.clone()
742 } else {
743 worktree.join(m)
744 };
745 if exists(&full)? {
746 targets.push(Target::Manifest(full));
747 } else {
748 tracing::warn!(
749 manifest = %m.display(),
750 "build cache: --manifest-path from a verify command does not exist; skipped"
751 );
752 }
753 }
754 Ok(targets)
755}
756
757fn file_exists(path: &Path) -> Result<bool> {
760 match std::fs::metadata(path) {
761 Ok(m) => Ok(m.is_file()),
762 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
763 Err(e) => Err(e).with_context(|| format!("stat {}", path.display())),
764 }
765}
766
767fn refresh_stale_packages(
775 worktree: &Path,
776 cache_dir: &Path,
777 manifests: &[PathBuf],
778) -> Result<Option<Vec<String>>> {
779 let targets = plan_targets(worktree, manifests, file_exists)?;
780 if targets.is_empty() {
781 return Ok(None);
782 }
783 let mut all = Vec::new();
784 let mut failed = Vec::new();
785 for target in &targets {
786 let manifest_arg = match target {
787 Target::Root => None,
788 Target::Manifest(m) => Some(m),
789 };
790 let mut cmd = std::process::Command::new("cargo");
791 cmd.args(["metadata", "--no-deps", "--format-version", "1"]);
792 if let Some(m) = manifest_arg {
793 cmd.arg("--manifest-path").arg(m);
794 }
795 let meta = cmd
796 .current_dir(worktree)
797 .quiet()
798 .output()
799 .context("run `cargo metadata`")?;
800 if !meta.status.success() {
801 bail!(
802 "cargo metadata failed: {}",
803 String::from_utf8_lossy(&meta.stderr)
804 );
805 }
806 let names = parse_workspace_package_names(&String::from_utf8_lossy(&meta.stdout));
807 for name in &names {
808 let mut clean = std::process::Command::new("cargo");
809 clean.arg("clean").arg("-p").arg(name);
810 if let Some(m) = manifest_arg {
811 clean.arg("--manifest-path").arg(m);
812 }
813 let out = clean
814 .arg("--target-dir")
815 .arg(cache_dir)
816 .current_dir(worktree)
817 .quiet()
818 .output()
819 .with_context(|| format!("cargo clean -p {name}"))?;
820 if !out.status.success() {
821 failed.push(format!(
830 "{name}: {}",
831 String::from_utf8_lossy(&out.stderr).trim()
832 ));
833 }
834 }
835 for n in names {
836 if !all.contains(&n) {
837 all.push(n);
838 }
839 }
840 }
841 if !failed.is_empty() {
842 bail!(
843 "cargo clean -p failed for {} package(s): {}",
844 failed.len(),
845 failed.join("; ")
846 );
847 }
848 Ok(Some(all))
849}
850
851const CACHEDIR_TAG: &str = "Signature: 8a477f597d28d172789f06886806bc55\n\
853# This file is a cache directory tag created by cargo.\n\
854# For information about cache directory tags see https://bford.info/cachedir/\n";
855
856fn restore_cachedir_tag(cache_dir: &Path) -> Result<bool> {
865 use std::io::Write as _;
866 if !cache_dir.is_dir() || !cache_dir.join(".rustc_info.json").is_file() {
867 return Ok(false);
868 }
869 let tag = cache_dir.join("CACHEDIR.TAG");
870 match std::fs::OpenOptions::new()
871 .write(true)
872 .create_new(true)
873 .open(&tag)
874 {
875 Ok(mut f) => {
876 f.write_all(CACHEDIR_TAG.as_bytes())
877 .with_context(|| format!("write {}", tag.display()))?;
878 Ok(true)
879 }
880 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
881 Err(e) => Err(e).with_context(|| format!("create {}", tag.display())),
882 }
883}
884
885pub fn ensure_fresh(
896 home: &Path,
897 cache_dir: &Path,
898 identity: &Identity,
899 manifests: &[PathBuf],
900) -> Result<()> {
901 match restore_cachedir_tag(cache_dir) {
902 Ok(true) => tracing::info!(
903 cache = %cache_dir.display(),
904 "build cache: restored a missing CACHEDIR.TAG"
905 ),
906 Ok(false) => {}
907 Err(e) => tracing::warn!(error = %e, "build cache: could not restore CACHEDIR.TAG"),
909 }
910 if needs_refresh(home, cache_dir, identity) {
911 let worktree = PathBuf::from(&identity.worktree);
912 match refresh_stale_packages(&worktree, cache_dir, manifests)? {
913 Some(cleaned) => tracing::info!(
914 ?cleaned,
915 cache = %cache_dir.display(),
916 "build cache: source identity changed; cleaned the workspace's own packages before reuse"
917 ),
918 None => {
919 tracing::warn!(
924 cache = %cache_dir.display(),
925 "build cache: no Cargo manifest found for the freshness check; skipping the package clean"
926 );
927 invalidate_identity(home, cache_dir);
928 return Ok(());
929 }
930 }
931 }
932 record_identity(home, cache_dir, identity)
933}
934
935#[cfg(test)]
936mod tests {
937 use super::*;
938
939 fn owner(pid: u32) -> Owner {
940 Owner {
941 run: "r1".to_owned(),
942 node: "gate".to_owned(),
943 seat: "gate".to_owned(),
944 pid,
945 worktree: "/w".to_owned(),
946 head: "deadbeef".to_owned(),
947 }
948 }
949
950 #[test]
951 fn a_missing_cachedir_tag_is_restored_only_for_a_cargo_target() {
952 let t = tempfile::TempDir::new().expect("temp");
953 let cache = t.path().join("cache");
954 assert!(!restore_cachedir_tag(&cache).expect("missing dir"));
955 assert!(!cache.exists(), "a missing directory is left for cargo");
956
957 std::fs::create_dir_all(&cache).expect("mkdir");
958 assert!(!restore_cachedir_tag(&cache).expect("no rustc info"));
959 assert!(!cache.join("CACHEDIR.TAG").exists());
960
961 std::fs::write(cache.join(".rustc_info.json"), "{}").expect("info");
962 assert!(restore_cachedir_tag(&cache).expect("restore"));
963 let body = std::fs::read_to_string(cache.join("CACHEDIR.TAG")).expect("tag");
964 assert_eq!(
965 body.lines().next(),
966 Some("Signature: 8a477f597d28d172789f06886806bc55")
967 );
968
969 std::fs::write(cache.join("CACHEDIR.TAG"), "custom").expect("custom");
970 assert!(!restore_cachedir_tag(&cache).expect("existing"));
971 assert_eq!(
972 std::fs::read_to_string(cache.join("CACHEDIR.TAG")).expect("tag"),
973 "custom"
974 );
975 }
976
977 #[test]
978 fn an_uncontended_lease_is_acquired_and_freed_on_release() {
979 let home = tempfile::TempDir::new().expect("temp");
980 let cache = home.path().join("cache");
981 let this = std::process::id();
982 match try_acquire(home.path(), &cache, &owner(this)).expect("acquire") {
983 AcquireOutcome::Acquired(g) => {
984 assert!(in_use(home.path(), &cache), "held while the guard lives");
985 g.release();
986 }
987 AcquireOutcome::Busy(b) => panic!("unexpectedly busy: {b:?}"),
988 }
989 assert!(!in_use(home.path(), &cache), "freed after release");
990 }
991
992 #[test]
993 fn a_lease_held_by_a_live_pid_is_reported_active_and_refuses_a_second_acquire() {
994 let home = tempfile::TempDir::new().expect("temp");
995 let cache = home.path().join("cache");
996 let this = std::process::id();
997 let _first =
1001 try_acquire(home.path(), &cache, &owner(this)).expect("first acquire succeeds");
1002 let mut second_owner = owner(this);
1003 second_owner.run = "r2".to_owned();
1004 match try_acquire(home.path(), &cache, &second_owner).expect("no io error") {
1005 AcquireOutcome::Busy(Busy::Active(held_by)) => assert_eq!(held_by.run, "r1"),
1006 other => panic!("expected Busy::Active, got a different outcome: {other:?}"),
1007 }
1008 assert!(in_use(home.path(), &cache));
1009 }
1010
1011 impl std::fmt::Debug for AcquireOutcome {
1012 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1013 match self {
1014 AcquireOutcome::Acquired(_) => write!(f, "Acquired"),
1015 AcquireOutcome::Busy(b) => write!(f, "Busy({b:?})"),
1016 }
1017 }
1018 }
1019
1020 #[test]
1021 fn a_lease_whose_pid_is_gone_is_stale_and_reclaimed_by_the_next_acquirer() {
1022 let home = tempfile::TempDir::new().expect("temp");
1023 let cache = home.path().join("cache");
1024 let dead_owner = owner(999_999);
1030 let path = lease_path(home.path(), &cache);
1031 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1032 write_new(&path, &cache, &dead_owner).expect("seed a stale lease");
1033 assert_eq!(
1034 classify_with(&path, |_| false),
1035 Status::Stale(dead_owner.clone())
1036 );
1037
1038 match try_acquire_with(home.path(), &cache, &owner(std::process::id()), |_| false)
1039 .expect("acquire")
1040 {
1041 AcquireOutcome::Acquired(_) => {}
1042 other => panic!("stale lease should have been reclaimed: {other:?}"),
1043 }
1044 }
1045
1046 #[test]
1047 fn an_unreadable_lease_is_unknown_and_never_reclaimed() {
1048 let home = tempfile::TempDir::new().expect("temp");
1049 let cache = home.path().join("cache");
1050 let path = lease_path(home.path(), &cache);
1051 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
1052 std::fs::write(&path, b"not json").unwrap();
1053 assert_eq!(classify(&path), Status::Unknown);
1054 assert!(in_use(home.path(), &cache), "unknown counts as in use");
1055 match try_acquire(home.path(), &cache, &owner(std::process::id())).expect("no io error") {
1056 AcquireOutcome::Busy(Busy::Unknown) => {}
1057 other => panic!("expected Busy::Unknown, got {other:?}"),
1058 }
1059 }
1060
1061 #[tokio::test]
1062 async fn waiting_for_a_busy_lease_times_out_within_its_own_budget() {
1063 let home = tempfile::TempDir::new().expect("temp");
1064 let cache = home.path().join("cache");
1065 let _held = try_acquire(home.path(), &cache, &owner(std::process::id()))
1066 .expect("acquire")
1067 .pipe();
1068 let mut waiter = owner(std::process::id());
1069 waiter.run = "r2".to_owned();
1070 let started = std::time::Instant::now();
1071 let err = wait_for(
1072 home.path(),
1073 &cache,
1074 &waiter,
1075 Duration::from_millis(150),
1076 Duration::from_millis(20),
1077 )
1078 .await
1079 .expect_err("still held, must time out");
1080 assert!(started.elapsed() < Duration::from_secs(2), "bounded wait");
1081 assert!(
1082 err.to_string().contains("r1"),
1083 "names the current holder: {err}"
1084 );
1085 }
1086
1087 #[tokio::test]
1088 async fn a_wait_succeeds_as_soon_as_the_lease_is_released() {
1089 let home = tempfile::TempDir::new().expect("temp");
1090 let cache = home.path().join("cache");
1091 let guard =
1092 match try_acquire(home.path(), &cache, &owner(std::process::id())).expect("acquire") {
1093 AcquireOutcome::Acquired(g) => g,
1094 AcquireOutcome::Busy(b) => panic!("unexpectedly busy: {b:?}"),
1095 };
1096 let home_path = home.path().to_path_buf();
1097 let cache_path = cache.clone();
1098 let mut waiter = owner(std::process::id());
1099 waiter.run = "r2".to_owned();
1100 let wait = tokio::spawn(async move {
1101 wait_for(
1102 &home_path,
1103 &cache_path,
1104 &waiter,
1105 Duration::from_secs(5),
1106 Duration::from_millis(10),
1107 )
1108 .await
1109 });
1110 tokio::time::sleep(Duration::from_millis(50)).await;
1111 guard.release();
1112 let acquired = wait.await.expect("task").expect("acquire after release");
1113 acquired.release();
1114 }
1115
1116 #[test]
1117 fn inventory_reports_active_stale_and_unknown_but_not_free() {
1118 let home = tempfile::TempDir::new().expect("temp");
1119 let active_cache = home.path().join("active");
1120 let stale_cache = home.path().join("stale");
1121 let unknown_cache = home.path().join("unknown");
1122
1123 let _held =
1124 try_acquire(home.path(), &active_cache, &owner(std::process::id())).expect("acquire");
1125 let stale_path = lease_path(home.path(), &stale_cache);
1126 std::fs::create_dir_all(stale_path.parent().unwrap()).unwrap();
1127 write_new(&stale_path, &stale_cache, &owner(999_999)).unwrap();
1128 let unknown_path = lease_path(home.path(), &unknown_cache);
1129 std::fs::write(&unknown_path, b"garbage").unwrap();
1130
1131 let entries = inventory_with(home.path(), |pid| pid == std::process::id());
1136 assert_eq!(entries.len(), 3, "{entries:?}");
1137 let by_dir = |dir: &Path| {
1138 entries
1139 .iter()
1140 .find(|e| e.cache_dir == dir.display().to_string())
1141 .unwrap_or_else(|| panic!("no entry for {}", dir.display()))
1142 };
1143 assert!(matches!(
1144 by_dir(&active_cache).status,
1145 EntryStatus::Active(_)
1146 ));
1147 assert!(matches!(by_dir(&stale_cache).status, EntryStatus::Stale(_)));
1148 let unknown = entries
1153 .iter()
1154 .find(|e| matches!(e.status, EntryStatus::Unknown))
1155 .unwrap_or_else(|| panic!("no Unknown entry: {entries:?}"));
1156 assert!(
1157 unknown
1158 .cache_dir
1159 .contains(&unknown_path.display().to_string()),
1160 "{unknown:?}"
1161 );
1162 }
1163
1164 #[test]
1165 fn a_released_lease_is_reported_idle_from_the_catalog_not_dropped_entirely() {
1166 let home = tempfile::TempDir::new().expect("temp");
1167 let cache = home.path().join("cache");
1168 let this = std::process::id();
1169
1170 match try_acquire(home.path(), &cache, &owner(this)).expect("acquire") {
1171 AcquireOutcome::Acquired(g) => g.release(),
1172 AcquireOutcome::Busy(b) => panic!("unexpectedly busy: {b:?}"),
1173 }
1174
1175 assert!(!in_use(home.path(), &cache));
1180 let entries = inventory_with(home.path(), |pid| pid == this);
1181 let entry = entries
1182 .iter()
1183 .find(|e| e.cache_dir == cache.display().to_string())
1184 .unwrap_or_else(|| panic!("no entry for a released cache: {entries:?}"));
1185 match &entry.status {
1186 EntryStatus::Idle(o) => assert_eq!(o.run, "r1"),
1187 other => panic!("expected Idle, got {other:?}"),
1188 }
1189 }
1190
1191 #[test]
1192 fn reacquiring_a_released_cache_reports_active_not_idle() {
1193 let home = tempfile::TempDir::new().expect("temp");
1194 let cache = home.path().join("cache");
1195 let this = std::process::id();
1196 match try_acquire(home.path(), &cache, &owner(this)).expect("acquire") {
1197 AcquireOutcome::Acquired(g) => g.release(),
1198 AcquireOutcome::Busy(b) => panic!("unexpectedly busy: {b:?}"),
1199 }
1200 let _held = try_acquire(home.path(), &cache, &owner(this)).expect("reacquire");
1201 let entries = inventory_with(home.path(), |pid| pid == this);
1202 assert_eq!(
1203 entries.len(),
1204 1,
1205 "the catalog row must not duplicate the live lease: {entries:?}"
1206 );
1207 assert!(matches!(entries[0].status, EntryStatus::Active(_)));
1208 }
1209
1210 #[test]
1211 fn maintenance_prune_refuses_a_cache_a_live_owner_holds() {
1212 let home = tempfile::TempDir::new().expect("temp");
1213 let cache = home.path().join("cache");
1214 std::fs::create_dir_all(&cache).unwrap();
1215 std::fs::write(cache.join("big"), vec![0u8; 100]).unwrap();
1216 let _held = try_acquire(home.path(), &cache, &owner(std::process::id())).expect("acquire");
1217
1218 let result = maintenance_prune(home.path(), &cache, 1).expect("no io error");
1219 assert!(
1220 result.is_none(),
1221 "must not prune while a live owner holds it"
1222 );
1223 assert!(cache.join("big").exists(), "nothing was deleted");
1224 }
1225
1226 #[test]
1227 fn maintenance_prune_acts_once_the_cache_is_free_and_releases_after() {
1228 let home = tempfile::TempDir::new().expect("temp");
1229 let cache = home.path().join("cache");
1230 std::fs::create_dir_all(&cache).unwrap();
1231 std::fs::write(cache.join("big"), vec![0u8; 100]).unwrap();
1232
1233 let pruned = maintenance_prune(home.path(), &cache, 1)
1234 .expect("no io error")
1235 .expect("cache was free");
1236 assert!(pruned.freed > 0);
1237 assert!(
1238 !in_use(home.path(), &cache),
1239 "the maintenance lease was released"
1240 );
1241 }
1242
1243 #[test]
1244 fn identity_drift_is_detected_once_and_then_settles() {
1245 let home = tempfile::TempDir::new().expect("temp");
1246 let cache = home.path().join("cache");
1247 let a = Identity {
1248 worktree: "/w/a".to_owned(),
1249 head: "aaaa".to_owned(),
1250 };
1251 let b = Identity {
1252 worktree: "/w/b".to_owned(),
1253 head: "bbbb".to_owned(),
1254 };
1255 assert!(
1256 needs_refresh(home.path(), &cache, &a),
1257 "nothing recorded yet"
1258 );
1259 record_identity(home.path(), &cache, &a).expect("record");
1260 assert!(
1261 !needs_refresh(home.path(), &cache, &a),
1262 "same identity, no refresh needed"
1263 );
1264 assert!(needs_refresh(home.path(), &cache, &b), "different source");
1265 record_identity(home.path(), &cache, &b).expect("record");
1266 assert!(!needs_refresh(home.path(), &cache, &b));
1267 }
1268
1269 #[test]
1270 fn invalidating_forgets_a_recorded_identity_so_the_next_check_refreshes() {
1271 let home = tempfile::TempDir::new().expect("temp");
1272 let cache = home.path().join("cache");
1273 let a = Identity {
1274 worktree: "/w/a".to_owned(),
1275 head: "aaaa".to_owned(),
1276 };
1277 record_identity(home.path(), &cache, &a).expect("record");
1278 assert!(!needs_refresh(home.path(), &cache, &a));
1279
1280 invalidate_identity(home.path(), &cache);
1285 assert!(
1286 needs_refresh(home.path(), &cache, &a),
1287 "invalidation must not be skippable by asking about the same identity again"
1288 );
1289
1290 invalidate_identity(home.path(), &home.path().join("never-recorded"));
1293 }
1294
1295 #[test]
1296 fn workspace_package_names_are_read_from_cargo_metadata_json() {
1297 let fixture = r#"{
1298 "packages": [
1299 {"name": "magi", "version": "0.1.0"},
1300 {"name": "magi-cli", "version": "0.1.0"}
1301 ],
1302 "workspace_members": []
1303 }"#;
1304 let mut names = parse_workspace_package_names(fixture);
1305 names.sort();
1306 assert_eq!(names, vec!["magi".to_owned(), "magi-cli".to_owned()]);
1307 assert_eq!(
1308 parse_workspace_package_names("not json"),
1309 Vec::<String>::new()
1310 );
1311 assert_eq!(parse_workspace_package_names("{}"), Vec::<String>::new());
1312 }
1313
1314 #[test]
1315 fn slugs_are_stable_and_filesystem_safe() {
1316 let a = slug(Path::new(r"C:\Users\op\Temp\magi-target"));
1317 let b = slug(Path::new(r"C:\Users\op\Temp\magi-target"));
1318 assert_eq!(a, b, "same input, same slug");
1319 assert!(
1320 a.chars()
1321 .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_'),
1322 "filesystem-safe: {a}"
1323 );
1324 }
1325
1326 #[test]
1327 fn busy_active_describes_the_holder() {
1328 let b = Busy::Active(owner(123));
1329 let s = b.describe();
1330 assert!(
1331 s.contains("r1") && s.contains("gate") && s.contains("123"),
1332 "{s}"
1333 );
1334 }
1335
1336 trait Pipe: Sized {
1337 fn pipe(self) -> Guard;
1338 }
1339 impl Pipe for AcquireOutcome {
1340 fn pipe(self) -> Guard {
1341 match self {
1342 AcquireOutcome::Acquired(g) => g,
1343 AcquireOutcome::Busy(b) => panic!("expected Acquired, got Busy({b:?})"),
1344 }
1345 }
1346 }
1347
1348 fn cmds(v: &[&str]) -> Vec<String> {
1349 v.iter().map(|s| (*s).to_owned()).collect()
1350 }
1351
1352 #[test]
1353 fn manifest_paths_are_read_in_every_spelling() {
1354 let got = extract_manifest_paths(&cmds(&[
1355 "cargo test --manifest-path tools/a/Cargo.toml",
1356 "cargo build --manifest-path=tools/b/Cargo.toml --release",
1357 "cargo fmt --manifest-path 'q q/Cargo.toml'",
1358 "cargo clippy --manifest-path \"d/Cargo.toml\" && cargo test --manifest-path tools/a/Cargo.toml",
1359 ]));
1360 let want: Vec<PathBuf> = [
1361 "tools/a/Cargo.toml",
1362 "tools/b/Cargo.toml",
1363 "q q/Cargo.toml",
1364 "d/Cargo.toml",
1365 ]
1366 .iter()
1367 .map(PathBuf::from)
1368 .collect();
1369 assert_eq!(got, want);
1370 assert!(extract_manifest_paths(&cmds(&["cargo test", "cd x && cargo test"])).is_empty());
1371 }
1372
1373 #[test]
1374 fn plan_targets_covers_blog_root_both_and_none() {
1375 let w = Path::new("/w");
1376 let only = |set: &'static [&'static str]| {
1377 move |p: &Path| Ok(set.iter().any(|s| p == Path::new(s)))
1378 };
1379 let m = vec![PathBuf::from("tools/postprocess/Cargo.toml")];
1380 let blog = plan_targets(w, &m, only(&["/w/tools/postprocess/Cargo.toml"])).unwrap();
1381 assert_eq!(
1382 blog,
1383 vec![Target::Manifest(PathBuf::from(
1384 "/w/tools/postprocess/Cargo.toml"
1385 ))]
1386 );
1387 let root = plan_targets(w, &[], only(&["/w/Cargo.toml"])).unwrap();
1388 assert_eq!(root, vec![Target::Root]);
1389 let both = plan_targets(
1390 w,
1391 &m,
1392 only(&["/w/Cargo.toml", "/w/tools/postprocess/Cargo.toml"]),
1393 )
1394 .unwrap();
1395 assert_eq!(both.len(), 2);
1396 assert_eq!(both[0], Target::Root);
1397 assert!(plan_targets(w, &m, only(&[])).unwrap().is_empty());
1398 assert!(plan_targets(w, &[], only(&[])).unwrap().is_empty());
1399 }
1400
1401 #[test]
1402 fn a_repo_without_a_reachable_manifest_passes_untouched() {
1403 let home = tempfile::TempDir::new().expect("temp");
1404 let wt = tempfile::TempDir::new().expect("temp");
1405 std::fs::create_dir_all(wt.path().join("sub")).unwrap();
1407 std::fs::write(wt.path().join("sub/Cargo.toml"), "").unwrap();
1408 let cache = home.path().join("cache");
1409 std::fs::create_dir_all(&cache).unwrap();
1410 std::fs::write(cache.join("artifact"), "x").unwrap();
1411 let old = Identity {
1412 worktree: "/w/old".to_owned(),
1413 head: "old".to_owned(),
1414 };
1415 record_identity(home.path(), &cache, &old).unwrap();
1416 let id = Identity::new(wt.path(), "head");
1417 ensure_fresh(home.path(), &cache, &id, &[]).expect("must not fail");
1418 assert_eq!(
1419 std::fs::read_to_string(cache.join("artifact")).unwrap(),
1420 "x"
1421 );
1422 assert!(
1423 needs_refresh(home.path(), &cache, &old),
1424 "identity was invalidated, not recorded"
1425 );
1426 assert!(needs_refresh(home.path(), &cache, &id));
1427 }
1428}