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