Skip to main content

ic_testkit/artifacts/
transaction.rs

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