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