1use std::{
2 collections::{BTreeMap, BTreeSet},
3 ffi::{OsStr, OsString},
4 fs::{self, File},
5 io,
6 path::{Path, PathBuf},
7 sync::atomic::{AtomicU64, Ordering},
8 time::{Duration, Instant},
9};
10
11use crate::timing::saturating_add_optional_duration;
12
13use super::{
14 cache_fs::{
15 ArtifactCacheMaintenance, ArtifactCachePrunePolicy, ArtifactCachePruneReport, CacheFsError,
16 LAST_USED_FILE, RetainedCacheEntry, canonicalize_allow_missing, directory_logical_size,
17 ensure_cache_directory_tag, is_sha256_directory, lock_cache_file,
18 perform_scheduled_cache_maintenance, prune_direct_child_directories,
19 record_cache_entry_use, remove_path_if_present, remove_unretained_entry,
20 try_lock_cache_file,
21 },
22 digest::{
23 InputDigest, InputHasher, copy_file_atomic, digest_bytes, digest_file,
24 digest_labeled_paths, os_bytes, write_atomic,
25 },
26 wasm_cache::{
27 ResolvedCargoBuildInputs, WasmBuildError, WasmBuildSpec, resolve_cargo_build_inputs,
28 },
29};
30
31const ARTIFACT_CACHE_FORMAT: &str = "ic-testkit-artifact-set-v1";
32const MANIFEST_FILE: &str = "manifest.ic-testkit";
33const OUTPUT_FILE_SUFFIX: &str = ".artifact";
34const MAX_PREPARATION_RETRIES: usize = 3;
35static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
36
37#[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)]
39pub enum ArtifactOutputValidation {
40 RegularFile,
42 #[default]
44 NonEmptyFile,
45}
46
47#[derive(Clone, Eq, PartialEq)]
55pub struct ArtifactCacheSpec {
56 cache_root: PathBuf,
57 namespace: String,
58 recipe_id: String,
59 coordination_scope: String,
60 inputs: Vec<LabeledPath>,
61 tools: Vec<LabeledPath>,
62 arguments: Vec<OsString>,
63 environment: BTreeMap<OsString, Option<OsString>>,
64 identities: Vec<LabeledIdentity>,
65 cargo_build_inputs: Vec<CargoBuildInputSet>,
66 outputs: Vec<OutputSpec>,
67 prune_policy: Option<ArtifactCachePrunePolicy>,
68 prune_interval: Option<Duration>,
69}
70
71pub enum ArtifactCachePreparation {
73 Reused(ArtifactCacheRecord),
75 Build(ArtifactBuildTransaction),
77}
78
79#[derive(Clone, Debug, Eq, PartialEq)]
81pub enum ArtifactCacheOutcome {
82 Built(ArtifactCacheRecord),
84 Reused(ArtifactCacheRecord),
86}
87
88#[derive(Clone, Debug, Eq, PartialEq)]
90pub struct ArtifactCacheRecord {
91 key: InputDigest,
92 input_digest: InputDigest,
93 artifacts: Vec<ArtifactCacheArtifact>,
94 _retention: RetainedCacheEntry,
95 timings: ArtifactCacheTimings,
96 maintenance: Option<ArtifactCacheMaintenance>,
97}
98
99#[derive(Clone, Debug, Eq, PartialEq)]
101pub struct ArtifactCacheArtifact {
102 name: String,
103 path: PathBuf,
104}
105
106#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
108pub struct ArtifactCacheTimings {
109 coordination_lock_wait: Duration,
110 content_lock_wait: Duration,
111 namespace_lock_wait: Duration,
112 input_capture: Duration,
113 cache_lookup: Duration,
114 caller_build: Option<Duration>,
115 output_validation: Duration,
116 publication: Duration,
117 materialization: Duration,
118 maintenance: Option<Duration>,
119 total: Duration,
120}
121
122pub struct ArtifactBuildTransaction {
124 spec: Box<ArtifactCacheSpec>,
125 resolved: ResolvedKey,
126 staging_directory: PathBuf,
127 entry_directory: PathBuf,
128 namespace_directory: PathBuf,
129 _coordination_lock: File,
130 _content_lock: File,
131 timings: ArtifactCacheTimings,
132 total_started: Instant,
133 caller_build_started: Instant,
134 staging_armed: bool,
135}
136
137#[non_exhaustive]
139#[derive(Debug)]
140pub enum ArtifactCacheError {
141 InvalidSpec { message: String },
143 Io {
145 operation: &'static str,
146 path: PathBuf,
147 source: io::Error,
148 },
149 InputsChangedDuringPreparation {
151 before: InputDigest,
152 after: InputDigest,
153 },
154 InputsChangedDuringAcquisition {
156 before: InputDigest,
157 after: InputDigest,
158 },
159 CargoBuildInputsChanged {
161 label: String,
163 before: InputDigest,
165 after: InputDigest,
167 },
168 CargoBuildInputRevalidation {
170 label: String,
172 source: WasmBuildError,
174 },
175 InvalidOutputs { outputs: Vec<(String, PathBuf)> },
177 UnknownOutput { name: String },
179 FailedTransactionCleanup {
181 transaction_error: Box<Self>,
182 path: PathBuf,
183 source: io::Error,
184 },
185}
186
187#[derive(Clone, Debug, Eq, PartialEq)]
188struct LabeledPath {
189 label: String,
190 path: PathBuf,
191}
192
193#[derive(Clone, Eq, PartialEq)]
194struct LabeledIdentity {
195 label: String,
196 value: Vec<u8>,
197}
198
199#[derive(Clone, Eq, PartialEq)]
200struct CargoBuildInputSet {
201 label: String,
202 build_spec: WasmBuildSpec,
203 resolved: ResolvedCargoBuildInputs,
204}
205
206#[derive(Clone, Debug, Eq, PartialEq)]
207struct OutputSpec {
208 name: String,
209 destination: PathBuf,
210 validation: ArtifactOutputValidation,
211}
212
213#[derive(Clone, Copy)]
214struct ResolvedKey {
215 key: InputDigest,
216 input_digest: InputDigest,
217}
218
219struct ArtifactInfo {
220 bytes: u64,
221 digest: InputDigest,
222}
223
224impl ArtifactCacheSpec {
225 #[must_use]
231 pub fn new(cache_root: &Path, namespace: &str, recipe_id: &str) -> Self {
232 Self {
233 cache_root: cache_root.to_owned(),
234 namespace: namespace.to_owned(),
235 recipe_id: recipe_id.to_owned(),
236 coordination_scope: namespace.to_owned(),
237 inputs: Vec::new(),
238 tools: Vec::new(),
239 arguments: Vec::new(),
240 environment: BTreeMap::new(),
241 identities: Vec::new(),
242 cargo_build_inputs: Vec::new(),
243 outputs: Vec::new(),
244 prune_policy: None,
245 prune_interval: None,
246 }
247 }
248
249 #[must_use]
251 pub fn with_coordination_scope(mut self, coordination_scope: &str) -> Self {
252 coordination_scope.clone_into(&mut self.coordination_scope);
253 self
254 }
255
256 #[must_use]
258 pub fn with_input(mut self, label: &str, path: &Path) -> Self {
259 self.inputs.push(LabeledPath {
260 label: label.to_owned(),
261 path: path.to_owned(),
262 });
263 self
264 }
265
266 #[must_use]
268 pub fn with_tool(mut self, label: &str, path: &Path) -> Self {
269 self.tools.push(LabeledPath {
270 label: label.to_owned(),
271 path: path.to_owned(),
272 });
273 self
274 }
275
276 #[must_use]
278 pub fn with_arguments<I, S>(mut self, arguments: I) -> Self
279 where
280 I: IntoIterator<Item = S>,
281 S: AsRef<OsStr>,
282 {
283 self.arguments = arguments
284 .into_iter()
285 .map(|argument| argument.as_ref().to_owned())
286 .collect();
287 self
288 }
289
290 #[must_use]
295 pub fn with_environment<I, K, V>(mut self, environment: I) -> Self
296 where
297 I: IntoIterator<Item = (K, V)>,
298 K: Into<OsString>,
299 V: Into<OsString>,
300 {
301 self.environment.extend(
302 environment
303 .into_iter()
304 .map(|(name, value)| (name.into(), Some(value.into()))),
305 );
306 self
307 }
308
309 #[must_use]
314 pub fn with_unset_environment<I, S>(mut self, names: I) -> Self
315 where
316 I: IntoIterator<Item = S>,
317 S: Into<OsString>,
318 {
319 self.environment
320 .extend(names.into_iter().map(|name| (name.into(), None)));
321 self
322 }
323
324 #[must_use]
326 pub fn with_identity_bytes(mut self, label: &str, value: &[u8]) -> Self {
327 self.identities.push(LabeledIdentity {
328 label: label.to_owned(),
329 value: value.to_vec(),
330 });
331 self
332 }
333
334 #[must_use]
341 pub fn with_cargo_build_inputs(
342 mut self,
343 label: &str,
344 build_spec: &WasmBuildSpec,
345 resolved: &ResolvedCargoBuildInputs,
346 ) -> Self {
347 self.cargo_build_inputs.push(CargoBuildInputSet {
348 label: label.to_owned(),
349 build_spec: build_spec.clone(),
350 resolved: resolved.clone(),
351 });
352 self
353 }
354
355 #[must_use]
357 pub fn with_output(self, name: &str, destination: &Path) -> Self {
358 self.with_output_validation(name, destination, ArtifactOutputValidation::NonEmptyFile)
359 }
360
361 #[must_use]
363 pub fn with_output_validation(
364 mut self,
365 name: &str,
366 destination: &Path,
367 validation: ArtifactOutputValidation,
368 ) -> Self {
369 self.outputs.push(OutputSpec {
370 name: name.to_owned(),
371 destination: destination.to_owned(),
372 validation,
373 });
374 self.outputs
375 .sort_by(|left, right| left.name.cmp(&right.name));
376 self
377 }
378
379 #[must_use]
381 pub const fn with_prune_policy(mut self, policy: ArtifactCachePrunePolicy) -> Self {
382 self.prune_policy = Some(policy);
383 self.prune_interval = None;
384 self
385 }
386
387 #[must_use]
395 pub const fn with_prune_policy_at_most_every(
396 mut self,
397 policy: ArtifactCachePrunePolicy,
398 minimum_interval: Duration,
399 ) -> Self {
400 self.prune_policy = Some(policy);
401 self.prune_interval = Some(minimum_interval);
402 self
403 }
404
405 #[must_use]
407 pub fn cache_root(&self) -> &Path {
408 &self.cache_root
409 }
410
411 #[must_use]
413 pub fn namespace(&self) -> &str {
414 &self.namespace
415 }
416
417 #[must_use]
419 pub fn recipe_id(&self) -> &str {
420 &self.recipe_id
421 }
422}
423
424impl std::fmt::Debug for ArtifactCacheSpec {
425 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
426 formatter
427 .debug_struct("ArtifactCacheSpec")
428 .field("cache_root", &self.cache_root)
429 .field("namespace", &self.namespace)
430 .field("recipe_id", &self.recipe_id)
431 .field("coordination_scope", &self.coordination_scope)
432 .field("inputs", &self.inputs)
433 .field("tools", &self.tools)
434 .field("argument_count", &self.arguments.len())
435 .field(
436 "environment_names",
437 &self.environment.keys().collect::<Vec<_>>(),
438 )
439 .field(
440 "identity_labels",
441 &self
442 .identities
443 .iter()
444 .map(|identity| identity.label.as_str())
445 .collect::<Vec<_>>(),
446 )
447 .field(
448 "cargo_build_input_labels",
449 &self
450 .cargo_build_inputs
451 .iter()
452 .map(|input| input.label.as_str())
453 .collect::<Vec<_>>(),
454 )
455 .field("outputs", &self.outputs)
456 .field("prune_policy", &self.prune_policy)
457 .field("prune_interval", &self.prune_interval)
458 .finish()
459 }
460}
461
462impl ArtifactCachePreparation {
463 #[must_use]
465 pub const fn reused_record(&self) -> Option<&ArtifactCacheRecord> {
466 match self {
467 Self::Reused(record) => Some(record),
468 Self::Build(_) => None,
469 }
470 }
471}
472
473impl ArtifactCacheOutcome {
474 #[must_use]
476 pub const fn record(&self) -> &ArtifactCacheRecord {
477 match self {
478 Self::Built(record) | Self::Reused(record) => record,
479 }
480 }
481
482 #[must_use]
484 pub const fn is_reused(&self) -> bool {
485 matches!(self, Self::Reused(_))
486 }
487}
488
489impl ArtifactCacheRecord {
490 #[must_use]
492 pub const fn key(&self) -> InputDigest {
493 self.key
494 }
495
496 #[must_use]
498 pub const fn input_digest(&self) -> InputDigest {
499 self.input_digest
500 }
501
502 #[must_use]
508 pub fn artifacts(&self) -> &[ArtifactCacheArtifact] {
509 &self.artifacts
510 }
511
512 #[must_use]
514 pub const fn timings(&self) -> ArtifactCacheTimings {
515 self.timings
516 }
517
518 #[must_use]
520 pub const fn maintenance(&self) -> Option<&ArtifactCacheMaintenance> {
521 self.maintenance.as_ref()
522 }
523}
524
525impl ArtifactCacheArtifact {
526 #[must_use]
528 pub fn name(&self) -> &str {
529 &self.name
530 }
531
532 #[must_use]
534 pub fn path(&self) -> &Path {
535 &self.path
536 }
537}
538
539impl ArtifactCacheTimings {
540 #[must_use]
542 pub const fn coordination_lock_wait(self) -> Duration {
543 self.coordination_lock_wait
544 }
545
546 #[must_use]
548 pub const fn content_lock_wait(self) -> Duration {
549 self.content_lock_wait
550 }
551
552 #[must_use]
554 pub const fn namespace_lock_wait(self) -> Duration {
555 self.namespace_lock_wait
556 }
557
558 #[must_use]
560 pub const fn input_capture(self) -> Duration {
561 self.input_capture
562 }
563
564 #[must_use]
566 pub const fn cache_lookup(self) -> Duration {
567 self.cache_lookup
568 }
569
570 #[must_use]
572 pub const fn caller_build(self) -> Option<Duration> {
573 self.caller_build
574 }
575
576 #[must_use]
578 pub const fn output_validation(self) -> Duration {
579 self.output_validation
580 }
581
582 #[must_use]
584 pub const fn publication(self) -> Duration {
585 self.publication
586 }
587
588 #[must_use]
590 pub const fn materialization(self) -> Duration {
591 self.materialization
592 }
593
594 #[must_use]
596 pub const fn maintenance(self) -> Option<Duration> {
597 self.maintenance
598 }
599
600 #[must_use]
602 pub const fn total(self) -> Duration {
603 self.total
604 }
605
606 pub(super) const fn saturating_add(self, other: Self) -> Self {
607 Self {
608 coordination_lock_wait: self
609 .coordination_lock_wait
610 .saturating_add(other.coordination_lock_wait),
611 content_lock_wait: self
612 .content_lock_wait
613 .saturating_add(other.content_lock_wait),
614 namespace_lock_wait: self
615 .namespace_lock_wait
616 .saturating_add(other.namespace_lock_wait),
617 input_capture: self.input_capture.saturating_add(other.input_capture),
618 cache_lookup: self.cache_lookup.saturating_add(other.cache_lookup),
619 caller_build: saturating_add_optional_duration(self.caller_build, other.caller_build),
620 output_validation: self
621 .output_validation
622 .saturating_add(other.output_validation),
623 publication: self.publication.saturating_add(other.publication),
624 materialization: self.materialization.saturating_add(other.materialization),
625 maintenance: saturating_add_optional_duration(self.maintenance, other.maintenance),
626 total: self.total.saturating_add(other.total),
627 }
628 }
629}
630
631impl std::fmt::Display for ArtifactCacheTimings {
632 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
633 write!(
634 formatter,
635 "total={:?} coordination_lock={:?} content_lock={:?} namespace_lock={:?} inputs={:?} lookup={:?} build={:?} validation={:?} publication={:?} materialization={:?} maintenance={:?}",
636 self.total,
637 self.coordination_lock_wait,
638 self.content_lock_wait,
639 self.namespace_lock_wait,
640 self.input_capture,
641 self.cache_lookup,
642 self.caller_build,
643 self.output_validation,
644 self.publication,
645 self.materialization,
646 self.maintenance,
647 )
648 }
649}
650
651impl std::fmt::Display for ArtifactCacheOutcome {
652 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
653 let state = if self.is_reused() { "reused" } else { "built" };
654 write!(
655 formatter,
656 "{state} key={} artifacts={} {}",
657 self.record().key,
658 self.record().artifacts.len(),
659 self.record().timings,
660 )
661 }
662}
663
664impl ArtifactBuildTransaction {
665 #[must_use]
670 pub fn staging_directory(&self) -> &Path {
671 &self.staging_directory
672 }
673
674 pub fn output_path(&self, name: &str) -> Result<PathBuf, ArtifactCacheError> {
676 self.output_index(name)
677 .map(|index| staged_output_path(&self.staging_directory, index))
678 .ok_or_else(|| ArtifactCacheError::UnknownOutput {
679 name: name.to_owned(),
680 })
681 }
682
683 pub fn import_output(&self, name: &str, source: &Path) -> Result<(), ArtifactCacheError> {
685 let destination = self.output_path(name)?;
686 copy_file_atomic(source, &destination).map_err(|source_error| ArtifactCacheError::Io {
687 operation: "import artifact output into staging",
688 path: destination,
689 source: source_error,
690 })?;
691 Ok(())
692 }
693
694 pub fn commit(mut self) -> Result<ArtifactCacheOutcome, ArtifactCacheError> {
696 let result = self.commit_inner();
697 match result {
698 Ok(outcome) => Ok(outcome),
699 Err(transaction_error) if self.staging_armed => {
700 let path = self.staging_directory.clone();
701 match remove_path_if_present(&path) {
702 Ok(()) => {
703 self.staging_armed = false;
704 Err(transaction_error)
705 }
706 Err(source) => Err(ArtifactCacheError::FailedTransactionCleanup {
707 transaction_error: Box::new(transaction_error),
708 path,
709 source,
710 }),
711 }
712 }
713 Err(transaction_error) => Err(transaction_error),
714 }
715 }
716
717 pub fn abort(mut self) -> Result<(), ArtifactCacheError> {
719 remove_path_if_present(&self.staging_directory).map_err(|source| {
720 ArtifactCacheError::Io {
721 operation: "abort artifact cache transaction",
722 path: self.staging_directory.clone(),
723 source,
724 }
725 })?;
726 self.staging_armed = false;
727 Ok(())
728 }
729
730 fn output_index(&self, name: &str) -> Option<usize> {
731 self.spec
732 .outputs
733 .iter()
734 .position(|output| output.name == name)
735 }
736
737 fn commit_inner(&mut self) -> Result<ArtifactCacheOutcome, ArtifactCacheError> {
738 self.timings.caller_build = Some(self.caller_build_started.elapsed());
739
740 let validation_started = Instant::now();
741 let output_info = inspect_complete_output_set(&self.spec, &self.staging_directory)?;
742 self.timings.output_validation = validation_started.elapsed();
743
744 let capture_started = Instant::now();
745 revalidate_cargo_build_input_fingerprints(&self.spec)?;
746 let verified = resolve_key(&self.spec)?;
747 self.timings.input_capture = self
748 .timings
749 .input_capture
750 .saturating_add(capture_started.elapsed());
751 if verified.input_digest != self.resolved.input_digest {
752 return Err(ArtifactCacheError::InputsChangedDuringAcquisition {
753 before: self.resolved.input_digest,
754 after: verified.input_digest,
755 });
756 }
757
758 let publication_started = Instant::now();
759 let manifest = manifest_contents(self.resolved.key, &self.spec, &output_info);
760 write_atomic(
761 &self.staging_directory.join(MANIFEST_FILE),
762 manifest.as_bytes(),
763 )
764 .map_err(|source| ArtifactCacheError::Io {
765 operation: "write artifact cache manifest",
766 path: self.staging_directory.join(MANIFEST_FILE),
767 source,
768 })?;
769 let namespace_lock_path = namespace_lock_path(&self.spec);
770 let (_namespace_lock, namespace_wait) =
771 lock_cache_file(&namespace_lock_path).map_err(artifact_cache_fs_error)?;
772 self.timings.namespace_lock_wait = self
773 .timings
774 .namespace_lock_wait
775 .saturating_add(namespace_wait);
776 remove_unretained_entry(&self.entry_directory).map_err(artifact_cache_fs_error)?;
777 fs::rename(&self.staging_directory, &self.entry_directory).map_err(|source| {
778 ArtifactCacheError::Io {
779 operation: "publish artifact cache entry",
780 path: self.entry_directory.clone(),
781 source,
782 }
783 })?;
784 self.staging_armed = false;
785 self.timings.publication = publication_started.elapsed();
786
787 let materialization_started = Instant::now();
788 materialize_outputs(&self.spec, &self.entry_directory)?;
789 self.timings.materialization = materialization_started.elapsed();
790 record_cache_entry_use(&self.entry_directory).map_err(artifact_cache_fs_error)?;
791 let (maintenance, maintenance_timing) = perform_maintenance_locked(
792 &self.spec,
793 &self.namespace_directory,
794 &self.entry_directory,
795 );
796 self.timings.maintenance = maintenance_timing;
797 self.timings.total = self.total_started.elapsed();
798
799 Ok(ArtifactCacheOutcome::Built(cache_record(
800 &self.spec,
801 self.resolved,
802 self.timings,
803 maintenance,
804 )?))
805 }
806}
807
808impl Drop for ArtifactBuildTransaction {
809 fn drop(&mut self) {
810 if self.staging_armed {
811 let _ = remove_path_if_present(&self.staging_directory);
812 }
813 }
814}
815
816pub fn prepare_artifact_cache(
818 spec: &ArtifactCacheSpec,
819) -> Result<ArtifactCachePreparation, ArtifactCacheError> {
820 let total_started = Instant::now();
821 validate_spec(spec)?;
822 validate_filesystem_boundaries(spec)?;
823 let namespace_directory = initialize_cache(spec)?;
824
825 let coordination_lock_path = coordination_lock_path(spec);
826 let (coordination_lock, coordination_wait) =
827 lock_cache_file(&coordination_lock_path).map_err(artifact_cache_fs_error)?;
828 let mut timings = ArtifactCacheTimings {
829 coordination_lock_wait: coordination_wait,
830 ..ArtifactCacheTimings::default()
831 };
832 let initial_cargo_started = Instant::now();
833 revalidate_cargo_build_input_fingerprints(spec)?;
834 timings.input_capture = initial_cargo_started.elapsed();
835 let mut last_change = None;
836
837 for _ in 0..MAX_PREPARATION_RETRIES {
838 let capture_started = Instant::now();
839 let resolved = resolve_key(spec)?;
840 timings.input_capture = timings
841 .input_capture
842 .saturating_add(capture_started.elapsed());
843
844 let content_lock_path = content_lock_path(spec, resolved.key);
845 let (content_lock, content_wait) =
846 lock_cache_file(&content_lock_path).map_err(artifact_cache_fs_error)?;
847 timings.content_lock_wait = timings.content_lock_wait.saturating_add(content_wait);
848
849 let verification_started = Instant::now();
850 let verified = resolve_key(spec)?;
851 timings.input_capture = timings
852 .input_capture
853 .saturating_add(verification_started.elapsed());
854 if resolved.input_digest != verified.input_digest {
855 last_change = Some((resolved.input_digest, verified.input_digest));
856 drop(content_lock);
857 continue;
858 }
859
860 let entry_directory = entry_directory(&namespace_directory, resolved.key);
861 let namespace_lock_path = namespace_lock_path(spec);
862 let (namespace_lock, namespace_wait) =
863 lock_cache_file(&namespace_lock_path).map_err(artifact_cache_fs_error)?;
864 timings.namespace_lock_wait = timings.namespace_lock_wait.saturating_add(namespace_wait);
865 let lookup_started = Instant::now();
866 let reusable = cache_entry_is_valid(spec, resolved.key, &entry_directory)?;
867 timings.cache_lookup = timings
868 .cache_lookup
869 .saturating_add(lookup_started.elapsed());
870
871 if reusable {
872 let materialization_started = Instant::now();
873 materialize_outputs(spec, &entry_directory)?;
874 timings.materialization = timings
875 .materialization
876 .saturating_add(materialization_started.elapsed());
877 let after_started = Instant::now();
878 revalidate_cargo_build_input_fingerprints(spec)?;
879 let after = resolve_key(spec)?;
880 timings.input_capture = timings
881 .input_capture
882 .saturating_add(after_started.elapsed());
883 if after.input_digest != resolved.input_digest {
884 last_change = Some((resolved.input_digest, after.input_digest));
885 drop(namespace_lock);
886 drop(content_lock);
887 continue;
888 }
889 record_cache_entry_use(&entry_directory).map_err(artifact_cache_fs_error)?;
890 let (maintenance, maintenance_timing) =
891 perform_maintenance_locked(spec, &namespace_directory, &entry_directory);
892 timings.maintenance = maintenance_timing;
893 timings.total = total_started.elapsed();
894 return Ok(ArtifactCachePreparation::Reused(cache_record(
895 spec,
896 resolved,
897 timings,
898 maintenance,
899 )?));
900 }
901
902 remove_unretained_entry(&entry_directory).map_err(artifact_cache_fs_error)?;
903 let staging_directory = create_staging_directory(&namespace_directory, resolved.key)?;
904 drop(namespace_lock);
905 return Ok(ArtifactCachePreparation::Build(ArtifactBuildTransaction {
906 spec: Box::new(spec.clone()),
907 resolved,
908 staging_directory,
909 entry_directory,
910 namespace_directory,
911 _coordination_lock: coordination_lock,
912 _content_lock: content_lock,
913 timings,
914 total_started,
915 caller_build_started: Instant::now(),
916 staging_armed: true,
917 }));
918 }
919
920 let (before, after) = last_change.expect("preparation retries require a recorded input change");
921 Err(ArtifactCacheError::InputsChangedDuringPreparation { before, after })
922}
923
924fn revalidate_cargo_build_input_fingerprints(
925 spec: &ArtifactCacheSpec,
926) -> Result<(), ArtifactCacheError> {
927 for cargo_input in &spec.cargo_build_inputs {
928 let current = resolve_cargo_build_inputs(&cargo_input.build_spec).map_err(|source| {
929 ArtifactCacheError::CargoBuildInputRevalidation {
930 label: cargo_input.label.clone(),
931 source,
932 }
933 })?;
934 let before_validation = cargo_input.resolved.validation_digest();
935 let after_validation = current.validation_digest();
936 if after_validation != before_validation {
937 return Err(ArtifactCacheError::CargoBuildInputsChanged {
938 label: cargo_input.label.clone(),
939 before: before_validation,
940 after: after_validation,
941 });
942 }
943 let before = cargo_input.resolved.fingerprint();
944 let after = current.fingerprint();
945 if after != before {
946 return Err(ArtifactCacheError::CargoBuildInputsChanged {
947 label: cargo_input.label.clone(),
948 before,
949 after,
950 });
951 }
952 }
953 Ok(())
954}
955
956pub fn prune_artifact_cache(
965 cache_root: &Path,
966 namespace: &str,
967 policy: ArtifactCachePrunePolicy,
968) -> Result<ArtifactCachePruneReport, ArtifactCacheError> {
969 validate_identifier("namespace", namespace)?;
970 if cache_root.as_os_str().is_empty() {
971 return invalid_spec("cache root must not be empty");
972 }
973 ensure_cache_directory_tag(cache_root).map_err(artifact_cache_fs_error)?;
974 let namespace_directory = namespace_directory_for(cache_root, namespace);
975 let lock_path = namespace_lock_path_for(cache_root, namespace);
976 let (_lock, _) = lock_cache_file(&lock_path).map_err(artifact_cache_fs_error)?;
977 prune_artifact_namespace_locked(cache_root, &namespace_directory, policy, None)
978}
979
980fn initialize_cache(spec: &ArtifactCacheSpec) -> Result<PathBuf, ArtifactCacheError> {
981 ensure_cache_directory_tag(&spec.cache_root).map_err(artifact_cache_fs_error)?;
982 let namespace = namespace_directory(spec);
983 fs::create_dir_all(entries_directory(&namespace)).map_err(|source| ArtifactCacheError::Io {
984 operation: "create artifact cache namespace",
985 path: namespace.clone(),
986 source,
987 })?;
988 Ok(namespace)
989}
990
991fn validate_spec(spec: &ArtifactCacheSpec) -> Result<(), ArtifactCacheError> {
992 validate_identifier("namespace", &spec.namespace)?;
993 validate_identifier("recipe identity", &spec.recipe_id)?;
994 validate_identifier("coordination scope", &spec.coordination_scope)?;
995 if spec.cache_root.as_os_str().is_empty() {
996 return invalid_spec("cache root must not be empty");
997 }
998 if spec.outputs.is_empty() {
999 return invalid_spec("at least one artifact output is required");
1000 }
1001
1002 let mut labels = BTreeSet::new();
1003 for (kind, paths) in [("input", &spec.inputs), ("tool", &spec.tools)] {
1004 for path in paths {
1005 validate_path_label(kind, &path.label)?;
1006 if !labels.insert((kind, path.label.as_str())) {
1007 return invalid_spec(&format!("duplicate {kind} label `{}`", path.label));
1008 }
1009 }
1010 }
1011 let mut identity_labels = BTreeSet::new();
1012 for identity in &spec.identities {
1013 validate_label("identity", &identity.label)?;
1014 if !identity_labels.insert(&identity.label) {
1015 return invalid_spec(&format!("duplicate identity label `{}`", identity.label));
1016 }
1017 }
1018 let mut cargo_input_labels = BTreeSet::new();
1019 for cargo_inputs in &spec.cargo_build_inputs {
1020 validate_label("Cargo build input", &cargo_inputs.label)?;
1021 if !cargo_input_labels.insert(&cargo_inputs.label) {
1022 return invalid_spec(&format!(
1023 "duplicate Cargo build input label `{}`",
1024 cargo_inputs.label
1025 ));
1026 }
1027 }
1028 if spec
1029 .environment
1030 .keys()
1031 .any(|name| name.as_os_str().is_empty())
1032 {
1033 return invalid_spec("environment names must not be empty");
1034 }
1035 let mut output_names = BTreeSet::new();
1036 let mut destinations = BTreeSet::new();
1037 for output in &spec.outputs {
1038 validate_output_name(&output.name)?;
1039 if !output_names.insert(&output.name) {
1040 return invalid_spec(&format!("duplicate output name `{}`", output.name));
1041 }
1042 if output.destination.as_os_str().is_empty() {
1043 return invalid_spec(&format!(
1044 "output `{}` destination must not be empty",
1045 output.name
1046 ));
1047 }
1048 if !destinations.insert(&output.destination) {
1049 return invalid_spec(&format!(
1050 "output `{}` shares a destination with another output",
1051 output.name
1052 ));
1053 }
1054 }
1055 Ok(())
1056}
1057
1058fn validate_filesystem_boundaries(spec: &ArtifactCacheSpec) -> Result<(), ArtifactCacheError> {
1059 let cache_root =
1060 canonicalize_allow_missing(&spec.cache_root).map_err(|source| ArtifactCacheError::Io {
1061 operation: "resolve artifact cache root",
1062 path: spec.cache_root.clone(),
1063 source,
1064 })?;
1065 let mut declared_paths = Vec::with_capacity(spec.inputs.len() + spec.tools.len());
1066 for (kind, labeled_paths) in [("input", &spec.inputs), ("tool", &spec.tools)] {
1067 for labeled in labeled_paths {
1068 let canonical =
1069 canonicalize_path(&labeled.path, "canonicalize declared artifact cache path")?;
1070 if canonical.starts_with(&cache_root) {
1071 return invalid_spec(&format!(
1072 "{kind} `{}` must not be located inside the artifact cache root",
1073 labeled.label
1074 ));
1075 }
1076 let is_directory = fs::metadata(&canonical)
1077 .map_err(|source| ArtifactCacheError::Io {
1078 operation: "inspect declared artifact cache path",
1079 path: canonical.clone(),
1080 source,
1081 })?
1082 .is_dir();
1083 let entry_path =
1084 resolve_file_entry(&labeled.path, "resolve declared artifact cache path entry")?;
1085 if entry_path != canonical {
1086 declared_paths.push((kind, labeled.label.as_str(), entry_path, false));
1087 }
1088 declared_paths.push((kind, labeled.label.as_str(), canonical, is_directory));
1089 }
1090 }
1091 for cargo_inputs in &spec.cargo_build_inputs {
1092 if resolved_cargo_inputs_watch_path(&cargo_inputs.resolved, &cache_root)? {
1093 return invalid_spec(&format!(
1094 "artifact cache root must be outside resolved Cargo build inputs `{}` or inside one of their generated-state exclusions",
1095 cargo_inputs.label
1096 ));
1097 }
1098 }
1099
1100 let mut destinations = BTreeSet::new();
1101 for output in &spec.outputs {
1102 let destination =
1103 resolve_file_entry(&output.destination, "resolve artifact output destination")?;
1104 if destination.starts_with(&cache_root) {
1105 return invalid_spec(&format!(
1106 "output `{}` destination must be outside the artifact cache root",
1107 output.name
1108 ));
1109 }
1110 if destinations.iter().any(|other: &PathBuf| {
1111 destination.starts_with(other) || other.starts_with(&destination)
1112 }) {
1113 return invalid_spec(&format!(
1114 "output `{}` destination overlaps another output",
1115 output.name
1116 ));
1117 }
1118 destinations.insert(destination.clone());
1119 for (kind, label, declared, is_directory) in &declared_paths {
1120 if destination == *declared || (*is_directory && destination.starts_with(declared)) {
1121 return invalid_spec(&format!(
1122 "output `{}` destination overlaps declared {kind} `{label}`",
1123 output.name
1124 ));
1125 }
1126 }
1127 for cargo_inputs in &spec.cargo_build_inputs {
1128 if resolved_cargo_inputs_watch_path(&cargo_inputs.resolved, &destination)? {
1129 return invalid_spec(&format!(
1130 "output `{}` destination must be outside resolved Cargo build inputs `{}` or inside one of their generated-state exclusions",
1131 output.name, cargo_inputs.label
1132 ));
1133 }
1134 }
1135 match fs::metadata(&destination) {
1136 Ok(metadata) if metadata.is_dir() => {
1137 return invalid_spec(&format!(
1138 "output `{}` destination must not be an existing directory",
1139 output.name
1140 ));
1141 }
1142 Ok(_) => {}
1143 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
1144 Err(source) => {
1145 return Err(ArtifactCacheError::Io {
1146 operation: "inspect artifact output destination",
1147 path: destination,
1148 source,
1149 });
1150 }
1151 }
1152 }
1153 Ok(())
1154}
1155
1156fn resolve_file_entry(path: &Path, operation: &'static str) -> Result<PathBuf, ArtifactCacheError> {
1159 let resolved = match (path.parent(), path.file_name()) {
1160 (Some(parent), Some(name)) => {
1161 canonicalize_allow_missing(parent).map(|parent| parent.join(name))
1162 }
1163 _ => canonicalize_allow_missing(path),
1164 };
1165 resolved.map_err(|source| ArtifactCacheError::Io {
1166 operation,
1167 path: path.to_owned(),
1168 source,
1169 })
1170}
1171
1172fn resolved_cargo_inputs_watch_path(
1173 resolved: &ResolvedCargoBuildInputs,
1174 candidate: &Path,
1175) -> Result<bool, ArtifactCacheError> {
1176 for exclusion in resolved.exclusions() {
1177 let exclusion =
1178 canonicalize_allow_missing(exclusion).map_err(|source| ArtifactCacheError::Io {
1179 operation: "resolve Cargo build input exclusion",
1180 path: exclusion.clone(),
1181 source,
1182 })?;
1183 if candidate.starts_with(exclusion) {
1184 return Ok(false);
1185 }
1186 }
1187
1188 for input in resolved.inputs() {
1189 let path = canonicalize_path(input.path(), "canonicalize resolved Cargo build input")?;
1190 let is_directory = fs::metadata(&path)
1191 .map_err(|source| ArtifactCacheError::Io {
1192 operation: "inspect resolved Cargo build input",
1193 path: path.clone(),
1194 source,
1195 })?
1196 .is_dir();
1197 if candidate == path || (is_directory && candidate.starts_with(path)) {
1198 return Ok(true);
1199 }
1200 }
1201 Ok(false)
1202}
1203
1204fn canonicalize_path(path: &Path, operation: &'static str) -> Result<PathBuf, ArtifactCacheError> {
1205 path.canonicalize()
1206 .map_err(|source| ArtifactCacheError::Io {
1207 operation,
1208 path: path.to_owned(),
1209 source,
1210 })
1211}
1212
1213fn validate_identifier(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1214 if value.is_empty() {
1215 return invalid_spec(&format!("{kind} must not be empty"));
1216 }
1217 if value.len() > 256 {
1218 return invalid_spec(&format!("{kind} must not exceed 256 bytes"));
1219 }
1220 Ok(())
1221}
1222
1223fn validate_label(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1224 if value.is_empty() {
1225 return invalid_spec(&format!("{kind} label must not be empty"));
1226 }
1227 if value.len() > 256 {
1228 return invalid_spec(&format!("{kind} label must not exceed 256 bytes"));
1229 }
1230 Ok(())
1231}
1232
1233fn validate_path_label(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1234 validate_label(kind, value)?;
1235 if value
1236 .bytes()
1237 .any(|byte| !(byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'/')))
1238 || value
1239 .split('/')
1240 .any(|component| component.is_empty() || matches!(component, "." | ".."))
1241 {
1242 return invalid_spec(&format!(
1243 "{kind} label `{value}` must be a portable relative logical path"
1244 ));
1245 }
1246 Ok(())
1247}
1248
1249fn validate_output_name(name: &str) -> Result<(), ArtifactCacheError> {
1250 if name.is_empty() || name.len() > 128 {
1251 return invalid_spec("output names must contain 1 to 128 bytes");
1252 }
1253 if !name
1254 .bytes()
1255 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
1256 {
1257 return invalid_spec(&format!(
1258 "output name `{name}` must use only ASCII letters, digits, dot, dash, or underscore"
1259 ));
1260 }
1261 if matches!(name, "." | "..") {
1262 return invalid_spec("output names must not be dot path components");
1263 }
1264 Ok(())
1265}
1266
1267fn invalid_spec<T>(message: &str) -> Result<T, ArtifactCacheError> {
1268 Err(ArtifactCacheError::InvalidSpec {
1269 message: message.to_owned(),
1270 })
1271}
1272
1273fn resolve_key(spec: &ArtifactCacheSpec) -> Result<ResolvedKey, ArtifactCacheError> {
1274 let paths = spec
1275 .inputs
1276 .iter()
1277 .map(|input| (PathBuf::from("input").join(&input.label), &input.path))
1278 .chain(
1279 spec.tools
1280 .iter()
1281 .map(|tool| (PathBuf::from("tool").join(&tool.label), &tool.path)),
1282 );
1283 let declared_input_digest = digest_labeled_paths(
1284 "artifact-set-inputs-v1",
1285 paths,
1286 std::slice::from_ref(&spec.cache_root),
1287 )
1288 .map_err(|source| ArtifactCacheError::Io {
1289 operation: "hash artifact cache inputs",
1290 path: spec.cache_root.clone(),
1291 source,
1292 })?;
1293 let input_digest = if spec.cargo_build_inputs.is_empty() {
1294 declared_input_digest
1295 } else {
1296 let mut cargo_inputs = spec.cargo_build_inputs.iter().collect::<Vec<_>>();
1297 cargo_inputs.sort_by(|left, right| left.label.cmp(&right.label));
1298 let mut inputs = InputHasher::new("artifact-set-inputs-with-cargo-v1");
1299 inputs.field("declared-input-digest", declared_input_digest.as_bytes());
1300 for cargo_input in cargo_inputs {
1301 let current = cargo_input
1302 .resolved
1303 .current_validation_digest()
1304 .map_err(|source| ArtifactCacheError::CargoBuildInputRevalidation {
1305 label: cargo_input.label.clone(),
1306 source,
1307 })?;
1308 let before = cargo_input.resolved.validation_digest();
1309 if current != before {
1310 return Err(ArtifactCacheError::CargoBuildInputsChanged {
1311 label: cargo_input.label.clone(),
1312 before,
1313 after: current,
1314 });
1315 }
1316 inputs.field("cargo-input-label", cargo_input.label.as_bytes());
1317 inputs.field(
1318 "cargo-build-fingerprint",
1319 cargo_input.resolved.fingerprint().as_bytes(),
1320 );
1321 inputs.field(
1322 "cargo-input-digest",
1323 cargo_input.resolved.input_digest().as_bytes(),
1324 );
1325 }
1326 inputs.finish()
1327 };
1328
1329 let mut hasher = InputHasher::new(ARTIFACT_CACHE_FORMAT);
1330 hasher.field("namespace", spec.namespace.as_bytes());
1331 hasher.field("recipe-id", spec.recipe_id.as_bytes());
1332 hasher.field("input-digest", input_digest.as_bytes());
1333 for argument in &spec.arguments {
1334 hasher.field("argument", &os_bytes(argument));
1335 }
1336 for (name, value) in &spec.environment {
1337 hasher.field("environment-name", &os_bytes(name));
1338 match value {
1339 Some(value) => hasher.field("environment-value", &os_bytes(value)),
1340 None => hasher.field("environment-unset", b""),
1341 }
1342 }
1343 let mut identities = spec.identities.iter().collect::<Vec<_>>();
1344 identities.sort_by(|left, right| left.label.cmp(&right.label));
1345 for identity in identities {
1346 hasher.field("identity-label", identity.label.as_bytes());
1347 hasher.field("identity-value", &identity.value);
1348 }
1349 for output in &spec.outputs {
1350 hasher.field("output-name", output.name.as_bytes());
1351 hasher.field(
1352 "output-validation",
1353 output.validation.cache_token().as_bytes(),
1354 );
1355 }
1356 Ok(ResolvedKey {
1357 key: hasher.finish(),
1358 input_digest,
1359 })
1360}
1361
1362fn cache_entry_is_valid(
1363 spec: &ArtifactCacheSpec,
1364 key: InputDigest,
1365 entry: &Path,
1366) -> Result<bool, ArtifactCacheError> {
1367 let entry_metadata = match fs::symlink_metadata(entry) {
1368 Ok(metadata) => metadata,
1369 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(false),
1370 Err(source) => {
1371 return Err(ArtifactCacheError::Io {
1372 operation: "inspect artifact cache entry",
1373 path: entry.to_owned(),
1374 source,
1375 });
1376 }
1377 };
1378 if !entry_metadata.file_type().is_dir() || !cache_entry_root_is_valid(entry)? {
1379 return Ok(false);
1380 }
1381 let manifest_path = entry.join(MANIFEST_FILE);
1382 let manifest_metadata = match fs::symlink_metadata(&manifest_path) {
1383 Ok(metadata) => metadata,
1384 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(false),
1385 Err(source) => {
1386 return Err(ArtifactCacheError::Io {
1387 operation: "inspect artifact cache manifest",
1388 path: manifest_path,
1389 source,
1390 });
1391 }
1392 };
1393 if !manifest_metadata.file_type().is_file() {
1394 return Ok(false);
1395 }
1396 let manifest = fs::read(&manifest_path).map_err(|source| ArtifactCacheError::Io {
1397 operation: "read artifact cache manifest",
1398 path: manifest_path,
1399 source,
1400 })?;
1401 if !manifest.starts_with(manifest_header(key).as_bytes()) {
1402 return Ok(false);
1403 }
1404 let Some(output_info) = inspect_cached_output_set(spec, entry)? else {
1405 return Ok(false);
1406 };
1407 Ok(manifest == manifest_contents(key, spec, &output_info).as_bytes())
1408}
1409
1410fn inspect_complete_output_set(
1411 spec: &ArtifactCacheSpec,
1412 root: &Path,
1413) -> Result<Vec<ArtifactInfo>, ArtifactCacheError> {
1414 let mut info = Vec::new();
1415 let mut invalid = Vec::new();
1416 let output_directory = root.join("outputs");
1417 if !is_plain_directory(
1418 &output_directory,
1419 "inspect artifact staging output directory",
1420 )? {
1421 invalid.push(("<outputs>".to_owned(), output_directory));
1422 return Err(ArtifactCacheError::InvalidOutputs { outputs: invalid });
1423 }
1424 let outputs = &spec.outputs;
1425 for (index, output) in outputs.iter().enumerate() {
1426 let path = staged_output_path(root, index);
1427 match inspect_artifact(&path, output.validation) {
1428 Ok(Some(artifact)) => info.push(artifact),
1429 Ok(None) => invalid.push((output.name.clone(), path)),
1430 Err(source) => {
1431 return Err(ArtifactCacheError::Io {
1432 operation: "inspect staged artifact output",
1433 path,
1434 source,
1435 });
1436 }
1437 }
1438 }
1439 invalid.extend(
1440 undeclared_output_paths(root, outputs.len())?
1441 .into_iter()
1442 .map(|path| ("<undeclared>".to_owned(), path)),
1443 );
1444 invalid.extend(
1445 undeclared_child_paths(
1446 root,
1447 |name| name == "outputs",
1448 "read artifact staging directory",
1449 )?
1450 .into_iter()
1451 .map(|path| ("<undeclared>".to_owned(), path)),
1452 );
1453 if invalid.is_empty() {
1454 Ok(info)
1455 } else {
1456 Err(ArtifactCacheError::InvalidOutputs { outputs: invalid })
1457 }
1458}
1459
1460fn inspect_cached_output_set(
1461 spec: &ArtifactCacheSpec,
1462 root: &Path,
1463) -> Result<Option<Vec<ArtifactInfo>>, ArtifactCacheError> {
1464 let mut info = Vec::new();
1465 let outputs = &spec.outputs;
1466 for (index, output) in outputs.iter().enumerate() {
1467 let path = staged_output_path(root, index);
1468 match inspect_artifact(&path, output.validation) {
1469 Ok(Some(artifact)) => info.push(artifact),
1470 Ok(None) => return Ok(None),
1471 Err(source) => {
1472 return Err(ArtifactCacheError::Io {
1473 operation: "inspect cached artifact output",
1474 path,
1475 source,
1476 });
1477 }
1478 }
1479 }
1480 if !undeclared_output_paths(root, outputs.len())?.is_empty() {
1481 return Ok(None);
1482 }
1483 Ok(Some(info))
1484}
1485
1486fn undeclared_output_paths(
1487 root: &Path,
1488 output_count: usize,
1489) -> Result<Vec<PathBuf>, ArtifactCacheError> {
1490 let output_directory = root.join("outputs");
1491 undeclared_child_paths(
1492 &output_directory,
1493 |name| {
1494 let Some(name) = name.to_str() else {
1495 return false;
1496 };
1497 name.strip_suffix(OUTPUT_FILE_SUFFIX)
1498 .and_then(|index| index.parse::<usize>().ok())
1499 .is_some_and(|index| index < output_count && name == format_output_index(index))
1500 },
1501 "read artifact output directory",
1502 )
1503}
1504
1505fn undeclared_child_paths(
1506 directory: &Path,
1507 is_expected: impl Fn(&OsStr) -> bool,
1508 operation: &'static str,
1509) -> Result<Vec<PathBuf>, ArtifactCacheError> {
1510 let entries = fs::read_dir(directory).map_err(|source| ArtifactCacheError::Io {
1511 operation,
1512 path: directory.to_owned(),
1513 source,
1514 })?;
1515 let mut undeclared = Vec::new();
1516 for entry in entries {
1517 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1518 operation,
1519 path: directory.to_owned(),
1520 source,
1521 })?;
1522 if !is_expected(&entry.file_name()) {
1523 undeclared.push(entry.path());
1524 }
1525 }
1526 Ok(undeclared)
1527}
1528
1529fn cache_entry_root_is_valid(root: &Path) -> Result<bool, ArtifactCacheError> {
1530 let expected = [
1531 "outputs",
1532 MANIFEST_FILE,
1533 LAST_USED_FILE,
1534 super::cache_fs::RETENTION_LOCK_FILE,
1535 ];
1536 if !undeclared_child_paths(
1537 root,
1538 |name| expected.iter().any(|expected| name == *expected),
1539 "read artifact cache entry",
1540 )?
1541 .is_empty()
1542 {
1543 return Ok(false);
1544 }
1545 if !is_plain_directory(
1546 &root.join("outputs"),
1547 "inspect artifact cache output directory",
1548 )? {
1549 return Ok(false);
1550 }
1551 let last_used = root.join(LAST_USED_FILE);
1552 match fs::symlink_metadata(&last_used) {
1553 Ok(metadata) => Ok(metadata.file_type().is_file()),
1554 Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(true),
1555 Err(source) => Err(ArtifactCacheError::Io {
1556 operation: "inspect artifact cache use marker",
1557 path: last_used,
1558 source,
1559 }),
1560 }
1561}
1562
1563fn is_plain_directory(path: &Path, operation: &'static str) -> Result<bool, ArtifactCacheError> {
1564 match fs::symlink_metadata(path) {
1565 Ok(metadata) => Ok(metadata.file_type().is_dir()),
1566 Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(false),
1567 Err(source) => Err(ArtifactCacheError::Io {
1568 operation,
1569 path: path.to_owned(),
1570 source,
1571 }),
1572 }
1573}
1574
1575fn inspect_artifact(
1576 path: &Path,
1577 validation: ArtifactOutputValidation,
1578) -> io::Result<Option<ArtifactInfo>> {
1579 let metadata = match fs::symlink_metadata(path) {
1580 Ok(metadata) => metadata,
1581 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1582 Err(error) => return Err(error),
1583 };
1584 if !metadata.file_type().is_file() {
1585 return Ok(None);
1586 }
1587 if validation == ArtifactOutputValidation::NonEmptyFile && metadata.len() == 0 {
1588 return Ok(None);
1589 }
1590 let (bytes, digest) = digest_file("artifact-set-output-v1", path)?;
1591 Ok(Some(ArtifactInfo { bytes, digest }))
1592}
1593
1594fn manifest_header(key: InputDigest) -> String {
1595 format!("{ARTIFACT_CACHE_FORMAT}\nkey:{key}\n")
1596}
1597
1598fn manifest_contents(
1599 key: InputDigest,
1600 spec: &ArtifactCacheSpec,
1601 output_info: &[ArtifactInfo],
1602) -> String {
1603 let mut manifest = manifest_header(key);
1604 for ((index, output), info) in spec.outputs.iter().enumerate().zip(output_info) {
1605 use std::fmt::Write as _;
1606 writeln!(
1607 manifest,
1608 "output:{index}:{}:{}:{}:{}",
1609 output.name,
1610 output.validation.cache_token(),
1611 info.bytes,
1612 info.digest,
1613 )
1614 .expect("writing an artifact manifest to a String cannot fail");
1615 }
1616 manifest
1617}
1618
1619fn materialize_outputs(spec: &ArtifactCacheSpec, entry: &Path) -> Result<(), ArtifactCacheError> {
1620 for (index, output) in spec.outputs.iter().enumerate() {
1621 let cached = staged_output_path(entry, index);
1622 copy_file_atomic(&cached, &output.destination).map_err(|source| {
1623 ArtifactCacheError::Io {
1624 operation: "materialize artifact output",
1625 path: output.destination.clone(),
1626 source,
1627 }
1628 })?;
1629 }
1630 Ok(())
1631}
1632
1633fn perform_maintenance_locked(
1634 spec: &ArtifactCacheSpec,
1635 namespace: &Path,
1636 protected_entry: &Path,
1637) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
1638 spec.prune_policy.map_or((None, None), |policy| {
1639 let identity = policy.maintenance_identity();
1640 perform_scheduled_cache_maintenance(namespace, spec.prune_interval, &identity, || {
1641 prune_artifact_namespace_locked(
1642 &spec.cache_root,
1643 namespace,
1644 policy,
1645 Some(protected_entry),
1646 )
1647 .map_err(|error| error.to_string())
1648 })
1649 })
1650}
1651
1652fn prune_artifact_namespace_locked(
1653 cache_root: &Path,
1654 namespace: &Path,
1655 policy: ArtifactCachePrunePolicy,
1656 protected_entry: Option<&Path>,
1657) -> Result<ArtifactCachePruneReport, ArtifactCacheError> {
1658 let mut report = prune_direct_child_directories(
1659 &entries_directory(namespace),
1660 policy,
1661 protected_entry,
1662 is_sha256_directory,
1663 )
1664 .map_err(artifact_cache_fs_error)?;
1665 remove_abandoned_staging(cache_root, namespace, &mut report)?;
1666 Ok(report)
1667}
1668
1669fn remove_abandoned_staging(
1670 cache_root: &Path,
1671 namespace: &Path,
1672 report: &mut ArtifactCachePruneReport,
1673) -> Result<(), ArtifactCacheError> {
1674 let staging_root = namespace.join("staging");
1675 let entries = match fs::read_dir(&staging_root) {
1676 Ok(entries) => entries,
1677 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
1678 Err(source) => {
1679 return Err(ArtifactCacheError::Io {
1680 operation: "read artifact staging root during pruning",
1681 path: staging_root,
1682 source,
1683 });
1684 }
1685 };
1686 for entry in entries {
1687 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1688 operation: "read artifact staging entry during pruning",
1689 path: staging_root.clone(),
1690 source,
1691 })?;
1692 let file_type = entry.file_type().map_err(|source| ArtifactCacheError::Io {
1693 operation: "inspect artifact staging entry during pruning",
1694 path: entry.path(),
1695 source,
1696 })?;
1697 if !file_type.is_dir() {
1698 continue;
1699 }
1700 let file_name = entry.file_name();
1701 let Some(key) = staging_content_key(&file_name) else {
1702 continue;
1703 };
1704 let lock_path = content_lock_path_for_key(cache_root, key);
1705 let Some(_content_lock) =
1706 try_lock_cache_file(&lock_path).map_err(artifact_cache_fs_error)?
1707 else {
1708 continue;
1709 };
1710 let path = entry.path();
1711 let bytes = directory_logical_size(&path).map_err(|source| ArtifactCacheError::Io {
1712 operation: "measure abandoned artifact staging directory",
1713 path: path.clone(),
1714 source,
1715 })?;
1716 remove_path_if_present(&path).map_err(|source| ArtifactCacheError::Io {
1717 operation: "remove abandoned artifact staging directory during pruning",
1718 path,
1719 source,
1720 })?;
1721 report.record_uncommitted_removal(bytes);
1722 }
1723 Ok(())
1724}
1725
1726fn staging_content_key(name: &OsStr) -> Option<&str> {
1727 let name = name.to_str()?;
1728 let (key, suffix) = name.split_once('-')?;
1729 (!suffix.is_empty() && key.len() == 64 && key.as_bytes().iter().all(u8::is_ascii_hexdigit))
1730 .then_some(key)
1731}
1732
1733fn cache_record(
1734 spec: &ArtifactCacheSpec,
1735 resolved: ResolvedKey,
1736 timings: ArtifactCacheTimings,
1737 maintenance: Option<ArtifactCacheMaintenance>,
1738) -> Result<ArtifactCacheRecord, ArtifactCacheError> {
1739 let entry = entry_directory(&namespace_directory(spec), resolved.key);
1740 let retention = RetainedCacheEntry::acquire(&entry).map_err(artifact_cache_fs_error)?;
1741 Ok(ArtifactCacheRecord {
1742 key: resolved.key,
1743 input_digest: resolved.input_digest,
1744 artifacts: spec
1745 .outputs
1746 .iter()
1747 .enumerate()
1748 .map(|(index, output)| ArtifactCacheArtifact {
1749 name: output.name.clone(),
1750 path: staged_output_path(&entry, index),
1751 })
1752 .collect(),
1753 timings,
1754 maintenance,
1755 _retention: retention,
1756 })
1757}
1758
1759fn create_staging_directory(
1760 namespace: &Path,
1761 key: InputDigest,
1762) -> Result<PathBuf, ArtifactCacheError> {
1763 let staging_root = namespace.join("staging");
1764 fs::create_dir_all(&staging_root).map_err(|source| ArtifactCacheError::Io {
1765 operation: "create artifact staging root",
1766 path: staging_root.clone(),
1767 source,
1768 })?;
1769 remove_same_key_staging(&staging_root, key)?;
1770 let sequence = STAGING_SEQUENCE.fetch_add(1, Ordering::Relaxed);
1771 let staging = staging_root.join(format!("{key}-{}-{sequence}", std::process::id()));
1772 fs::create_dir_all(staging.join("outputs")).map_err(|source| ArtifactCacheError::Io {
1773 operation: "create artifact transaction staging directory",
1774 path: staging.clone(),
1775 source,
1776 })?;
1777 Ok(staging)
1778}
1779
1780fn remove_same_key_staging(root: &Path, key: InputDigest) -> Result<(), ArtifactCacheError> {
1781 let prefix = format!("{key}-");
1782 let entries = fs::read_dir(root).map_err(|source| ArtifactCacheError::Io {
1783 operation: "read artifact staging root",
1784 path: root.to_owned(),
1785 source,
1786 })?;
1787 for entry in entries {
1788 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1789 operation: "read artifact staging entry",
1790 path: root.to_owned(),
1791 source,
1792 })?;
1793 if entry.file_name().to_string_lossy().starts_with(&prefix) {
1794 remove_path_if_present(&entry.path()).map_err(|source| ArtifactCacheError::Io {
1795 operation: "remove abandoned artifact staging directory",
1796 path: entry.path(),
1797 source,
1798 })?;
1799 }
1800 }
1801 Ok(())
1802}
1803
1804fn staged_output_path(root: &Path, index: usize) -> PathBuf {
1805 root.join("outputs").join(format_output_index(index))
1806}
1807
1808fn format_output_index(index: usize) -> String {
1809 format!("{index:04}{OUTPUT_FILE_SUFFIX}")
1810}
1811
1812fn namespace_directory(spec: &ArtifactCacheSpec) -> PathBuf {
1813 namespace_directory_for(&spec.cache_root, &spec.namespace)
1814}
1815
1816fn namespace_directory_for(cache_root: &Path, namespace: &str) -> PathBuf {
1817 cache_root
1818 .join(".ic-testkit/artifact-sets/namespaces")
1819 .join(identifier_digest("artifact-cache-namespace-v1", namespace))
1820}
1821
1822fn entries_directory(namespace: &Path) -> PathBuf {
1823 namespace.join("entries")
1824}
1825
1826fn entry_directory(namespace: &Path, key: InputDigest) -> PathBuf {
1827 entries_directory(namespace).join(key.to_hex())
1828}
1829
1830fn coordination_lock_path(spec: &ArtifactCacheSpec) -> PathBuf {
1831 spec.cache_root
1832 .join(".ic-testkit/artifact-sets/locks/coordination")
1833 .join(format!(
1834 "{}.lock",
1835 identifier_digest("artifact-cache-coordination-v1", &spec.coordination_scope,)
1836 ))
1837}
1838
1839fn content_lock_path(spec: &ArtifactCacheSpec, key: InputDigest) -> PathBuf {
1840 content_lock_path_for_key(&spec.cache_root, &key.to_hex())
1841}
1842
1843fn content_lock_path_for_key(cache_root: &Path, key: &str) -> PathBuf {
1844 cache_root
1845 .join(".ic-testkit/artifact-sets/locks/content")
1846 .join(format!("{key}.lock"))
1847}
1848
1849fn namespace_lock_path(spec: &ArtifactCacheSpec) -> PathBuf {
1850 namespace_lock_path_for(&spec.cache_root, &spec.namespace)
1851}
1852
1853fn namespace_lock_path_for(cache_root: &Path, namespace: &str) -> PathBuf {
1854 cache_root
1855 .join(".ic-testkit/artifact-sets/locks/namespaces")
1856 .join(format!(
1857 "{}.lock",
1858 identifier_digest("artifact-cache-namespace-v1", namespace)
1859 ))
1860}
1861
1862fn identifier_digest(domain: &str, identifier: &str) -> String {
1863 digest_bytes(domain, identifier.as_bytes()).to_hex()
1864}
1865
1866fn artifact_cache_fs_error(error: CacheFsError) -> ArtifactCacheError {
1867 ArtifactCacheError::Io {
1868 operation: error.operation,
1869 path: error.path,
1870 source: error.source,
1871 }
1872}
1873
1874impl ArtifactOutputValidation {
1875 const fn cache_token(self) -> &'static str {
1876 match self {
1877 Self::RegularFile => "regular-file",
1878 Self::NonEmptyFile => "nonempty-file",
1879 }
1880 }
1881}
1882
1883impl std::fmt::Display for ArtifactCacheError {
1884 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1885 match self {
1886 Self::InvalidSpec { message } => {
1887 write!(formatter, "invalid artifact cache spec: {message}")
1888 }
1889 Self::Io {
1890 operation,
1891 path,
1892 source,
1893 } => write!(
1894 formatter,
1895 "failed to {operation} at {}: {source}",
1896 path.display()
1897 ),
1898 Self::InputsChangedDuringPreparation { before, after } => write!(
1899 formatter,
1900 "artifact inputs repeatedly changed during cache preparation: {before} -> {after}",
1901 ),
1902 Self::InputsChangedDuringAcquisition { before, after } => write!(
1903 formatter,
1904 "artifact inputs changed while the caller was building: {before} -> {after}",
1905 ),
1906 Self::CargoBuildInputsChanged {
1907 label,
1908 before,
1909 after,
1910 } => write!(
1911 formatter,
1912 "resolved Cargo build inputs `{label}` changed: {before} -> {after}",
1913 ),
1914 Self::CargoBuildInputRevalidation { label, source } => write!(
1915 formatter,
1916 "failed to revalidate resolved Cargo build inputs `{label}`: {source}",
1917 ),
1918 Self::InvalidOutputs { outputs } => write!(
1919 formatter,
1920 "artifact transaction has missing or invalid outputs: {}",
1921 outputs
1922 .iter()
1923 .map(|(name, path)| format!("{name} ({})", path.display()))
1924 .collect::<Vec<_>>()
1925 .join(", "),
1926 ),
1927 Self::UnknownOutput { name } => {
1928 write!(
1929 formatter,
1930 "artifact transaction has no output named `{name}`"
1931 )
1932 }
1933 Self::FailedTransactionCleanup {
1934 transaction_error,
1935 path,
1936 source,
1937 } => write!(
1938 formatter,
1939 "artifact transaction failed ({transaction_error}) and staging cleanup at {} also failed: {source}",
1940 path.display(),
1941 ),
1942 }
1943 }
1944}
1945
1946impl std::error::Error for ArtifactCacheError {
1947 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
1948 match self {
1949 Self::Io { source, .. } | Self::FailedTransactionCleanup { source, .. } => Some(source),
1950 Self::CargoBuildInputRevalidation { source, .. } => Some(source),
1951 _ => None,
1952 }
1953 }
1954}
1955
1956#[cfg(test)]
1957#[path = "transaction/tests.rs"]
1958mod tests;