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#[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
53pub(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#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
97pub struct ArtifactCachePrunePolicy {
98 max_age: Option<Duration>,
99 max_size_bytes: Option<u64>,
100}
101
102#[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#[non_exhaustive]
115#[derive(Clone, Debug, Eq, PartialEq)]
116pub enum ArtifactCacheMaintenance {
117 Pruned(ArtifactCachePruneReport),
119 PruneFailed {
121 message: String,
123 },
124}
125
126impl ArtifactCachePrunePolicy {
127 #[must_use]
129 pub const fn new() -> Self {
130 Self {
131 max_age: None,
132 max_size_bytes: None,
133 }
134 }
135
136 #[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 #[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 #[must_use]
152 pub const fn max_age(self) -> Option<Duration> {
153 self.max_age
154 }
155
156 #[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 #[must_use]
174 pub const fn entries_scanned(self) -> usize {
175 self.entries_scanned
176 }
177
178 #[must_use]
180 pub const fn entries_removed(self) -> usize {
181 self.entries_removed
182 }
183
184 #[must_use]
186 pub const fn entries_retained(self) -> usize {
187 self.entries_scanned.saturating_sub(self.entries_removed)
188 }
189
190 #[must_use]
192 pub const fn bytes_before(self) -> u64 {
193 self.bytes_before
194 }
195
196 #[must_use]
198 pub const fn bytes_removed(self) -> u64 {
199 self.bytes_removed
200 }
201
202 #[must_use]
204 pub const fn bytes_retained(self) -> u64 {
205 self.bytes_before.saturating_sub(self.bytes_removed)
206 }
207
208 #[must_use]
210 pub const fn uncommitted_directories_removed(self) -> usize {
211 self.uncommitted_directories_removed
212 }
213
214 #[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 #[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 #[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(¤t)?;
520 if metadata.is_dir() {
521 for entry in fs::read_dir(¤t)? {
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}