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