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 write_atomic,
27 },
28 wasm_cache::{
29 ResolvedCargoBuildInputs, WasmBuildError, WasmBuildSpec, resolve_cargo_build_inputs,
30 },
31};
32
33const ARTIFACT_CACHE_FORMAT: &str = "ic-testkit-artifact-set-v1";
34const MANIFEST_FILE: &str = "manifest.ic-testkit";
35const OUTPUT_FILE_SUFFIX: &str = ".artifact";
36const MAX_PREPARATION_RETRIES: usize = 3;
37static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
38
39#[derive(Clone, Copy, Debug, Default, Eq, Ord, PartialEq, PartialOrd)]
41pub enum ArtifactOutputValidation {
42 RegularFile,
44 #[default]
46 NonEmptyFile,
47}
48
49#[derive(Clone, Eq, PartialEq)]
57pub struct ArtifactCacheSpec {
58 cache_root: PathBuf,
59 namespace: String,
60 recipe_id: String,
61 coordination_scope: String,
62 inputs: Vec<LabeledPath>,
63 tools: Vec<LabeledPath>,
64 arguments: Vec<OsString>,
65 environment: BTreeMap<OsString, Option<OsString>>,
66 identities: Vec<LabeledIdentity>,
68 cargo_build_inputs: Vec<CargoBuildInputSet>,
69 outputs: Vec<OutputSpec>,
70 prune_policy: Option<ArtifactCachePrunePolicy>,
71 prune_interval: Option<Duration>,
72}
73
74pub enum ArtifactCachePreparation {
76 Reused(ArtifactCacheRecord),
78 Build(ArtifactBuildTransaction),
80}
81
82#[derive(Clone, Debug, Eq, PartialEq)]
84pub enum ArtifactCacheOutcome {
85 Built(ArtifactCacheRecord),
87 Reused(ArtifactCacheRecord),
89}
90
91#[derive(Clone, Debug, Eq, PartialEq)]
93pub struct ArtifactCacheRecord {
94 key: InputDigest,
95 input_digest: InputDigest,
96 artifacts: Vec<ArtifactCacheArtifact>,
97 _retention: RetainedCacheEntry,
98 timings: ArtifactCacheTimings,
99 maintenance: Option<ArtifactCacheMaintenance>,
100}
101
102#[derive(Clone, Debug, Eq, PartialEq)]
104pub struct ArtifactCacheArtifact {
105 name: String,
106 path: PathBuf,
107}
108
109#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
111pub struct ArtifactCacheTimings {
112 coordination_lock_wait: Duration,
113 content_lock_wait: Duration,
114 namespace_lock_wait: Duration,
115 input_capture: Duration,
116 cache_lookup: Duration,
117 caller_build: Option<Duration>,
118 output_validation: Duration,
119 publication: Duration,
120 materialization: Duration,
121 maintenance: Option<Duration>,
122 total: Duration,
123}
124
125pub struct ArtifactBuildTransaction {
127 spec: Box<ArtifactCacheSpec>,
128 resolved: ResolvedKey,
129 staging_directory: PathBuf,
130 entry_directory: PathBuf,
131 namespace_directory: PathBuf,
132 _coordination_lock: File,
133 _content_lock: File,
134 timings: ArtifactCacheTimings,
135 total_started: Instant,
136 caller_build_started: Instant,
137 staging_armed: bool,
138}
139
140#[non_exhaustive]
142#[derive(Debug)]
143pub enum ArtifactCacheError {
144 InvalidSpec { message: String },
146 Io {
148 operation: &'static str,
149 path: PathBuf,
150 source: io::Error,
151 },
152 InputsChangedDuringPreparation {
154 before: InputDigest,
155 after: InputDigest,
156 },
157 InputsChangedDuringAcquisition {
159 before: InputDigest,
160 after: InputDigest,
161 },
162 CargoBuildInputsChanged {
164 label: String,
166 before: InputDigest,
168 after: InputDigest,
170 },
171 CargoBuildInputRevalidation {
173 label: String,
175 source: WasmBuildError,
177 },
178 InvalidOutputs { outputs: Vec<(String, PathBuf)> },
180 UnknownOutput { name: String },
182 FailedTransactionCleanup {
184 transaction_error: Box<Self>,
185 path: PathBuf,
186 source: io::Error,
187 },
188}
189
190#[derive(Clone, Debug, Eq, PartialEq)]
191struct LabeledPath {
192 label: String,
193 path: PathBuf,
194}
195
196#[derive(Clone, Eq, PartialEq)]
197struct LabeledIdentity {
198 label: String,
199 value: Vec<u8>,
200}
201
202#[derive(Clone, Eq, PartialEq)]
203struct CargoBuildInputSet {
204 label: String,
205 build_spec: WasmBuildSpec,
206 resolved: ResolvedCargoBuildInputs,
207}
208
209#[derive(Clone, Debug, Eq, PartialEq)]
210struct OutputSpec {
211 name: String,
212 destination: PathBuf,
213 validation: ArtifactOutputValidation,
214}
215
216#[derive(Clone, Copy)]
217struct ResolvedKey {
218 key: InputDigest,
219 input_digest: InputDigest,
220}
221
222impl ArtifactCacheSpec {
223 #[must_use]
229 pub fn new(cache_root: &Path, namespace: &str, recipe_id: &str) -> Self {
230 Self {
231 cache_root: cache_root.to_owned(),
232 namespace: namespace.to_owned(),
233 recipe_id: recipe_id.to_owned(),
234 coordination_scope: namespace.to_owned(),
235 inputs: Vec::new(),
236 tools: Vec::new(),
237 arguments: Vec::new(),
238 environment: BTreeMap::new(),
239 identities: Vec::new(),
240 cargo_build_inputs: Vec::new(),
241 outputs: Vec::new(),
242 prune_policy: None,
243 prune_interval: None,
244 }
245 }
246
247 #[must_use]
249 pub fn with_coordination_scope(mut self, coordination_scope: &str) -> Self {
250 coordination_scope.clone_into(&mut self.coordination_scope);
251 self
252 }
253
254 #[must_use]
256 pub fn with_input(mut self, label: &str, path: &Path) -> Self {
257 self.inputs.push(LabeledPath {
258 label: label.to_owned(),
259 path: path.to_owned(),
260 });
261 self
262 }
263
264 #[must_use]
266 pub fn with_tool(mut self, label: &str, path: &Path) -> Self {
267 self.tools.push(LabeledPath {
268 label: label.to_owned(),
269 path: path.to_owned(),
270 });
271 self
272 }
273
274 #[must_use]
276 pub fn with_arguments<I, S>(mut self, arguments: I) -> Self
277 where
278 I: IntoIterator<Item = S>,
279 S: AsRef<OsStr>,
280 {
281 self.arguments = arguments
282 .into_iter()
283 .map(|argument| argument.as_ref().to_owned())
284 .collect();
285 self
286 }
287
288 #[must_use]
293 pub fn with_environment<I, K, V>(mut self, environment: I) -> Self
294 where
295 I: IntoIterator<Item = (K, V)>,
296 K: Into<OsString>,
297 V: Into<OsString>,
298 {
299 self.environment.extend(
300 environment
301 .into_iter()
302 .map(|(name, value)| (name.into(), Some(value.into()))),
303 );
304 self
305 }
306
307 #[must_use]
312 pub fn with_unset_environment<I, S>(mut self, names: I) -> Self
313 where
314 I: IntoIterator<Item = S>,
315 S: Into<OsString>,
316 {
317 self.environment
318 .extend(names.into_iter().map(|name| (name.into(), None)));
319 self
320 }
321
322 #[must_use]
326 pub fn with_identity_bytes(mut self, label: &str, value: &[u8]) -> Self {
327 self.identities.push(LabeledIdentity {
328 label: label.to_owned(),
329 value: value.to_vec(),
330 });
331 self.identities
332 .sort_by(|left, right| left.label.cmp(&right.label));
333 self
334 }
335
336 #[must_use]
344 pub fn with_cargo_build_inputs(
345 mut self,
346 label: &str,
347 build_spec: &WasmBuildSpec,
348 resolved: &ResolvedCargoBuildInputs,
349 ) -> Self {
350 self.cargo_build_inputs.push(CargoBuildInputSet {
351 label: label.to_owned(),
352 build_spec: build_spec.clone(),
353 resolved: resolved.clone(),
354 });
355 self.cargo_build_inputs
356 .sort_by(|left, right| left.label.cmp(&right.label));
357 self
358 }
359
360 #[must_use]
362 pub fn with_output(self, name: &str, destination: &Path) -> Self {
363 self.with_output_validation(name, destination, ArtifactOutputValidation::NonEmptyFile)
364 }
365
366 #[must_use]
368 pub fn with_output_validation(
369 mut self,
370 name: &str,
371 destination: &Path,
372 validation: ArtifactOutputValidation,
373 ) -> Self {
374 self.outputs.push(OutputSpec {
375 name: name.to_owned(),
376 destination: destination.to_owned(),
377 validation,
378 });
379 self.outputs
380 .sort_by(|left, right| left.name.cmp(&right.name));
381 self
382 }
383
384 #[must_use]
386 pub const fn with_prune_policy(mut self, policy: ArtifactCachePrunePolicy) -> Self {
387 self.prune_policy = Some(policy);
388 self.prune_interval = None;
389 self
390 }
391
392 #[must_use]
400 pub const fn with_prune_policy_at_most_every(
401 mut self,
402 policy: ArtifactCachePrunePolicy,
403 minimum_interval: Duration,
404 ) -> Self {
405 self.prune_policy = Some(policy);
406 self.prune_interval = Some(minimum_interval);
407 self
408 }
409
410 #[must_use]
412 pub fn cache_root(&self) -> &Path {
413 &self.cache_root
414 }
415
416 #[must_use]
418 pub fn namespace(&self) -> &str {
419 &self.namespace
420 }
421
422 #[must_use]
424 pub fn recipe_id(&self) -> &str {
425 &self.recipe_id
426 }
427}
428
429impl std::fmt::Debug for ArtifactCacheSpec {
430 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
431 formatter
432 .debug_struct("ArtifactCacheSpec")
433 .field("cache_root", &self.cache_root)
434 .field("namespace", &self.namespace)
435 .field("recipe_id", &self.recipe_id)
436 .field("coordination_scope", &self.coordination_scope)
437 .field("inputs", &self.inputs)
438 .field("tools", &self.tools)
439 .field("argument_count", &self.arguments.len())
440 .field(
441 "environment_names",
442 &self.environment.keys().collect::<Vec<_>>(),
443 )
444 .field(
445 "identity_labels",
446 &self
447 .identities
448 .iter()
449 .map(|identity| identity.label.as_str())
450 .collect::<Vec<_>>(),
451 )
452 .field(
453 "cargo_build_input_labels",
454 &self
455 .cargo_build_inputs
456 .iter()
457 .map(|input| input.label.as_str())
458 .collect::<Vec<_>>(),
459 )
460 .field("outputs", &self.outputs)
461 .field("prune_policy", &self.prune_policy)
462 .field("prune_interval", &self.prune_interval)
463 .finish()
464 }
465}
466
467impl ArtifactCachePreparation {
468 #[must_use]
470 pub const fn reused_record(&self) -> Option<&ArtifactCacheRecord> {
471 match self {
472 Self::Reused(record) => Some(record),
473 Self::Build(_) => None,
474 }
475 }
476}
477
478impl ArtifactCacheOutcome {
479 #[must_use]
481 pub const fn record(&self) -> &ArtifactCacheRecord {
482 match self {
483 Self::Built(record) | Self::Reused(record) => record,
484 }
485 }
486
487 #[must_use]
489 pub const fn is_reused(&self) -> bool {
490 matches!(self, Self::Reused(_))
491 }
492}
493
494impl ArtifactCacheRecord {
495 #[must_use]
497 pub const fn key(&self) -> InputDigest {
498 self.key
499 }
500
501 #[must_use]
503 pub const fn input_digest(&self) -> InputDigest {
504 self.input_digest
505 }
506
507 #[must_use]
513 pub fn artifacts(&self) -> &[ArtifactCacheArtifact] {
514 &self.artifacts
515 }
516
517 #[must_use]
519 pub const fn timings(&self) -> ArtifactCacheTimings {
520 self.timings
521 }
522
523 #[must_use]
525 pub const fn maintenance(&self) -> Option<&ArtifactCacheMaintenance> {
526 self.maintenance.as_ref()
527 }
528}
529
530impl ArtifactCacheArtifact {
531 #[must_use]
533 pub fn name(&self) -> &str {
534 &self.name
535 }
536
537 #[must_use]
539 pub fn path(&self) -> &Path {
540 &self.path
541 }
542}
543
544impl ArtifactCacheTimings {
545 #[must_use]
547 pub const fn coordination_lock_wait(self) -> Duration {
548 self.coordination_lock_wait
549 }
550
551 #[must_use]
553 pub const fn content_lock_wait(self) -> Duration {
554 self.content_lock_wait
555 }
556
557 #[must_use]
559 pub const fn namespace_lock_wait(self) -> Duration {
560 self.namespace_lock_wait
561 }
562
563 #[must_use]
565 pub const fn input_capture(self) -> Duration {
566 self.input_capture
567 }
568
569 #[must_use]
571 pub const fn cache_lookup(self) -> Duration {
572 self.cache_lookup
573 }
574
575 #[must_use]
577 pub const fn caller_build(self) -> Option<Duration> {
578 self.caller_build
579 }
580
581 #[must_use]
583 pub const fn output_validation(self) -> Duration {
584 self.output_validation
585 }
586
587 #[must_use]
589 pub const fn publication(self) -> Duration {
590 self.publication
591 }
592
593 #[must_use]
595 pub const fn materialization(self) -> Duration {
596 self.materialization
597 }
598
599 #[must_use]
601 pub const fn maintenance(self) -> Option<Duration> {
602 self.maintenance
603 }
604
605 #[must_use]
607 pub const fn total(self) -> Duration {
608 self.total
609 }
610
611 pub(super) const fn saturating_add(self, other: Self) -> Self {
612 Self {
613 coordination_lock_wait: self
614 .coordination_lock_wait
615 .saturating_add(other.coordination_lock_wait),
616 content_lock_wait: self
617 .content_lock_wait
618 .saturating_add(other.content_lock_wait),
619 namespace_lock_wait: self
620 .namespace_lock_wait
621 .saturating_add(other.namespace_lock_wait),
622 input_capture: self.input_capture.saturating_add(other.input_capture),
623 cache_lookup: self.cache_lookup.saturating_add(other.cache_lookup),
624 caller_build: saturating_add_optional_duration(self.caller_build, other.caller_build),
625 output_validation: self
626 .output_validation
627 .saturating_add(other.output_validation),
628 publication: self.publication.saturating_add(other.publication),
629 materialization: self.materialization.saturating_add(other.materialization),
630 maintenance: saturating_add_optional_duration(self.maintenance, other.maintenance),
631 total: self.total.saturating_add(other.total),
632 }
633 }
634}
635
636impl std::fmt::Display for ArtifactCacheTimings {
637 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
638 write!(
639 formatter,
640 "total={:?} coordination_lock={:?} content_lock={:?} namespace_lock={:?} inputs={:?} lookup={:?} build={:?} validation={:?} publication={:?} materialization={:?} maintenance={:?}",
641 self.total,
642 self.coordination_lock_wait,
643 self.content_lock_wait,
644 self.namespace_lock_wait,
645 self.input_capture,
646 self.cache_lookup,
647 self.caller_build,
648 self.output_validation,
649 self.publication,
650 self.materialization,
651 self.maintenance,
652 )
653 }
654}
655
656impl std::fmt::Display for ArtifactCacheOutcome {
657 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
658 let state = if self.is_reused() { "reused" } else { "built" };
659 write!(
660 formatter,
661 "{state} key={} artifacts={} {}",
662 self.record().key,
663 self.record().artifacts.len(),
664 self.record().timings,
665 )
666 }
667}
668
669impl ArtifactBuildTransaction {
670 #[must_use]
675 pub fn staging_directory(&self) -> &Path {
676 &self.staging_directory
677 }
678
679 pub fn output_path(&self, name: &str) -> Result<PathBuf, ArtifactCacheError> {
681 self.output_index(name)
682 .map(|index| staged_output_path(&self.staging_directory, index))
683 .ok_or_else(|| ArtifactCacheError::UnknownOutput {
684 name: name.to_owned(),
685 })
686 }
687
688 pub fn import_output(&self, name: &str, source: &Path) -> Result<(), ArtifactCacheError> {
690 let destination = self.output_path(name)?;
691 copy_file_atomic(source, &destination).map_err(|source_error| ArtifactCacheError::Io {
692 operation: "import artifact output into staging",
693 path: destination,
694 source: source_error,
695 })?;
696 Ok(())
697 }
698
699 pub fn commit(mut self) -> Result<ArtifactCacheOutcome, ArtifactCacheError> {
701 let result = self.commit_inner();
702 match result {
703 Ok(outcome) => Ok(outcome),
704 Err(transaction_error) if self.staging_armed => {
705 let path = self.staging_directory.clone();
706 match remove_path_if_present(&path) {
707 Ok(()) => {
708 self.staging_armed = false;
709 Err(transaction_error)
710 }
711 Err(source) => Err(ArtifactCacheError::FailedTransactionCleanup {
712 transaction_error: Box::new(transaction_error),
713 path,
714 source,
715 }),
716 }
717 }
718 Err(transaction_error) => Err(transaction_error),
719 }
720 }
721
722 pub fn abort(mut self) -> Result<(), ArtifactCacheError> {
724 remove_path_if_present(&self.staging_directory).map_err(|source| {
725 ArtifactCacheError::Io {
726 operation: "abort artifact cache transaction",
727 path: self.staging_directory.clone(),
728 source,
729 }
730 })?;
731 self.staging_armed = false;
732 Ok(())
733 }
734
735 fn output_index(&self, name: &str) -> Option<usize> {
736 self.spec
737 .outputs
738 .iter()
739 .position(|output| output.name == name)
740 }
741
742 fn commit_inner(&mut self) -> Result<ArtifactCacheOutcome, ArtifactCacheError> {
743 self.timings.caller_build = Some(self.caller_build_started.elapsed());
744
745 let validation_started = Instant::now();
746 let output_info = inspect_complete_output_set(&self.spec, &self.staging_directory)?;
747 self.timings.output_validation = validation_started.elapsed();
748
749 let capture_started = Instant::now();
750 revalidate_cargo_build_input_fingerprints(&self.spec)?;
751 let verified = resolve_key(&self.spec)?;
752 self.timings.input_capture = self
753 .timings
754 .input_capture
755 .saturating_add(capture_started.elapsed());
756 if verified.input_digest != self.resolved.input_digest {
757 return Err(ArtifactCacheError::InputsChangedDuringAcquisition {
758 before: self.resolved.input_digest,
759 after: verified.input_digest,
760 });
761 }
762
763 let publication_started = Instant::now();
764 let manifest =
765 manifest_contents(self.resolved.key, &self.spec, output_info.iter().copied());
766 write_atomic(
767 &self.staging_directory.join(MANIFEST_FILE),
768 manifest.as_bytes(),
769 )
770 .map_err(|source| ArtifactCacheError::Io {
771 operation: "write artifact cache manifest",
772 path: self.staging_directory.join(MANIFEST_FILE),
773 source,
774 })?;
775 let namespace_lock_path = namespace_lock_path(&self.spec);
776 let (_namespace_lock, namespace_wait) =
777 lock_cache_file(&namespace_lock_path).map_err(artifact_cache_fs_error)?;
778 self.timings.namespace_lock_wait = self
779 .timings
780 .namespace_lock_wait
781 .saturating_add(namespace_wait);
782 remove_unretained_entry(&self.entry_directory).map_err(artifact_cache_fs_error)?;
783 fs::rename(&self.staging_directory, &self.entry_directory).map_err(|source| {
784 ArtifactCacheError::Io {
785 operation: "publish artifact cache entry",
786 path: self.entry_directory.clone(),
787 source,
788 }
789 })?;
790 self.staging_armed = false;
791 self.timings.publication = publication_started.elapsed();
792
793 let materialization_started = Instant::now();
794 materialize_outputs(&self.spec, &self.entry_directory, None)?;
795 self.timings.materialization = materialization_started.elapsed();
796 record_cache_entry_use(&self.entry_directory).map_err(artifact_cache_fs_error)?;
797 let (maintenance, maintenance_timing) = perform_maintenance_locked(
798 &self.spec,
799 &self.namespace_directory,
800 &self.entry_directory,
801 );
802 self.timings.maintenance = maintenance_timing;
803 self.timings.total = self.total_started.elapsed();
804
805 Ok(ArtifactCacheOutcome::Built(cache_record(
806 &self.spec,
807 self.resolved,
808 self.timings,
809 maintenance,
810 )?))
811 }
812}
813
814impl Drop for ArtifactBuildTransaction {
815 fn drop(&mut self) {
816 if self.staging_armed {
817 let _ = remove_path_if_present(&self.staging_directory);
818 }
819 }
820}
821
822pub fn prepare_artifact_cache(
834 spec: &ArtifactCacheSpec,
835) -> Result<ArtifactCachePreparation, ArtifactCacheError> {
836 let total_started = Instant::now();
837 validate_spec(spec)?;
838 validate_filesystem_boundaries(spec)?;
839 let namespace_directory = initialize_cache(spec)?;
840
841 let coordination_lock_path = coordination_lock_path(spec);
842 let (coordination_lock, coordination_wait) =
843 lock_cache_file(&coordination_lock_path).map_err(artifact_cache_fs_error)?;
844 let mut timings = ArtifactCacheTimings {
845 coordination_lock_wait: coordination_wait,
846 ..ArtifactCacheTimings::default()
847 };
848 let initial_cargo_started = Instant::now();
849 revalidate_cargo_build_input_fingerprints(spec)?;
850 timings.input_capture = initial_cargo_started.elapsed();
851 let mut last_change = None;
852
853 for _ in 0..MAX_PREPARATION_RETRIES {
854 let capture_started = Instant::now();
855 let resolved = resolve_key(spec)?;
856 timings.input_capture = timings
857 .input_capture
858 .saturating_add(capture_started.elapsed());
859
860 let content_lock_path = content_lock_path(spec, resolved.key);
861 let (content_lock, content_wait) =
862 lock_cache_file(&content_lock_path).map_err(artifact_cache_fs_error)?;
863 timings.content_lock_wait = timings.content_lock_wait.saturating_add(content_wait);
864
865 let verification_started = Instant::now();
866 let verified = resolve_key(spec)?;
867 timings.input_capture = timings
868 .input_capture
869 .saturating_add(verification_started.elapsed());
870 if resolved.input_digest != verified.input_digest {
871 last_change = Some((resolved.input_digest, verified.input_digest));
872 drop(content_lock);
873 continue;
874 }
875
876 let entry_directory = entry_directory(&namespace_directory, resolved.key);
877 let namespace_lock_path = namespace_lock_path(spec);
878 let (namespace_lock, namespace_wait) =
879 lock_cache_file(&namespace_lock_path).map_err(artifact_cache_fs_error)?;
880 timings.namespace_lock_wait = timings.namespace_lock_wait.saturating_add(namespace_wait);
881 let lookup_started = Instant::now();
882 let verified_outputs = validated_cache_outputs(spec, resolved.key, &entry_directory)?;
883 timings.cache_lookup = timings
884 .cache_lookup
885 .saturating_add(lookup_started.elapsed());
886
887 if let Some(output_info) = verified_outputs {
888 let materialization_started = Instant::now();
889 materialize_outputs(spec, &entry_directory, Some(&output_info))?;
890 timings.materialization = timings
891 .materialization
892 .saturating_add(materialization_started.elapsed());
893 let after_started = Instant::now();
894 revalidate_cargo_build_input_fingerprints(spec)?;
895 let after = resolve_key(spec)?;
896 timings.input_capture = timings
897 .input_capture
898 .saturating_add(after_started.elapsed());
899 if after.input_digest != resolved.input_digest {
900 last_change = Some((resolved.input_digest, after.input_digest));
901 drop(namespace_lock);
902 drop(content_lock);
903 continue;
904 }
905 record_cache_entry_use(&entry_directory).map_err(artifact_cache_fs_error)?;
906 let (maintenance, maintenance_timing) =
907 perform_maintenance_locked(spec, &namespace_directory, &entry_directory);
908 timings.maintenance = maintenance_timing;
909 timings.total = total_started.elapsed();
910 return Ok(ArtifactCachePreparation::Reused(cache_record(
911 spec,
912 resolved,
913 timings,
914 maintenance,
915 )?));
916 }
917
918 remove_unretained_entry(&entry_directory).map_err(artifact_cache_fs_error)?;
919 let staging_directory = create_staging_directory(&namespace_directory, resolved.key)?;
920 drop(namespace_lock);
921 return Ok(ArtifactCachePreparation::Build(ArtifactBuildTransaction {
922 spec: Box::new(spec.clone()),
923 resolved,
924 staging_directory,
925 entry_directory,
926 namespace_directory,
927 _coordination_lock: coordination_lock,
928 _content_lock: content_lock,
929 timings,
930 total_started,
931 caller_build_started: Instant::now(),
932 staging_armed: true,
933 }));
934 }
935
936 let (before, after) = last_change.expect("preparation retries require a recorded input change");
937 Err(ArtifactCacheError::InputsChangedDuringPreparation { before, after })
938}
939
940fn revalidate_cargo_build_input_fingerprints(
941 spec: &ArtifactCacheSpec,
942) -> Result<(), ArtifactCacheError> {
943 for cargo_input in &spec.cargo_build_inputs {
944 let current = resolve_cargo_build_inputs(&cargo_input.build_spec).map_err(|source| {
945 ArtifactCacheError::CargoBuildInputRevalidation {
946 label: cargo_input.label.clone(),
947 source,
948 }
949 })?;
950 let before_validation = cargo_input.resolved.validation_digest();
951 let after_validation = current.validation_digest();
952 if after_validation != before_validation {
953 return Err(ArtifactCacheError::CargoBuildInputsChanged {
954 label: cargo_input.label.clone(),
955 before: before_validation,
956 after: after_validation,
957 });
958 }
959 let before = cargo_input.resolved.fingerprint();
960 let after = current.fingerprint();
961 if after != before {
962 return Err(ArtifactCacheError::CargoBuildInputsChanged {
963 label: cargo_input.label.clone(),
964 before,
965 after,
966 });
967 }
968 }
969 Ok(())
970}
971
972pub fn prune_artifact_cache(
981 cache_root: &Path,
982 namespace: &str,
983 policy: ArtifactCachePrunePolicy,
984) -> Result<ArtifactCachePruneReport, ArtifactCacheError> {
985 validate_identifier("namespace", namespace)?;
986 if cache_root.as_os_str().is_empty() {
987 return invalid_spec("cache root must not be empty");
988 }
989 ensure_cache_directory_tag(cache_root).map_err(artifact_cache_fs_error)?;
990 let namespace_directory = namespace_directory_for(cache_root, namespace);
991 let lock_path = namespace_lock_path_for(cache_root, namespace);
992 let (_lock, _) = lock_cache_file(&lock_path).map_err(artifact_cache_fs_error)?;
993 prune_artifact_namespace_locked(cache_root, &namespace_directory, policy, None)
994}
995
996fn initialize_cache(spec: &ArtifactCacheSpec) -> Result<PathBuf, ArtifactCacheError> {
997 ensure_cache_directory_tag(&spec.cache_root).map_err(artifact_cache_fs_error)?;
998 let namespace = namespace_directory(spec);
999 fs::create_dir_all(entries_directory(&namespace)).map_err(|source| ArtifactCacheError::Io {
1000 operation: "create artifact cache namespace",
1001 path: namespace.clone(),
1002 source,
1003 })?;
1004 Ok(namespace)
1005}
1006
1007fn validate_spec(spec: &ArtifactCacheSpec) -> Result<(), ArtifactCacheError> {
1008 validate_identifier("namespace", &spec.namespace)?;
1009 validate_identifier("recipe identity", &spec.recipe_id)?;
1010 validate_identifier("coordination scope", &spec.coordination_scope)?;
1011 if spec.cache_root.as_os_str().is_empty() {
1012 return invalid_spec("cache root must not be empty");
1013 }
1014 if spec.outputs.is_empty() {
1015 return invalid_spec("at least one artifact output is required");
1016 }
1017
1018 let mut labels = BTreeSet::new();
1019 for (kind, paths) in [("input", &spec.inputs), ("tool", &spec.tools)] {
1020 for path in paths {
1021 validate_path_label(kind, &path.label)?;
1022 if !labels.insert((kind, path.label.as_str())) {
1023 return invalid_spec(&format!("duplicate {kind} label `{}`", path.label));
1024 }
1025 }
1026 }
1027 let mut identity_labels = BTreeSet::new();
1028 for identity in &spec.identities {
1029 validate_label("identity", &identity.label)?;
1030 if !identity_labels.insert(&identity.label) {
1031 return invalid_spec(&format!("duplicate identity label `{}`", identity.label));
1032 }
1033 }
1034 let mut cargo_input_labels = BTreeSet::new();
1035 for cargo_inputs in &spec.cargo_build_inputs {
1036 validate_label("Cargo build input", &cargo_inputs.label)?;
1037 if !cargo_input_labels.insert(&cargo_inputs.label) {
1038 return invalid_spec(&format!(
1039 "duplicate Cargo build input label `{}`",
1040 cargo_inputs.label
1041 ));
1042 }
1043 }
1044 if spec
1045 .environment
1046 .keys()
1047 .any(|name| name.as_os_str().is_empty())
1048 {
1049 return invalid_spec("environment names must not be empty");
1050 }
1051 let mut output_names = BTreeSet::new();
1052 let mut destinations = BTreeSet::new();
1053 for output in &spec.outputs {
1054 validate_output_name(&output.name)?;
1055 if !output_names.insert(&output.name) {
1056 return invalid_spec(&format!("duplicate output name `{}`", output.name));
1057 }
1058 if output.destination.as_os_str().is_empty() {
1059 return invalid_spec(&format!(
1060 "output `{}` destination must not be empty",
1061 output.name
1062 ));
1063 }
1064 if !destinations.insert(&output.destination) {
1065 return invalid_spec(&format!(
1066 "output `{}` shares a destination with another output",
1067 output.name
1068 ));
1069 }
1070 }
1071 Ok(())
1072}
1073
1074fn validate_filesystem_boundaries(spec: &ArtifactCacheSpec) -> Result<(), ArtifactCacheError> {
1075 let cache_root =
1076 canonicalize_allow_missing(&spec.cache_root).map_err(|source| ArtifactCacheError::Io {
1077 operation: "resolve artifact cache root",
1078 path: spec.cache_root.clone(),
1079 source,
1080 })?;
1081 let mut declared_paths = Vec::with_capacity(spec.inputs.len() + spec.tools.len());
1082 for (kind, labeled_paths) in [("input", &spec.inputs), ("tool", &spec.tools)] {
1083 for labeled in labeled_paths {
1084 let canonical =
1085 canonicalize_path(&labeled.path, "canonicalize declared artifact cache path")?;
1086 if canonical.starts_with(&cache_root) {
1087 return invalid_spec(&format!(
1088 "{kind} `{}` must not be located inside the artifact cache root",
1089 labeled.label
1090 ));
1091 }
1092 let is_directory = fs::metadata(&canonical)
1093 .map_err(|source| ArtifactCacheError::Io {
1094 operation: "inspect declared artifact cache path",
1095 path: canonical.clone(),
1096 source,
1097 })?
1098 .is_dir();
1099 let entry_path =
1100 resolve_file_entry(&labeled.path, "resolve declared artifact cache path entry")?;
1101 if entry_path != canonical {
1102 declared_paths.push((kind, labeled.label.as_str(), entry_path, false));
1103 }
1104 declared_paths.push((kind, labeled.label.as_str(), canonical, is_directory));
1105 }
1106 }
1107 for cargo_inputs in &spec.cargo_build_inputs {
1108 if resolved_cargo_inputs_watch_path(&cargo_inputs.resolved, &cache_root)? {
1109 return invalid_spec(&format!(
1110 "artifact cache root must be outside resolved Cargo build inputs `{}` or inside one of their generated-state exclusions",
1111 cargo_inputs.label
1112 ));
1113 }
1114 }
1115
1116 let mut destinations = BTreeSet::new();
1117 for output in &spec.outputs {
1118 let destination =
1119 resolve_file_entry(&output.destination, "resolve artifact output destination")?;
1120 if destination.starts_with(&cache_root) {
1121 return invalid_spec(&format!(
1122 "output `{}` destination must be outside the artifact cache root",
1123 output.name
1124 ));
1125 }
1126 if destinations.iter().any(|other: &PathBuf| {
1127 destination.starts_with(other) || other.starts_with(&destination)
1128 }) {
1129 return invalid_spec(&format!(
1130 "output `{}` destination overlaps another output",
1131 output.name
1132 ));
1133 }
1134 destinations.insert(destination.clone());
1135 for (kind, label, declared, is_directory) in &declared_paths {
1136 if destination == *declared || (*is_directory && destination.starts_with(declared)) {
1137 return invalid_spec(&format!(
1138 "output `{}` destination overlaps declared {kind} `{label}`",
1139 output.name
1140 ));
1141 }
1142 }
1143 for cargo_inputs in &spec.cargo_build_inputs {
1144 if resolved_cargo_inputs_watch_path(&cargo_inputs.resolved, &destination)? {
1145 return invalid_spec(&format!(
1146 "output `{}` destination must be outside resolved Cargo build inputs `{}` or inside one of their generated-state exclusions",
1147 output.name, cargo_inputs.label
1148 ));
1149 }
1150 }
1151 match fs::symlink_metadata(&destination) {
1155 Ok(metadata)
1156 if metadata.is_dir()
1157 || (metadata.file_type().is_symlink()
1158 && fs::metadata(&destination).is_ok_and(|target| target.is_dir())) =>
1159 {
1160 return invalid_spec(&format!(
1161 "output `{}` destination must not be an existing directory",
1162 output.name
1163 ));
1164 }
1165 Ok(_) => {}
1166 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
1167 Err(source) => {
1168 return Err(ArtifactCacheError::Io {
1169 operation: "inspect artifact output destination",
1170 path: destination,
1171 source,
1172 });
1173 }
1174 }
1175 }
1176 Ok(())
1177}
1178
1179fn resolve_file_entry(path: &Path, operation: &'static str) -> Result<PathBuf, ArtifactCacheError> {
1182 let resolved = match (path.parent(), path.file_name()) {
1183 (Some(parent), Some(name)) => {
1184 canonicalize_allow_missing(parent).map(|parent| parent.join(name))
1185 }
1186 _ => canonicalize_allow_missing(path),
1187 };
1188 resolved.map_err(|source| ArtifactCacheError::Io {
1189 operation,
1190 path: path.to_owned(),
1191 source,
1192 })
1193}
1194
1195fn resolved_cargo_inputs_watch_path(
1196 resolved: &ResolvedCargoBuildInputs,
1197 candidate: &Path,
1198) -> Result<bool, ArtifactCacheError> {
1199 for exclusion in resolved.exclusions() {
1200 let exclusion =
1201 canonicalize_allow_missing(exclusion).map_err(|source| ArtifactCacheError::Io {
1202 operation: "resolve Cargo build input exclusion",
1203 path: exclusion.clone(),
1204 source,
1205 })?;
1206 if candidate.starts_with(exclusion) {
1207 return Ok(false);
1208 }
1209 }
1210
1211 for input in resolved.inputs() {
1212 let path = canonicalize_path(input.path(), "canonicalize resolved Cargo build input")?;
1213 let is_directory = fs::metadata(&path)
1214 .map_err(|source| ArtifactCacheError::Io {
1215 operation: "inspect resolved Cargo build input",
1216 path: path.clone(),
1217 source,
1218 })?
1219 .is_dir();
1220 if candidate == path || (is_directory && candidate.starts_with(path)) {
1221 return Ok(true);
1222 }
1223 }
1224 Ok(false)
1225}
1226
1227fn canonicalize_path(path: &Path, operation: &'static str) -> Result<PathBuf, ArtifactCacheError> {
1228 path.canonicalize()
1229 .map_err(|source| ArtifactCacheError::Io {
1230 operation,
1231 path: path.to_owned(),
1232 source,
1233 })
1234}
1235
1236fn validate_identifier(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1237 if value.is_empty() {
1238 return invalid_spec(&format!("{kind} must not be empty"));
1239 }
1240 if value.len() > 256 {
1241 return invalid_spec(&format!("{kind} must not exceed 256 bytes"));
1242 }
1243 Ok(())
1244}
1245
1246fn validate_label(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1247 if value.is_empty() {
1248 return invalid_spec(&format!("{kind} label must not be empty"));
1249 }
1250 if value.len() > 256 {
1251 return invalid_spec(&format!("{kind} label must not exceed 256 bytes"));
1252 }
1253 Ok(())
1254}
1255
1256fn validate_path_label(kind: &str, value: &str) -> Result<(), ArtifactCacheError> {
1257 validate_label(kind, value)?;
1258 if value
1259 .bytes()
1260 .any(|byte| !(byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'/')))
1261 || value
1262 .split('/')
1263 .any(|component| component.is_empty() || matches!(component, "." | ".."))
1264 {
1265 return invalid_spec(&format!(
1266 "{kind} label `{value}` must be a portable relative logical path"
1267 ));
1268 }
1269 Ok(())
1270}
1271
1272fn validate_output_name(name: &str) -> Result<(), ArtifactCacheError> {
1273 if name.is_empty() || name.len() > 128 {
1274 return invalid_spec("output names must contain 1 to 128 bytes");
1275 }
1276 if !name
1277 .bytes()
1278 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
1279 {
1280 return invalid_spec(&format!(
1281 "output name `{name}` must use only ASCII letters, digits, dot, dash, or underscore"
1282 ));
1283 }
1284 if matches!(name, "." | "..") {
1285 return invalid_spec("output names must not be dot path components");
1286 }
1287 Ok(())
1288}
1289
1290fn invalid_spec<T>(message: &str) -> Result<T, ArtifactCacheError> {
1291 Err(ArtifactCacheError::InvalidSpec {
1292 message: message.to_owned(),
1293 })
1294}
1295
1296fn resolve_key(spec: &ArtifactCacheSpec) -> Result<ResolvedKey, ArtifactCacheError> {
1297 let paths = spec
1298 .inputs
1299 .iter()
1300 .map(|input| (PathBuf::from("input").join(&input.label), &input.path))
1301 .chain(
1302 spec.tools
1303 .iter()
1304 .map(|tool| (PathBuf::from("tool").join(&tool.label), &tool.path)),
1305 );
1306 let declared_input_digest = digest_labeled_paths(
1307 "artifact-set-inputs-v1",
1308 paths,
1309 std::slice::from_ref(&spec.cache_root),
1310 )
1311 .map_err(|source| ArtifactCacheError::Io {
1312 operation: "hash artifact cache inputs",
1313 path: spec.cache_root.clone(),
1314 source,
1315 })?;
1316 let input_digest = if spec.cargo_build_inputs.is_empty() {
1317 declared_input_digest
1318 } else {
1319 let mut inputs = InputHasher::new("artifact-set-inputs-with-cargo-v1");
1320 inputs.field("declared-input-digest", declared_input_digest.as_bytes());
1321 for cargo_input in &spec.cargo_build_inputs {
1322 let current = cargo_input
1323 .resolved
1324 .current_validation_digest()
1325 .map_err(|source| ArtifactCacheError::CargoBuildInputRevalidation {
1326 label: cargo_input.label.clone(),
1327 source,
1328 })?;
1329 let before = cargo_input.resolved.validation_digest();
1330 if current != before {
1331 return Err(ArtifactCacheError::CargoBuildInputsChanged {
1332 label: cargo_input.label.clone(),
1333 before,
1334 after: current,
1335 });
1336 }
1337 inputs.field("cargo-input-label", cargo_input.label.as_bytes());
1338 inputs.field(
1339 "cargo-build-fingerprint",
1340 cargo_input.resolved.fingerprint().as_bytes(),
1341 );
1342 inputs.field(
1343 "cargo-input-digest",
1344 cargo_input.resolved.input_digest().as_bytes(),
1345 );
1346 }
1347 inputs.finish()
1348 };
1349
1350 let mut hasher = InputHasher::new(ARTIFACT_CACHE_FORMAT);
1351 hasher.field("namespace", spec.namespace.as_bytes());
1352 hasher.field("recipe-id", spec.recipe_id.as_bytes());
1353 hasher.field("input-digest", input_digest.as_bytes());
1354 for argument in &spec.arguments {
1355 hasher.field("argument", &os_bytes(argument));
1356 }
1357 for (name, value) in &spec.environment {
1358 hasher.field("environment-name", &os_bytes(name));
1359 match value {
1360 Some(value) => hasher.field("environment-value", &os_bytes(value)),
1361 None => hasher.field("environment-unset", b""),
1362 }
1363 }
1364 for identity in &spec.identities {
1365 hasher.field("identity-label", identity.label.as_bytes());
1366 hasher.field("identity-value", &identity.value);
1367 }
1368 for output in &spec.outputs {
1369 hasher.field("output-name", output.name.as_bytes());
1370 hasher.field(
1371 "output-validation",
1372 output.validation.cache_token().as_bytes(),
1373 );
1374 }
1375 Ok(ResolvedKey {
1376 key: hasher.finish(),
1377 input_digest,
1378 })
1379}
1380
1381fn validated_cache_outputs(
1382 spec: &ArtifactCacheSpec,
1383 key: InputDigest,
1384 entry: &Path,
1385) -> Result<Option<Vec<FileDigest>>, ArtifactCacheError> {
1386 let entry_metadata = match fs::symlink_metadata(entry) {
1387 Ok(metadata) => metadata,
1388 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1389 Err(source) => {
1390 return Err(ArtifactCacheError::Io {
1391 operation: "inspect artifact cache entry",
1392 path: entry.to_owned(),
1393 source,
1394 });
1395 }
1396 };
1397 if !entry_metadata.file_type().is_dir() || !cache_entry_root_is_valid(entry)? {
1398 return Ok(None);
1399 }
1400 let manifest_path = entry.join(MANIFEST_FILE);
1401 let manifest_metadata = match fs::symlink_metadata(&manifest_path) {
1402 Ok(metadata) => metadata,
1403 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1404 Err(source) => {
1405 return Err(ArtifactCacheError::Io {
1406 operation: "inspect artifact cache manifest",
1407 path: manifest_path,
1408 source,
1409 });
1410 }
1411 };
1412 if !manifest_metadata.file_type().is_file() {
1413 return Ok(None);
1414 }
1415 let maximum_len = manifest_contents(
1419 key,
1420 spec,
1421 std::iter::repeat(FileDigest {
1422 bytes: u64::MAX,
1423 digest: key,
1424 }),
1425 )
1426 .len();
1427 let Some(manifest) = read_file_with_limit(&manifest_path, maximum_len).map_err(|source| {
1428 ArtifactCacheError::Io {
1429 operation: "read artifact cache manifest",
1430 path: manifest_path,
1431 source,
1432 }
1433 })?
1434 else {
1435 return Ok(None);
1436 };
1437 if !manifest.starts_with(manifest_header(key).as_bytes()) {
1438 return Ok(None);
1439 }
1440 let Some(output_info) = inspect_cached_output_set(spec, entry)? else {
1441 return Ok(None);
1442 };
1443 Ok(
1444 (manifest == manifest_contents(key, spec, output_info.iter().copied()).as_bytes())
1445 .then_some(output_info),
1446 )
1447}
1448
1449fn inspect_complete_output_set(
1450 spec: &ArtifactCacheSpec,
1451 root: &Path,
1452) -> Result<Vec<FileDigest>, ArtifactCacheError> {
1453 let mut info = Vec::new();
1454 let mut invalid = Vec::new();
1455 let output_directory = root.join("outputs");
1456 if !is_plain_directory(
1457 &output_directory,
1458 "inspect artifact staging output directory",
1459 )? {
1460 invalid.push(("<outputs>".to_owned(), output_directory));
1461 return Err(ArtifactCacheError::InvalidOutputs { outputs: invalid });
1462 }
1463 let outputs = &spec.outputs;
1464 for (index, output) in outputs.iter().enumerate() {
1465 let path = staged_output_path(root, index);
1466 match inspect_artifact(&path, output.validation) {
1467 Ok(Some(artifact)) => info.push(artifact),
1468 Ok(None) => invalid.push((output.name.clone(), path)),
1469 Err(source) => {
1470 return Err(ArtifactCacheError::Io {
1471 operation: "inspect staged artifact output",
1472 path,
1473 source,
1474 });
1475 }
1476 }
1477 }
1478 invalid.extend(
1479 undeclared_output_paths(root, outputs.len())?
1480 .into_iter()
1481 .map(|path| ("<undeclared>".to_owned(), path)),
1482 );
1483 invalid.extend(
1484 undeclared_child_paths(
1485 root,
1486 |name| name == "outputs",
1487 "read artifact staging directory",
1488 )?
1489 .into_iter()
1490 .map(|path| ("<undeclared>".to_owned(), path)),
1491 );
1492 if invalid.is_empty() {
1493 Ok(info)
1494 } else {
1495 Err(ArtifactCacheError::InvalidOutputs { outputs: invalid })
1496 }
1497}
1498
1499fn inspect_cached_output_set(
1500 spec: &ArtifactCacheSpec,
1501 root: &Path,
1502) -> Result<Option<Vec<FileDigest>>, ArtifactCacheError> {
1503 let mut info = Vec::new();
1504 let outputs = &spec.outputs;
1505 for (index, output) in outputs.iter().enumerate() {
1506 let path = staged_output_path(root, index);
1507 match inspect_artifact(&path, output.validation) {
1508 Ok(Some(artifact)) => info.push(artifact),
1509 Ok(None) => return Ok(None),
1510 Err(source) => {
1511 return Err(ArtifactCacheError::Io {
1512 operation: "inspect cached artifact output",
1513 path,
1514 source,
1515 });
1516 }
1517 }
1518 }
1519 if !undeclared_output_paths(root, outputs.len())?.is_empty() {
1520 return Ok(None);
1521 }
1522 Ok(Some(info))
1523}
1524
1525fn undeclared_output_paths(
1526 root: &Path,
1527 output_count: usize,
1528) -> Result<Vec<PathBuf>, ArtifactCacheError> {
1529 let output_directory = root.join("outputs");
1530 undeclared_child_paths(
1531 &output_directory,
1532 |name| {
1533 let Some(name) = name.to_str() else {
1534 return false;
1535 };
1536 name.strip_suffix(OUTPUT_FILE_SUFFIX)
1537 .and_then(|index| index.parse::<usize>().ok())
1538 .is_some_and(|index| index < output_count && name == format_output_index(index))
1539 },
1540 "read artifact output directory",
1541 )
1542}
1543
1544fn undeclared_child_paths(
1545 directory: &Path,
1546 is_expected: impl Fn(&OsStr) -> bool,
1547 operation: &'static str,
1548) -> Result<Vec<PathBuf>, ArtifactCacheError> {
1549 let entries = fs::read_dir(directory).map_err(|source| ArtifactCacheError::Io {
1550 operation,
1551 path: directory.to_owned(),
1552 source,
1553 })?;
1554 let mut undeclared = Vec::new();
1555 for entry in entries {
1556 let entry = entry.map_err(|source| ArtifactCacheError::Io {
1557 operation,
1558 path: directory.to_owned(),
1559 source,
1560 })?;
1561 if !is_expected(&entry.file_name()) {
1562 undeclared.push(entry.path());
1563 }
1564 }
1565 Ok(undeclared)
1566}
1567
1568fn cache_entry_root_is_valid(root: &Path) -> Result<bool, ArtifactCacheError> {
1569 let expected = [
1570 "outputs",
1571 MANIFEST_FILE,
1572 LAST_USED_FILE,
1573 super::cache_fs::RETENTION_LOCK_FILE,
1574 ];
1575 if !undeclared_child_paths(
1576 root,
1577 |name| expected.iter().any(|expected| name == *expected),
1578 "read artifact cache entry",
1579 )?
1580 .is_empty()
1581 {
1582 return Ok(false);
1583 }
1584 if !is_plain_directory(
1585 &root.join("outputs"),
1586 "inspect artifact cache output directory",
1587 )? {
1588 return Ok(false);
1589 }
1590 let last_used = root.join(LAST_USED_FILE);
1591 match fs::symlink_metadata(&last_used) {
1592 Ok(metadata) => Ok(metadata.file_type().is_file()),
1593 Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(true),
1594 Err(source) => Err(ArtifactCacheError::Io {
1595 operation: "inspect artifact cache use marker",
1596 path: last_used,
1597 source,
1598 }),
1599 }
1600}
1601
1602fn is_plain_directory(path: &Path, operation: &'static str) -> Result<bool, ArtifactCacheError> {
1603 match fs::symlink_metadata(path) {
1604 Ok(metadata) => Ok(metadata.file_type().is_dir()),
1605 Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(false),
1606 Err(source) => Err(ArtifactCacheError::Io {
1607 operation,
1608 path: path.to_owned(),
1609 source,
1610 }),
1611 }
1612}
1613
1614fn inspect_artifact(
1615 path: &Path,
1616 validation: ArtifactOutputValidation,
1617) -> io::Result<Option<FileDigest>> {
1618 let metadata = match fs::symlink_metadata(path) {
1619 Ok(metadata) => metadata,
1620 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1621 Err(error) => return Err(error),
1622 };
1623 if !metadata.file_type().is_file() {
1624 return Ok(None);
1625 }
1626 if validation == ArtifactOutputValidation::NonEmptyFile && metadata.len() == 0 {
1627 return Ok(None);
1628 }
1629 digest_file("artifact-set-output-v1", path).map(Some)
1630}
1631
1632fn manifest_header(key: InputDigest) -> String {
1633 format!("{ARTIFACT_CACHE_FORMAT}\nkey:{key}\n")
1634}
1635
1636fn manifest_contents(
1637 key: InputDigest,
1638 spec: &ArtifactCacheSpec,
1639 output_info: impl IntoIterator<Item = FileDigest>,
1640) -> String {
1641 let mut manifest = manifest_header(key);
1642 for ((index, output), info) in spec.outputs.iter().enumerate().zip(output_info) {
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)]
2005mod tests;