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 FileDigest, InputDigest, InputHasher, copy_file_atomic, destination_matches_digest,
24 digest_bytes, digest_file, digest_labeled_paths, os_bytes, read_file_with_limit,
25 write_atomic,
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 write_atomic(
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 use std::fmt::Write as _;
1643 writeln!(
1644 manifest,
1645 "output:{index}:{}:{}:{}:{}",
1646 output.name,
1647 output.validation.cache_token(),
1648 info.bytes,
1649 info.digest,
1650 )
1651 .expect("writing an artifact manifest to a String cannot fail");
1652 }
1653 manifest
1654}
1655
1656fn materialize_outputs(
1657 spec: &ArtifactCacheSpec,
1658 entry: &Path,
1659 verified_outputs: Option<&[FileDigest]>,
1660) -> Result<(), ArtifactCacheError> {
1661 for (index, output) in spec.outputs.iter().enumerate() {
1662 if verified_outputs.is_some_and(|info| {
1665 destination_matches_digest("artifact-set-output-v1", &output.destination, &info[index])
1666 }) {
1667 continue;
1668 }
1669 let cached = staged_output_path(entry, index);
1670 copy_file_atomic(&cached, &output.destination).map_err(|source| {
1671 ArtifactCacheError::Io {
1672 operation: "materialize artifact output",
1673 path: output.destination.clone(),
1674 source,
1675 }
1676 })?;
1677 }
1678 Ok(())
1679}
1680
1681fn perform_maintenance_locked(
1682 spec: &ArtifactCacheSpec,
1683 namespace: &Path,
1684 protected_entry: &Path,
1685) -> (Option<ArtifactCacheMaintenance>, Option<Duration>) {
1686 spec.prune_policy.map_or((None, None), |policy| {
1687 let identity = policy.maintenance_identity();
1688 perform_scheduled_cache_maintenance(namespace, spec.prune_interval, &identity, || {
1689 prune_artifact_namespace_locked(
1690 &spec.cache_root,
1691 namespace,
1692 policy,
1693 Some(protected_entry),
1694 )
1695 .map_err(|error| error.to_string())
1696 })
1697 })
1698}
1699
1700fn prune_artifact_namespace_locked(
1701 cache_root: &Path,
1702 namespace: &Path,
1703 policy: ArtifactCachePrunePolicy,
1704 protected_entry: Option<&Path>,
1705) -> Result<ArtifactCachePruneReport, ArtifactCacheError> {
1706 let mut report = prune_direct_child_directories(
1707 &entries_directory(namespace),
1708 policy,
1709 protected_entry,
1710 is_sha256_directory,
1711 )
1712 .map_err(artifact_cache_fs_error)?;
1713 remove_abandoned_staging(cache_root, namespace, &mut report)?;
1714 Ok(report)
1715}
1716
1717fn remove_abandoned_staging(
1718 cache_root: &Path,
1719 namespace: &Path,
1720 report: &mut ArtifactCachePruneReport,
1721) -> Result<(), ArtifactCacheError> {
1722 let staging_root = namespace.join("staging");
1723 let entries = match fs::read_dir(&staging_root) {
1724 Ok(entries) => entries,
1725 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(()),
1726 Err(source) => {
1727 return Err(ArtifactCacheError::Io {
1728 operation: "read artifact staging root during pruning",
1729 path: staging_root,
1730 source,
1731 });
1732 }
1733 };
1734 for entry in entries {
1735 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1736 operation: "read artifact staging entry during pruning",
1737 path: staging_root.clone(),
1738 source,
1739 })?;
1740 let file_type = entry.file_type().map_err(|source| ArtifactCacheError::Io {
1741 operation: "inspect artifact staging entry during pruning",
1742 path: entry.path(),
1743 source,
1744 })?;
1745 if !file_type.is_dir() {
1746 continue;
1747 }
1748 let file_name = entry.file_name();
1749 let Some(key) = staging_content_key(&file_name) else {
1750 continue;
1751 };
1752 let lock_path = content_lock_path_for_key(cache_root, key);
1753 let Some(_content_lock) =
1754 try_lock_cache_file(&lock_path).map_err(artifact_cache_fs_error)?
1755 else {
1756 continue;
1757 };
1758 let path = entry.path();
1759 let bytes = directory_logical_size(&path).map_err(|source| ArtifactCacheError::Io {
1760 operation: "measure abandoned artifact staging directory",
1761 path: path.clone(),
1762 source,
1763 })?;
1764 remove_path_if_present(&path).map_err(|source| ArtifactCacheError::Io {
1765 operation: "remove abandoned artifact staging directory during pruning",
1766 path,
1767 source,
1768 })?;
1769 report.record_uncommitted_removal(bytes);
1770 }
1771 Ok(())
1772}
1773
1774fn staging_content_key(name: &OsStr) -> Option<&str> {
1775 let name = name.to_str()?;
1776 let (key, suffix) = name.split_once('-')?;
1777 (!suffix.is_empty() && key.len() == 64 && key.as_bytes().iter().all(u8::is_ascii_hexdigit))
1778 .then_some(key)
1779}
1780
1781fn cache_record(
1782 spec: &ArtifactCacheSpec,
1783 resolved: ResolvedKey,
1784 timings: ArtifactCacheTimings,
1785 maintenance: Option<ArtifactCacheMaintenance>,
1786) -> Result<ArtifactCacheRecord, ArtifactCacheError> {
1787 let entry = entry_directory(&namespace_directory(spec), resolved.key);
1788 let retention = RetainedCacheEntry::acquire(&entry).map_err(artifact_cache_fs_error)?;
1789 Ok(ArtifactCacheRecord {
1790 key: resolved.key,
1791 input_digest: resolved.input_digest,
1792 artifacts: spec
1793 .outputs
1794 .iter()
1795 .enumerate()
1796 .map(|(index, output)| ArtifactCacheArtifact {
1797 name: output.name.clone(),
1798 path: staged_output_path(&entry, index),
1799 })
1800 .collect(),
1801 timings,
1802 maintenance,
1803 _retention: retention,
1804 })
1805}
1806
1807fn create_staging_directory(
1808 namespace: &Path,
1809 key: InputDigest,
1810) -> Result<PathBuf, ArtifactCacheError> {
1811 let staging_root = namespace.join("staging");
1812 fs::create_dir_all(&staging_root).map_err(|source| ArtifactCacheError::Io {
1813 operation: "create artifact staging root",
1814 path: staging_root.clone(),
1815 source,
1816 })?;
1817 remove_same_key_staging(&staging_root, key)?;
1818 let sequence = STAGING_SEQUENCE.fetch_add(1, Ordering::Relaxed);
1819 let staging = staging_root.join(format!("{key}-{}-{sequence}", std::process::id()));
1820 fs::create_dir_all(staging.join("outputs")).map_err(|source| ArtifactCacheError::Io {
1821 operation: "create artifact transaction staging directory",
1822 path: staging.clone(),
1823 source,
1824 })?;
1825 Ok(staging)
1826}
1827
1828fn remove_same_key_staging(root: &Path, key: InputDigest) -> Result<(), ArtifactCacheError> {
1829 let prefix = format!("{key}-");
1830 let entries = fs::read_dir(root).map_err(|source| ArtifactCacheError::Io {
1831 operation: "read artifact staging root",
1832 path: root.to_owned(),
1833 source,
1834 })?;
1835 for entry in entries {
1836 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1837 operation: "read artifact staging entry",
1838 path: root.to_owned(),
1839 source,
1840 })?;
1841 if entry.file_name().to_string_lossy().starts_with(&prefix) {
1842 remove_path_if_present(&entry.path()).map_err(|source| ArtifactCacheError::Io {
1843 operation: "remove abandoned artifact staging directory",
1844 path: entry.path(),
1845 source,
1846 })?;
1847 }
1848 }
1849 Ok(())
1850}
1851
1852fn staged_output_path(root: &Path, index: usize) -> PathBuf {
1853 root.join("outputs").join(format_output_index(index))
1854}
1855
1856fn format_output_index(index: usize) -> String {
1857 format!("{index:04}{OUTPUT_FILE_SUFFIX}")
1858}
1859
1860fn namespace_directory(spec: &ArtifactCacheSpec) -> PathBuf {
1861 namespace_directory_for(&spec.cache_root, &spec.namespace)
1862}
1863
1864fn namespace_directory_for(cache_root: &Path, namespace: &str) -> PathBuf {
1865 cache_root
1866 .join(".ic-testkit/artifact-sets/namespaces")
1867 .join(identifier_digest("artifact-cache-namespace-v1", namespace))
1868}
1869
1870fn entries_directory(namespace: &Path) -> PathBuf {
1871 namespace.join("entries")
1872}
1873
1874fn entry_directory(namespace: &Path, key: InputDigest) -> PathBuf {
1875 entries_directory(namespace).join(key.to_hex())
1876}
1877
1878fn coordination_lock_path(spec: &ArtifactCacheSpec) -> PathBuf {
1879 spec.cache_root
1880 .join(".ic-testkit/artifact-sets/locks/coordination")
1881 .join(format!(
1882 "{}.lock",
1883 identifier_digest("artifact-cache-coordination-v1", &spec.coordination_scope,)
1884 ))
1885}
1886
1887fn content_lock_path(spec: &ArtifactCacheSpec, key: InputDigest) -> PathBuf {
1888 content_lock_path_for_key(&spec.cache_root, &key.to_hex())
1889}
1890
1891fn content_lock_path_for_key(cache_root: &Path, key: &str) -> PathBuf {
1892 cache_root
1893 .join(".ic-testkit/artifact-sets/locks/content")
1894 .join(format!("{key}.lock"))
1895}
1896
1897fn namespace_lock_path(spec: &ArtifactCacheSpec) -> PathBuf {
1898 namespace_lock_path_for(&spec.cache_root, &spec.namespace)
1899}
1900
1901fn namespace_lock_path_for(cache_root: &Path, namespace: &str) -> PathBuf {
1902 cache_root
1903 .join(".ic-testkit/artifact-sets/locks/namespaces")
1904 .join(format!(
1905 "{}.lock",
1906 identifier_digest("artifact-cache-namespace-v1", namespace)
1907 ))
1908}
1909
1910fn identifier_digest(domain: &str, identifier: &str) -> String {
1911 digest_bytes(domain, identifier.as_bytes()).to_hex()
1912}
1913
1914fn artifact_cache_fs_error(error: CacheFsError) -> ArtifactCacheError {
1915 ArtifactCacheError::Io {
1916 operation: error.operation,
1917 path: error.path,
1918 source: error.source,
1919 }
1920}
1921
1922impl ArtifactOutputValidation {
1923 const fn cache_token(self) -> &'static str {
1924 match self {
1925 Self::RegularFile => "regular-file",
1926 Self::NonEmptyFile => "nonempty-file",
1927 }
1928 }
1929}
1930
1931impl std::fmt::Display for ArtifactCacheError {
1932 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1933 match self {
1934 Self::InvalidSpec { message } => {
1935 write!(formatter, "invalid artifact cache spec: {message}")
1936 }
1937 Self::Io {
1938 operation,
1939 path,
1940 source,
1941 } => write!(
1942 formatter,
1943 "failed to {operation} at {}: {source}",
1944 path.display()
1945 ),
1946 Self::InputsChangedDuringPreparation { before, after } => write!(
1947 formatter,
1948 "artifact inputs repeatedly changed during cache preparation: {before} -> {after}",
1949 ),
1950 Self::InputsChangedDuringAcquisition { before, after } => write!(
1951 formatter,
1952 "artifact inputs changed while the caller was building: {before} -> {after}",
1953 ),
1954 Self::CargoBuildInputsChanged {
1955 label,
1956 before,
1957 after,
1958 } => write!(
1959 formatter,
1960 "resolved Cargo build inputs `{label}` changed: {before} -> {after}",
1961 ),
1962 Self::CargoBuildInputRevalidation { label, source } => write!(
1963 formatter,
1964 "failed to revalidate resolved Cargo build inputs `{label}`: {source}",
1965 ),
1966 Self::InvalidOutputs { outputs } => write!(
1967 formatter,
1968 "artifact transaction has missing or invalid outputs: {}",
1969 outputs
1970 .iter()
1971 .map(|(name, path)| format!("{name} ({})", path.display()))
1972 .collect::<Vec<_>>()
1973 .join(", "),
1974 ),
1975 Self::UnknownOutput { name } => {
1976 write!(
1977 formatter,
1978 "artifact transaction has no output named `{name}`"
1979 )
1980 }
1981 Self::FailedTransactionCleanup {
1982 transaction_error,
1983 path,
1984 source,
1985 } => write!(
1986 formatter,
1987 "artifact transaction failed ({transaction_error}) and staging cleanup at {} also failed: {source}",
1988 path.display(),
1989 ),
1990 }
1991 }
1992}
1993
1994impl std::error::Error for ArtifactCacheError {
1995 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
1996 match self {
1997 Self::Io { source, .. } | Self::FailedTransactionCleanup { source, .. } => Some(source),
1998 Self::CargoBuildInputRevalidation { source, .. } => Some(source),
1999 _ => None,
2000 }
2001 }
2002}
2003
2004#[cfg(test)]
2005#[path = "transaction/tests.rs"]
2006mod tests;