1use fs2::FileExt as _;
2use std::{
3 fs::{self, File, OpenOptions},
4 io,
5 path::{Component, Path, PathBuf},
6 sync::Arc,
7 thread,
8 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
9};
10
11use super::digest::write_atomic;
12
13pub(super) fn canonicalize_allow_missing(path: &Path) -> io::Result<PathBuf> {
16 let absolute = if path.is_absolute() {
17 path.to_owned()
18 } else {
19 std::env::current_dir()?.join(path)
20 };
21 let mut resolved = PathBuf::new();
22 let mut missing_depth = 0_usize;
23 for component in absolute.components() {
24 match component {
25 Component::Prefix(_) | Component::RootDir | Component::Normal(_) => {
26 let candidate = resolved.join(component.as_os_str());
27 if missing_depth == 0 && matches!(component, Component::Normal(_)) {
28 match candidate.canonicalize() {
29 Ok(canonical) => resolved = canonical,
30 Err(error) if error.kind() == io::ErrorKind::NotFound => {
31 resolved = candidate;
32 missing_depth = 1;
33 }
34 Err(error) => return Err(error),
35 }
36 } else {
37 resolved = candidate;
38 if matches!(component, Component::Normal(_)) && missing_depth > 0 {
39 missing_depth += 1;
40 }
41 }
42 }
43 Component::CurDir => {}
44 Component::ParentDir => {
45 resolved.pop();
46 missing_depth = missing_depth.saturating_sub(1);
47 }
48 }
49 }
50 Ok(resolved)
51}
52
53const CACHE_DIRECTORY_TAG: &str = "Signature: 8a477f597d28d172789f06886806bc55\n\
54# This file is a cache directory tag created by ic-testkit.\n\
55# For information about cache directory tags see https://bford.info/cachedir/\n";
56pub(super) const CACHE_DIRECTORY_TAG_SIGNATURE: &str =
57 "Signature: 8a477f597d28d172789f06886806bc55\n";
58pub(super) const LAST_USED_FILE: &str = ".ic-testkit-last-used";
59const LAST_MAINTENANCE_FILE: &str = ".ic-testkit-last-maintenance";
60pub(super) const RETENTION_LOCK_FILE: &str = ".ic-testkit-retention-v1";
61
62#[derive(Clone, Debug)]
65pub(super) struct RetainedCacheEntry {
66 path: PathBuf,
67 _lock: Arc<File>,
68}
69
70impl PartialEq for RetainedCacheEntry {
71 fn eq(&self, other: &Self) -> bool {
72 self.path == other.path
73 }
74}
75
76impl Eq for RetainedCacheEntry {}
77
78impl RetainedCacheEntry {
79 pub(super) fn acquire(path: &Path) -> Result<Self, CacheFsError> {
80 let file = open_cache_lock_file(&path.join(RETENTION_LOCK_FILE))?;
81 fs2::FileExt::lock_shared(&file).map_err(|source| CacheFsError {
82 operation: "retain cache entry",
83 path: path.to_owned(),
84 source,
85 })?;
86 Ok(Self {
87 path: path.to_owned(),
88 _lock: Arc::new(file),
89 })
90 }
91}
92
93pub(super) fn remove_unretained_entry(path: &Path) -> Result<(), CacheFsError> {
95 let metadata = match fs::symlink_metadata(path) {
96 Ok(metadata) => metadata,
97 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
98 Err(source) => {
99 return Err(CacheFsError {
100 operation: "inspect cache entry",
101 path: path.to_owned(),
102 source,
103 });
104 }
105 };
106 if !metadata.is_dir() {
107 return remove_path_if_present(path).map_err(|source| CacheFsError {
108 operation: "remove invalid cache entry",
109 path: path.to_owned(),
110 source,
111 });
112 }
113 let _lock =
114 try_lock_cache_file(&path.join(RETENTION_LOCK_FILE))?.ok_or_else(|| CacheFsError {
115 operation: "replace retained cache entry",
116 path: path.to_owned(),
117 source: io::Error::new(
118 io::ErrorKind::WouldBlock,
119 "cache entry is retained by a consumer",
120 ),
121 })?;
122 remove_path_if_present(path).map_err(|source| CacheFsError {
123 operation: "remove cache entry",
124 path: path.to_owned(),
125 source,
126 })
127}
128
129#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
137pub struct ArtifactCachePrunePolicy {
138 max_age: Option<Duration>,
139 max_size_bytes: Option<u64>,
140}
141
142#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
144pub struct ArtifactCachePruneReport {
145 entries_scanned: usize,
146 entries_removed: usize,
147 bytes_before: u64,
148 bytes_removed: u64,
149 uncommitted_directories_removed: usize,
150 uncommitted_bytes_removed: u64,
151}
152
153#[non_exhaustive]
155#[derive(Clone, Debug, Eq, PartialEq)]
156pub enum ArtifactCacheMaintenance {
157 Pruned(ArtifactCachePruneReport),
159 PruneFailed {
161 message: String,
163 },
164}
165
166impl ArtifactCachePrunePolicy {
167 #[must_use]
169 pub const fn new() -> Self {
170 Self {
171 max_age: None,
172 max_size_bytes: None,
173 }
174 }
175
176 #[must_use]
178 pub const fn with_max_age(mut self, max_age: Duration) -> Self {
179 self.max_age = Some(max_age);
180 self
181 }
182
183 #[must_use]
185 pub const fn with_max_size_bytes(mut self, bytes: u64) -> Self {
186 self.max_size_bytes = Some(bytes);
187 self
188 }
189
190 #[must_use]
192 pub const fn max_age(self) -> Option<Duration> {
193 self.max_age
194 }
195
196 #[must_use]
198 pub const fn max_size_bytes(self) -> Option<u64> {
199 self.max_size_bytes
200 }
201
202 pub(super) fn maintenance_identity(self) -> String {
203 format!(
204 "age={:?};size={:?}",
205 self.max_age.map(|duration| duration.as_nanos()),
206 self.max_size_bytes
207 )
208 }
209}
210
211impl ArtifactCachePruneReport {
212 #[must_use]
214 pub const fn entries_scanned(self) -> usize {
215 self.entries_scanned
216 }
217
218 #[must_use]
220 pub const fn entries_removed(self) -> usize {
221 self.entries_removed
222 }
223
224 #[must_use]
226 pub const fn entries_retained(self) -> usize {
227 self.entries_scanned.saturating_sub(self.entries_removed)
228 }
229
230 #[must_use]
232 pub const fn bytes_before(self) -> u64 {
233 self.bytes_before
234 }
235
236 #[must_use]
238 pub const fn bytes_removed(self) -> u64 {
239 self.bytes_removed
240 }
241
242 #[must_use]
244 pub const fn bytes_retained(self) -> u64 {
245 self.bytes_before.saturating_sub(self.bytes_removed)
246 }
247
248 #[must_use]
250 pub const fn uncommitted_directories_removed(self) -> usize {
251 self.uncommitted_directories_removed
252 }
253
254 #[must_use]
256 pub const fn uncommitted_bytes_removed(self) -> u64 {
257 self.uncommitted_bytes_removed
258 }
259
260 pub(super) const fn record_uncommitted_removal(&mut self, bytes: u64) {
261 self.uncommitted_directories_removed += 1;
262 self.uncommitted_bytes_removed = self.uncommitted_bytes_removed.saturating_add(bytes);
263 }
264}
265
266impl ArtifactCacheMaintenance {
267 #[must_use]
269 pub const fn prune_report(&self) -> Option<ArtifactCachePruneReport> {
270 match self {
271 Self::Pruned(report) => Some(*report),
272 Self::PruneFailed { .. } => None,
273 }
274 }
275
276 #[must_use]
278 pub fn failure_message(&self) -> Option<&str> {
279 match self {
280 Self::Pruned(_) => None,
281 Self::PruneFailed { message } => Some(message),
282 }
283 }
284}
285
286#[derive(Debug)]
287pub(super) struct CacheFsError {
288 pub(super) operation: &'static str,
289 pub(super) path: PathBuf,
290 pub(super) source: io::Error,
291}
292
293pub(super) fn ensure_cache_directory_tag(cache_root: &Path) -> Result<(), CacheFsError> {
294 let path = cache_root.join("CACHEDIR.TAG");
295 if fs::read_to_string(&path)
296 .is_ok_and(|contents| contents.starts_with(CACHE_DIRECTORY_TAG_SIGNATURE))
297 {
298 return Ok(());
299 }
300 write_atomic(&path, CACHE_DIRECTORY_TAG.as_bytes()).map_err(|source| CacheFsError {
301 operation: "write cache directory tag",
302 path,
303 source,
304 })
305}
306
307pub(super) fn lock_cache_file(path: &Path) -> Result<(File, Duration), CacheFsError> {
308 let file = open_cache_lock_file(path)?;
309 let started = Instant::now();
310 file.lock_exclusive().map_err(|source| CacheFsError {
311 operation: "lock cache",
312 path: path.to_owned(),
313 source,
314 })?;
315 Ok((file, started.elapsed()))
316}
317
318pub(super) fn lock_cache_file_with_wait_observer(
319 path: &Path,
320 poll_interval: Duration,
321 mut observer: impl FnMut(Duration),
322) -> Result<(File, Duration), CacheFsError> {
323 let file = open_cache_lock_file(path)?;
324 let started = Instant::now();
325 loop {
326 match file.try_lock_exclusive() {
327 Ok(()) => return Ok((file, started.elapsed())),
328 Err(error) if error.kind() == io::ErrorKind::WouldBlock => {
329 observer(started.elapsed());
330 thread::sleep(poll_interval.min(Duration::from_millis(25)));
331 }
332 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
333 Err(source) => {
334 return Err(CacheFsError {
335 operation: "try lock cache",
336 path: path.to_owned(),
337 source,
338 });
339 }
340 }
341 }
342}
343
344pub(super) fn try_lock_cache_file(path: &Path) -> Result<Option<File>, CacheFsError> {
345 let file = open_cache_lock_file(path)?;
346 match file.try_lock_exclusive() {
347 Ok(()) => Ok(Some(file)),
348 Err(error) if error.kind() == io::ErrorKind::WouldBlock => Ok(None),
349 Err(source) => Err(CacheFsError {
350 operation: "try lock cache",
351 path: path.to_owned(),
352 source,
353 }),
354 }
355}
356
357fn open_cache_lock_file(path: &Path) -> Result<File, CacheFsError> {
358 if let Some(parent) = path.parent() {
359 fs::create_dir_all(parent).map_err(|source| CacheFsError {
360 operation: "create cache lock directory",
361 path: parent.to_owned(),
362 source,
363 })?;
364 }
365 OpenOptions::new()
366 .create(true)
367 .read(true)
368 .write(true)
369 .truncate(false)
370 .open(path)
371 .map_err(|source| CacheFsError {
372 operation: "open cache lock",
373 path: path.to_owned(),
374 source,
375 })
376}
377
378pub(super) fn record_cache_entry_use(path: &Path) -> Result<(), CacheFsError> {
379 write_last_used(path, SystemTime::now())
380}
381
382pub(super) fn cache_maintenance_due(
383 path: &Path,
384 minimum_interval: Option<Duration>,
385 maintenance_identity: &str,
386) -> Result<bool, CacheFsError> {
387 let Some(minimum_interval) = minimum_interval else {
388 return Ok(true);
389 };
390 let marker = path.join(LAST_MAINTENANCE_FILE);
391 let contents = match fs::read_to_string(&marker) {
392 Ok(contents) => contents,
393 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(true),
394 Err(source) => {
395 return Err(CacheFsError {
396 operation: "read cache maintenance time",
397 path: marker,
398 source,
399 });
400 }
401 };
402 let mut lines = contents.lines();
403 let Some(last_maintenance) = lines.next().and_then(decode_system_time) else {
404 return Ok(true);
405 };
406 if lines.next() != Some(maintenance_identity) {
407 return Ok(true);
408 }
409 Ok(match SystemTime::now().duration_since(last_maintenance) {
410 Ok(elapsed) => elapsed >= minimum_interval,
411 Err(_) => true,
412 })
413}
414
415pub(super) fn record_cache_maintenance(
416 path: &Path,
417 maintenance_identity: &str,
418) -> Result<(), CacheFsError> {
419 fs::create_dir_all(path).map_err(|source| CacheFsError {
420 operation: "create cache maintenance directory",
421 path: path.to_owned(),
422 source,
423 })?;
424 let marker = path.join(LAST_MAINTENANCE_FILE);
425 let elapsed = encode_system_time(&marker, SystemTime::now())?;
426 let contents = format!("{}\n{maintenance_identity}\n", elapsed.as_nanos());
427 write_atomic(&marker, contents.as_bytes()).map_err(|source| CacheFsError {
428 operation: "record cache maintenance time",
429 path: marker,
430 source,
431 })
432}
433
434pub(super) fn perform_scheduled_cache_maintenance(
435 path: &Path,
436 minimum_interval: Option<Duration>,
437 maintenance_identity: &str,
438 maintenance: impl FnOnce() -> Result<ArtifactCachePruneReport, String>,
439) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
440 let started = Instant::now();
441 match cache_maintenance_due(path, minimum_interval, maintenance_identity) {
442 Ok(false) => return (None, Some(started.elapsed())),
443 Ok(true) => {}
444 Err(error) => {
445 return (
446 Some(ArtifactCacheMaintenance::PruneFailed {
447 message: error.to_string(),
448 }),
449 Some(started.elapsed()),
450 );
451 }
452 }
453
454 let result = maintenance();
455 let marker = record_cache_maintenance(path, maintenance_identity);
456 let outcome = match (result, marker) {
457 (Ok(report), Ok(())) => ArtifactCacheMaintenance::Pruned(report),
458 (Err(message), Ok(())) => ArtifactCacheMaintenance::PruneFailed { message },
459 (Ok(_), Err(error)) => ArtifactCacheMaintenance::PruneFailed {
460 message: error.to_string(),
461 },
462 (Err(message), Err(marker)) => ArtifactCacheMaintenance::PruneFailed {
463 message: format!(
464 "{message}; additionally failed to record the maintenance attempt: {marker}"
465 ),
466 },
467 };
468 (Some(outcome), Some(started.elapsed()))
469}
470
471pub(super) fn write_last_used(path: &Path, last_used: SystemTime) -> Result<(), CacheFsError> {
472 let marker = path.join(LAST_USED_FILE);
473 write_system_time(&marker, last_used, "record cache use time")
474}
475
476fn write_system_time(
477 path: &Path,
478 timestamp: SystemTime,
479 operation: &'static str,
480) -> Result<(), CacheFsError> {
481 let elapsed = encode_system_time(path, timestamp)?;
482 write_atomic(path, elapsed.as_nanos().to_string().as_bytes()).map_err(|source| CacheFsError {
483 operation,
484 path: path.to_owned(),
485 source,
486 })
487}
488
489fn encode_system_time(path: &Path, timestamp: SystemTime) -> Result<Duration, CacheFsError> {
490 timestamp
491 .duration_since(UNIX_EPOCH)
492 .map_err(|source| CacheFsError {
493 operation: "encode cache time",
494 path: path.to_owned(),
495 source: io::Error::new(io::ErrorKind::InvalidInput, source),
496 })
497}
498
499fn decode_system_time(contents: &str) -> Option<SystemTime> {
500 let nanoseconds = contents.parse::<u128>().ok()?;
501 let seconds = u64::try_from(nanoseconds / 1_000_000_000).ok()?;
502 let subsecond_nanos = (nanoseconds % 1_000_000_000) as u32;
503 UNIX_EPOCH.checked_add(Duration::new(seconds, subsecond_nanos))
504}
505
506pub(super) fn prune_direct_child_directories(
507 cache_root: &Path,
508 policy: ArtifactCachePrunePolicy,
509 protected_entry: Option<&Path>,
510 is_eligible: impl Fn(&Path) -> bool,
511) -> Result<ArtifactCachePruneReport, CacheFsError> {
512 let mut entries = cache_entries(cache_root, is_eligible)?;
513 let bytes_before = entries
514 .iter()
515 .fold(0_u64, |total, entry| total.saturating_add(entry.bytes));
516 let mut report = ArtifactCachePruneReport {
517 entries_scanned: entries.len(),
518 entries_removed: 0,
519 bytes_before,
520 bytes_removed: 0,
521 uncommitted_directories_removed: 0,
522 uncommitted_bytes_removed: 0,
523 };
524 let now = SystemTime::now();
525
526 if let Some(max_age) = policy.max_age() {
527 for entry in &mut entries {
528 let age = now.duration_since(entry.last_used).unwrap_or_default();
529 if protected_entry != Some(entry.path.as_path()) && age > max_age {
530 remove_cache_entry(entry, &mut report)?;
531 }
532 }
533 }
534
535 if let Some(max_size_bytes) = policy.max_size_bytes() {
536 entries.sort_by(|left, right| {
537 left.last_used
538 .cmp(&right.last_used)
539 .then_with(|| left.path.cmp(&right.path))
540 });
541 for entry in &mut entries {
542 if report.bytes_retained() <= max_size_bytes {
543 break;
544 }
545 if protected_entry == Some(entry.path.as_path()) {
546 continue;
547 }
548 remove_cache_entry(entry, &mut report)?;
549 }
550 }
551
552 Ok(report)
553}
554
555pub(super) fn directory_logical_size(path: &Path) -> io::Result<u64> {
556 let mut total = 0_u64;
557 let mut pending = vec![path.to_owned()];
558 while let Some(current) = pending.pop() {
559 let metadata = fs::symlink_metadata(¤t)?;
560 if metadata.is_dir() {
561 for entry in fs::read_dir(¤t)? {
562 pending.push(entry?.path());
563 }
564 } else {
565 total = total.saturating_add(metadata.len());
566 }
567 }
568 Ok(total)
569}
570
571pub(super) fn is_sha256_directory(path: &Path) -> bool {
572 path.file_name().is_some_and(|name| {
573 let bytes = name.as_encoded_bytes();
574 bytes.len() == 64 && bytes.iter().all(u8::is_ascii_hexdigit)
575 })
576}
577
578pub(super) fn remove_path_if_present(path: &Path) -> io::Result<()> {
579 let metadata = match fs::symlink_metadata(path) {
580 Ok(metadata) => metadata,
581 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
582 Err(error) => return Err(error),
583 };
584 if metadata.file_type().is_dir() {
585 fs::remove_dir_all(path)
586 } else {
587 fs::remove_file(path)
588 }
589}
590
591struct CacheEntry {
592 path: PathBuf,
593 bytes: u64,
594 last_used: SystemTime,
595 removed: bool,
596}
597
598impl std::fmt::Display for CacheFsError {
599 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
600 write!(
601 formatter,
602 "failed to {} at {}: {}",
603 self.operation,
604 self.path.display(),
605 self.source
606 )
607 }
608}
609
610impl std::error::Error for CacheFsError {
611 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
612 Some(&self.source)
613 }
614}
615
616fn cache_entries(
617 cache_root: &Path,
618 is_eligible: impl Fn(&Path) -> bool,
619) -> Result<Vec<CacheEntry>, CacheFsError> {
620 let read_dir = match fs::read_dir(cache_root) {
621 Ok(read_dir) => read_dir,
622 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
623 Err(source) => {
624 return Err(CacheFsError {
625 operation: "read cache directory",
626 path: cache_root.to_owned(),
627 source,
628 });
629 }
630 };
631 let mut entries = Vec::new();
632 for directory_entry in read_dir {
633 let directory_entry = directory_entry.map_err(|source| CacheFsError {
634 operation: "read cache entry",
635 path: cache_root.to_owned(),
636 source,
637 })?;
638 let path = directory_entry.path();
639 let file_type = directory_entry.file_type().map_err(|source| CacheFsError {
640 operation: "inspect cache entry",
641 path: path.clone(),
642 source,
643 })?;
644 if !file_type.is_dir() || !is_eligible(&path) {
645 continue;
646 }
647 let bytes = directory_logical_size(&path).map_err(|source| CacheFsError {
648 operation: "measure cache entry",
649 path: path.clone(),
650 source,
651 })?;
652 let last_used = cache_entry_last_used(&path).map_err(|source| CacheFsError {
653 operation: "read cache use time",
654 path: path.clone(),
655 source,
656 })?;
657 entries.push(CacheEntry {
658 path,
659 bytes,
660 last_used,
661 removed: false,
662 });
663 }
664 Ok(entries)
665}
666
667pub(super) fn cache_entry_last_used(path: &Path) -> io::Result<SystemTime> {
668 let marker = path.join(LAST_USED_FILE);
669 if let Ok(contents) = fs::read_to_string(&marker)
670 && let Some(timestamp) = decode_system_time(&contents)
671 {
672 return Ok(timestamp);
673 }
674 fs::metadata(path)?.modified()
675}
676
677fn remove_cache_entry(
678 entry: &mut CacheEntry,
679 report: &mut ArtifactCachePruneReport,
680) -> Result<(), CacheFsError> {
681 if entry.removed {
682 return Ok(());
683 }
684 let Some(_retention_lock) = try_lock_cache_file(&entry.path.join(RETENTION_LOCK_FILE))? else {
685 return Ok(());
686 };
687 remove_path_if_present(&entry.path).map_err(|source| CacheFsError {
688 operation: "prune cache entry",
689 path: entry.path.clone(),
690 source,
691 })?;
692 entry.removed = true;
693 report.entries_removed += 1;
694 report.bytes_removed = report.bytes_removed.saturating_add(entry.bytes);
695 Ok(())
696}
697
698#[cfg(all(test, unix))]
699mod tests {
700 use super::canonicalize_allow_missing;
701 use crate::artifacts::test_support::unique_temp_directory;
702 use std::fs;
703
704 #[test]
705 fn missing_parent_traversal_resumes_existing_symlink_resolution() {
706 let root = unique_temp_directory("canonical-missing-parent");
707 let target = root.join("real");
708 fs::create_dir_all(&target).unwrap();
709 std::os::unix::fs::symlink(&target, root.join("alias")).unwrap();
710 let path = root.join("missing/../alias/generated/nested/../output");
711 assert_eq!(
712 canonicalize_allow_missing(&path).unwrap(),
713 target.canonicalize().unwrap().join("generated/output")
714 );
715 assert_eq!(
716 canonicalize_allow_missing(&root.join("alias/../other/output")).unwrap(),
717 root.canonicalize().unwrap().join("other/output")
718 );
719 fs::remove_dir_all(root).unwrap();
720 }
721}