Skip to main content

ic_testkit/artifacts/
transaction.rs

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