Skip to main content

ic_testkit/artifacts/
cache_fs.rs

1use fs2::FileExt as _;
2use std::{
3    fs::{self, File},
4    io::{self, Read as _},
5    path::{Path, PathBuf},
6    sync::Arc,
7    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
8};
9
10use super::digest::read_stamp_with_limit;
11
12/// Resolve existing components through symlinks and normalize a missing suffix.
13/// Parent traversal can return from a missing suffix to existing components.
14pub(super) fn canonicalize_allow_missing(path: &Path) -> io::Result<PathBuf> {
15    let base = if path.is_absolute() {
16        PathBuf::new()
17    } else {
18        std::env::current_dir()?
19    };
20    ic_host_fs::path::canonicalize_allow_missing(path, &base)
21}
22
23const CACHE_DIRECTORY_TAG: &str = "Signature: 8a477f597d28d172789f06886806bc55\n\
24# This file is a cache directory tag created by ic-testkit.\n\
25# For information about cache directory tags see https://bford.info/cachedir/\n";
26pub(super) const CACHE_DIRECTORY_TAG_SIGNATURE: &str =
27    "Signature: 8a477f597d28d172789f06886806bc55";
28pub(super) const LAST_USED_FILE: &str = ".ic-testkit-last-used";
29const LAST_MAINTENANCE_FILE: &str = ".ic-testkit-last-maintenance";
30// Nanoseconds are written as decimal u128 values, which need at most 39 bytes.
31const MAX_TIMESTAMP_BYTES: usize = 39;
32pub(super) const RETENTION_LOCK_FILE: &str = ".ic-testkit-retention-v1";
33
34/// Acquired under the producer/namespace lock before handing an entry to a
35/// consumer. Clones share ownership; the OS releases locks on process exit.
36#[derive(Clone, Debug)]
37pub(super) struct RetainedCacheEntry {
38    path: PathBuf,
39    _lock: Arc<RetentionLock>,
40}
41
42#[derive(Debug)]
43struct RetentionLock(File);
44
45impl Drop for RetentionLock {
46    fn drop(&mut self) {
47        // Closing alone can leave a flock held by a descriptor inherited during
48        // a concurrent spawn before exec. Release it when the final record owner
49        // drops, rather than waiting for unrelated child descriptors to close.
50        // If unlocking fails, closing the owned file remains the fallback.
51        let _ = fs2::FileExt::unlock(&self.0);
52    }
53}
54
55impl PartialEq for RetainedCacheEntry {
56    fn eq(&self, other: &Self) -> bool {
57        self.path == other.path
58    }
59}
60
61impl Eq for RetainedCacheEntry {}
62
63impl RetainedCacheEntry {
64    pub(super) fn path(&self) -> &Path {
65        &self.path
66    }
67
68    pub(super) fn acquire(path: &Path) -> Result<Self, CacheFsError> {
69        let file = open_cache_lock_file(&path.join(RETENTION_LOCK_FILE))?;
70        fs2::FileExt::lock_shared(&file).map_err(|source| CacheFsError {
71            operation: "retain cache entry",
72            path: path.to_owned(),
73            source,
74        })?;
75        Ok(Self {
76            path: path.to_owned(),
77            _lock: Arc::new(RetentionLock(file)),
78        })
79    }
80}
81
82/// The caller must hold the producer/namespace lock throughout this operation.
83pub(super) fn remove_unretained_entry(path: &Path) -> Result<(), CacheFsError> {
84    let metadata = match fs::symlink_metadata(path) {
85        Ok(metadata) => metadata,
86        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
87        Err(source) => {
88            return Err(CacheFsError {
89                operation: "inspect cache entry",
90                path: path.to_owned(),
91                source,
92            });
93        }
94    };
95    if !metadata.is_dir() {
96        return remove_path_if_present(path).map_err(|source| CacheFsError {
97            operation: "remove invalid cache entry",
98            path: path.to_owned(),
99            source,
100        });
101    }
102    let _lock =
103        try_lock_cache_file(&path.join(RETENTION_LOCK_FILE))?.ok_or_else(|| CacheFsError {
104            operation: "replace retained cache entry",
105            path: path.to_owned(),
106            source: io::Error::new(
107                io::ErrorKind::WouldBlock,
108                "cache entry is retained by a consumer",
109            ),
110        })?;
111    remove_path_if_present(path).map_err(|source| CacheFsError {
112        operation: "remove cache entry",
113        path: path.to_owned(),
114        source,
115    })
116}
117
118/// Caller-selected retention limits for content-addressed artifact entries.
119///
120/// Age pruning runs before size pruning. A policy without either limit scans
121/// the selected cache namespace and updates its cache metadata without
122/// removing entries. Entries retained by live acquisition records are skipped,
123/// even when this temporarily exceeds the limits. They become eligible for the
124/// next maintenance pass after their final owner drops or its process exits.
125#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
126pub struct ArtifactCachePrunePolicy {
127    max_age: Option<Duration>,
128    max_size_bytes: Option<u64>,
129}
130
131/// Summary of one lock-coordinated artifact-cache pruning pass.
132#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
133pub struct ArtifactCachePruneReport {
134    entries_scanned: usize,
135    entries_removed: usize,
136    bytes_before: u64,
137    bytes_removed: u64,
138    uncommitted_directories_removed: usize,
139    uncommitted_bytes_removed: u64,
140}
141
142/// Nonfatal retention attempted as part of a successful cache acquisition.
143#[non_exhaustive]
144#[derive(Clone, Debug, Eq, PartialEq)]
145pub enum ArtifactCacheMaintenance {
146    /// Configured retention completed under the cache lock.
147    Pruned(ArtifactCachePruneReport),
148    /// Configured retention failed after the requested artifacts were ready.
149    PruneFailed {
150        /// Cache error rendered without invalidating the successful acquisition.
151        message: String,
152    },
153}
154
155impl ArtifactCachePrunePolicy {
156    /// Create a policy that records cache metadata without removing entries.
157    #[must_use]
158    pub const fn new() -> Self {
159        Self {
160            max_age: None,
161            max_size_bytes: None,
162        }
163    }
164
165    /// Remove entries older than `max_age` before applying the size limit.
166    #[must_use]
167    pub const fn with_max_age(mut self, max_age: Duration) -> Self {
168        self.max_age = Some(max_age);
169        self
170    }
171
172    /// Remove least-recently-used entries until retained logical size is at most `bytes`.
173    #[must_use]
174    pub const fn with_max_size_bytes(mut self, bytes: u64) -> Self {
175        self.max_size_bytes = Some(bytes);
176        self
177    }
178
179    /// Configured maximum entry age, if any.
180    #[must_use]
181    pub const fn max_age(self) -> Option<Duration> {
182        self.max_age
183    }
184
185    /// Configured maximum logical cache size in bytes, if any.
186    #[must_use]
187    pub const fn max_size_bytes(self) -> Option<u64> {
188        self.max_size_bytes
189    }
190
191    pub(super) fn maintenance_identity(self) -> String {
192        format!(
193            "age={:?};size={:?}",
194            self.max_age.map(|duration| duration.as_nanos()),
195            self.max_size_bytes
196        )
197    }
198}
199
200impl ArtifactCachePruneReport {
201    /// Number of content-addressed directories considered for pruning.
202    #[must_use]
203    pub const fn entries_scanned(self) -> usize {
204        self.entries_scanned
205    }
206
207    /// Number of content-addressed directories removed.
208    #[must_use]
209    pub const fn entries_removed(self) -> usize {
210        self.entries_removed
211    }
212
213    /// Number of content-addressed directories retained.
214    #[must_use]
215    pub const fn entries_retained(self) -> usize {
216        self.entries_scanned.saturating_sub(self.entries_removed)
217    }
218
219    /// Logical bytes occupied by scanned entries before pruning.
220    #[must_use]
221    pub const fn bytes_before(self) -> u64 {
222        self.bytes_before
223    }
224
225    /// Logical bytes removed by pruning.
226    #[must_use]
227    pub const fn bytes_removed(self) -> u64 {
228        self.bytes_removed
229    }
230
231    /// Logical bytes occupied by retained entries after pruning.
232    #[must_use]
233    pub const fn bytes_retained(self) -> u64 {
234        self.bytes_before.saturating_sub(self.bytes_removed)
235    }
236
237    /// Abandoned transaction directories removed outside the committed-entry totals.
238    #[must_use]
239    pub const fn uncommitted_directories_removed(self) -> usize {
240        self.uncommitted_directories_removed
241    }
242
243    /// Logical bytes removed from abandoned transaction directories.
244    #[must_use]
245    pub const fn uncommitted_bytes_removed(self) -> u64 {
246        self.uncommitted_bytes_removed
247    }
248
249    pub(super) const fn record_uncommitted_removal(&mut self, bytes: u64) {
250        self.uncommitted_directories_removed += 1;
251        self.uncommitted_bytes_removed = self.uncommitted_bytes_removed.saturating_add(bytes);
252    }
253}
254
255impl ArtifactCacheMaintenance {
256    /// Successful pruning report, or `None` when maintenance failed.
257    #[must_use]
258    pub const fn prune_report(&self) -> Option<ArtifactCachePruneReport> {
259        match self {
260            Self::Pruned(report) => Some(*report),
261            Self::PruneFailed { .. } => None,
262        }
263    }
264
265    /// Rendered maintenance failure, or `None` when pruning succeeded.
266    #[must_use]
267    pub fn failure_message(&self) -> Option<&str> {
268        match self {
269            Self::Pruned(_) => None,
270            Self::PruneFailed { message } => Some(message),
271        }
272    }
273}
274
275#[derive(Debug)]
276pub(super) struct CacheFsError {
277    pub(super) operation: &'static str,
278    pub(super) path: PathBuf,
279    pub(super) source: io::Error,
280}
281
282pub(super) fn ensure_cache_directory_tag(cache_root: &Path) -> Result<(), CacheFsError> {
283    let path = cache_root.join("CACHEDIR.TAG");
284    // The standard recognizes the first 43 bytes, without requiring a newline
285    // or interpreting the remaining text. Symlinks are not valid tag files.
286    let mut signature = [0_u8; CACHE_DIRECTORY_TAG_SIGNATURE.len()];
287    if fs::symlink_metadata(&path).is_ok_and(|metadata| metadata.file_type().is_file())
288        && File::open(&path)
289            .and_then(|mut file| file.read_exact(&mut signature))
290            .is_ok()
291        && signature == CACHE_DIRECTORY_TAG_SIGNATURE.as_bytes()
292    {
293        return Ok(());
294    }
295    ic_host_fs::durable::write_bytes(&path, CACHE_DIRECTORY_TAG.as_bytes()).map_err(|source| {
296        CacheFsError {
297            operation: "write cache directory tag",
298            path,
299            source,
300        }
301    })
302}
303
304pub(super) fn lock_cache_file(path: &Path) -> Result<(File, Duration), CacheFsError> {
305    let file = open_cache_lock_file(path)?;
306    let started = Instant::now();
307    file.lock_exclusive().map_err(|source| CacheFsError {
308        operation: "lock cache",
309        path: path.to_owned(),
310        source,
311    })?;
312    Ok((file, started.elapsed()))
313}
314
315pub(super) fn lock_cache_file_with_wait_observer(
316    path: &Path,
317    poll_interval: Duration,
318    mut observer: impl FnMut(Duration),
319) -> Result<(File, Duration), CacheFsError> {
320    let file = open_cache_lock_file(path)?;
321    let wait = ic_host_fs::durable::lock_exclusive_with_wait(
322        &file,
323        poll_interval.min(Duration::from_millis(25)),
324        |elapsed| {
325            observer(elapsed);
326            Ok(())
327        },
328    )
329    .map_err(|source| CacheFsError {
330        operation: "try lock cache",
331        path: path.to_owned(),
332        source,
333    })?;
334    Ok((file, wait))
335}
336
337pub(super) fn try_lock_cache_file(path: &Path) -> Result<Option<File>, CacheFsError> {
338    match ic_host_fs::durable::try_lock_regular_file_with_parents(path).map_err(io::Error::from) {
339        Ok(file) => Ok(Some(file)),
340        Err(error) if error.kind() == io::ErrorKind::WouldBlock => Ok(None),
341        Err(source) => Err(CacheFsError {
342            operation: "try lock cache",
343            path: path.to_owned(),
344            source,
345        }),
346    }
347}
348
349fn open_cache_lock_file(path: &Path) -> Result<File, CacheFsError> {
350    ic_host_fs::durable::open_regular_lock_file_with_parents(path)
351        .map_err(io::Error::from)
352        .map_err(|source| CacheFsError {
353            operation: "open cache lock",
354            path: path.to_owned(),
355            source,
356        })
357}
358
359pub(super) fn record_cache_entry_use(path: &Path) -> Result<(), CacheFsError> {
360    write_last_used(path, SystemTime::now())
361}
362
363pub(super) fn cache_maintenance_due(
364    path: &Path,
365    minimum_interval: Option<Duration>,
366    maintenance_identity: &str,
367) -> Result<bool, CacheFsError> {
368    let Some(minimum_interval) = minimum_interval else {
369        return Ok(true);
370    };
371    let marker = path.join(LAST_MAINTENANCE_FILE);
372    // Allow both LF and CRLF for the timestamp and policy-identity lines.
373    let maximum_len = MAX_TIMESTAMP_BYTES + maintenance_identity.len() + 4;
374    let contents = match read_stamp_with_limit(&marker, maximum_len) {
375        Ok(Some(contents)) => contents,
376        Ok(None) => return Ok(true),
377        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(true),
378        Err(source) => {
379            return Err(CacheFsError {
380                operation: "read cache maintenance time",
381                path: marker,
382                source,
383            });
384        }
385    };
386    let mut lines = contents.lines();
387    let Some(last_maintenance) = lines.next().and_then(decode_system_time) else {
388        return Ok(true);
389    };
390    if lines.next() != Some(maintenance_identity) {
391        return Ok(true);
392    }
393    Ok(match SystemTime::now().duration_since(last_maintenance) {
394        Ok(elapsed) => elapsed >= minimum_interval,
395        Err(_) => true,
396    })
397}
398
399pub(super) fn record_cache_maintenance(
400    path: &Path,
401    maintenance_identity: &str,
402) -> Result<(), CacheFsError> {
403    let marker = path.join(LAST_MAINTENANCE_FILE);
404    let elapsed = encode_system_time(&marker, SystemTime::now())?;
405    let contents = format!("{}\n{maintenance_identity}\n", elapsed.as_nanos());
406    ic_host_fs::durable::write_bytes(&marker, contents.as_bytes()).map_err(|source| CacheFsError {
407        operation: "record cache maintenance time",
408        path: marker,
409        source,
410    })
411}
412
413pub(super) fn perform_scheduled_cache_maintenance(
414    path: &Path,
415    minimum_interval: Option<Duration>,
416    maintenance_identity: &str,
417    maintenance: impl FnOnce() -> Result<ArtifactCachePruneReport, String>,
418) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
419    let started = Instant::now();
420    match cache_maintenance_due(path, minimum_interval, maintenance_identity) {
421        Ok(false) => return (None, Some(started.elapsed())),
422        Ok(true) => {}
423        Err(error) => {
424            return (
425                Some(ArtifactCacheMaintenance::PruneFailed {
426                    message: error.to_string(),
427                }),
428                Some(started.elapsed()),
429            );
430        }
431    }
432
433    let result = maintenance();
434    let marker = record_cache_maintenance(path, maintenance_identity);
435    let outcome = match (result, marker) {
436        (Ok(report), Ok(())) => ArtifactCacheMaintenance::Pruned(report),
437        (Err(message), Ok(())) => ArtifactCacheMaintenance::PruneFailed { message },
438        (Ok(_), Err(error)) => ArtifactCacheMaintenance::PruneFailed {
439            message: error.to_string(),
440        },
441        (Err(message), Err(marker)) => ArtifactCacheMaintenance::PruneFailed {
442            message: format!(
443                "{message}; additionally failed to record the maintenance attempt: {marker}"
444            ),
445        },
446    };
447    (Some(outcome), Some(started.elapsed()))
448}
449
450pub(super) fn write_last_used(path: &Path, last_used: SystemTime) -> Result<(), CacheFsError> {
451    let marker = path.join(LAST_USED_FILE);
452    write_system_time(&marker, last_used, "record cache use time")
453}
454
455fn write_system_time(
456    path: &Path,
457    timestamp: SystemTime,
458    operation: &'static str,
459) -> Result<(), CacheFsError> {
460    let elapsed = encode_system_time(path, timestamp)?;
461    ic_host_fs::durable::write_bytes(path, elapsed.as_nanos().to_string().as_bytes()).map_err(
462        |source| CacheFsError {
463            operation,
464            path: path.to_owned(),
465            source,
466        },
467    )
468}
469
470fn encode_system_time(path: &Path, timestamp: SystemTime) -> Result<Duration, CacheFsError> {
471    timestamp
472        .duration_since(UNIX_EPOCH)
473        .map_err(|source| CacheFsError {
474            operation: "encode cache time",
475            path: path.to_owned(),
476            source: io::Error::new(io::ErrorKind::InvalidInput, source),
477        })
478}
479
480fn decode_system_time(contents: &str) -> Option<SystemTime> {
481    let nanoseconds = contents.parse::<u128>().ok()?;
482    let seconds = u64::try_from(nanoseconds / 1_000_000_000).ok()?;
483    let subsecond_nanos = (nanoseconds % 1_000_000_000) as u32;
484    UNIX_EPOCH.checked_add(Duration::new(seconds, subsecond_nanos))
485}
486
487pub(super) fn prune_direct_child_directories(
488    cache_root: &Path,
489    policy: ArtifactCachePrunePolicy,
490    protected_entry: Option<&Path>,
491    is_eligible: impl Fn(&Path) -> bool,
492) -> Result<ArtifactCachePruneReport, CacheFsError> {
493    let mut entries = cache_entries(cache_root, is_eligible)?;
494    let bytes_before = entries
495        .iter()
496        .fold(0_u64, |total, entry| total.saturating_add(entry.bytes));
497    let mut report = ArtifactCachePruneReport {
498        entries_scanned: entries.len(),
499        entries_removed: 0,
500        bytes_before,
501        bytes_removed: 0,
502        uncommitted_directories_removed: 0,
503        uncommitted_bytes_removed: 0,
504    };
505    let now = SystemTime::now();
506
507    if let Some(max_age) = policy.max_age() {
508        for entry in &mut entries {
509            let age = now.duration_since(entry.last_used).unwrap_or_default();
510            if protected_entry != Some(entry.path.as_path()) && age > max_age {
511                remove_cache_entry(entry, &mut report)?;
512            }
513        }
514    }
515
516    if let Some(max_size_bytes) = policy.max_size_bytes()
517        && report.bytes_retained() > max_size_bytes
518    {
519        entries.sort_by(|left, right| {
520            left.last_used
521                .cmp(&right.last_used)
522                .then_with(|| left.path.cmp(&right.path))
523        });
524        for entry in &mut entries {
525            if report.bytes_retained() <= max_size_bytes {
526                break;
527            }
528            if protected_entry == Some(entry.path.as_path()) {
529                continue;
530            }
531            remove_cache_entry(entry, &mut report)?;
532        }
533    }
534
535    Ok(report)
536}
537
538pub(super) fn directory_logical_size(path: &Path) -> io::Result<u64> {
539    let mut total = 0_u64;
540    let mut pending = vec![path.to_owned()];
541    while let Some(current) = pending.pop() {
542        let metadata = fs::symlink_metadata(&current)?;
543        if metadata.is_dir() {
544            for entry in fs::read_dir(&current)? {
545                let path = entry?.path();
546                let metadata = fs::symlink_metadata(&path)?;
547                if metadata.is_dir() {
548                    pending.push(path);
549                } else {
550                    total = total.saturating_add(metadata.len());
551                }
552            }
553        } else {
554            total = total.saturating_add(metadata.len());
555        }
556    }
557    Ok(total)
558}
559
560pub(super) fn is_sha256_directory(path: &Path) -> bool {
561    path.file_name().is_some_and(|name| {
562        let bytes = name.as_encoded_bytes();
563        bytes.len() == 64 && bytes.iter().all(u8::is_ascii_hexdigit)
564    })
565}
566
567pub(super) fn remove_path_if_present(path: &Path) -> io::Result<()> {
568    let metadata = match fs::symlink_metadata(path) {
569        Ok(metadata) => metadata,
570        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
571        Err(error) => return Err(error),
572    };
573    if metadata.file_type().is_dir() {
574        fs::remove_dir_all(path)
575    } else {
576        fs::remove_file(path)
577    }
578}
579
580struct CacheEntry {
581    path: PathBuf,
582    bytes: u64,
583    last_used: SystemTime,
584    removed: bool,
585}
586
587impl std::fmt::Display for CacheFsError {
588    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
589        write!(
590            formatter,
591            "failed to {} at {}: {}",
592            self.operation,
593            self.path.display(),
594            self.source
595        )
596    }
597}
598
599impl std::error::Error for CacheFsError {
600    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
601        Some(&self.source)
602    }
603}
604
605fn cache_entries(
606    cache_root: &Path,
607    is_eligible: impl Fn(&Path) -> bool,
608) -> Result<Vec<CacheEntry>, CacheFsError> {
609    let read_dir = match fs::read_dir(cache_root) {
610        Ok(read_dir) => read_dir,
611        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
612        Err(source) => {
613            return Err(CacheFsError {
614                operation: "read cache directory",
615                path: cache_root.to_owned(),
616                source,
617            });
618        }
619    };
620    let mut entries = Vec::new();
621    for directory_entry in read_dir {
622        let directory_entry = directory_entry.map_err(|source| CacheFsError {
623            operation: "read cache entry",
624            path: cache_root.to_owned(),
625            source,
626        })?;
627        let path = directory_entry.path();
628        let file_type = directory_entry.file_type().map_err(|source| CacheFsError {
629            operation: "inspect cache entry",
630            path: path.clone(),
631            source,
632        })?;
633        if !file_type.is_dir() || !is_eligible(&path) {
634            continue;
635        }
636        let bytes = directory_logical_size(&path).map_err(|source| CacheFsError {
637            operation: "measure cache entry",
638            path: path.clone(),
639            source,
640        })?;
641        let last_used = cache_entry_last_used(&path).map_err(|source| CacheFsError {
642            operation: "read cache use time",
643            path: path.clone(),
644            source,
645        })?;
646        entries.push(CacheEntry {
647            path,
648            bytes,
649            last_used,
650            removed: false,
651        });
652    }
653    Ok(entries)
654}
655
656pub(super) fn cache_entry_last_used(path: &Path) -> io::Result<SystemTime> {
657    let marker = path.join(LAST_USED_FILE);
658    if let Ok(Some(contents)) = read_stamp_with_limit(&marker, MAX_TIMESTAMP_BYTES)
659        && let Some(timestamp) = decode_system_time(&contents)
660    {
661        return Ok(timestamp);
662    }
663    fs::metadata(path)?.modified()
664}
665
666fn remove_cache_entry(
667    entry: &mut CacheEntry,
668    report: &mut ArtifactCachePruneReport,
669) -> Result<(), CacheFsError> {
670    if entry.removed {
671        return Ok(());
672    }
673    let Some(_retention_lock) = try_lock_cache_file(&entry.path.join(RETENTION_LOCK_FILE))? else {
674        return Ok(());
675    };
676    remove_path_if_present(&entry.path).map_err(|source| CacheFsError {
677        operation: "prune cache entry",
678        path: entry.path.clone(),
679        source,
680    })?;
681    entry.removed = true;
682    report.entries_removed += 1;
683    report.bytes_removed = report.bytes_removed.saturating_add(entry.bytes);
684    Ok(())
685}
686
687#[cfg(test)]
688mod tests {
689    use super::directory_logical_size;
690    use crate::artifacts::test_support::unique_temp_directory;
691    use std::{
692        fs,
693        time::{Duration, SystemTime, UNIX_EPOCH},
694    };
695
696    #[cfg(unix)]
697    #[test]
698    fn cache_lock_admission_preserves_bytes_and_returns_an_unlocked_close_on_exec_file() {
699        use std::os::fd::AsRawFd as _;
700
701        let root = unique_temp_directory("cache-lock-admission");
702        let path = root.join("nested/lock");
703        let file = super::open_cache_lock_file(&path).unwrap();
704        assert!(super::try_lock_cache_file(&path).unwrap().is_some());
705        // SAFETY: inspect descriptor flags on the live file without changing them.
706        let flags = unsafe { libc::fcntl(file.as_raw_fd(), libc::F_GETFD) };
707        assert!(flags >= 0);
708        assert_ne!(flags & libc::FD_CLOEXEC, 0);
709        drop(file);
710        fs::write(&path, b"retained lock bytes").unwrap();
711        drop(super::lock_cache_file(&path).unwrap());
712        assert_eq!(fs::read(&path).unwrap(), b"retained lock bytes");
713        fs::remove_dir_all(root).unwrap();
714    }
715
716    #[cfg(unix)]
717    #[test]
718    fn cache_lock_callers_refuse_redirected_final_entries() {
719        use std::os::unix::fs::symlink;
720
721        let root = unique_temp_directory("cache-lock-redirected");
722        let target = root.join("target");
723        fs::write(&target, b"target bytes").unwrap();
724        let redirected = root.join("redirected");
725        symlink(&target, &redirected).unwrap();
726        for result in [
727            super::open_cache_lock_file(&redirected),
728            super::lock_cache_file(&redirected).map(|(file, _)| file),
729            super::lock_cache_file_with_wait_observer(
730                &redirected,
731                Duration::from_millis(5),
732                |_| panic!("redirected lock reached acquisition"),
733            )
734            .map(|(file, _)| file),
735        ] {
736            let error = result.unwrap_err();
737            assert_eq!(error.path, redirected);
738            assert!(matches!(
739                error.source.get_ref().and_then(|cause| cause
740                    .downcast_ref::<ic_host_fs::durable::RegularFileLockError>(
741                )),
742                Some(ic_host_fs::durable::RegularFileLockError::NotRegular)
743            ));
744        }
745        symlink(&target, root.join(super::RETENTION_LOCK_FILE)).unwrap();
746        assert!(super::RetainedCacheEntry::acquire(&root).is_err());
747        assert_eq!(fs::read(&target).unwrap(), b"target bytes");
748        fs::remove_dir_all(root).unwrap();
749    }
750
751    #[cfg(unix)]
752    #[test]
753    fn nonblocking_cache_locks_preserve_contention_and_reject_redirected_files() {
754        use std::os::unix::fs::symlink;
755
756        let root = unique_temp_directory("nonblocking-cache-lock-admission");
757        let lock_path = root.join("lock");
758        fs::write(&lock_path, b"existing lock bytes").unwrap();
759        let held = super::try_lock_cache_file(&lock_path).unwrap().unwrap();
760        assert!(super::try_lock_cache_file(&lock_path).unwrap().is_none());
761        assert_eq!(fs::read(&lock_path).unwrap(), b"existing lock bytes");
762        drop(held);
763        assert!(super::try_lock_cache_file(&lock_path).unwrap().is_some());
764
765        let redirected = root.join("redirected-lock");
766        symlink(&lock_path, &redirected).unwrap();
767        let error = super::try_lock_cache_file(&redirected).unwrap_err();
768        assert_eq!(error.path, redirected);
769        assert!(matches!(
770                error.source.get_ref().and_then(|cause| cause
771                    .downcast_ref::<ic_host_fs::durable::RegularFileLockError>(
772                )),
773                Some(ic_host_fs::durable::RegularFileLockError::NotRegular)
774            ));
775        assert_eq!(fs::read(&lock_path).unwrap(), b"existing lock bytes");
776        fs::remove_dir_all(root).unwrap();
777    }
778
779    #[cfg(unix)]
780    #[test]
781    fn final_retention_owner_releases_lock_with_a_duplicate_descriptor_open() {
782        let root = unique_temp_directory("retention-duplicate-descriptor");
783        let retained = super::RetainedCacheEntry::acquire(&root).unwrap();
784        // A process spawn can duplicate this descriptor before close-on-exec.
785        let inherited = retained._lock.0.try_clone().unwrap();
786        let clone = retained.clone();
787        let independently_retained = super::RetainedCacheEntry::acquire(&root).unwrap();
788        let lock_path = root.join(super::RETENTION_LOCK_FILE);
789        let available = || super::try_lock_cache_file(&lock_path).unwrap().is_some();
790        drop(retained);
791        assert!(!available(), "a record clone still retains the entry");
792        drop(clone);
793        assert!(
794            !available(),
795            "an independent acquisition still retains the entry"
796        );
797        drop(independently_retained);
798        assert!(
799            available(),
800            "descriptor duplication must not extend record ownership"
801        );
802        drop(inherited);
803        fs::remove_dir_all(root).unwrap();
804    }
805
806    #[test]
807    fn last_use_markers_preserve_timestamps_and_bounded_fallbacks() {
808        let root = unique_temp_directory("bounded-last-use-marker");
809        let marker = root.join(super::LAST_USED_FILE);
810        let modified = || fs::metadata(&root).unwrap().modified().unwrap();
811        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
812        for timestamp in [
813            UNIX_EPOCH + Duration::from_nanos(123),
814            SystemTime::now() + Duration::from_secs(3600),
815        ] {
816            super::write_last_used(&root, timestamp).unwrap();
817            assert_eq!(super::cache_entry_last_used(&root).unwrap(), timestamp);
818        }
819        let overflowing_timestamp = u128::MAX.to_string();
820        for invalid in [
821            b"invalid".as_slice(),
822            overflowing_timestamp.as_bytes(),
823            &[0xff],
824        ] {
825            fs::write(&marker, invalid).unwrap();
826            assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
827        }
828        fs::File::create(&marker)
829            .unwrap()
830            .set_len(1024 * 1024 * 1024)
831            .unwrap();
832        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
833        fs::remove_file(&marker).unwrap();
834        fs::create_dir(&marker).unwrap();
835        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
836        fs::remove_dir_all(root).unwrap();
837    }
838
839    #[test]
840    fn maintenance_markers_preserve_policy_intervals_and_read_errors() {
841        let root = unique_temp_directory("bounded-maintenance-marker");
842        let marker = root.join(super::LAST_MAINTENANCE_FILE);
843        let identity = super::ArtifactCachePrunePolicy::new().maintenance_identity();
844        let interval = Some(Duration::from_secs(3600));
845        let due = || super::cache_maintenance_due(&root, interval, &identity).unwrap();
846        assert!(due());
847        super::record_cache_maintenance(&root, &identity).unwrap();
848        assert!(!due());
849        assert!(super::cache_maintenance_due(&root, interval, "other-policy").unwrap());
850        assert!(super::cache_maintenance_due(&root, Some(Duration::ZERO), &identity).unwrap());
851
852        let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap();
853        for suffix in ["", "\r\n"] {
854            fs::write(&marker, format!("{}\r\n{identity}{suffix}", now.as_nanos())).unwrap();
855            assert!(!due());
856        }
857        for timestamp in [
858            "invalid".to_owned(),
859            "0".to_owned(),
860            (now + Duration::from_secs(3600)).as_nanos().to_string(),
861        ] {
862            fs::write(&marker, format!("{timestamp}\n{identity}\n")).unwrap();
863            assert!(due());
864        }
865        super::record_cache_maintenance(&root, &identity).unwrap();
866        fs::OpenOptions::new()
867            .write(true)
868            .open(&marker)
869            .unwrap()
870            .set_len(1024 * 1024 * 1024)
871            .unwrap();
872        assert!(due());
873
874        fs::write(&marker, [0xff]).unwrap();
875        let error = super::cache_maintenance_due(&root, interval, &identity).unwrap_err();
876        assert_eq!(error.operation, "read cache maintenance time");
877        assert_eq!(error.path, marker);
878        assert_eq!(error.source.kind(), std::io::ErrorKind::InvalidData);
879        fs::remove_file(&marker).unwrap();
880        fs::create_dir(&marker).unwrap();
881        assert!(super::cache_maintenance_due(&root, interval, &identity).is_err());
882        assert!(super::cache_maintenance_due(&root, None, &identity).unwrap());
883        fs::remove_dir_all(root).unwrap();
884    }
885
886    #[test]
887    fn cache_directory_tags_preserve_valid_standard_signatures() {
888        let root = unique_temp_directory("cache-tag-signatures");
889        let tag = root.join("CACHEDIR.TAG");
890        let signature = "Signature: 8a477f597d28d172789f06886806bc55";
891        for contents in [
892            signature.to_owned(),
893            format!("{signature}\r\n# Created by another application\r\n"),
894            format!("{signature}\n# {}\n", "comment".repeat(100_000)),
895        ] {
896            fs::write(&tag, &contents).unwrap();
897            super::ensure_cache_directory_tag(&root).unwrap();
898            assert_eq!(fs::read_to_string(&tag).unwrap(), contents);
899        }
900        for invalid in ["", &signature[..42], "Signature: incorrect"] {
901            fs::write(&tag, invalid).unwrap();
902            super::ensure_cache_directory_tag(&root).unwrap();
903            assert_eq!(
904                fs::read_to_string(&tag).unwrap(),
905                super::CACHE_DIRECTORY_TAG
906            );
907        }
908        fs::remove_dir_all(root).unwrap();
909    }
910
911    #[test]
912    #[cfg(unix)]
913    fn cache_directory_tag_replaces_symlinks_without_changing_referents() {
914        let root = unique_temp_directory("cache-tag-symlink");
915        let referent = root.join("other-application-tag");
916        let contents = "Signature: 8a477f597d28d172789f06886806bc55\n# Preserve this file\n";
917        fs::write(&referent, contents).unwrap();
918        let tag = root.join("CACHEDIR.TAG");
919        std::os::unix::fs::symlink(&referent, &tag).unwrap();
920        super::ensure_cache_directory_tag(&root).unwrap();
921        assert!(fs::symlink_metadata(&tag).unwrap().file_type().is_file());
922        assert_eq!(fs::read_to_string(&referent).unwrap(), contents);
923        assert_eq!(
924            fs::read_to_string(&tag).unwrap(),
925            super::CACHE_DIRECTORY_TAG
926        );
927
928        fs::remove_file(&tag).unwrap();
929        fs::create_dir(&tag).unwrap();
930        let error = super::ensure_cache_directory_tag(&root).unwrap_err();
931        assert_eq!(error.operation, "write cache directory tag");
932        assert!(tag.is_dir());
933        fs::remove_dir_all(root).unwrap();
934    }
935
936    #[test]
937    fn directory_size_sums_wide_and_nested_files() {
938        let root = unique_temp_directory("directory-logical-size");
939        assert_eq!(directory_logical_size(&root).unwrap(), 0);
940        fs::create_dir_all(root.join("wide")).unwrap();
941        fs::create_dir_all(root.join("nested/deep/empty")).unwrap();
942        let mut expected = 0;
943        for index in 0..128 {
944            let bytes = vec![42; index % 13];
945            fs::write(root.join("wide").join(index.to_string()), &bytes).unwrap();
946            expected += bytes.len() as u64;
947        }
948        let sparse = root.join("nested/deep/sparse");
949        fs::File::create(&sparse)
950            .unwrap()
951            .set_len(1024 * 1024)
952            .unwrap();
953        assert_eq!(directory_logical_size(&sparse).unwrap(), 1024 * 1024);
954        assert_eq!(
955            directory_logical_size(&root).unwrap(),
956            expected + 1024 * 1024
957        );
958        assert_eq!(
959            directory_logical_size(&root.join("missing"))
960                .unwrap_err()
961                .kind(),
962            std::io::ErrorKind::NotFound,
963        );
964        fs::remove_dir_all(root).unwrap();
965    }
966
967    #[test]
968    #[cfg(unix)]
969    fn directory_size_counts_symlinks_without_following_them() {
970        let root = unique_temp_directory("directory-size-symlinks");
971        let walked = root.join("walked");
972        fs::create_dir_all(&walked).unwrap();
973        fs::create_dir_all(root.join("external")).unwrap();
974        fs::write(root.join("external/payload"), vec![42; 4096]).unwrap();
975        fs::write(root.join("outside-file"), vec![42; 4096]).unwrap();
976        fs::write(walked.join("payload"), b"abc").unwrap();
977        let targets = ["../external", "../outside-file", "missing", "."];
978        for (index, target) in targets.iter().enumerate() {
979            std::os::unix::fs::symlink(target, walked.join(index.to_string())).unwrap();
980        }
981        let expected = 3 + targets
982            .iter()
983            .map(|target| target.len() as u64)
984            .sum::<u64>();
985        assert_eq!(directory_logical_size(&walked).unwrap(), expected);
986        assert_eq!(
987            directory_logical_size(&walked.join("0")).unwrap(),
988            targets[0].len() as u64,
989        );
990        fs::remove_dir_all(root).unwrap();
991    }
992}