Skip to main content

ic_testkit/artifacts/
cache_fs.rs

1use fs2::FileExt as _;
2use std::{
3    fs::{self, File, OpenOptions},
4    io,
5    path::{Path, PathBuf},
6    sync::Arc,
7    thread,
8    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
9};
10
11use super::digest::write_atomic;
12
13const CACHE_DIRECTORY_TAG: &str = "Signature: 8a477f597d28d172789f06886806bc55\n\
14# This file is a cache directory tag created by ic-testkit.\n\
15# For information about cache directory tags see https://bford.info/cachedir/\n";
16pub(super) const CACHE_DIRECTORY_TAG_SIGNATURE: &str =
17    "Signature: 8a477f597d28d172789f06886806bc55\n";
18pub(super) const LAST_USED_FILE: &str = ".ic-testkit-last-used";
19const LAST_MAINTENANCE_FILE: &str = ".ic-testkit-last-maintenance";
20pub(super) const RETENTION_LOCK_FILE: &str = ".ic-testkit-retention-v1";
21
22/// Acquired under the producer/namespace lock before handing an entry to a
23/// consumer. Clones share ownership; the OS releases locks on process exit.
24#[derive(Clone, Debug)]
25pub(super) struct RetainedCacheEntry {
26    path: PathBuf,
27    _lock: Arc<File>,
28}
29
30impl PartialEq for RetainedCacheEntry {
31    fn eq(&self, other: &Self) -> bool {
32        self.path == other.path
33    }
34}
35
36impl Eq for RetainedCacheEntry {}
37
38impl RetainedCacheEntry {
39    pub(super) fn acquire(path: &Path) -> Result<Self, CacheFsError> {
40        let file = open_cache_lock_file(&path.join(RETENTION_LOCK_FILE))?;
41        fs2::FileExt::lock_shared(&file).map_err(|source| CacheFsError {
42            operation: "retain cache entry",
43            path: path.to_owned(),
44            source,
45        })?;
46        Ok(Self {
47            path: path.to_owned(),
48            _lock: Arc::new(file),
49        })
50    }
51}
52
53/// The caller must hold the producer/namespace lock throughout this operation.
54pub(super) fn remove_unretained_entry(path: &Path) -> Result<(), CacheFsError> {
55    let metadata = match fs::symlink_metadata(path) {
56        Ok(metadata) => metadata,
57        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
58        Err(source) => {
59            return Err(CacheFsError {
60                operation: "inspect cache entry",
61                path: path.to_owned(),
62                source,
63            });
64        }
65    };
66    if !metadata.is_dir() {
67        return remove_path_if_present(path).map_err(|source| CacheFsError {
68            operation: "remove invalid cache entry",
69            path: path.to_owned(),
70            source,
71        });
72    }
73    let _lock =
74        try_lock_cache_file(&path.join(RETENTION_LOCK_FILE))?.ok_or_else(|| CacheFsError {
75            operation: "replace retained cache entry",
76            path: path.to_owned(),
77            source: io::Error::new(
78                io::ErrorKind::WouldBlock,
79                "cache entry is retained by a consumer",
80            ),
81        })?;
82    remove_path_if_present(path).map_err(|source| CacheFsError {
83        operation: "remove cache entry",
84        path: path.to_owned(),
85        source,
86    })
87}
88
89/// Caller-selected retention limits for content-addressed artifact entries.
90///
91/// Age pruning runs before size pruning. A policy without either limit scans
92/// the selected cache namespace and updates its cache metadata without
93/// removing entries. Entries retained by live acquisition records are skipped,
94/// even when this temporarily exceeds the limits. They become eligible for the
95/// next maintenance pass after their final owner drops or its process exits.
96#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
97pub struct ArtifactCachePrunePolicy {
98    max_age: Option<Duration>,
99    max_size_bytes: Option<u64>,
100}
101
102/// Summary of one lock-coordinated artifact-cache pruning pass.
103#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
104pub struct ArtifactCachePruneReport {
105    entries_scanned: usize,
106    entries_removed: usize,
107    bytes_before: u64,
108    bytes_removed: u64,
109    uncommitted_directories_removed: usize,
110    uncommitted_bytes_removed: u64,
111}
112
113/// Nonfatal retention attempted as part of a successful cache acquisition.
114#[non_exhaustive]
115#[derive(Clone, Debug, Eq, PartialEq)]
116pub enum ArtifactCacheMaintenance {
117    /// Configured retention completed under the cache lock.
118    Pruned(ArtifactCachePruneReport),
119    /// Configured retention failed after the requested artifacts were ready.
120    PruneFailed {
121        /// Cache error rendered without invalidating the successful acquisition.
122        message: String,
123    },
124}
125
126impl ArtifactCachePrunePolicy {
127    /// Create a policy that records cache metadata without removing entries.
128    #[must_use]
129    pub const fn new() -> Self {
130        Self {
131            max_age: None,
132            max_size_bytes: None,
133        }
134    }
135
136    /// Remove entries older than `max_age` before applying the size limit.
137    #[must_use]
138    pub const fn with_max_age(mut self, max_age: Duration) -> Self {
139        self.max_age = Some(max_age);
140        self
141    }
142
143    /// Remove least-recently-used entries until retained logical size is at most `bytes`.
144    #[must_use]
145    pub const fn with_max_size_bytes(mut self, bytes: u64) -> Self {
146        self.max_size_bytes = Some(bytes);
147        self
148    }
149
150    /// Configured maximum entry age, if any.
151    #[must_use]
152    pub const fn max_age(self) -> Option<Duration> {
153        self.max_age
154    }
155
156    /// Configured maximum logical cache size in bytes, if any.
157    #[must_use]
158    pub const fn max_size_bytes(self) -> Option<u64> {
159        self.max_size_bytes
160    }
161
162    pub(super) fn maintenance_identity(self) -> String {
163        format!(
164            "age={:?};size={:?}",
165            self.max_age.map(|duration| duration.as_nanos()),
166            self.max_size_bytes
167        )
168    }
169}
170
171impl ArtifactCachePruneReport {
172    /// Number of content-addressed directories considered for pruning.
173    #[must_use]
174    pub const fn entries_scanned(self) -> usize {
175        self.entries_scanned
176    }
177
178    /// Number of content-addressed directories removed.
179    #[must_use]
180    pub const fn entries_removed(self) -> usize {
181        self.entries_removed
182    }
183
184    /// Number of content-addressed directories retained.
185    #[must_use]
186    pub const fn entries_retained(self) -> usize {
187        self.entries_scanned.saturating_sub(self.entries_removed)
188    }
189
190    /// Logical bytes occupied by scanned entries before pruning.
191    #[must_use]
192    pub const fn bytes_before(self) -> u64 {
193        self.bytes_before
194    }
195
196    /// Logical bytes removed by pruning.
197    #[must_use]
198    pub const fn bytes_removed(self) -> u64 {
199        self.bytes_removed
200    }
201
202    /// Logical bytes occupied by retained entries after pruning.
203    #[must_use]
204    pub const fn bytes_retained(self) -> u64 {
205        self.bytes_before.saturating_sub(self.bytes_removed)
206    }
207
208    /// Abandoned transaction directories removed outside the committed-entry totals.
209    #[must_use]
210    pub const fn uncommitted_directories_removed(self) -> usize {
211        self.uncommitted_directories_removed
212    }
213
214    /// Logical bytes removed from abandoned transaction directories.
215    #[must_use]
216    pub const fn uncommitted_bytes_removed(self) -> u64 {
217        self.uncommitted_bytes_removed
218    }
219
220    pub(super) const fn record_uncommitted_removal(&mut self, bytes: u64) {
221        self.uncommitted_directories_removed += 1;
222        self.uncommitted_bytes_removed = self.uncommitted_bytes_removed.saturating_add(bytes);
223    }
224}
225
226impl ArtifactCacheMaintenance {
227    /// Successful pruning report, or `None` when maintenance failed.
228    #[must_use]
229    pub const fn prune_report(&self) -> Option<ArtifactCachePruneReport> {
230        match self {
231            Self::Pruned(report) => Some(*report),
232            Self::PruneFailed { .. } => None,
233        }
234    }
235
236    /// Rendered maintenance failure, or `None` when pruning succeeded.
237    #[must_use]
238    pub fn failure_message(&self) -> Option<&str> {
239        match self {
240            Self::Pruned(_) => None,
241            Self::PruneFailed { message } => Some(message),
242        }
243    }
244}
245
246#[derive(Debug)]
247pub(super) struct CacheFsError {
248    pub(super) operation: &'static str,
249    pub(super) path: PathBuf,
250    pub(super) source: io::Error,
251}
252
253pub(super) fn ensure_cache_directory_tag(cache_root: &Path) -> Result<(), CacheFsError> {
254    let path = cache_root.join("CACHEDIR.TAG");
255    if fs::read_to_string(&path)
256        .is_ok_and(|contents| contents.starts_with(CACHE_DIRECTORY_TAG_SIGNATURE))
257    {
258        return Ok(());
259    }
260    write_atomic(&path, CACHE_DIRECTORY_TAG.as_bytes()).map_err(|source| CacheFsError {
261        operation: "write cache directory tag",
262        path,
263        source,
264    })
265}
266
267pub(super) fn lock_cache_file(path: &Path) -> Result<(File, Duration), CacheFsError> {
268    let file = open_cache_lock_file(path)?;
269    let started = Instant::now();
270    file.lock_exclusive().map_err(|source| CacheFsError {
271        operation: "lock cache",
272        path: path.to_owned(),
273        source,
274    })?;
275    Ok((file, started.elapsed()))
276}
277
278pub(super) fn lock_cache_file_with_wait_observer(
279    path: &Path,
280    poll_interval: Duration,
281    mut observer: impl FnMut(Duration),
282) -> Result<(File, Duration), CacheFsError> {
283    let file = open_cache_lock_file(path)?;
284    let started = Instant::now();
285    loop {
286        match file.try_lock_exclusive() {
287            Ok(()) => return Ok((file, started.elapsed())),
288            Err(error) if error.kind() == io::ErrorKind::WouldBlock => {
289                observer(started.elapsed());
290                thread::sleep(poll_interval.min(Duration::from_millis(25)));
291            }
292            Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
293            Err(source) => {
294                return Err(CacheFsError {
295                    operation: "try lock cache",
296                    path: path.to_owned(),
297                    source,
298                });
299            }
300        }
301    }
302}
303
304pub(super) fn try_lock_cache_file(path: &Path) -> Result<Option<File>, CacheFsError> {
305    let file = open_cache_lock_file(path)?;
306    match file.try_lock_exclusive() {
307        Ok(()) => Ok(Some(file)),
308        Err(error) if error.kind() == io::ErrorKind::WouldBlock => Ok(None),
309        Err(source) => Err(CacheFsError {
310            operation: "try lock cache",
311            path: path.to_owned(),
312            source,
313        }),
314    }
315}
316
317fn open_cache_lock_file(path: &Path) -> Result<File, CacheFsError> {
318    if let Some(parent) = path.parent() {
319        fs::create_dir_all(parent).map_err(|source| CacheFsError {
320            operation: "create cache lock directory",
321            path: parent.to_owned(),
322            source,
323        })?;
324    }
325    OpenOptions::new()
326        .create(true)
327        .read(true)
328        .write(true)
329        .truncate(false)
330        .open(path)
331        .map_err(|source| CacheFsError {
332            operation: "open cache lock",
333            path: path.to_owned(),
334            source,
335        })
336}
337
338pub(super) fn record_cache_entry_use(path: &Path) -> Result<(), CacheFsError> {
339    write_last_used(path, SystemTime::now())
340}
341
342pub(super) fn cache_maintenance_due(
343    path: &Path,
344    minimum_interval: Option<Duration>,
345    maintenance_identity: &str,
346) -> Result<bool, CacheFsError> {
347    let Some(minimum_interval) = minimum_interval else {
348        return Ok(true);
349    };
350    let marker = path.join(LAST_MAINTENANCE_FILE);
351    let contents = match fs::read_to_string(&marker) {
352        Ok(contents) => contents,
353        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(true),
354        Err(source) => {
355            return Err(CacheFsError {
356                operation: "read cache maintenance time",
357                path: marker,
358                source,
359            });
360        }
361    };
362    let mut lines = contents.lines();
363    let Some(last_maintenance) = lines.next().and_then(decode_system_time) else {
364        return Ok(true);
365    };
366    if lines.next() != Some(maintenance_identity) {
367        return Ok(true);
368    }
369    Ok(match SystemTime::now().duration_since(last_maintenance) {
370        Ok(elapsed) => elapsed >= minimum_interval,
371        Err(_) => true,
372    })
373}
374
375pub(super) fn record_cache_maintenance(
376    path: &Path,
377    maintenance_identity: &str,
378) -> Result<(), CacheFsError> {
379    fs::create_dir_all(path).map_err(|source| CacheFsError {
380        operation: "create cache maintenance directory",
381        path: path.to_owned(),
382        source,
383    })?;
384    let marker = path.join(LAST_MAINTENANCE_FILE);
385    let elapsed = encode_system_time(&marker, SystemTime::now())?;
386    let contents = format!("{}\n{maintenance_identity}\n", elapsed.as_nanos());
387    write_atomic(&marker, contents.as_bytes()).map_err(|source| CacheFsError {
388        operation: "record cache maintenance time",
389        path: marker,
390        source,
391    })
392}
393
394pub(super) fn perform_scheduled_cache_maintenance(
395    path: &Path,
396    minimum_interval: Option<Duration>,
397    maintenance_identity: &str,
398    maintenance: impl FnOnce() -> Result<ArtifactCachePruneReport, String>,
399) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
400    let started = Instant::now();
401    match cache_maintenance_due(path, minimum_interval, maintenance_identity) {
402        Ok(false) => return (None, Some(started.elapsed())),
403        Ok(true) => {}
404        Err(error) => {
405            return (
406                Some(ArtifactCacheMaintenance::PruneFailed {
407                    message: error.to_string(),
408                }),
409                Some(started.elapsed()),
410            );
411        }
412    }
413
414    let result = maintenance();
415    let marker = record_cache_maintenance(path, maintenance_identity);
416    let outcome = match (result, marker) {
417        (Ok(report), Ok(())) => ArtifactCacheMaintenance::Pruned(report),
418        (Err(message), Ok(())) => ArtifactCacheMaintenance::PruneFailed { message },
419        (Ok(_), Err(error)) => ArtifactCacheMaintenance::PruneFailed {
420            message: error.to_string(),
421        },
422        (Err(message), Err(marker)) => ArtifactCacheMaintenance::PruneFailed {
423            message: format!(
424                "{message}; additionally failed to record the maintenance attempt: {marker}"
425            ),
426        },
427    };
428    (Some(outcome), Some(started.elapsed()))
429}
430
431pub(super) fn write_last_used(path: &Path, last_used: SystemTime) -> Result<(), CacheFsError> {
432    let marker = path.join(LAST_USED_FILE);
433    write_system_time(&marker, last_used, "record cache use time")
434}
435
436fn write_system_time(
437    path: &Path,
438    timestamp: SystemTime,
439    operation: &'static str,
440) -> Result<(), CacheFsError> {
441    let elapsed = encode_system_time(path, timestamp)?;
442    write_atomic(path, elapsed.as_nanos().to_string().as_bytes()).map_err(|source| CacheFsError {
443        operation,
444        path: path.to_owned(),
445        source,
446    })
447}
448
449fn encode_system_time(path: &Path, timestamp: SystemTime) -> Result<Duration, CacheFsError> {
450    timestamp
451        .duration_since(UNIX_EPOCH)
452        .map_err(|source| CacheFsError {
453            operation: "encode cache time",
454            path: path.to_owned(),
455            source: io::Error::new(io::ErrorKind::InvalidInput, source),
456        })
457}
458
459fn decode_system_time(contents: &str) -> Option<SystemTime> {
460    let nanoseconds = contents.parse::<u128>().ok()?;
461    let seconds = u64::try_from(nanoseconds / 1_000_000_000).ok()?;
462    let subsecond_nanos = (nanoseconds % 1_000_000_000) as u32;
463    UNIX_EPOCH.checked_add(Duration::new(seconds, subsecond_nanos))
464}
465
466pub(super) fn prune_direct_child_directories(
467    cache_root: &Path,
468    policy: ArtifactCachePrunePolicy,
469    protected_entry: Option<&Path>,
470    is_eligible: impl Fn(&Path) -> bool,
471) -> Result<ArtifactCachePruneReport, CacheFsError> {
472    let mut entries = cache_entries(cache_root, is_eligible)?;
473    let bytes_before = entries
474        .iter()
475        .fold(0_u64, |total, entry| total.saturating_add(entry.bytes));
476    let mut report = ArtifactCachePruneReport {
477        entries_scanned: entries.len(),
478        entries_removed: 0,
479        bytes_before,
480        bytes_removed: 0,
481        uncommitted_directories_removed: 0,
482        uncommitted_bytes_removed: 0,
483    };
484    let now = SystemTime::now();
485
486    if let Some(max_age) = policy.max_age() {
487        for entry in &mut entries {
488            let age = now.duration_since(entry.last_used).unwrap_or_default();
489            if protected_entry != Some(entry.path.as_path()) && age > max_age {
490                remove_cache_entry(entry, &mut report)?;
491            }
492        }
493    }
494
495    if let Some(max_size_bytes) = policy.max_size_bytes() {
496        entries.sort_by(|left, right| {
497            left.last_used
498                .cmp(&right.last_used)
499                .then_with(|| left.path.cmp(&right.path))
500        });
501        for entry in &mut entries {
502            if report.bytes_retained() <= max_size_bytes {
503                break;
504            }
505            if protected_entry == Some(entry.path.as_path()) {
506                continue;
507            }
508            remove_cache_entry(entry, &mut report)?;
509        }
510    }
511
512    Ok(report)
513}
514
515pub(super) fn directory_logical_size(path: &Path) -> io::Result<u64> {
516    let mut total = 0_u64;
517    let mut pending = vec![path.to_owned()];
518    while let Some(current) = pending.pop() {
519        let metadata = fs::symlink_metadata(&current)?;
520        if metadata.is_dir() {
521            for entry in fs::read_dir(&current)? {
522                pending.push(entry?.path());
523            }
524        } else {
525            total = total.saturating_add(metadata.len());
526        }
527    }
528    Ok(total)
529}
530
531pub(super) fn is_sha256_directory(path: &Path) -> bool {
532    path.file_name().is_some_and(|name| {
533        let bytes = name.as_encoded_bytes();
534        bytes.len() == 64 && bytes.iter().all(u8::is_ascii_hexdigit)
535    })
536}
537
538pub(super) fn remove_path_if_present(path: &Path) -> io::Result<()> {
539    let metadata = match fs::symlink_metadata(path) {
540        Ok(metadata) => metadata,
541        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
542        Err(error) => return Err(error),
543    };
544    if metadata.file_type().is_dir() {
545        fs::remove_dir_all(path)
546    } else {
547        fs::remove_file(path)
548    }
549}
550
551struct CacheEntry {
552    path: PathBuf,
553    bytes: u64,
554    last_used: SystemTime,
555    removed: bool,
556}
557
558impl std::fmt::Display for CacheFsError {
559    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
560        write!(
561            formatter,
562            "failed to {} at {}: {}",
563            self.operation,
564            self.path.display(),
565            self.source
566        )
567    }
568}
569
570impl std::error::Error for CacheFsError {
571    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
572        Some(&self.source)
573    }
574}
575
576fn cache_entries(
577    cache_root: &Path,
578    is_eligible: impl Fn(&Path) -> bool,
579) -> Result<Vec<CacheEntry>, CacheFsError> {
580    let read_dir = match fs::read_dir(cache_root) {
581        Ok(read_dir) => read_dir,
582        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
583        Err(source) => {
584            return Err(CacheFsError {
585                operation: "read cache directory",
586                path: cache_root.to_owned(),
587                source,
588            });
589        }
590    };
591    let mut entries = Vec::new();
592    for directory_entry in read_dir {
593        let directory_entry = directory_entry.map_err(|source| CacheFsError {
594            operation: "read cache entry",
595            path: cache_root.to_owned(),
596            source,
597        })?;
598        let path = directory_entry.path();
599        let file_type = directory_entry.file_type().map_err(|source| CacheFsError {
600            operation: "inspect cache entry",
601            path: path.clone(),
602            source,
603        })?;
604        if !file_type.is_dir() || !is_eligible(&path) {
605            continue;
606        }
607        let bytes = directory_logical_size(&path).map_err(|source| CacheFsError {
608            operation: "measure cache entry",
609            path: path.clone(),
610            source,
611        })?;
612        let last_used = cache_entry_last_used(&path).map_err(|source| CacheFsError {
613            operation: "read cache use time",
614            path: path.clone(),
615            source,
616        })?;
617        entries.push(CacheEntry {
618            path,
619            bytes,
620            last_used,
621            removed: false,
622        });
623    }
624    Ok(entries)
625}
626
627pub(super) fn cache_entry_last_used(path: &Path) -> io::Result<SystemTime> {
628    let marker = path.join(LAST_USED_FILE);
629    if let Ok(contents) = fs::read_to_string(&marker)
630        && let Some(timestamp) = decode_system_time(&contents)
631    {
632        return Ok(timestamp);
633    }
634    fs::metadata(path)?.modified()
635}
636
637fn remove_cache_entry(
638    entry: &mut CacheEntry,
639    report: &mut ArtifactCachePruneReport,
640) -> Result<(), CacheFsError> {
641    if entry.removed {
642        return Ok(());
643    }
644    let Some(_retention_lock) = try_lock_cache_file(&entry.path.join(RETENTION_LOCK_FILE))? else {
645        return Ok(());
646    };
647    remove_path_if_present(&entry.path).map_err(|source| CacheFsError {
648        operation: "prune cache entry",
649        path: entry.path.clone(),
650        source,
651    })?;
652    entry.removed = true;
653    report.entries_removed += 1;
654    report.bytes_removed = report.bytes_removed.saturating_add(entry.bytes);
655    Ok(())
656}