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