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