Skip to main content

ic_testkit/artifacts/
cache_fs.rs

1use fs2::FileExt as _;
2use std::{
3    fs::{self, File, OpenOptions},
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    if let Some(parent) = path.parent() {
351        fs::create_dir_all(parent).map_err(|source| CacheFsError {
352            operation: "create cache lock directory",
353            path: parent.to_owned(),
354            source,
355        })?;
356    }
357    OpenOptions::new()
358        .create(true)
359        .read(true)
360        .write(true)
361        .truncate(false)
362        .open(path)
363        .map_err(|source| CacheFsError {
364            operation: "open cache lock",
365            path: path.to_owned(),
366            source,
367        })
368}
369
370pub(super) fn record_cache_entry_use(path: &Path) -> Result<(), CacheFsError> {
371    write_last_used(path, SystemTime::now())
372}
373
374pub(super) fn cache_maintenance_due(
375    path: &Path,
376    minimum_interval: Option<Duration>,
377    maintenance_identity: &str,
378) -> Result<bool, CacheFsError> {
379    let Some(minimum_interval) = minimum_interval else {
380        return Ok(true);
381    };
382    let marker = path.join(LAST_MAINTENANCE_FILE);
383    // Allow both LF and CRLF for the timestamp and policy-identity lines.
384    let maximum_len = MAX_TIMESTAMP_BYTES + maintenance_identity.len() + 4;
385    let contents = match read_stamp_with_limit(&marker, maximum_len) {
386        Ok(Some(contents)) => contents,
387        Ok(None) => return Ok(true),
388        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(true),
389        Err(source) => {
390            return Err(CacheFsError {
391                operation: "read cache maintenance time",
392                path: marker,
393                source,
394            });
395        }
396    };
397    let mut lines = contents.lines();
398    let Some(last_maintenance) = lines.next().and_then(decode_system_time) else {
399        return Ok(true);
400    };
401    if lines.next() != Some(maintenance_identity) {
402        return Ok(true);
403    }
404    Ok(match SystemTime::now().duration_since(last_maintenance) {
405        Ok(elapsed) => elapsed >= minimum_interval,
406        Err(_) => true,
407    })
408}
409
410pub(super) fn record_cache_maintenance(
411    path: &Path,
412    maintenance_identity: &str,
413) -> Result<(), CacheFsError> {
414    let marker = path.join(LAST_MAINTENANCE_FILE);
415    let elapsed = encode_system_time(&marker, SystemTime::now())?;
416    let contents = format!("{}\n{maintenance_identity}\n", elapsed.as_nanos());
417    ic_host_fs::durable::write_bytes(&marker, contents.as_bytes()).map_err(|source| CacheFsError {
418        operation: "record cache maintenance time",
419        path: marker,
420        source,
421    })
422}
423
424pub(super) fn perform_scheduled_cache_maintenance(
425    path: &Path,
426    minimum_interval: Option<Duration>,
427    maintenance_identity: &str,
428    maintenance: impl FnOnce() -> Result<ArtifactCachePruneReport, String>,
429) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
430    let started = Instant::now();
431    match cache_maintenance_due(path, minimum_interval, maintenance_identity) {
432        Ok(false) => return (None, Some(started.elapsed())),
433        Ok(true) => {}
434        Err(error) => {
435            return (
436                Some(ArtifactCacheMaintenance::PruneFailed {
437                    message: error.to_string(),
438                }),
439                Some(started.elapsed()),
440            );
441        }
442    }
443
444    let result = maintenance();
445    let marker = record_cache_maintenance(path, maintenance_identity);
446    let outcome = match (result, marker) {
447        (Ok(report), Ok(())) => ArtifactCacheMaintenance::Pruned(report),
448        (Err(message), Ok(())) => ArtifactCacheMaintenance::PruneFailed { message },
449        (Ok(_), Err(error)) => ArtifactCacheMaintenance::PruneFailed {
450            message: error.to_string(),
451        },
452        (Err(message), Err(marker)) => ArtifactCacheMaintenance::PruneFailed {
453            message: format!(
454                "{message}; additionally failed to record the maintenance attempt: {marker}"
455            ),
456        },
457    };
458    (Some(outcome), Some(started.elapsed()))
459}
460
461pub(super) fn write_last_used(path: &Path, last_used: SystemTime) -> Result<(), CacheFsError> {
462    let marker = path.join(LAST_USED_FILE);
463    write_system_time(&marker, last_used, "record cache use time")
464}
465
466fn write_system_time(
467    path: &Path,
468    timestamp: SystemTime,
469    operation: &'static str,
470) -> Result<(), CacheFsError> {
471    let elapsed = encode_system_time(path, timestamp)?;
472    ic_host_fs::durable::write_bytes(path, elapsed.as_nanos().to_string().as_bytes()).map_err(
473        |source| CacheFsError {
474            operation,
475            path: path.to_owned(),
476            source,
477        },
478    )
479}
480
481fn encode_system_time(path: &Path, timestamp: SystemTime) -> Result<Duration, CacheFsError> {
482    timestamp
483        .duration_since(UNIX_EPOCH)
484        .map_err(|source| CacheFsError {
485            operation: "encode cache time",
486            path: path.to_owned(),
487            source: io::Error::new(io::ErrorKind::InvalidInput, source),
488        })
489}
490
491fn decode_system_time(contents: &str) -> Option<SystemTime> {
492    let nanoseconds = contents.parse::<u128>().ok()?;
493    let seconds = u64::try_from(nanoseconds / 1_000_000_000).ok()?;
494    let subsecond_nanos = (nanoseconds % 1_000_000_000) as u32;
495    UNIX_EPOCH.checked_add(Duration::new(seconds, subsecond_nanos))
496}
497
498pub(super) fn prune_direct_child_directories(
499    cache_root: &Path,
500    policy: ArtifactCachePrunePolicy,
501    protected_entry: Option<&Path>,
502    is_eligible: impl Fn(&Path) -> bool,
503) -> Result<ArtifactCachePruneReport, CacheFsError> {
504    let mut entries = cache_entries(cache_root, is_eligible)?;
505    let bytes_before = entries
506        .iter()
507        .fold(0_u64, |total, entry| total.saturating_add(entry.bytes));
508    let mut report = ArtifactCachePruneReport {
509        entries_scanned: entries.len(),
510        entries_removed: 0,
511        bytes_before,
512        bytes_removed: 0,
513        uncommitted_directories_removed: 0,
514        uncommitted_bytes_removed: 0,
515    };
516    let now = SystemTime::now();
517
518    if let Some(max_age) = policy.max_age() {
519        for entry in &mut entries {
520            let age = now.duration_since(entry.last_used).unwrap_or_default();
521            if protected_entry != Some(entry.path.as_path()) && age > max_age {
522                remove_cache_entry(entry, &mut report)?;
523            }
524        }
525    }
526
527    if let Some(max_size_bytes) = policy.max_size_bytes()
528        && report.bytes_retained() > max_size_bytes
529    {
530        entries.sort_by(|left, right| {
531            left.last_used
532                .cmp(&right.last_used)
533                .then_with(|| left.path.cmp(&right.path))
534        });
535        for entry in &mut entries {
536            if report.bytes_retained() <= max_size_bytes {
537                break;
538            }
539            if protected_entry == Some(entry.path.as_path()) {
540                continue;
541            }
542            remove_cache_entry(entry, &mut report)?;
543        }
544    }
545
546    Ok(report)
547}
548
549pub(super) fn directory_logical_size(path: &Path) -> io::Result<u64> {
550    let mut total = 0_u64;
551    let mut pending = vec![path.to_owned()];
552    while let Some(current) = pending.pop() {
553        let metadata = fs::symlink_metadata(&current)?;
554        if metadata.is_dir() {
555            for entry in fs::read_dir(&current)? {
556                let path = entry?.path();
557                let metadata = fs::symlink_metadata(&path)?;
558                if metadata.is_dir() {
559                    pending.push(path);
560                } else {
561                    total = total.saturating_add(metadata.len());
562                }
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(Some(contents)) = read_stamp_with_limit(&marker, MAX_TIMESTAMP_BYTES)
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(test)]
699mod tests {
700    use super::directory_logical_size;
701    use crate::artifacts::test_support::unique_temp_directory;
702    use std::{
703        fs,
704        time::{Duration, SystemTime, UNIX_EPOCH},
705    };
706
707    #[cfg(unix)]
708    #[test]
709    fn nonblocking_cache_locks_preserve_contention_and_reject_redirected_files() {
710        use std::os::unix::fs::symlink;
711
712        let root = unique_temp_directory("nonblocking-cache-lock-admission");
713        let lock_path = root.join("lock");
714        fs::write(&lock_path, b"existing lock bytes").unwrap();
715        let held = super::try_lock_cache_file(&lock_path).unwrap().unwrap();
716        assert!(super::try_lock_cache_file(&lock_path).unwrap().is_none());
717        assert_eq!(fs::read(&lock_path).unwrap(), b"existing lock bytes");
718        drop(held);
719        assert!(super::try_lock_cache_file(&lock_path).unwrap().is_some());
720
721        let redirected = root.join("redirected-lock");
722        symlink(&lock_path, &redirected).unwrap();
723        let error = super::try_lock_cache_file(&redirected).unwrap_err();
724        assert_eq!(error.path, redirected);
725        assert!(matches!(
726                error.source.get_ref().and_then(|cause| cause
727                    .downcast_ref::<ic_host_fs::durable::RegularFileLockError>(
728                )),
729                Some(ic_host_fs::durable::RegularFileLockError::NotRegular)
730            ));
731        assert_eq!(fs::read(&lock_path).unwrap(), b"existing lock bytes");
732        fs::remove_dir_all(root).unwrap();
733    }
734
735    #[cfg(unix)]
736    #[test]
737    fn final_retention_owner_releases_lock_with_a_duplicate_descriptor_open() {
738        let root = unique_temp_directory("retention-duplicate-descriptor");
739        let retained = super::RetainedCacheEntry::acquire(&root).unwrap();
740        // A process spawn can duplicate this descriptor before close-on-exec.
741        let inherited = retained._lock.0.try_clone().unwrap();
742        let clone = retained.clone();
743        let independently_retained = super::RetainedCacheEntry::acquire(&root).unwrap();
744        let lock_path = root.join(super::RETENTION_LOCK_FILE);
745        let available = || super::try_lock_cache_file(&lock_path).unwrap().is_some();
746        drop(retained);
747        assert!(!available(), "a record clone still retains the entry");
748        drop(clone);
749        assert!(
750            !available(),
751            "an independent acquisition still retains the entry"
752        );
753        drop(independently_retained);
754        assert!(
755            available(),
756            "descriptor duplication must not extend record ownership"
757        );
758        drop(inherited);
759        fs::remove_dir_all(root).unwrap();
760    }
761
762    #[test]
763    fn last_use_markers_preserve_timestamps_and_bounded_fallbacks() {
764        let root = unique_temp_directory("bounded-last-use-marker");
765        let marker = root.join(super::LAST_USED_FILE);
766        let modified = || fs::metadata(&root).unwrap().modified().unwrap();
767        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
768        for timestamp in [
769            UNIX_EPOCH + Duration::from_nanos(123),
770            SystemTime::now() + Duration::from_secs(3600),
771        ] {
772            super::write_last_used(&root, timestamp).unwrap();
773            assert_eq!(super::cache_entry_last_used(&root).unwrap(), timestamp);
774        }
775        let overflowing_timestamp = u128::MAX.to_string();
776        for invalid in [
777            b"invalid".as_slice(),
778            overflowing_timestamp.as_bytes(),
779            &[0xff],
780        ] {
781            fs::write(&marker, invalid).unwrap();
782            assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
783        }
784        fs::File::create(&marker)
785            .unwrap()
786            .set_len(1024 * 1024 * 1024)
787            .unwrap();
788        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
789        fs::remove_file(&marker).unwrap();
790        fs::create_dir(&marker).unwrap();
791        assert_eq!(super::cache_entry_last_used(&root).unwrap(), modified());
792        fs::remove_dir_all(root).unwrap();
793    }
794
795    #[test]
796    fn maintenance_markers_preserve_policy_intervals_and_read_errors() {
797        let root = unique_temp_directory("bounded-maintenance-marker");
798        let marker = root.join(super::LAST_MAINTENANCE_FILE);
799        let identity = super::ArtifactCachePrunePolicy::new().maintenance_identity();
800        let interval = Some(Duration::from_secs(3600));
801        let due = || super::cache_maintenance_due(&root, interval, &identity).unwrap();
802        assert!(due());
803        super::record_cache_maintenance(&root, &identity).unwrap();
804        assert!(!due());
805        assert!(super::cache_maintenance_due(&root, interval, "other-policy").unwrap());
806        assert!(super::cache_maintenance_due(&root, Some(Duration::ZERO), &identity).unwrap());
807
808        let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap();
809        for suffix in ["", "\r\n"] {
810            fs::write(&marker, format!("{}\r\n{identity}{suffix}", now.as_nanos())).unwrap();
811            assert!(!due());
812        }
813        for timestamp in [
814            "invalid".to_owned(),
815            "0".to_owned(),
816            (now + Duration::from_secs(3600)).as_nanos().to_string(),
817        ] {
818            fs::write(&marker, format!("{timestamp}\n{identity}\n")).unwrap();
819            assert!(due());
820        }
821        super::record_cache_maintenance(&root, &identity).unwrap();
822        fs::OpenOptions::new()
823            .write(true)
824            .open(&marker)
825            .unwrap()
826            .set_len(1024 * 1024 * 1024)
827            .unwrap();
828        assert!(due());
829
830        fs::write(&marker, [0xff]).unwrap();
831        let error = super::cache_maintenance_due(&root, interval, &identity).unwrap_err();
832        assert_eq!(error.operation, "read cache maintenance time");
833        assert_eq!(error.path, marker);
834        assert_eq!(error.source.kind(), std::io::ErrorKind::InvalidData);
835        fs::remove_file(&marker).unwrap();
836        fs::create_dir(&marker).unwrap();
837        assert!(super::cache_maintenance_due(&root, interval, &identity).is_err());
838        assert!(super::cache_maintenance_due(&root, None, &identity).unwrap());
839        fs::remove_dir_all(root).unwrap();
840    }
841
842    #[test]
843    fn cache_directory_tags_preserve_valid_standard_signatures() {
844        let root = unique_temp_directory("cache-tag-signatures");
845        let tag = root.join("CACHEDIR.TAG");
846        let signature = "Signature: 8a477f597d28d172789f06886806bc55";
847        for contents in [
848            signature.to_owned(),
849            format!("{signature}\r\n# Created by another application\r\n"),
850            format!("{signature}\n# {}\n", "comment".repeat(100_000)),
851        ] {
852            fs::write(&tag, &contents).unwrap();
853            super::ensure_cache_directory_tag(&root).unwrap();
854            assert_eq!(fs::read_to_string(&tag).unwrap(), contents);
855        }
856        for invalid in ["", &signature[..42], "Signature: incorrect"] {
857            fs::write(&tag, invalid).unwrap();
858            super::ensure_cache_directory_tag(&root).unwrap();
859            assert_eq!(
860                fs::read_to_string(&tag).unwrap(),
861                super::CACHE_DIRECTORY_TAG
862            );
863        }
864        fs::remove_dir_all(root).unwrap();
865    }
866
867    #[test]
868    #[cfg(unix)]
869    fn cache_directory_tag_replaces_symlinks_without_changing_referents() {
870        let root = unique_temp_directory("cache-tag-symlink");
871        let referent = root.join("other-application-tag");
872        let contents = "Signature: 8a477f597d28d172789f06886806bc55\n# Preserve this file\n";
873        fs::write(&referent, contents).unwrap();
874        let tag = root.join("CACHEDIR.TAG");
875        std::os::unix::fs::symlink(&referent, &tag).unwrap();
876        super::ensure_cache_directory_tag(&root).unwrap();
877        assert!(fs::symlink_metadata(&tag).unwrap().file_type().is_file());
878        assert_eq!(fs::read_to_string(&referent).unwrap(), contents);
879        assert_eq!(
880            fs::read_to_string(&tag).unwrap(),
881            super::CACHE_DIRECTORY_TAG
882        );
883
884        fs::remove_file(&tag).unwrap();
885        fs::create_dir(&tag).unwrap();
886        let error = super::ensure_cache_directory_tag(&root).unwrap_err();
887        assert_eq!(error.operation, "write cache directory tag");
888        assert!(tag.is_dir());
889        fs::remove_dir_all(root).unwrap();
890    }
891
892    #[test]
893    fn directory_size_sums_wide_and_nested_files() {
894        let root = unique_temp_directory("directory-logical-size");
895        assert_eq!(directory_logical_size(&root).unwrap(), 0);
896        fs::create_dir_all(root.join("wide")).unwrap();
897        fs::create_dir_all(root.join("nested/deep/empty")).unwrap();
898        let mut expected = 0;
899        for index in 0..128 {
900            let bytes = vec![42; index % 13];
901            fs::write(root.join("wide").join(index.to_string()), &bytes).unwrap();
902            expected += bytes.len() as u64;
903        }
904        let sparse = root.join("nested/deep/sparse");
905        fs::File::create(&sparse)
906            .unwrap()
907            .set_len(1024 * 1024)
908            .unwrap();
909        assert_eq!(directory_logical_size(&sparse).unwrap(), 1024 * 1024);
910        assert_eq!(
911            directory_logical_size(&root).unwrap(),
912            expected + 1024 * 1024
913        );
914        assert_eq!(
915            directory_logical_size(&root.join("missing"))
916                .unwrap_err()
917                .kind(),
918            std::io::ErrorKind::NotFound,
919        );
920        fs::remove_dir_all(root).unwrap();
921    }
922
923    #[test]
924    #[cfg(unix)]
925    fn directory_size_counts_symlinks_without_following_them() {
926        let root = unique_temp_directory("directory-size-symlinks");
927        let walked = root.join("walked");
928        fs::create_dir_all(&walked).unwrap();
929        fs::create_dir_all(root.join("external")).unwrap();
930        fs::write(root.join("external/payload"), vec![42; 4096]).unwrap();
931        fs::write(root.join("outside-file"), vec![42; 4096]).unwrap();
932        fs::write(walked.join("payload"), b"abc").unwrap();
933        let targets = ["../external", "../outside-file", "missing", "."];
934        for (index, target) in targets.iter().enumerate() {
935            std::os::unix::fs::symlink(target, walked.join(index.to_string())).unwrap();
936        }
937        let expected = 3 + targets
938            .iter()
939            .map(|target| target.len() as u64)
940            .sum::<u64>();
941        assert_eq!(directory_logical_size(&walked).unwrap(), expected);
942        assert_eq!(
943            directory_logical_size(&walked.join("0")).unwrap(),
944            targets[0].len() as u64,
945        );
946        fs::remove_dir_all(root).unwrap();
947    }
948}