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