Skip to main content

mobench_artifacts/
lib.rs

1//! Contained artifact roots and atomic run workspaces.
2//!
3//! [`ApprovedRoot`] validates paths immediately before use. It prevents
4//! traversal through pre-existing symlinks, but it does not yet provide the
5//! descriptor-relative operations needed to close hostile concurrent-swap
6//! races on every supported platform.
7
8use std::collections::{BTreeMap, BTreeSet};
9use std::fmt;
10use std::fs::{self, File, OpenOptions};
11use std::io::{BufReader, Read, Write};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicU64, Ordering};
14use std::sync::{Arc, Mutex, OnceLock};
15use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
16
17use fs2::FileExt;
18use semver::Version;
19use serde::{Deserialize, Serialize};
20use sha2::{Digest, Sha256};
21use thiserror::Error;
22
23/// A validated artifact identifier that is safe to use as one path component.
24#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
25pub struct ArtifactId(String);
26
27impl ArtifactId {
28    /// Validate and retain a single artifact path component.
29    pub fn new(value: impl Into<String>) -> Result<Self, ArtifactPathError> {
30        let value = value.into();
31        let path = Path::new(&value);
32        let is_one_normal_component = {
33            let mut components = path.components();
34            matches!(components.next(), Some(std::path::Component::Normal(_)))
35                && components.next().is_none()
36        };
37        if value.is_empty()
38            || value.contains('/')
39            || value.contains('\\')
40            || path.is_absolute()
41            || !is_one_normal_component
42        {
43            return Err(ArtifactPathError::InvalidArtifactId { value });
44        }
45        Ok(Self(value))
46    }
47
48    /// Return the validated component.
49    pub fn as_str(&self) -> &str {
50        &self.0
51    }
52}
53
54impl fmt::Display for ArtifactId {
55    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
56        formatter.write_str(self.as_str())
57    }
58}
59
60const STAGING_DIRECTORY: &str = ".mobench-staging";
61const STAGING_QUARANTINE_DIRECTORY: &str = ".mobench-staging-quarantine";
62const STAGING_LOCK_FILE: &str = ".mobench-staging.lock";
63const WORKSPACE_LOCK_FILE: &str = ".mobench-workspace.lock";
64const RUN_READER_LOCK_FILE: &str = ".mobench-run-reader.lock";
65const RUN_RETENTION_QUARANTINE_DIRECTORY: &str = ".mobench-run-quarantine";
66const LATEST_DIRECTORY: &str = ".mobench-latest";
67const LATEST_STAGING_DIRECTORY: &str = "staging";
68const LATEST_GENERATIONS_DIRECTORY: &str = "generations";
69const LATEST_QUARANTINE_DIRECTORY: &str = "quarantine";
70const LATEST_CURRENT_FILE: &str = "current";
71const LATEST_LOCK_FILE: &str = ".mobench-latest.lock";
72const LATEST_READER_LOCK_FILE: &str = ".mobench-reader.lock";
73const LATEST_RETENTION_BOUNDARY_FILE: &str = ".mobench-retention-boundary.json";
74const RETAIN_LATEST_GENERATIONS: usize = 8;
75const RETAIN_PUBLISHED_RUNS: usize = 32;
76const RETAIN_QUARANTINE_ENTRIES: usize = 16;
77/// File name of the versioned manifest stored in every published run.
78pub const RUN_MANIFEST_FILE: &str = "mobench-run-manifest.json";
79/// File name of the versioned manifest stored in every latest generation.
80pub const LATEST_MANIFEST_FILE: &str = "manifest.json";
81const RUN_MANIFEST_VERSION: u32 = 1;
82const LATEST_MANIFEST_VERSION: u32 = 1;
83const LATEST_LOCK_TIMEOUT: Duration = Duration::from_secs(5);
84const LATEST_LOCK_POLL_INTERVAL: Duration = Duration::from_millis(5);
85const MAX_ALLOCATION_ATTEMPTS: usize = 1_024;
86static WORKSPACE_SEQUENCE: AtomicU64 = AtomicU64::new(0);
87static ACTIVE_WORKSPACES: OnceLock<Mutex<BTreeSet<PathBuf>>> = OnceLock::new();
88
89/// A unique private directory for building one complete run before publication.
90#[derive(Debug)]
91pub struct RunWorkspace {
92    root: ApprovedRoot,
93    logical_id: ArtifactId,
94    expected_latest_generation: Option<String>,
95    workspace_lock: Option<File>,
96    staging_path: Option<PathBuf>,
97    published_path: PathBuf,
98}
99
100/// A completed run directory that has been atomically renamed into visibility.
101#[derive(Debug, Clone)]
102pub struct PublishedRun {
103    root: ApprovedRoot,
104    path: PathBuf,
105    _reader_lease: Arc<File>,
106}
107
108/// One regular file recorded in a durable artifact manifest.
109#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
110pub struct ManifestArtifact {
111    /// Slash-separated path relative to the manifest's directory.
112    pub relative_path: String,
113    /// Exact file length in bytes.
114    pub size: u64,
115    /// Lower-case hexadecimal SHA-256 digest.
116    pub sha256: String,
117}
118
119/// Durable identity and integrity record for one published run.
120#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121pub struct RunManifest {
122    /// Manifest format version. Unknown versions are rejected.
123    pub format_version: u32,
124    /// Version of the mobench-artifacts producer.
125    pub producer_version: String,
126    /// Caller-provided logical run identity.
127    pub logical_id: String,
128    /// Collision-resistant published directory identity.
129    pub publication_id: String,
130    /// Latest generation observed when this run workspace was allocated.
131    /// Refresh uses this as a compare-and-swap token.
132    pub expected_latest_generation: Option<String>,
133    /// Every regular run artifact except this manifest, sorted by path.
134    pub artifacts: Vec<ManifestArtifact>,
135}
136
137/// Mapping from a published run artifact into one latest-generation alias.
138#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
139pub struct LatestManifestArtifact {
140    /// Source path relative to the published run.
141    pub source_relative_path: String,
142    /// Stable public alias and path relative to the generation directory.
143    pub destination_relative_path: String,
144    /// Exact file length in bytes.
145    pub size: u64,
146    /// Lower-case hexadecimal SHA-256 digest.
147    pub sha256: String,
148}
149
150/// Durable identity and integrity record for one immutable latest generation.
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub struct LatestManifest {
153    /// Manifest format version. Unknown versions are rejected.
154    pub format_version: u32,
155    /// Version of the mobench-artifacts producer.
156    pub producer_version: String,
157    /// Identity named by the atomic `current` pointer.
158    pub generation: String,
159    /// Previously committed generation, when one existed.
160    pub predecessor_generation: Option<String>,
161    /// Logical identity of the source run.
162    pub source_logical_id: String,
163    /// Publication identity of the source run.
164    pub source_publication_id: String,
165    /// Aliases in this generation, sorted by destination path.
166    pub artifacts: Vec<LatestManifestArtifact>,
167}
168
169#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
170struct RetentionBoundary {
171    format_version: u32,
172    predecessor_generation: String,
173}
174
175/// A reader pinned to one immutable latest generation.
176///
177/// Opening resolves the atomic pointer once and verifies the generation
178/// manifest and every listed artifact. Subsequent reads never re-resolve
179/// `current`, so a concurrent writer cannot mix generations for this reader.
180#[derive(Debug, Clone)]
181pub struct LatestSnapshot {
182    root: ApprovedRoot,
183    generation_path: PathBuf,
184    manifest: LatestManifest,
185    _reader_lease: Arc<File>,
186    _source_run_lease: Arc<File>,
187    #[cfg(unix)]
188    generation_directory: Arc<File>,
189}
190
191#[derive(Debug)]
192struct LatestUpdateLock {
193    file: File,
194}
195
196#[derive(Debug)]
197struct StableAliasTransaction {
198    path: PathBuf,
199    root_path: PathBuf,
200    installed: Vec<PathBuf>,
201    backups: Vec<(PathBuf, PathBuf)>,
202}
203
204#[derive(Debug)]
205enum PointerCommitFailure {
206    BeforeCommit(ArtifactPathError),
207    AfterCommit(ArtifactPathError),
208}
209
210impl StableAliasTransaction {
211    fn rollback(self, original: ArtifactPathError) -> ArtifactPathError {
212        let error =
213            latest_update_failure(original, &self.installed, &self.backups, self.path.clone());
214        let rollback_complete = !matches!(&error, ArtifactPathError::LatestRollback { .. });
215        if rollback_complete {
216            let _ = fs::remove_dir_all(&self.path);
217        }
218        let _ = sync_directory(&self.root_path);
219        error
220    }
221
222    fn finish(self) {
223        let _ = fs::remove_dir_all(self.path);
224    }
225}
226
227impl Drop for LatestUpdateLock {
228    fn drop(&mut self) {
229        let _ = FileExt::unlock(&self.file);
230    }
231}
232
233/// One published artifact to copy to a stable convenience path.
234#[derive(Debug, Clone, PartialEq, Eq)]
235pub struct LatestArtifact {
236    source: ArtifactId,
237    destination: ArtifactId,
238}
239
240impl LatestArtifact {
241    /// Refresh a stable path with a published artifact of the same name.
242    pub fn same(id: ArtifactId) -> Self {
243        Self {
244            source: id.clone(),
245            destination: id,
246        }
247    }
248
249    /// Refresh a stable destination from a differently named published file.
250    pub fn new(source: ArtifactId, destination: ArtifactId) -> Self {
251        Self {
252            source,
253            destination,
254        }
255    }
256}
257
258impl RunWorkspace {
259    /// Allocate a unique staging directory and reserve a collision-resistant
260    /// publication path derived from the validated logical run ID.
261    pub fn allocate(
262        root: impl AsRef<Path>,
263        logical_id: &ArtifactId,
264    ) -> Result<Self, ArtifactPathError> {
265        let root_path = root.as_ref();
266        fs::create_dir_all(root_path).map_err(|source| ArtifactPathError::Io {
267            operation: "create artifact root",
268            path: root_path.to_path_buf(),
269            source,
270        })?;
271        let root = ApprovedRoot::existing(root_path)?;
272        let expected_latest_generation =
273            recover_latest_for_writer(&root)?.map(|snapshot| snapshot.manifest.generation);
274        let _staging_manager_lock = acquire_staging_manager_lock(&root)?;
275        let staging_root = root.prepare_dir(STAGING_DIRECTORY)?;
276        quarantine_abandoned_run_staging(&root, &staging_root)?;
277
278        for _ in 0..MAX_ALLOCATION_ATTEMPTS {
279            let nonce = workspace_nonce();
280            let staging_path = staging_root.join(&nonce);
281            let published_path = root.path().join(format!("{logical_id}--{nonce}"));
282            if fs::symlink_metadata(&published_path).is_ok() {
283                continue;
284            }
285
286            match fs::create_dir(&staging_path) {
287                Ok(()) => {
288                    if fs::symlink_metadata(&published_path).is_ok() {
289                        let _ = fs::remove_dir(&staging_path);
290                        continue;
291                    }
292                    let workspace_lock_path = staging_path.join(WORKSPACE_LOCK_FILE);
293                    let workspace_lock = match OpenOptions::new()
294                        .read(true)
295                        .write(true)
296                        .create_new(true)
297                        .open(&workspace_lock_path)
298                    {
299                        Ok(file) => file,
300                        Err(source) => {
301                            let _ = fs::remove_dir_all(&staging_path);
302                            return Err(ArtifactPathError::Io {
303                                operation: "create run-workspace lock",
304                                path: workspace_lock_path,
305                                source,
306                            });
307                        }
308                    };
309                    if let Err(source) = FileExt::lock_exclusive(&workspace_lock) {
310                        let _ = fs::remove_dir_all(&staging_path);
311                        return Err(ArtifactPathError::Io {
312                            operation: "lock run workspace",
313                            path: workspace_lock_path,
314                            source,
315                        });
316                    }
317                    register_active_workspace(&staging_path);
318                    return Ok(Self {
319                        root,
320                        logical_id: logical_id.clone(),
321                        expected_latest_generation,
322                        workspace_lock: Some(workspace_lock),
323                        staging_path: Some(staging_path),
324                        published_path,
325                    });
326                }
327                Err(source) if source.kind() == std::io::ErrorKind::AlreadyExists => continue,
328                Err(source) => {
329                    return Err(ArtifactPathError::Io {
330                        operation: "create run staging directory",
331                        path: staging_path,
332                        source,
333                    });
334                }
335            }
336        }
337
338        Err(ArtifactPathError::WorkspaceAllocationExhausted {
339            logical_id: logical_id.clone(),
340        })
341    }
342
343    /// Return the private directory where run outputs must be written.
344    pub fn staging_path(&self) -> &Path {
345        self.staging_path
346            .as_deref()
347            .expect("unpublished workspace retains its staging path")
348    }
349
350    /// Return the unique path where the completed run will be published.
351    pub fn published_path(&self) -> &Path {
352        &self.published_path
353    }
354
355    /// Return the approved output root shared by published and latest artifacts.
356    pub fn root(&self) -> &ApprovedRoot {
357        &self.root
358    }
359
360    /// Publish the complete staging directory without replacing any existing
361    /// file, directory, or symbolic link at the destination.
362    pub fn publish(
363        mut self,
364        required_files: &[ArtifactId],
365    ) -> Result<PublishedRun, ArtifactPathError> {
366        let staging_path = self
367            .staging_path
368            .as_ref()
369            .expect("unpublished workspace retains its staging path")
370            .clone();
371        let published_path = self.published_path.clone();
372        match fs::symlink_metadata(&staging_path) {
373            Ok(metadata) if metadata.file_type().is_symlink() => {
374                return Err(ArtifactPathError::SymlinkComponent { path: staging_path });
375            }
376            Ok(metadata) if !metadata.is_dir() => {
377                return Err(ArtifactPathError::DirectoryComponentNotDirectory {
378                    path: staging_path,
379                });
380            }
381            Ok(_) => {}
382            Err(source) => {
383                return Err(ArtifactPathError::Io {
384                    operation: "inspect run staging directory",
385                    path: staging_path,
386                    source,
387                });
388            }
389        }
390        let staging_relative = staging_path.strip_prefix(self.root.path()).map_err(|_| {
391            ArtifactPathError::InvalidRelativePath {
392                path: staging_path.clone(),
393            }
394        })?;
395        self.root.prepare_dir(staging_relative)?;
396
397        let staging_root = ApprovedRoot::existing(&staging_path)?;
398        for required_file in required_files {
399            let required_path = staging_root.prepare_file(required_file.as_str())?;
400            match fs::symlink_metadata(&required_path) {
401                Ok(metadata) if metadata.file_type().is_symlink() => {
402                    return Err(ArtifactPathError::SymlinkComponent {
403                        path: required_path,
404                    });
405                }
406                Ok(metadata) if !metadata.is_file() => {
407                    return Err(ArtifactPathError::FileDestinationNotFile {
408                        path: required_path,
409                    });
410                }
411                Ok(_) => {}
412                Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
413                    return Err(ArtifactPathError::RequiredArtifactMissing {
414                        path: required_path,
415                    });
416                }
417                Err(source) => {
418                    return Err(ArtifactPathError::Io {
419                        operation: "inspect required staged artifact",
420                        path: required_path,
421                        source,
422                    });
423                }
424            }
425        }
426
427        let manifest_path = staging_path.join(RUN_MANIFEST_FILE);
428        match fs::symlink_metadata(&manifest_path) {
429            Ok(_) => {
430                return Err(ArtifactPathError::ReservedManifestPath {
431                    path: manifest_path,
432                });
433            }
434            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
435            Err(source) => {
436                return Err(ArtifactPathError::Io {
437                    operation: "inspect reserved run manifest path",
438                    path: manifest_path,
439                    source,
440                });
441            }
442        }
443
444        create_run_reader_lease_file(&staging_path)?;
445        let artifacts =
446            inspect_artifact_tree(&staging_path, &[WORKSPACE_LOCK_FILE, RUN_READER_LOCK_FILE])?;
447        let publication_id = published_path
448            .file_name()
449            .and_then(|name| name.to_str())
450            .ok_or_else(|| ArtifactPathError::NonUtf8ArtifactPath {
451                path: published_path.clone(),
452            })?
453            .to_owned();
454        let manifest = RunManifest {
455            format_version: RUN_MANIFEST_VERSION,
456            producer_version: env!("CARGO_PKG_VERSION").to_owned(),
457            logical_id: self.logical_id.as_str().to_owned(),
458            publication_id,
459            expected_latest_generation: self.expected_latest_generation.clone(),
460            artifacts,
461        };
462        write_json_file(&manifest_path, &manifest, "write run manifest")?;
463        sync_artifact_tree(&staging_path)?;
464
465        match fs::symlink_metadata(&published_path) {
466            Ok(metadata) if metadata.file_type().is_symlink() => {
467                return Err(ArtifactPathError::SymlinkComponent {
468                    path: published_path,
469                });
470            }
471            Ok(_) => {
472                return Err(ArtifactPathError::PublicationDestinationExists {
473                    path: published_path,
474                });
475            }
476            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
477            Err(source) => {
478                return Err(ArtifactPathError::Io {
479                    operation: "inspect run publication destination",
480                    path: published_path,
481                    source,
482                });
483            }
484        }
485
486        let _staging_manager_lock = acquire_staging_manager_lock(&self.root)?;
487        let workspace_lock_path = staging_path.join(WORKSPACE_LOCK_FILE);
488        let workspace_lock = self.workspace_lock.take();
489        if let Some(workspace_lock) = workspace_lock.as_ref() {
490            let _ = FileExt::unlock(workspace_lock);
491        }
492        // Close before unlinking so publication also works on platforms that
493        // do not allow deleting an open file. The staging-manager lock keeps
494        // recovery from racing this short handoff window.
495        drop(workspace_lock);
496        fs::remove_file(&workspace_lock_path).map_err(|source| ArtifactPathError::Io {
497            operation: "remove run-workspace lock before publication",
498            path: workspace_lock_path,
499            source,
500        })?;
501        sync_directory(&staging_path)?;
502
503        rename_directory_noreplace(&staging_path, &published_path).map_err(|source| {
504            ArtifactPathError::Io {
505                operation: "publish completed run",
506                path: published_path.clone(),
507                source,
508            }
509        })?;
510        record_durability_event("publish_run");
511        let published_sync = sync_directory(
512            published_path
513                .parent()
514                .expect("published run always has an artifact-root parent"),
515        );
516        let staging_sync = sync_directory(
517            staging_path
518                .parent()
519                .expect("staging run always has a staging-root parent"),
520        );
521        unregister_active_workspace(&staging_path);
522        self.staging_path = None;
523
524        if let Err(source) = published_sync.and(staging_sync) {
525            return Err(ArtifactPathError::PublicationDurabilityUncertain {
526                path: published_path,
527                source: Box::new(source),
528            });
529        }
530
531        let reader_lease = open_shared_lease(
532            &published_path.join(RUN_READER_LOCK_FILE),
533            "open published-run reader lease",
534            "lock published-run reader lease",
535        )
536        .map_err(
537            |source| ArtifactPathError::PublicationPostCommitMaintenance {
538                path: published_path.clone(),
539                source: Box::new(source),
540            },
541        )?;
542        let published = PublishedRun {
543            root: self.root.clone(),
544            path: published_path,
545            _reader_lease: reader_lease,
546        };
547        prune_published_runs(&self.root).map_err(|source| {
548            ArtifactPathError::PublicationPostCommitMaintenance {
549                path: published.path.clone(),
550                source: Box::new(source),
551            }
552        })?;
553        Ok(published)
554    }
555}
556
557impl Drop for RunWorkspace {
558    fn drop(&mut self) {
559        if let Some(staging_path) = self.staging_path.as_ref() {
560            if let Ok(_manager_lock) = acquire_staging_manager_lock(&self.root) {
561                let _ = fs::remove_dir_all(staging_path);
562            }
563            unregister_active_workspace(staging_path);
564        }
565    }
566}
567
568impl PublishedRun {
569    /// Return the unique published run directory.
570    pub fn path(&self) -> &Path {
571        &self.path
572    }
573
574    /// Return the approved output root containing this run.
575    pub fn root(&self) -> &ApprovedRoot {
576        &self.root
577    }
578
579    /// Read and verify this run's versioned manifest and artifact digests.
580    pub fn manifest(&self) -> Result<RunManifest, ArtifactPathError> {
581        validate_run_manifest(&self.path)
582    }
583
584    /// Commit an immutable latest generation and refresh legacy stable copies.
585    ///
586    /// The generation becomes authoritative through one atomic `current`
587    /// pointer rename. Root-level stable copies are retained for compatibility;
588    /// because that legacy file set cannot be swapped atomically, the next
589    /// refresh first repairs it from the already committed generation.
590    pub fn refresh_latest(&self, artifacts: &[LatestArtifact]) -> Result<(), ArtifactPathError> {
591        let published_root = ApprovedRoot::existing(&self.path)?;
592        let mut prepared: Vec<(PathBuf, PathBuf, ArtifactId, ArtifactId)> =
593            Vec::with_capacity(artifacts.len());
594        for artifact in artifacts {
595            if prepared
596                .iter()
597                .any(|(_, _, _, destination)| destination == &artifact.destination)
598            {
599                return Err(ArtifactPathError::DuplicateLatestDestination {
600                    id: artifact.destination.clone(),
601                });
602            }
603            let source = published_root.prepare_file(artifact.source.as_str())?;
604            match fs::symlink_metadata(&source) {
605                Ok(metadata) if metadata.file_type().is_symlink() => {
606                    return Err(ArtifactPathError::SymlinkComponent { path: source });
607                }
608                Ok(metadata) if !metadata.is_file() => {
609                    return Err(ArtifactPathError::FileDestinationNotFile { path: source });
610                }
611                Ok(_) => {}
612                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
613                    return Err(ArtifactPathError::RequiredArtifactMissing { path: source });
614                }
615                Err(source_error) => {
616                    return Err(ArtifactPathError::Io {
617                        operation: "inspect published latest source",
618                        path: source,
619                        source: source_error,
620                    });
621                }
622            }
623            let destination = self.root.prepare_file(artifact.destination.as_str())?;
624            prepared.push((
625                source,
626                destination,
627                artifact.source.clone(),
628                artifact.destination.clone(),
629            ));
630        }
631        let run_manifest = self.manifest()?;
632        maybe_mutate_latest_source_after_validation();
633
634        let latest_root_path = self.root.prepare_dir(LATEST_DIRECTORY)?;
635        let latest_root = ApprovedRoot::existing(&latest_root_path)?;
636        latest_root.prepare_dir(LATEST_STAGING_DIRECTORY)?;
637        latest_root.prepare_dir(LATEST_GENERATIONS_DIRECTORY)?;
638        latest_root.prepare_dir(LATEST_QUARANTINE_DIRECTORY)?;
639        // Allocation and lease acquisition are one short serialized operation.
640        // Without this fence, another writer could observe the directory in the
641        // interval before its lease exists and quarantine an active staging tree.
642        let (generation, staging_path, staging_lease) = {
643            let _latest_update_lock = acquire_latest_update_lock(&self.root)?;
644            let (generation, staging_path) = allocate_latest_generation(&latest_root)?;
645            let lease_result = (|| {
646                create_reader_lease_file(&staging_path)?;
647                let staging_lease_path = staging_path.join(LATEST_READER_LOCK_FILE);
648                let staging_lease = OpenOptions::new()
649                    .read(true)
650                    .write(true)
651                    .open(&staging_lease_path)
652                    .map_err(|source| ArtifactPathError::Io {
653                        operation: "open active latest-generation staging lease",
654                        path: staging_lease_path.clone(),
655                        source,
656                    })?;
657                FileExt::lock_exclusive(&staging_lease).map_err(|source| {
658                    ArtifactPathError::Io {
659                        operation: "lock active latest-generation staging lease",
660                        path: staging_lease_path,
661                        source,
662                    }
663                })?;
664                Ok(staging_lease)
665            })();
666            match lease_result {
667                Ok(staging_lease) => (generation, staging_path, staging_lease),
668                Err(error) => {
669                    let _ = quarantine_latest_path(&latest_root, &staging_path);
670                    return Err(error);
671                }
672            }
673        };
674        let mut staging_lease = Some(staging_lease);
675
676        // Copy, hash, and sync potentially large artifacts before entering the
677        // serialized commit section. The staging lease prevents another
678        // writer's recovery pass from quarantining this active generation.
679        let preparation_result = (|| {
680            let mut manifest_artifacts = Vec::with_capacity(prepared.len());
681            for (source, _, source_id, destination_id) in &prepared {
682                let staged = staging_path.join(destination_id.as_str());
683                copy_new_synced(source, &staged, "stage latest generation artifact")?;
684                let (size, sha256) = digest_file(&staged)?;
685                let expected = run_manifest
686                    .artifacts
687                    .iter()
688                    .find(|artifact| artifact.relative_path == source_id.as_str())
689                    .ok_or_else(|| ArtifactPathError::ManifestIntegrity {
690                        path: self.path.join(RUN_MANIFEST_FILE),
691                        detail: format!(
692                            "latest source `{source_id}` is not recorded in the run manifest"
693                        ),
694                    })?;
695                if expected.size != size || expected.sha256 != sha256 {
696                    return Err(ArtifactPathError::ManifestIntegrity {
697                        path: source.clone(),
698                        detail: format!(
699                            "latest source `{source_id}` changed after run-manifest validation"
700                        ),
701                    });
702                }
703                manifest_artifacts.push(LatestManifestArtifact {
704                    source_relative_path: source_id.as_str().to_owned(),
705                    destination_relative_path: destination_id.as_str().to_owned(),
706                    size,
707                    sha256,
708                });
709            }
710            manifest_artifacts.sort_by(|left, right| {
711                left.destination_relative_path
712                    .cmp(&right.destination_relative_path)
713            });
714            sync_artifact_tree(&staging_path)?;
715            Ok(manifest_artifacts)
716        })();
717        let manifest_artifacts = match preparation_result {
718            Ok(manifest_artifacts) => manifest_artifacts,
719            Err(error) => {
720                if let Some(lease) = staging_lease.take() {
721                    let _ = FileExt::unlock(&lease);
722                    drop(lease);
723                }
724                if let Err(cleanup) = quarantine_latest_path(&latest_root, &staging_path) {
725                    return Err(ArtifactPathError::LatestStagingCleanup {
726                        original: Box::new(error),
727                        cleanup: Box::new(cleanup),
728                        staging_path,
729                    });
730                }
731                return Err(error);
732            }
733        };
734
735        let generation_path = latest_root
736            .path()
737            .join(LATEST_GENERATIONS_DIRECTORY)
738            .join(&generation);
739        let mut generation_committed = false;
740        let generation_result = (|| {
741            let _latest_update_lock = acquire_latest_update_lock(&self.root)?;
742            let predecessor = recover_latest_state_locked(&self.root, &latest_root)?;
743            let observed_generation = predecessor
744                .as_ref()
745                .map(|snapshot| snapshot.manifest.generation.clone());
746            if run_manifest.expected_latest_generation != observed_generation {
747                return Err(ArtifactPathError::StaleLatestGeneration {
748                    expected: run_manifest.expected_latest_generation.clone(),
749                    observed: observed_generation,
750                });
751            }
752
753            let manifest = LatestManifest {
754                format_version: LATEST_MANIFEST_VERSION,
755                producer_version: env!("CARGO_PKG_VERSION").to_owned(),
756                generation: generation.clone(),
757                predecessor_generation: predecessor
758                    .as_ref()
759                    .map(|snapshot| snapshot.manifest.generation.clone()),
760                source_logical_id: run_manifest.logical_id.clone(),
761                source_publication_id: run_manifest.publication_id.clone(),
762                artifacts: manifest_artifacts.clone(),
763            };
764            write_json_file(
765                &staging_path.join(LATEST_MANIFEST_FILE),
766                &manifest,
767                "write latest-generation manifest",
768            )?;
769            sync_directory(&staging_path)?;
770
771            fs::rename(&staging_path, &generation_path).map_err(|source| {
772                ArtifactPathError::Io {
773                    operation: "publish latest generation",
774                    path: generation_path.clone(),
775                    source,
776                }
777            })?;
778            record_durability_event("publish_generation");
779            sync_directory(&latest_root.path().join(LATEST_STAGING_DIRECTORY))?;
780            sync_directory(&latest_root.path().join(LATEST_GENERATIONS_DIRECTORY))?;
781            if let Some(lease) = staging_lease.take() {
782                let _ = FileExt::unlock(&lease);
783                drop(lease);
784            }
785
786            let generation_manifest = validate_latest_manifest(&generation_path, &generation)?;
787            let candidate = make_latest_snapshot(
788                self.root.clone(),
789                generation_path.clone(),
790                generation_manifest,
791            )?;
792            let alias_transaction =
793                install_stable_aliases_transactional(&self.root, &latest_root, &candidate)?;
794            match commit_current_pointer(&latest_root, &generation) {
795                Ok(()) => {
796                    generation_committed = true;
797                    alias_transaction.finish();
798                }
799                Err(PointerCommitFailure::BeforeCommit(pointer_error)) => {
800                    return Err(alias_transaction.rollback(pointer_error));
801                }
802                Err(PointerCommitFailure::AfterCommit(durability_error)) => {
803                    generation_committed = true;
804                    alias_transaction.finish();
805                    return Err(ArtifactPathError::LatestCommitDurabilityUncertain {
806                        generation: generation.clone(),
807                        source: Box::new(durability_error),
808                    });
809                }
810            }
811            prune_latest_generations(&latest_root, candidate.manifest()).map_err(|source| {
812                ArtifactPathError::LatestRetentionAfterCommit {
813                    generation: generation.clone(),
814                    source: Box::new(source),
815                }
816            })?;
817            Ok(())
818        })();
819
820        if generation_result.is_err() {
821            if let Some(lease) = staging_lease.take() {
822                let _ = FileExt::unlock(&lease);
823                drop(lease);
824            }
825            if staging_path.exists() {
826                quarantine_latest_path(&latest_root, &staging_path)?;
827            } else if generation_path.exists() && !generation_committed {
828                quarantine_latest_path(&latest_root, &generation_path)?;
829            }
830        }
831        generation_result
832    }
833}
834
835impl LatestSnapshot {
836    /// Resolve and verify the committed generation without writing or locking.
837    pub fn open(root: impl AsRef<Path>) -> Result<Self, ArtifactPathError> {
838        let root = ApprovedRoot::existing(root)?;
839        load_current_snapshot(&root, false)?.ok_or_else(|| {
840            ArtifactPathError::LatestSnapshotUnavailable {
841                path: root.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE),
842            }
843        })
844    }
845
846    /// Repair legacy root-level aliases from the committed generation.
847    ///
848    /// This is idempotent and is also run automatically before every writer
849    /// creates a new generation. The returned reader remains pinned to the
850    /// generation used for recovery.
851    pub fn recover_stable_aliases(root: impl AsRef<Path>) -> Result<Self, ArtifactPathError> {
852        let root = ApprovedRoot::existing(root)?;
853        let snapshot = recover_latest_for_writer(&root)?.ok_or_else(|| {
854            ArtifactPathError::LatestSnapshotUnavailable {
855                path: root.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE),
856            }
857        })?;
858        Ok(snapshot)
859    }
860
861    /// Return the pinned generation identity.
862    pub fn generation(&self) -> &str {
863        &self.manifest.generation
864    }
865
866    /// Return the immutable generation directory.
867    pub fn path(&self) -> &Path {
868        &self.generation_path
869    }
870
871    /// Return the approved artifact root containing this snapshot.
872    pub fn root(&self) -> &ApprovedRoot {
873        &self.root
874    }
875
876    /// Return the verified versioned manifest.
877    pub fn manifest(&self) -> &LatestManifest {
878        &self.manifest
879    }
880
881    /// Open one alias listed by this pinned generation.
882    ///
883    /// Callers that need integrity verification should prefer
884    /// [`LatestSnapshot::read_artifact`], which verifies the exact bytes read
885    /// against the pinned manifest.
886    pub fn open_artifact(&self, id: &ArtifactId) -> Result<File, ArtifactPathError> {
887        if !self
888            .manifest
889            .artifacts
890            .iter()
891            .any(|artifact| artifact.destination_relative_path == id.as_str())
892        {
893            return Err(ArtifactPathError::LatestArtifactNotInSnapshot { id: id.clone() });
894        }
895        let path = self.generation_path.join(id.as_str());
896        open_snapshot_artifact(self, id).map_err(|source| ArtifactPathError::Io {
897            operation: "open latest snapshot artifact",
898            path,
899            source,
900        })
901    }
902
903    /// Read one alias from this pinned generation.
904    pub fn read_artifact(&self, id: &ArtifactId) -> Result<Vec<u8>, ArtifactPathError> {
905        let expected = self
906            .manifest
907            .artifacts
908            .iter()
909            .find(|artifact| artifact.destination_relative_path == id.as_str())
910            .ok_or_else(|| ArtifactPathError::LatestArtifactNotInSnapshot { id: id.clone() })?;
911        let mut file = self.open_artifact(id)?;
912        let mut contents = Vec::new();
913        file.read_to_end(&mut contents)
914            .map_err(|source| ArtifactPathError::Io {
915                operation: "read latest snapshot artifact",
916                path: self.generation_path.join(id.as_str()),
917                source,
918            })?;
919        let actual_digest = sha256_bytes(&contents);
920        if contents.len() as u64 != expected.size || actual_digest != expected.sha256 {
921            return Err(ArtifactPathError::CorruptGeneration {
922                path: self.generation_path.join(id.as_str()),
923                detail: "artifact bytes changed after snapshot validation".to_owned(),
924            });
925        }
926        Ok(contents)
927    }
928}
929
930fn acquire_latest_update_lock(root: &ApprovedRoot) -> Result<LatestUpdateLock, ArtifactPathError> {
931    acquire_named_lock(
932        root,
933        LATEST_LOCK_FILE,
934        "open latest-artifact lock",
935        "lock latest artifacts",
936    )
937}
938
939fn acquire_staging_manager_lock(
940    root: &ApprovedRoot,
941) -> Result<LatestUpdateLock, ArtifactPathError> {
942    acquire_named_lock(
943        root,
944        STAGING_LOCK_FILE,
945        "open run-staging lock",
946        "lock run staging",
947    )
948}
949
950fn acquire_named_lock(
951    root: &ApprovedRoot,
952    name: &str,
953    open_operation: &'static str,
954    lock_operation: &'static str,
955) -> Result<LatestUpdateLock, ArtifactPathError> {
956    let path = root.prepare_file(name)?;
957    let file = OpenOptions::new()
958        .read(true)
959        .write(true)
960        .create(true)
961        .truncate(false)
962        .open(&path)
963        .map_err(|source| ArtifactPathError::Io {
964            operation: open_operation,
965            path: path.clone(),
966            source,
967        })?;
968    let started = Instant::now();
969
970    loop {
971        match FileExt::try_lock_exclusive(&file) {
972            Ok(()) => return Ok(LatestUpdateLock { file }),
973            Err(source)
974                if source.kind() == std::io::ErrorKind::WouldBlock
975                    && started.elapsed() < LATEST_LOCK_TIMEOUT =>
976            {
977                std::thread::sleep(LATEST_LOCK_POLL_INTERVAL);
978            }
979            Err(source) if source.kind() == std::io::ErrorKind::WouldBlock => {
980                return Err(ArtifactPathError::LatestLockTimeout { path });
981            }
982            Err(source) => {
983                return Err(ArtifactPathError::Io {
984                    operation: lock_operation,
985                    path,
986                    source,
987                });
988            }
989        }
990    }
991}
992
993fn quarantine_abandoned_run_staging(
994    root: &ApprovedRoot,
995    staging_root: &Path,
996) -> Result<(), ArtifactPathError> {
997    let quarantine_root = root.prepare_dir(STAGING_QUARANTINE_DIRECTORY)?;
998    let entries = fs::read_dir(staging_root).map_err(|source| ArtifactPathError::Io {
999        operation: "list run staging for recovery",
1000        path: staging_root.to_path_buf(),
1001        source,
1002    })?;
1003    for entry in entries {
1004        let entry = entry.map_err(|source| ArtifactPathError::Io {
1005            operation: "inspect run staging entry",
1006            path: staging_root.to_path_buf(),
1007            source,
1008        })?;
1009        let staging_path = entry.path();
1010        if is_active_workspace(&staging_path) {
1011            continue;
1012        }
1013        let metadata =
1014            fs::symlink_metadata(&staging_path).map_err(|source| ArtifactPathError::Io {
1015                operation: "inspect run staging workspace",
1016                path: staging_path.clone(),
1017                source,
1018            })?;
1019        if metadata.file_type().is_symlink() || !metadata.is_dir() {
1020            quarantine_run_staging_path(staging_root, &quarantine_root, &staging_path)?;
1021            continue;
1022        }
1023
1024        let workspace_lock_path = staging_path.join(WORKSPACE_LOCK_FILE);
1025        match fs::symlink_metadata(&workspace_lock_path) {
1026            Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_file() => {
1027                quarantine_run_staging_path(staging_root, &quarantine_root, &staging_path)?;
1028                continue;
1029            }
1030            Ok(_) => {}
1031            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
1032            Err(source) => {
1033                return Err(ArtifactPathError::Io {
1034                    operation: "inspect run-workspace lock",
1035                    path: workspace_lock_path,
1036                    source,
1037                });
1038            }
1039        }
1040        let workspace_lock = OpenOptions::new()
1041            .read(true)
1042            .write(true)
1043            .create(true)
1044            .truncate(false)
1045            .open(&workspace_lock_path)
1046            .map_err(|source| ArtifactPathError::Io {
1047                operation: "open staged run-workspace lock",
1048                path: workspace_lock_path.clone(),
1049                source,
1050            })?;
1051        match FileExt::try_lock_exclusive(&workspace_lock) {
1052            Ok(()) => {
1053                quarantine_run_staging_path(staging_root, &quarantine_root, &staging_path)?;
1054            }
1055            Err(source) if source.kind() == std::io::ErrorKind::WouldBlock => {}
1056            Err(source) => {
1057                return Err(ArtifactPathError::Io {
1058                    operation: "lock staged run workspace for recovery",
1059                    path: workspace_lock_path,
1060                    source,
1061                });
1062            }
1063        }
1064    }
1065    prune_quarantine_entries(&quarantine_root, RETAIN_QUARANTINE_ENTRIES)
1066}
1067
1068fn quarantine_run_staging_path(
1069    staging_root: &Path,
1070    quarantine_root: &Path,
1071    staging_path: &Path,
1072) -> Result<(), ArtifactPathError> {
1073    let name = staging_path
1074        .file_name()
1075        .and_then(|name| name.to_str())
1076        .unwrap_or("unreadable");
1077    let quarantine_path = quarantine_root.join(format!("{name}--{}", workspace_nonce()));
1078    fs::rename(staging_path, &quarantine_path).map_err(|source| ArtifactPathError::Io {
1079        operation: "quarantine abandoned run workspace",
1080        path: staging_path.to_path_buf(),
1081        source,
1082    })?;
1083    sync_directory(staging_root)?;
1084    sync_directory(quarantine_root)
1085}
1086
1087fn allocate_latest_generation(
1088    latest_root: &ApprovedRoot,
1089) -> Result<(String, PathBuf), ArtifactPathError> {
1090    let staging_root = latest_root.prepare_dir(LATEST_STAGING_DIRECTORY)?;
1091    for _ in 0..MAX_ALLOCATION_ATTEMPTS {
1092        let generation = format!("generation-{}", workspace_nonce());
1093        let path = staging_root.join(&generation);
1094        match fs::create_dir(&path) {
1095            Ok(()) => return Ok((generation, path)),
1096            Err(source) if source.kind() == std::io::ErrorKind::AlreadyExists => continue,
1097            Err(source) => {
1098                return Err(ArtifactPathError::Io {
1099                    operation: "create latest staging generation",
1100                    path,
1101                    source,
1102                });
1103            }
1104        }
1105    }
1106    Err(ArtifactPathError::LatestTransactionAllocationExhausted)
1107}
1108
1109fn quarantine_abandoned_latest_staging(
1110    latest_root: &ApprovedRoot,
1111) -> Result<(), ArtifactPathError> {
1112    let staging = latest_root.prepare_dir(LATEST_STAGING_DIRECTORY)?;
1113    let entries = fs::read_dir(&staging).map_err(|source| ArtifactPathError::Io {
1114        operation: "list abandoned latest staging",
1115        path: staging.clone(),
1116        source,
1117    })?;
1118    for entry in entries {
1119        let entry = entry.map_err(|source| ArtifactPathError::Io {
1120            operation: "inspect abandoned latest staging entry",
1121            path: staging.clone(),
1122            source,
1123        })?;
1124        let staging_path = entry.path();
1125        let metadata =
1126            fs::symlink_metadata(&staging_path).map_err(|source| ArtifactPathError::Io {
1127                operation: "inspect abandoned latest staging entry",
1128                path: staging_path.clone(),
1129                source,
1130            })?;
1131        if metadata.file_type().is_symlink() || !metadata.is_dir() {
1132            quarantine_latest_path(latest_root, &staging_path)?;
1133            continue;
1134        }
1135
1136        let lease_path = staging_path.join(LATEST_READER_LOCK_FILE);
1137        let lease_metadata = match fs::symlink_metadata(&lease_path) {
1138            Ok(metadata) => metadata,
1139            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
1140                quarantine_latest_path(latest_root, &staging_path)?;
1141                continue;
1142            }
1143            Err(source) => {
1144                return Err(ArtifactPathError::Io {
1145                    operation: "inspect latest-generation staging lease",
1146                    path: lease_path,
1147                    source,
1148                });
1149            }
1150        };
1151        if lease_metadata.file_type().is_symlink() || !lease_metadata.is_file() {
1152            quarantine_latest_path(latest_root, &staging_path)?;
1153            continue;
1154        }
1155        let lease = OpenOptions::new()
1156            .read(true)
1157            .write(true)
1158            .open(&lease_path)
1159            .map_err(|source| ArtifactPathError::Io {
1160                operation: "open latest-generation staging lease for recovery",
1161                path: lease_path.clone(),
1162                source,
1163            })?;
1164        match FileExt::try_lock_exclusive(&lease) {
1165            Ok(()) => quarantine_latest_path(latest_root, &staging_path)?,
1166            Err(source) if source.kind() == std::io::ErrorKind::WouldBlock => {}
1167            Err(source) => {
1168                return Err(ArtifactPathError::Io {
1169                    operation: "lock latest-generation staging lease for recovery",
1170                    path: lease_path,
1171                    source,
1172                });
1173            }
1174        }
1175    }
1176    Ok(())
1177}
1178
1179fn quarantine_latest_path(
1180    latest_root: &ApprovedRoot,
1181    source_path: &Path,
1182) -> Result<(), ArtifactPathError> {
1183    let quarantine = latest_root.prepare_dir(LATEST_QUARANTINE_DIRECTORY)?;
1184    let source_name = source_path
1185        .file_name()
1186        .and_then(|name| name.to_str())
1187        .unwrap_or("unreadable");
1188    let destination = quarantine.join(format!("{source_name}--{}", workspace_nonce()));
1189    fs::rename(source_path, &destination).map_err(|source| ArtifactPathError::Io {
1190        operation: "quarantine abandoned latest staging",
1191        path: source_path.to_path_buf(),
1192        source,
1193    })?;
1194    sync_directory(
1195        source_path
1196            .parent()
1197            .expect("latest staging entry always has a parent"),
1198    )?;
1199    sync_directory(&quarantine)?;
1200    prune_quarantine_entries(&quarantine, RETAIN_QUARANTINE_ENTRIES)
1201}
1202
1203fn commit_current_pointer(
1204    latest_root: &ApprovedRoot,
1205    generation: &str,
1206) -> Result<(), PointerCommitFailure> {
1207    ArtifactId::new(generation.to_owned()).map_err(PointerCommitFailure::BeforeCommit)?;
1208    let temporary = latest_root
1209        .path()
1210        .join(format!(".current-{}", workspace_nonce()));
1211    let mut file = OpenOptions::new()
1212        .write(true)
1213        .create_new(true)
1214        .open(&temporary)
1215        .map_err(|source| {
1216            PointerCommitFailure::BeforeCommit(ArtifactPathError::Io {
1217                operation: "create latest pointer candidate",
1218                path: temporary.clone(),
1219                source,
1220            })
1221        })?;
1222    file.write_all(generation.as_bytes())
1223        .and_then(|()| file.write_all(b"\n"))
1224        .map_err(|source| {
1225            PointerCommitFailure::BeforeCommit(ArtifactPathError::Io {
1226                operation: "write latest pointer candidate",
1227                path: temporary.clone(),
1228                source,
1229            })
1230        })?;
1231    file.sync_all().map_err(|source| {
1232        PointerCommitFailure::BeforeCommit(ArtifactPathError::Io {
1233            operation: "sync latest pointer candidate",
1234            path: temporary.clone(),
1235            source,
1236        })
1237    })?;
1238    record_durability_event("sync_pointer_file");
1239
1240    let current = latest_root.path().join(LATEST_CURRENT_FILE);
1241    match fs::symlink_metadata(&current) {
1242        Ok(metadata) if metadata.is_dir() => {
1243            let _ = fs::remove_file(&temporary);
1244            return Err(PointerCommitFailure::BeforeCommit(
1245                ArtifactPathError::FileDestinationNotFile { path: current },
1246            ));
1247        }
1248        Ok(_) => {}
1249        Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
1250        Err(source) => {
1251            let _ = fs::remove_file(&temporary);
1252            return Err(PointerCommitFailure::BeforeCommit(ArtifactPathError::Io {
1253                operation: "inspect current latest pointer",
1254                path: current,
1255                source,
1256            }));
1257        }
1258    }
1259    replace_file_atomically(&temporary, &current).map_err(|source| {
1260        PointerCommitFailure::BeforeCommit(ArtifactPathError::Io {
1261            operation: "commit latest pointer",
1262            path: current,
1263            source,
1264        })
1265    })?;
1266    record_durability_event("commit_pointer");
1267    sync_directory(latest_root.path()).map_err(PointerCommitFailure::AfterCommit)
1268}
1269
1270fn recover_latest_for_writer(
1271    root: &ApprovedRoot,
1272) -> Result<Option<LatestSnapshot>, ArtifactPathError> {
1273    let _lock = acquire_latest_update_lock(root)?;
1274    let latest_path = root.prepare_dir(LATEST_DIRECTORY)?;
1275    let latest_root = ApprovedRoot::existing(latest_path)?;
1276    latest_root.prepare_dir(LATEST_STAGING_DIRECTORY)?;
1277    latest_root.prepare_dir(LATEST_GENERATIONS_DIRECTORY)?;
1278    latest_root.prepare_dir(LATEST_QUARANTINE_DIRECTORY)?;
1279    recover_latest_state_locked(root, &latest_root)
1280}
1281
1282fn recover_latest_state_locked(
1283    root: &ApprovedRoot,
1284    latest_root: &ApprovedRoot,
1285) -> Result<Option<LatestSnapshot>, ArtifactPathError> {
1286    let (snapshot, needs_alias_refresh) = match load_current_snapshot(root, true) {
1287        Ok(Some(snapshot)) => {
1288            quarantine_generations_outside_committed_chain(latest_root, &snapshot.manifest)?;
1289            (Some(snapshot), true)
1290        }
1291        Ok(None) => (
1292            recover_current_from_generation_chain(root, latest_root)?,
1293            false,
1294        ),
1295        Err(error) if is_recoverable_pointer_error(&error) => (
1296            recover_current_from_generation_chain(root, latest_root)?,
1297            false,
1298        ),
1299        Err(error) => return Err(error),
1300    };
1301    if let Some(snapshot) = snapshot.as_ref() {
1302        reject_producer_downgrade(&snapshot.manifest.producer_version)?;
1303        if needs_alias_refresh {
1304            install_stable_aliases_transactional(root, latest_root, snapshot)?.finish();
1305        }
1306    }
1307    quarantine_abandoned_latest_staging(latest_root)?;
1308    Ok(snapshot)
1309}
1310
1311fn quarantine_generations_outside_committed_chain(
1312    latest_root: &ApprovedRoot,
1313    committed: &LatestManifest,
1314) -> Result<(), ArtifactPathError> {
1315    let mut committed_chain = BTreeSet::new();
1316    let mut cursor = committed.clone();
1317    loop {
1318        committed_chain.insert(cursor.generation.clone());
1319        let Some(predecessor) = cursor.predecessor_generation.as_ref() else {
1320            break;
1321        };
1322        let cursor_path = latest_root
1323            .path()
1324            .join(LATEST_GENERATIONS_DIRECTORY)
1325            .join(&cursor.generation);
1326        if load_retention_boundary(&cursor_path, &cursor)?.is_some() {
1327            break;
1328        }
1329        let path = latest_root
1330            .path()
1331            .join(LATEST_GENERATIONS_DIRECTORY)
1332            .join(predecessor);
1333        cursor = validate_latest_manifest(&path, predecessor)?;
1334    }
1335
1336    let generations = latest_root.path().join(LATEST_GENERATIONS_DIRECTORY);
1337    let entries = fs::read_dir(&generations).map_err(|source| ArtifactPathError::Io {
1338        operation: "list uncommitted latest generations",
1339        path: generations,
1340        source,
1341    })?;
1342    for entry in entries {
1343        let entry = entry.map_err(|source| ArtifactPathError::Io {
1344            operation: "inspect uncommitted latest generation",
1345            path: latest_root.path().join(LATEST_GENERATIONS_DIRECTORY),
1346            source,
1347        })?;
1348        let path = entry.path();
1349        let generation = entry
1350            .file_name()
1351            .to_str()
1352            .ok_or_else(|| ArtifactPathError::NonUtf8ArtifactPath { path: path.clone() })?
1353            .to_owned();
1354        if committed_chain.contains(&generation) {
1355            continue;
1356        }
1357        // The current pointer already anchors and validates the committed
1358        // chain. Everything else is uncommitted recovery material, so move it
1359        // without parsing attacker-controlled or crash-torn contents.
1360        quarantine_latest_path(latest_root, &path)?;
1361    }
1362    Ok(())
1363}
1364
1365fn is_recoverable_pointer_error(error: &ArtifactPathError) -> bool {
1366    fn is_current(path: &Path) -> bool {
1367        path.file_name().and_then(|name| name.to_str()) == Some(LATEST_CURRENT_FILE)
1368    }
1369
1370    match error {
1371        ArtifactPathError::InvalidLatestPointer { .. } => true,
1372        ArtifactPathError::LatestSnapshotUnavailable { path } => is_current(path),
1373        ArtifactPathError::SymlinkComponent { path }
1374        | ArtifactPathError::FileDestinationNotFile { path } => is_current(path),
1375        ArtifactPathError::Io {
1376            operation, path, ..
1377        } => operation.contains("latest snapshot pointer") && is_current(path),
1378        _ => false,
1379    }
1380}
1381
1382fn recover_current_from_generation_chain(
1383    root: &ApprovedRoot,
1384    latest_root: &ApprovedRoot,
1385) -> Result<Option<LatestSnapshot>, ArtifactPathError> {
1386    let candidate = select_unique_generation_tip(root, latest_root)?;
1387    // A directory at `current` blocks replacement even when there is no
1388    // generation to recover. Isolate it before returning or committing.
1389    quarantine_directory_pointer_for_recovery(latest_root)?;
1390    let Some(candidate) = candidate else {
1391        return Ok(None);
1392    };
1393    reject_producer_downgrade(&candidate.manifest.producer_version)?;
1394    let alias_transaction = install_stable_aliases_transactional(root, latest_root, &candidate)?;
1395    match commit_current_pointer(latest_root, candidate.generation()) {
1396        Ok(()) => alias_transaction.finish(),
1397        Err(PointerCommitFailure::BeforeCommit(error)) => {
1398            return Err(alias_transaction.rollback(error));
1399        }
1400        Err(PointerCommitFailure::AfterCommit(error)) => {
1401            alias_transaction.finish();
1402            return Err(ArtifactPathError::LatestCommitDurabilityUncertain {
1403                generation: candidate.generation().to_owned(),
1404                source: Box::new(error),
1405            });
1406        }
1407    }
1408    Ok(Some(candidate))
1409}
1410
1411fn quarantine_directory_pointer_for_recovery(
1412    latest_root: &ApprovedRoot,
1413) -> Result<(), ArtifactPathError> {
1414    let current = latest_root.path().join(LATEST_CURRENT_FILE);
1415    let metadata = match fs::symlink_metadata(&current) {
1416        Ok(metadata) => metadata,
1417        Err(source) if source.kind() == std::io::ErrorKind::NotFound => return Ok(()),
1418        Err(source) => {
1419            return Err(ArtifactPathError::Io {
1420                operation: "inspect corrupt latest pointer for recovery",
1421                path: current,
1422                source,
1423            });
1424        }
1425    };
1426    if !metadata.is_dir() || metadata.file_type().is_symlink() {
1427        return Ok(());
1428    }
1429    let quarantine = latest_root.prepare_dir(LATEST_QUARANTINE_DIRECTORY)?;
1430    let destination = quarantine.join(format!("current-corrupt--{}", workspace_nonce()));
1431    fs::rename(&current, &destination).map_err(|source| ArtifactPathError::Io {
1432        operation: "quarantine corrupt latest pointer directory",
1433        path: current,
1434        source,
1435    })?;
1436    sync_directory(latest_root.path())?;
1437    sync_directory(&quarantine)
1438}
1439
1440fn select_unique_generation_tip(
1441    root: &ApprovedRoot,
1442    latest_root: &ApprovedRoot,
1443) -> Result<Option<LatestSnapshot>, ArtifactPathError> {
1444    let generations_path = latest_root.path().join(LATEST_GENERATIONS_DIRECTORY);
1445    let mut entries = fs::read_dir(&generations_path)
1446        .map_err(|source| ArtifactPathError::Io {
1447            operation: "list latest generations for recovery",
1448            path: generations_path.clone(),
1449            source,
1450        })?
1451        .collect::<Result<Vec<_>, _>>()
1452        .map_err(|source| ArtifactPathError::Io {
1453            operation: "inspect latest generation for recovery",
1454            path: generations_path.clone(),
1455            source,
1456        })?;
1457    entries.sort_by_key(|entry| entry.file_name());
1458
1459    let mut manifests = BTreeMap::new();
1460    let mut retention_boundaries = BTreeSet::new();
1461    for entry in entries {
1462        let path = entry.path();
1463        let metadata = fs::symlink_metadata(&path).map_err(|source| ArtifactPathError::Io {
1464            operation: "inspect latest generation for recovery",
1465            path: path.clone(),
1466            source,
1467        })?;
1468        if metadata.file_type().is_symlink() {
1469            return Err(ArtifactPathError::SymlinkComponent { path });
1470        }
1471        if !metadata.is_dir() {
1472            return Err(ArtifactPathError::DirectoryComponentNotDirectory { path });
1473        }
1474        let generation = entry
1475            .file_name()
1476            .to_str()
1477            .ok_or_else(|| ArtifactPathError::NonUtf8ArtifactPath { path: path.clone() })?
1478            .to_owned();
1479        ArtifactId::new(generation.clone())?;
1480        let manifest = validate_latest_manifest(&path, &generation)?;
1481        if load_retention_boundary(&path, &manifest)?.is_some() {
1482            retention_boundaries.insert(generation.clone());
1483        }
1484        manifests.insert(generation, (path, manifest));
1485    }
1486    if manifests.is_empty() {
1487        return Ok(None);
1488    }
1489
1490    let mut children: BTreeMap<String, Vec<String>> = BTreeMap::new();
1491    let mut predecessors = BTreeSet::new();
1492    for (generation, (_, manifest)) in &manifests {
1493        if retention_boundaries.contains(generation) {
1494            continue;
1495        }
1496        if let Some(predecessor) = manifest.predecessor_generation.as_ref() {
1497            let Some((_, predecessor_manifest)) = manifests.get(predecessor) else {
1498                return Err(ArtifactPathError::BrokenGenerationChain {
1499                    generation: generation.clone(),
1500                    predecessor: predecessor.clone(),
1501                });
1502            };
1503            reject_chain_downgrade(predecessor_manifest, manifest)?;
1504            children
1505                .entry(predecessor.clone())
1506                .or_default()
1507                .push(generation.clone());
1508            predecessors.insert(predecessor.clone());
1509        }
1510    }
1511    for (predecessor, children) in &children {
1512        if children.len() > 1 {
1513            return Err(ArtifactPathError::AmbiguousGenerationFork {
1514                predecessor: predecessor.clone(),
1515                children: children.clone(),
1516            });
1517        }
1518    }
1519    validate_generation_map_acyclic(&manifests, &retention_boundaries)?;
1520
1521    let tips: Vec<_> = manifests
1522        .keys()
1523        .filter(|generation| !predecessors.contains(*generation))
1524        .cloned()
1525        .collect();
1526    if tips.len() != 1 {
1527        return Err(ArtifactPathError::AmbiguousGenerationTips { tips });
1528    }
1529    let tip = &tips[0];
1530    let (generation_path, manifest) = manifests
1531        .remove(tip)
1532        .expect("selected generation tip came from the manifest map");
1533    Ok(Some(make_latest_snapshot(
1534        root.clone(),
1535        generation_path,
1536        manifest,
1537    )?))
1538}
1539
1540fn validate_generation_map_acyclic(
1541    manifests: &BTreeMap<String, (PathBuf, LatestManifest)>,
1542    retention_boundaries: &BTreeSet<String>,
1543) -> Result<(), ArtifactPathError> {
1544    let mut complete = BTreeSet::new();
1545    for generation in manifests.keys() {
1546        let mut chain = BTreeSet::new();
1547        let mut cursor = generation.as_str();
1548        while !complete.contains(cursor) {
1549            if !chain.insert(cursor.to_owned()) {
1550                return Err(ArtifactPathError::GenerationChainCycle {
1551                    generation: cursor.to_owned(),
1552                });
1553            }
1554            if retention_boundaries.contains(cursor) {
1555                break;
1556            }
1557            let Some(predecessor) = manifests
1558                .get(cursor)
1559                .and_then(|(_, manifest)| manifest.predecessor_generation.as_deref())
1560            else {
1561                break;
1562            };
1563            cursor = predecessor;
1564        }
1565        complete.extend(chain);
1566    }
1567    Ok(())
1568}
1569
1570fn load_current_snapshot(
1571    root: &ApprovedRoot,
1572    allow_missing: bool,
1573) -> Result<Option<LatestSnapshot>, ArtifactPathError> {
1574    let latest_path = root.path().join(LATEST_DIRECTORY);
1575    match fs::symlink_metadata(&latest_path) {
1576        Ok(metadata) if metadata.file_type().is_symlink() => {
1577            return Err(ArtifactPathError::SymlinkComponent { path: latest_path });
1578        }
1579        Ok(metadata) if !metadata.is_dir() => {
1580            return Err(ArtifactPathError::DirectoryComponentNotDirectory { path: latest_path });
1581        }
1582        Ok(_) => {}
1583        Err(source) if source.kind() == std::io::ErrorKind::NotFound && allow_missing => {
1584            return Ok(None);
1585        }
1586        Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
1587            return Err(ArtifactPathError::LatestSnapshotUnavailable { path: latest_path });
1588        }
1589        Err(source) => {
1590            return Err(ArtifactPathError::Io {
1591                operation: "inspect latest snapshot root",
1592                path: latest_path,
1593                source,
1594            });
1595        }
1596    }
1597    let latest_root = ApprovedRoot::existing(&latest_path)?;
1598    let current = latest_root.path().join(LATEST_CURRENT_FILE);
1599    let metadata = match fs::symlink_metadata(&current) {
1600        Ok(metadata) if metadata.file_type().is_symlink() => {
1601            return Err(ArtifactPathError::SymlinkComponent { path: current });
1602        }
1603        Ok(metadata) if !metadata.is_file() => {
1604            return Err(ArtifactPathError::FileDestinationNotFile { path: current });
1605        }
1606        Ok(metadata) => metadata,
1607        Err(source) if source.kind() == std::io::ErrorKind::NotFound && allow_missing => {
1608            return Ok(None);
1609        }
1610        Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
1611            return Err(ArtifactPathError::LatestSnapshotUnavailable { path: current });
1612        }
1613        Err(source) => {
1614            return Err(ArtifactPathError::Io {
1615                operation: "inspect latest snapshot pointer",
1616                path: current,
1617                source,
1618            });
1619        }
1620    };
1621    if metadata.len() > 512 {
1622        return Err(ArtifactPathError::InvalidLatestPointer { path: current });
1623    }
1624    let generation = fs::read_to_string(&current)
1625        .map_err(|source| ArtifactPathError::Io {
1626            operation: "read latest snapshot pointer",
1627            path: current.clone(),
1628            source,
1629        })?
1630        .trim()
1631        .to_owned();
1632    ArtifactId::new(generation.clone()).map_err(|_| ArtifactPathError::InvalidLatestPointer {
1633        path: current.clone(),
1634    })?;
1635
1636    let generation_path = latest_root
1637        .path()
1638        .join(LATEST_GENERATIONS_DIRECTORY)
1639        .join(&generation);
1640    match fs::symlink_metadata(&generation_path) {
1641        Ok(metadata) if metadata.file_type().is_symlink() => {
1642            return Err(ArtifactPathError::SymlinkComponent {
1643                path: generation_path,
1644            });
1645        }
1646        Ok(metadata) if !metadata.is_dir() => {
1647            return Err(ArtifactPathError::DirectoryComponentNotDirectory {
1648                path: generation_path,
1649            });
1650        }
1651        Ok(_) => {}
1652        Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
1653            return Err(ArtifactPathError::InvalidLatestPointer { path: current });
1654        }
1655        Err(source) => {
1656            return Err(ArtifactPathError::Io {
1657                operation: "inspect latest generation",
1658                path: generation_path,
1659                source,
1660            });
1661        }
1662    }
1663    let manifest = validate_latest_manifest(&generation_path, &generation)?;
1664    validate_predecessor_chain(&latest_root, &manifest)?;
1665    Ok(Some(make_latest_snapshot(
1666        root.clone(),
1667        generation_path,
1668        manifest,
1669    )?))
1670}
1671
1672fn make_latest_snapshot(
1673    root: ApprovedRoot,
1674    generation_path: PathBuf,
1675    manifest: LatestManifest,
1676) -> Result<LatestSnapshot, ArtifactPathError> {
1677    let lease_path = generation_path.join(LATEST_READER_LOCK_FILE);
1678    let reader_lease = Arc::new(OpenOptions::new().read(true).open(&lease_path).map_err(
1679        |source| ArtifactPathError::Io {
1680            operation: "open latest-generation reader lease",
1681            path: lease_path.clone(),
1682            source,
1683        },
1684    )?);
1685    FileExt::lock_shared(reader_lease.as_ref()).map_err(|source| ArtifactPathError::Io {
1686        operation: "lock latest-generation reader lease",
1687        path: lease_path,
1688        source,
1689    })?;
1690    let source_run_lease = open_shared_lease(
1691        &root
1692            .path()
1693            .join(&manifest.source_publication_id)
1694            .join(RUN_READER_LOCK_FILE),
1695        "open latest source-run reader lease",
1696        "lock latest source-run reader lease",
1697    )?;
1698    #[cfg(unix)]
1699    let generation_directory = Arc::new(open_directory_no_follow(&generation_path).map_err(
1700        |source| ArtifactPathError::Io {
1701            operation: "pin latest generation directory",
1702            path: generation_path.clone(),
1703            source,
1704        },
1705    )?);
1706    Ok(LatestSnapshot {
1707        root,
1708        generation_path,
1709        manifest,
1710        _reader_lease: reader_lease,
1711        _source_run_lease: source_run_lease,
1712        #[cfg(unix)]
1713        generation_directory,
1714    })
1715}
1716
1717fn open_shared_lease(
1718    path: &Path,
1719    open_operation: &'static str,
1720    lock_operation: &'static str,
1721) -> Result<Arc<File>, ArtifactPathError> {
1722    let lease = Arc::new(OpenOptions::new().read(true).open(path).map_err(|source| {
1723        ArtifactPathError::Io {
1724            operation: open_operation,
1725            path: path.to_path_buf(),
1726            source,
1727        }
1728    })?);
1729    FileExt::lock_shared(lease.as_ref()).map_err(|source| ArtifactPathError::Io {
1730        operation: lock_operation,
1731        path: path.to_path_buf(),
1732        source,
1733    })?;
1734    Ok(lease)
1735}
1736
1737fn create_run_reader_lease_file(directory: &Path) -> Result<(), ArtifactPathError> {
1738    create_lease_file(
1739        directory,
1740        RUN_READER_LOCK_FILE,
1741        "create published-run reader lease",
1742    )
1743}
1744
1745fn create_reader_lease_file(directory: &Path) -> Result<(), ArtifactPathError> {
1746    create_lease_file(
1747        directory,
1748        LATEST_READER_LOCK_FILE,
1749        "create latest-generation reader lease",
1750    )
1751}
1752
1753fn create_lease_file(
1754    directory: &Path,
1755    name: &str,
1756    operation: &'static str,
1757) -> Result<(), ArtifactPathError> {
1758    let path = directory.join(name);
1759    let file = OpenOptions::new()
1760        .read(true)
1761        .write(true)
1762        .create_new(true)
1763        .open(&path)
1764        .map_err(|source| ArtifactPathError::Io {
1765            operation,
1766            path: path.clone(),
1767            source,
1768        })?;
1769    file.sync_all().map_err(|source| ArtifactPathError::Io {
1770        operation: "sync latest-generation reader lease",
1771        path,
1772        source,
1773    })
1774}
1775
1776fn prune_latest_generations(
1777    latest_root: &ApprovedRoot,
1778    current: &LatestManifest,
1779) -> Result<(), ArtifactPathError> {
1780    let generations_root = latest_root.path().join(LATEST_GENERATIONS_DIRECTORY);
1781    let mut chain = Vec::new();
1782    let mut cursor = current.clone();
1783    loop {
1784        let path = generations_root.join(&cursor.generation);
1785        let boundary = load_retention_boundary(&path, &cursor)?;
1786        chain.push((path, cursor.clone()));
1787        if boundary.is_some() {
1788            break;
1789        }
1790        let Some(predecessor) = cursor.predecessor_generation.as_ref() else {
1791            break;
1792        };
1793        let predecessor_path = generations_root.join(predecessor);
1794        cursor = validate_latest_manifest(&predecessor_path, predecessor)?;
1795    }
1796    if chain.len() <= RETAIN_LATEST_GENERATIONS {
1797        return Ok(());
1798    }
1799
1800    // Retention is lease-aware and all-or-nothing. Keeping the complete tail
1801    // temporarily is preferable to stranding a leased generation outside the
1802    // retained chain.
1803    let mut leases = Vec::new();
1804    for (path, _) in &chain[RETAIN_LATEST_GENERATIONS..] {
1805        let lease_path = path.join(LATEST_READER_LOCK_FILE);
1806        let lease = OpenOptions::new()
1807            .read(true)
1808            .write(true)
1809            .open(&lease_path)
1810            .map_err(|source| ArtifactPathError::Io {
1811                operation: "open latest-generation retention lease",
1812                path: lease_path.clone(),
1813                source,
1814            })?;
1815        match FileExt::try_lock_exclusive(&lease) {
1816            Ok(()) => leases.push(lease),
1817            Err(source) if source.kind() == std::io::ErrorKind::WouldBlock => return Ok(()),
1818            Err(source) => {
1819                return Err(ArtifactPathError::Io {
1820                    operation: "lock latest generation for retention",
1821                    path: lease_path,
1822                    source,
1823                });
1824            }
1825        }
1826    }
1827
1828    let (oldest_retained_path, oldest_retained) = &chain[RETAIN_LATEST_GENERATIONS - 1];
1829    let predecessor = oldest_retained
1830        .predecessor_generation
1831        .clone()
1832        .expect("a prunable chain has a predecessor after the retention boundary");
1833    let boundary = RetentionBoundary {
1834        format_version: 1,
1835        predecessor_generation: predecessor,
1836    };
1837    let boundary_path = oldest_retained_path.join(LATEST_RETENTION_BOUNDARY_FILE);
1838    let boundary_temp =
1839        oldest_retained_path.join(format!(".retention-boundary-{}", workspace_nonce()));
1840    write_json_file(
1841        &boundary_temp,
1842        &boundary,
1843        "write latest-generation retention boundary",
1844    )?;
1845    replace_file_atomically(&boundary_temp, &boundary_path).map_err(|source| {
1846        ArtifactPathError::Io {
1847            operation: "commit latest-generation retention boundary",
1848            path: boundary_path.clone(),
1849            source,
1850        }
1851    })?;
1852    sync_directory(oldest_retained_path)?;
1853
1854    let quarantine = latest_root.prepare_dir(LATEST_QUARANTINE_DIRECTORY)?;
1855    let mut retired = Vec::new();
1856    for (path, _) in &chain[RETAIN_LATEST_GENERATIONS..] {
1857        let name = path
1858            .file_name()
1859            .and_then(|name| name.to_str())
1860            .unwrap_or("generation");
1861        let destination = quarantine.join(format!("retired-{name}--{}", workspace_nonce()));
1862        fs::rename(path, &destination).map_err(|source| ArtifactPathError::Io {
1863            operation: "retire latest generation",
1864            path: path.clone(),
1865            source,
1866        })?;
1867        retired.push(destination);
1868    }
1869    sync_directory(&generations_root)?;
1870    sync_directory(&quarantine)?;
1871    drop(leases);
1872
1873    for path in retired {
1874        fs::remove_dir_all(&path).map_err(|source| ArtifactPathError::Io {
1875            operation: "delete retired latest generation",
1876            path,
1877            source,
1878        })?;
1879    }
1880    sync_directory(&quarantine)
1881}
1882
1883fn prune_published_runs(root: &ApprovedRoot) -> Result<(), ArtifactPathError> {
1884    let mut runs = Vec::new();
1885    for entry in fs::read_dir(root.path()).map_err(|source| ArtifactPathError::Io {
1886        operation: "list published runs for retention",
1887        path: root.path().to_path_buf(),
1888        source,
1889    })? {
1890        let entry = entry.map_err(|source| ArtifactPathError::Io {
1891            operation: "inspect published run for retention",
1892            path: root.path().to_path_buf(),
1893            source,
1894        })?;
1895        let path = entry.path();
1896        let name = entry.file_name().to_string_lossy().into_owned();
1897        if name.starts_with('.') {
1898            continue;
1899        }
1900        let metadata = fs::symlink_metadata(&path).map_err(|source| ArtifactPathError::Io {
1901            operation: "inspect published run type for retention",
1902            path: path.clone(),
1903            source,
1904        })?;
1905        if metadata.file_type().is_symlink() || !metadata.is_dir() {
1906            continue;
1907        }
1908        if !path.join(RUN_MANIFEST_FILE).is_file() || !path.join(RUN_READER_LOCK_FILE).is_file() {
1909            continue;
1910        }
1911        // Never delete a directory merely because its name resembles a run.
1912        // Full manifest verification proves Mobench owns the complete tree.
1913        if validate_run_manifest(&path).is_err() {
1914            continue;
1915        }
1916        let order_key = name
1917            .rsplit_once("--")
1918            .map(|(_, nonce)| nonce.to_owned())
1919            .unwrap_or_else(|| name.clone());
1920        runs.push((order_key, path));
1921    }
1922    runs.sort_by(|left, right| right.0.cmp(&left.0));
1923    if runs.len() <= RETAIN_PUBLISHED_RUNS {
1924        return Ok(());
1925    }
1926
1927    let quarantine = root.prepare_dir(RUN_RETENTION_QUARANTINE_DIRECTORY)?;
1928    for (_, path) in &runs[RETAIN_PUBLISHED_RUNS..] {
1929        let lease_path = path.join(RUN_READER_LOCK_FILE);
1930        let lease = OpenOptions::new()
1931            .read(true)
1932            .write(true)
1933            .open(&lease_path)
1934            .map_err(|source| ArtifactPathError::Io {
1935                operation: "open published-run retention lease",
1936                path: lease_path.clone(),
1937                source,
1938            })?;
1939        match FileExt::try_lock_exclusive(&lease) {
1940            Ok(()) => {}
1941            Err(source) if source.kind() == std::io::ErrorKind::WouldBlock => continue,
1942            Err(source) => {
1943                return Err(ArtifactPathError::Io {
1944                    operation: "lock published run for retention",
1945                    path: lease_path,
1946                    source,
1947                });
1948            }
1949        }
1950        let name = path
1951            .file_name()
1952            .and_then(|name| name.to_str())
1953            .unwrap_or("run");
1954        let retired = quarantine.join(format!("retired-{name}--{}", workspace_nonce()));
1955        fs::rename(path, &retired).map_err(|source| ArtifactPathError::Io {
1956            operation: "retire published run",
1957            path: path.clone(),
1958            source,
1959        })?;
1960        drop(lease);
1961        fs::remove_dir_all(&retired).map_err(|source| ArtifactPathError::Io {
1962            operation: "delete retired published run",
1963            path: retired,
1964            source,
1965        })?;
1966    }
1967    sync_directory(root.path())?;
1968    sync_directory(&quarantine)?;
1969    prune_quarantine_entries(&quarantine, RETAIN_QUARANTINE_ENTRIES)
1970}
1971
1972fn prune_quarantine_entries(directory: &Path, retain: usize) -> Result<(), ArtifactPathError> {
1973    let mut entries = fs::read_dir(directory)
1974        .map_err(|source| ArtifactPathError::Io {
1975            operation: "list artifact quarantine for retention",
1976            path: directory.to_path_buf(),
1977            source,
1978        })?
1979        .collect::<Result<Vec<_>, _>>()
1980        .map_err(|source| ArtifactPathError::Io {
1981            operation: "inspect artifact quarantine for retention",
1982            path: directory.to_path_buf(),
1983            source,
1984        })?;
1985    entries.sort_by_key(|entry| std::cmp::Reverse(entry.file_name()));
1986    for entry in entries.into_iter().skip(retain) {
1987        let path = entry.path();
1988        let metadata = fs::symlink_metadata(&path).map_err(|source| ArtifactPathError::Io {
1989            operation: "inspect quarantined artifact for deletion",
1990            path: path.clone(),
1991            source,
1992        })?;
1993        let result = if metadata.file_type().is_symlink() || metadata.is_file() {
1994            fs::remove_file(&path)
1995        } else if metadata.is_dir() {
1996            fs::remove_dir_all(&path)
1997        } else {
1998            fs::remove_file(&path)
1999        };
2000        result.map_err(|source| ArtifactPathError::Io {
2001            operation: "delete expired quarantined artifact",
2002            path,
2003            source,
2004        })?;
2005    }
2006    sync_directory(directory)
2007}
2008
2009fn validate_predecessor_chain(
2010    latest_root: &ApprovedRoot,
2011    tip: &LatestManifest,
2012) -> Result<(), ArtifactPathError> {
2013    let mut visited = BTreeSet::new();
2014    let mut child = tip.clone();
2015    loop {
2016        if !visited.insert(child.generation.clone()) {
2017            return Err(ArtifactPathError::GenerationChainCycle {
2018                generation: child.generation,
2019            });
2020        }
2021        let Some(predecessor) = child.predecessor_generation.as_ref() else {
2022            return Ok(());
2023        };
2024        let child_path = latest_root
2025            .path()
2026            .join(LATEST_GENERATIONS_DIRECTORY)
2027            .join(&child.generation);
2028        if load_retention_boundary(&child_path, &child)?.is_some() {
2029            return Ok(());
2030        }
2031        let predecessor_path = latest_root
2032            .path()
2033            .join(LATEST_GENERATIONS_DIRECTORY)
2034            .join(predecessor);
2035        let predecessor_manifest = match fs::symlink_metadata(&predecessor_path) {
2036            Ok(metadata) if metadata.file_type().is_symlink() => {
2037                return Err(ArtifactPathError::SymlinkComponent {
2038                    path: predecessor_path,
2039                });
2040            }
2041            Ok(metadata) if !metadata.is_dir() => {
2042                return Err(ArtifactPathError::DirectoryComponentNotDirectory {
2043                    path: predecessor_path,
2044                });
2045            }
2046            Ok(_) => validate_latest_manifest(&predecessor_path, predecessor)?,
2047            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
2048                return Err(ArtifactPathError::BrokenGenerationChain {
2049                    generation: child.generation,
2050                    predecessor: predecessor.clone(),
2051                });
2052            }
2053            Err(source) => {
2054                return Err(ArtifactPathError::Io {
2055                    operation: "inspect predecessor generation",
2056                    path: predecessor_path,
2057                    source,
2058                });
2059            }
2060        };
2061        reject_chain_downgrade(&predecessor_manifest, &child)?;
2062        child = predecessor_manifest;
2063    }
2064}
2065
2066fn load_retention_boundary(
2067    generation_path: &Path,
2068    manifest: &LatestManifest,
2069) -> Result<Option<RetentionBoundary>, ArtifactPathError> {
2070    let path = generation_path.join(LATEST_RETENTION_BOUNDARY_FILE);
2071    let metadata = match fs::symlink_metadata(&path) {
2072        Ok(metadata) => metadata,
2073        Err(source) if source.kind() == std::io::ErrorKind::NotFound => return Ok(None),
2074        Err(source) => {
2075            return Err(ArtifactPathError::Io {
2076                operation: "inspect latest-generation retention boundary",
2077                path,
2078                source,
2079            });
2080        }
2081    };
2082    if metadata.file_type().is_symlink() {
2083        return Err(ArtifactPathError::SymlinkComponent { path });
2084    }
2085    if !metadata.is_file() {
2086        return Err(ArtifactPathError::FileDestinationNotFile { path });
2087    }
2088    let body = fs::read(&path).map_err(|source| ArtifactPathError::Io {
2089        operation: "read latest-generation retention boundary",
2090        path: path.clone(),
2091        source,
2092    })?;
2093    let boundary: RetentionBoundary =
2094        serde_json::from_slice(&body).map_err(|source| ArtifactPathError::ManifestJson {
2095            path: path.clone(),
2096            source,
2097        })?;
2098    if boundary.format_version != 1
2099        || manifest.predecessor_generation.as_deref()
2100            != Some(boundary.predecessor_generation.as_str())
2101    {
2102        return Err(ArtifactPathError::CorruptGeneration {
2103            path,
2104            detail: "retention boundary does not match the manifest predecessor".to_owned(),
2105        });
2106    }
2107    Ok(Some(boundary))
2108}
2109
2110fn reject_chain_downgrade(
2111    predecessor: &LatestManifest,
2112    child: &LatestManifest,
2113) -> Result<(), ArtifactPathError> {
2114    if parse_producer_version(&predecessor.producer_version)?
2115        > parse_producer_version(&child.producer_version)?
2116    {
2117        return Err(ArtifactPathError::ProducerDowngrade {
2118            existing: predecessor.producer_version.clone(),
2119            attempted: child.producer_version.clone(),
2120        });
2121    }
2122    Ok(())
2123}
2124
2125fn reject_producer_downgrade(existing: &str) -> Result<(), ArtifactPathError> {
2126    let existing_version = parse_producer_version(existing)?;
2127    let producer_version = parse_producer_version(env!("CARGO_PKG_VERSION"))?;
2128    if existing_version > producer_version {
2129        return Err(ArtifactPathError::ProducerDowngrade {
2130            existing: existing.to_owned(),
2131            attempted: env!("CARGO_PKG_VERSION").to_owned(),
2132        });
2133    }
2134    Ok(())
2135}
2136
2137fn install_stable_aliases_transactional(
2138    root: &ApprovedRoot,
2139    latest_root: &ApprovedRoot,
2140    snapshot: &LatestSnapshot,
2141) -> Result<StableAliasTransaction, ArtifactPathError> {
2142    let current_aliases: BTreeSet<ArtifactId> = match load_current_snapshot(root, true) {
2143        Ok(Some(current)) => current
2144            .manifest
2145            .artifacts
2146            .iter()
2147            .map(|artifact| ArtifactId::new(artifact.destination_relative_path.clone()))
2148            .collect::<Result<_, _>>()?,
2149        Ok(None) => BTreeSet::new(),
2150        Err(error) if is_recoverable_pointer_error(&error) => BTreeSet::new(),
2151        Err(error) => return Err(error),
2152    };
2153    let next_aliases: BTreeSet<ArtifactId> = snapshot
2154        .manifest
2155        .artifacts
2156        .iter()
2157        .map(|artifact| ArtifactId::new(artifact.destination_relative_path.clone()))
2158        .collect::<Result<_, _>>()?;
2159
2160    let staging_root = latest_root.prepare_dir(LATEST_STAGING_DIRECTORY)?;
2161    let transaction_path = staging_root.join(format!("aliases-{}", workspace_nonce()));
2162    fs::create_dir(&transaction_path).map_err(|source| ArtifactPathError::Io {
2163        operation: "create stable-alias transaction",
2164        path: transaction_path.clone(),
2165        source,
2166    })?;
2167    let transaction_root = ApprovedRoot::existing(&transaction_path)?;
2168    let new_root = transaction_root.prepare_dir("new")?;
2169    let backup_root = transaction_root.prepare_dir("backup")?;
2170
2171    let mut prepared = Vec::with_capacity(snapshot.manifest.artifacts.len());
2172    for artifact in &snapshot.manifest.artifacts {
2173        let id = ArtifactId::new(artifact.destination_relative_path.clone())?;
2174        let source = snapshot.generation_path.join(id.as_str());
2175        let destination = root.prepare_file(id.as_str())?;
2176        let staged = new_root.join(id.as_str());
2177        if let Err(error) = copy_new_synced(&source, &staged, "stage stable latest alias") {
2178            let _ = fs::remove_dir_all(&transaction_path);
2179            return Err(error);
2180        }
2181        prepared.push((staged, destination, id));
2182    }
2183    sync_artifact_tree(&transaction_path)?;
2184
2185    let mut transaction = StableAliasTransaction {
2186        path: transaction_path,
2187        root_path: root.path().to_path_buf(),
2188        installed: Vec::new(),
2189        backups: Vec::new(),
2190    };
2191
2192    // Aliases owned by the previous committed generation but omitted by the
2193    // candidate must disappear in the same rollback-capable transaction. Their
2194    // backups restore the old alias set if the pointer does not commit.
2195    for id in current_aliases.difference(&next_aliases) {
2196        let destination = root.prepare_file(id.as_str())?;
2197        match fs::symlink_metadata(&destination) {
2198            Ok(metadata) if metadata.file_type().is_symlink() => {
2199                let original = ArtifactPathError::SymlinkComponent { path: destination };
2200                return Err(transaction.rollback(original));
2201            }
2202            Ok(metadata) if !metadata.is_file() => {
2203                let original = ArtifactPathError::FileDestinationNotFile { path: destination };
2204                return Err(transaction.rollback(original));
2205            }
2206            Ok(_) => {
2207                let backup = backup_root.join(id.as_str());
2208                if let Err(source) = fs::rename(&destination, &backup) {
2209                    let original = ArtifactPathError::Io {
2210                        operation: "retire obsolete stable latest alias",
2211                        path: destination,
2212                        source,
2213                    };
2214                    return Err(transaction.rollback(original));
2215                }
2216                transaction.backups.push((destination, backup));
2217                if let Err(error) =
2218                    sync_directory(root.path()).and_then(|()| sync_directory(&backup_root))
2219                {
2220                    return Err(transaction.rollback(error));
2221                }
2222            }
2223            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
2224            Err(source) => {
2225                let original = ArtifactPathError::Io {
2226                    operation: "inspect obsolete stable latest alias",
2227                    path: destination,
2228                    source,
2229                };
2230                return Err(transaction.rollback(original));
2231            }
2232        }
2233    }
2234
2235    for (staged, destination, id) in prepared {
2236        match fs::symlink_metadata(&destination) {
2237            Ok(metadata) if metadata.file_type().is_symlink() => {
2238                let original = ArtifactPathError::SymlinkComponent { path: destination };
2239                return Err(transaction.rollback(original));
2240            }
2241            Ok(metadata) if !metadata.is_file() => {
2242                let original = ArtifactPathError::FileDestinationNotFile { path: destination };
2243                return Err(transaction.rollback(original));
2244            }
2245            Ok(_) => {
2246                let backup = backup_root.join(id.as_str());
2247                if let Err(source) = fs::rename(&destination, &backup) {
2248                    let original = ArtifactPathError::Io {
2249                        operation: "back up stable latest alias",
2250                        path: destination,
2251                        source,
2252                    };
2253                    return Err(transaction.rollback(original));
2254                }
2255                transaction.backups.push((destination.clone(), backup));
2256                if let Err(error) =
2257                    sync_directory(root.path()).and_then(|()| sync_directory(&backup_root))
2258                {
2259                    return Err(transaction.rollback(error));
2260                }
2261            }
2262            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
2263            Err(source) => {
2264                let original = ArtifactPathError::Io {
2265                    operation: "inspect stable latest alias",
2266                    path: destination,
2267                    source,
2268                };
2269                return Err(transaction.rollback(original));
2270            }
2271        }
2272
2273        if let Err(original) = maybe_fail_alias_install(transaction.installed.len(), &destination) {
2274            return Err(transaction.rollback(original));
2275        }
2276        if let Err(source) = fs::rename(&staged, &destination) {
2277            let original = ArtifactPathError::Io {
2278                operation: "install stable latest alias",
2279                path: destination,
2280                source,
2281            };
2282            return Err(transaction.rollback(original));
2283        }
2284        transaction.installed.push(destination);
2285        if let Err(error) = sync_directory(root.path()) {
2286            return Err(transaction.rollback(error));
2287        }
2288    }
2289    record_durability_event("install_aliases");
2290    Ok(transaction)
2291}
2292
2293fn remove_installed_latest(installed: &[PathBuf]) -> Result<(), ArtifactPathError> {
2294    let mut first_error = None;
2295    for destination in installed.iter().rev() {
2296        match fs::remove_file(destination) {
2297            Ok(()) => {}
2298            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {}
2299            Err(source) => {
2300                first_error.get_or_insert_with(|| ArtifactPathError::Io {
2301                    operation: "roll back installed latest artifact",
2302                    path: destination.clone(),
2303                    source,
2304                });
2305            }
2306        }
2307    }
2308    first_error.map_or(Ok(()), Err)
2309}
2310
2311fn restore_latest_backups(backups: &[(PathBuf, PathBuf)]) -> Result<(), ArtifactPathError> {
2312    let mut first_error = None;
2313    for (destination, backup) in backups.iter().rev() {
2314        if let Err(source) = fs::rename(backup, destination) {
2315            first_error.get_or_insert_with(|| ArtifactPathError::Io {
2316                operation: "restore previous latest artifact",
2317                path: destination.clone(),
2318                source,
2319            });
2320        }
2321    }
2322    first_error.map_or(Ok(()), Err)
2323}
2324
2325fn latest_update_failure(
2326    original: ArtifactPathError,
2327    installed: &[PathBuf],
2328    backups: &[(PathBuf, PathBuf)],
2329    recovery_path: PathBuf,
2330) -> ArtifactPathError {
2331    let remove_error = remove_installed_latest(installed).err();
2332    let restore_error = restore_latest_backups(backups).err();
2333    let rollback = remove_error.or(restore_error);
2334
2335    match rollback {
2336        Some(rollback) => ArtifactPathError::LatestRollback {
2337            original: Box::new(original),
2338            rollback: Box::new(rollback),
2339            recovery_path,
2340        },
2341        None => original,
2342    }
2343}
2344
2345fn validate_run_manifest(run_path: &Path) -> Result<RunManifest, ArtifactPathError> {
2346    let manifest_path = run_path.join(RUN_MANIFEST_FILE);
2347    let manifest: RunManifest = read_json_file(&manifest_path, "read run manifest")?;
2348    if manifest.format_version != RUN_MANIFEST_VERSION {
2349        return Err(ArtifactPathError::UnknownManifestVersion {
2350            path: manifest_path,
2351            found: manifest.format_version,
2352            supported: RUN_MANIFEST_VERSION,
2353        });
2354    }
2355    parse_producer_version(&manifest.producer_version)?;
2356    ArtifactId::new(manifest.logical_id.clone())?;
2357    ArtifactId::new(manifest.publication_id.clone())?;
2358    if let Some(expected) = manifest.expected_latest_generation.as_ref() {
2359        ArtifactId::new(expected.clone())?;
2360    }
2361    let expected_publication = run_path
2362        .file_name()
2363        .and_then(|name| name.to_str())
2364        .ok_or_else(|| ArtifactPathError::NonUtf8ArtifactPath {
2365            path: run_path.to_path_buf(),
2366        })?;
2367    if manifest.publication_id != expected_publication {
2368        return Err(ArtifactPathError::ManifestIntegrity {
2369            path: manifest_path,
2370            detail: "publication identity does not match its directory".to_owned(),
2371        });
2372    }
2373
2374    let actual = inspect_artifact_tree(run_path, &[RUN_MANIFEST_FILE, RUN_READER_LOCK_FILE])?;
2375    if manifest.artifacts != actual {
2376        return Err(ArtifactPathError::ManifestIntegrity {
2377            path: manifest_path,
2378            detail: "run artifact paths, sizes, or digests do not match".to_owned(),
2379        });
2380    }
2381    Ok(manifest)
2382}
2383
2384fn validate_latest_manifest(
2385    generation_path: &Path,
2386    expected_generation: &str,
2387) -> Result<LatestManifest, ArtifactPathError> {
2388    let manifest_path = generation_path.join(LATEST_MANIFEST_FILE);
2389    let manifest: LatestManifest = read_json_file(&manifest_path, "read latest manifest")?;
2390    if manifest.format_version != LATEST_MANIFEST_VERSION {
2391        return Err(ArtifactPathError::UnknownManifestVersion {
2392            path: manifest_path,
2393            found: manifest.format_version,
2394            supported: LATEST_MANIFEST_VERSION,
2395        });
2396    }
2397    parse_producer_version(&manifest.producer_version)?;
2398    ArtifactId::new(manifest.generation.clone())?;
2399    ArtifactId::new(manifest.source_logical_id.clone())?;
2400    ArtifactId::new(manifest.source_publication_id.clone())?;
2401    if let Some(predecessor) = manifest.predecessor_generation.as_ref() {
2402        ArtifactId::new(predecessor.clone())?;
2403    }
2404    if manifest.generation != expected_generation {
2405        return Err(ArtifactPathError::CorruptGeneration {
2406            path: generation_path.to_path_buf(),
2407            detail: "manifest generation does not match current pointer".to_owned(),
2408        });
2409    }
2410
2411    let actual = inspect_artifact_tree(
2412        generation_path,
2413        &[
2414            LATEST_MANIFEST_FILE,
2415            LATEST_READER_LOCK_FILE,
2416            LATEST_RETENTION_BOUNDARY_FILE,
2417        ],
2418    )?;
2419    let mut expected = Vec::with_capacity(manifest.artifacts.len());
2420    for artifact in &manifest.artifacts {
2421        ArtifactId::new(artifact.source_relative_path.clone())?;
2422        ArtifactId::new(artifact.destination_relative_path.clone())?;
2423        expected.push(ManifestArtifact {
2424            relative_path: artifact.destination_relative_path.clone(),
2425            size: artifact.size,
2426            sha256: artifact.sha256.clone(),
2427        });
2428    }
2429    expected.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
2430    if expected != actual {
2431        return Err(ArtifactPathError::CorruptGeneration {
2432            path: generation_path.to_path_buf(),
2433            detail: "generation artifact paths, sizes, or digests do not match".to_owned(),
2434        });
2435    }
2436    Ok(manifest)
2437}
2438
2439fn parse_producer_version(version: &str) -> Result<Version, ArtifactPathError> {
2440    Version::parse(version).map_err(|source| ArtifactPathError::InvalidProducerVersion {
2441        version: version.to_owned(),
2442        source,
2443    })
2444}
2445
2446fn read_json_file<T: for<'de> Deserialize<'de>>(
2447    path: &Path,
2448    operation: &'static str,
2449) -> Result<T, ArtifactPathError> {
2450    let metadata = fs::symlink_metadata(path).map_err(|source| ArtifactPathError::Io {
2451        operation,
2452        path: path.to_path_buf(),
2453        source,
2454    })?;
2455    if metadata.file_type().is_symlink() {
2456        return Err(ArtifactPathError::SymlinkComponent {
2457            path: path.to_path_buf(),
2458        });
2459    }
2460    if !metadata.is_file() {
2461        return Err(ArtifactPathError::FileDestinationNotFile {
2462            path: path.to_path_buf(),
2463        });
2464    }
2465    let file = File::open(path).map_err(|source| ArtifactPathError::Io {
2466        operation,
2467        path: path.to_path_buf(),
2468        source,
2469    })?;
2470    serde_json::from_reader(BufReader::new(file)).map_err(|source| {
2471        ArtifactPathError::ManifestJson {
2472            path: path.to_path_buf(),
2473            source,
2474        }
2475    })
2476}
2477
2478fn write_json_file<T: Serialize>(
2479    path: &Path,
2480    value: &T,
2481    operation: &'static str,
2482) -> Result<(), ArtifactPathError> {
2483    let mut encoded =
2484        serde_json::to_vec_pretty(value).map_err(|source| ArtifactPathError::ManifestJson {
2485            path: path.to_path_buf(),
2486            source,
2487        })?;
2488    encoded.push(b'\n');
2489    let mut file = OpenOptions::new()
2490        .write(true)
2491        .create_new(true)
2492        .open(path)
2493        .map_err(|source| ArtifactPathError::Io {
2494            operation,
2495            path: path.to_path_buf(),
2496            source,
2497        })?;
2498    file.write_all(&encoded)
2499        .map_err(|source| ArtifactPathError::Io {
2500            operation,
2501            path: path.to_path_buf(),
2502            source,
2503        })?;
2504    file.sync_all().map_err(|source| ArtifactPathError::Io {
2505        operation: "sync manifest",
2506        path: path.to_path_buf(),
2507        source,
2508    })?;
2509    record_durability_event("sync_manifest");
2510    Ok(())
2511}
2512
2513fn inspect_artifact_tree(
2514    root: &Path,
2515    excluded_relative_paths: &[&str],
2516) -> Result<Vec<ManifestArtifact>, ArtifactPathError> {
2517    let mut artifacts = Vec::new();
2518    inspect_artifact_tree_recursive(root, root, excluded_relative_paths, &mut artifacts)?;
2519    artifacts.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
2520    Ok(artifacts)
2521}
2522
2523fn inspect_artifact_tree_recursive(
2524    root: &Path,
2525    directory: &Path,
2526    excluded_relative_paths: &[&str],
2527    artifacts: &mut Vec<ManifestArtifact>,
2528) -> Result<(), ArtifactPathError> {
2529    let mut entries = fs::read_dir(directory)
2530        .map_err(|source| ArtifactPathError::Io {
2531            operation: "list artifact directory",
2532            path: directory.to_path_buf(),
2533            source,
2534        })?
2535        .collect::<Result<Vec<_>, _>>()
2536        .map_err(|source| ArtifactPathError::Io {
2537            operation: "inspect artifact directory entry",
2538            path: directory.to_path_buf(),
2539            source,
2540        })?;
2541    entries.sort_by_key(|entry| entry.file_name());
2542    for entry in entries {
2543        let path = entry.path();
2544        let metadata = fs::symlink_metadata(&path).map_err(|source| ArtifactPathError::Io {
2545            operation: "inspect artifact tree entry",
2546            path: path.clone(),
2547            source,
2548        })?;
2549        if metadata.file_type().is_symlink() {
2550            return Err(ArtifactPathError::SymlinkComponent { path });
2551        }
2552        if metadata.is_dir() {
2553            inspect_artifact_tree_recursive(root, &path, excluded_relative_paths, artifacts)?;
2554            continue;
2555        }
2556        if !metadata.is_file() {
2557            return Err(ArtifactPathError::FileDestinationNotFile { path });
2558        }
2559        reject_hardlinked_regular_file(&metadata, &path)?;
2560        let relative_path = manifest_relative_path(root, &path)?;
2561        if excluded_relative_paths.contains(&relative_path.as_str()) {
2562            continue;
2563        }
2564        let (size, sha256) = digest_file(&path)?;
2565        artifacts.push(ManifestArtifact {
2566            relative_path,
2567            size,
2568            sha256,
2569        });
2570    }
2571    Ok(())
2572}
2573
2574fn manifest_relative_path(root: &Path, path: &Path) -> Result<String, ArtifactPathError> {
2575    let relative = path
2576        .strip_prefix(root)
2577        .map_err(|_| ArtifactPathError::InvalidRelativePath {
2578            path: path.to_path_buf(),
2579        })?;
2580    let mut components = Vec::new();
2581    for component in relative.components() {
2582        let std::path::Component::Normal(component) = component else {
2583            return Err(ArtifactPathError::InvalidRelativePath {
2584                path: relative.to_path_buf(),
2585            });
2586        };
2587        components.push(component.to_str().ok_or_else(|| {
2588            ArtifactPathError::NonUtf8ArtifactPath {
2589                path: path.to_path_buf(),
2590            }
2591        })?);
2592    }
2593    Ok(components.join("/"))
2594}
2595
2596fn digest_file(path: &Path) -> Result<(u64, String), ArtifactPathError> {
2597    let metadata = fs::symlink_metadata(path).map_err(|source| ArtifactPathError::Io {
2598        operation: "inspect artifact for digest",
2599        path: path.to_path_buf(),
2600        source,
2601    })?;
2602    if metadata.file_type().is_symlink() {
2603        return Err(ArtifactPathError::SymlinkComponent {
2604            path: path.to_path_buf(),
2605        });
2606    }
2607    if !metadata.is_file() {
2608        return Err(ArtifactPathError::FileDestinationNotFile {
2609            path: path.to_path_buf(),
2610        });
2611    }
2612    reject_hardlinked_regular_file(&metadata, path)?;
2613    let mut file = File::open(path).map_err(|source| ArtifactPathError::Io {
2614        operation: "open artifact for digest",
2615        path: path.to_path_buf(),
2616        source,
2617    })?;
2618    let mut digest = Sha256::new();
2619    let mut buffer = [0_u8; 64 * 1024];
2620    loop {
2621        let read = file
2622            .read(&mut buffer)
2623            .map_err(|source| ArtifactPathError::Io {
2624                operation: "read artifact for digest",
2625                path: path.to_path_buf(),
2626                source,
2627            })?;
2628        if read == 0 {
2629            break;
2630        }
2631        digest.update(&buffer[..read]);
2632    }
2633    Ok((metadata.len(), format!("{:x}", digest.finalize())))
2634}
2635
2636fn sha256_bytes(bytes: &[u8]) -> String {
2637    let mut digest = Sha256::new();
2638    digest.update(bytes);
2639    format!("{:x}", digest.finalize())
2640}
2641
2642#[cfg(unix)]
2643fn open_directory_no_follow(path: &Path) -> std::io::Result<File> {
2644    use std::os::unix::fs::OpenOptionsExt;
2645
2646    OpenOptions::new()
2647        .read(true)
2648        .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC)
2649        .open(path)
2650}
2651
2652#[cfg(unix)]
2653fn open_snapshot_artifact(snapshot: &LatestSnapshot, id: &ArtifactId) -> std::io::Result<File> {
2654    use std::ffi::CString;
2655    use std::os::fd::{AsRawFd, FromRawFd};
2656
2657    let name = CString::new(id.as_str())
2658        .map_err(|_| std::io::Error::new(std::io::ErrorKind::InvalidInput, "name contains NUL"))?;
2659    // SAFETY: `generation_directory` pins a real directory descriptor and the
2660    // validated artifact ID is exactly one child component. O_NOFOLLOW rejects
2661    // a raced symlink at the leaf.
2662    let fd = unsafe {
2663        libc::openat(
2664            snapshot.generation_directory.as_raw_fd(),
2665            name.as_ptr(),
2666            libc::O_RDONLY | libc::O_CLOEXEC | libc::O_NOFOLLOW,
2667        )
2668    };
2669    if fd < 0 {
2670        return Err(std::io::Error::last_os_error());
2671    }
2672    // SAFETY: `openat` returned a new owned descriptor.
2673    let file = unsafe { File::from_raw_fd(fd) };
2674    if !file.metadata()?.is_file() {
2675        return Err(std::io::Error::new(
2676            std::io::ErrorKind::InvalidData,
2677            "snapshot artifact is not a regular file",
2678        ));
2679    }
2680    Ok(file)
2681}
2682
2683#[cfg(not(unix))]
2684fn open_snapshot_artifact(snapshot: &LatestSnapshot, id: &ArtifactId) -> std::io::Result<File> {
2685    File::open(snapshot.generation_path.join(id.as_str()))
2686}
2687
2688#[cfg(any(target_os = "macos", target_os = "ios"))]
2689fn rename_directory_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2690    use std::ffi::CString;
2691    use std::os::unix::ffi::OsStrExt;
2692
2693    let source = CString::new(source.as_os_str().as_bytes()).map_err(|_| {
2694        std::io::Error::new(std::io::ErrorKind::InvalidInput, "source contains NUL")
2695    })?;
2696    let destination = CString::new(destination.as_os_str().as_bytes()).map_err(|_| {
2697        std::io::Error::new(std::io::ErrorKind::InvalidInput, "destination contains NUL")
2698    })?;
2699    // SAFETY: both C strings remain alive for the duration of the call. The
2700    // platform flag provides an atomic no-replace directory rename.
2701    let result =
2702        unsafe { libc::renamex_np(source.as_ptr(), destination.as_ptr(), libc::RENAME_EXCL) };
2703    if result == 0 {
2704        Ok(())
2705    } else {
2706        Err(std::io::Error::last_os_error())
2707    }
2708}
2709
2710#[cfg(any(target_os = "linux", target_os = "android"))]
2711fn rename_directory_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2712    use std::ffi::CString;
2713    use std::os::unix::ffi::OsStrExt;
2714
2715    let source = CString::new(source.as_os_str().as_bytes()).map_err(|_| {
2716        std::io::Error::new(std::io::ErrorKind::InvalidInput, "source contains NUL")
2717    })?;
2718    let destination = CString::new(destination.as_os_str().as_bytes()).map_err(|_| {
2719        std::io::Error::new(std::io::ErrorKind::InvalidInput, "destination contains NUL")
2720    })?;
2721    // SAFETY: the arguments are valid C strings and AT_FDCWD selects their
2722    // absolute paths. RENAME_NOREPLACE makes collision handling atomic.
2723    let result = unsafe {
2724        libc::syscall(
2725            libc::SYS_renameat2,
2726            libc::AT_FDCWD,
2727            source.as_ptr(),
2728            libc::AT_FDCWD,
2729            destination.as_ptr(),
2730            libc::RENAME_NOREPLACE,
2731        )
2732    };
2733    if result == 0 {
2734        Ok(())
2735    } else {
2736        Err(std::io::Error::last_os_error())
2737    }
2738}
2739
2740#[cfg(not(any(
2741    target_os = "macos",
2742    target_os = "ios",
2743    target_os = "linux",
2744    target_os = "android",
2745    windows
2746)))]
2747fn rename_directory_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2748    // These targets retain the preflight check but cannot promise an atomic
2749    // no-replace rename through Rust's portable API.
2750    fs::rename(source, destination)
2751}
2752
2753#[cfg(windows)]
2754fn rename_directory_noreplace(source: &Path, destination: &Path) -> std::io::Result<()> {
2755    use std::os::windows::ffi::OsStrExt;
2756    use windows_sys::Win32::Storage::FileSystem::{MOVEFILE_WRITE_THROUGH, MoveFileExW};
2757
2758    let source: Vec<u16> = source
2759        .as_os_str()
2760        .encode_wide()
2761        .chain(std::iter::once(0))
2762        .collect();
2763    let destination: Vec<u16> = destination
2764        .as_os_str()
2765        .encode_wide()
2766        .chain(std::iter::once(0))
2767        .collect();
2768    // Without REPLACE_EXISTING, this is an atomic no-replace move. The
2769    // write-through flag requests durable completion before returning.
2770    let result = unsafe {
2771        MoveFileExW(
2772            source.as_ptr(),
2773            destination.as_ptr(),
2774            MOVEFILE_WRITE_THROUGH,
2775        )
2776    };
2777    if result != 0 {
2778        Ok(())
2779    } else {
2780        Err(std::io::Error::last_os_error())
2781    }
2782}
2783
2784#[cfg(windows)]
2785fn replace_file_atomically(source: &Path, destination: &Path) -> std::io::Result<()> {
2786    use std::os::windows::ffi::OsStrExt;
2787    use windows_sys::Win32::Storage::FileSystem::{
2788        MOVEFILE_REPLACE_EXISTING, MOVEFILE_WRITE_THROUGH, MoveFileExW,
2789    };
2790
2791    let source: Vec<u16> = source
2792        .as_os_str()
2793        .encode_wide()
2794        .chain(std::iter::once(0))
2795        .collect();
2796    let destination: Vec<u16> = destination
2797        .as_os_str()
2798        .encode_wide()
2799        .chain(std::iter::once(0))
2800        .collect();
2801    // SAFETY: both UTF-16 buffers are NUL-terminated and alive for the call.
2802    let result = unsafe {
2803        MoveFileExW(
2804            source.as_ptr(),
2805            destination.as_ptr(),
2806            MOVEFILE_REPLACE_EXISTING | MOVEFILE_WRITE_THROUGH,
2807        )
2808    };
2809    if result != 0 {
2810        Ok(())
2811    } else {
2812        Err(std::io::Error::last_os_error())
2813    }
2814}
2815
2816#[cfg(not(windows))]
2817fn replace_file_atomically(source: &Path, destination: &Path) -> std::io::Result<()> {
2818    fs::rename(source, destination)
2819}
2820
2821fn copy_new_synced(
2822    source: &Path,
2823    destination: &Path,
2824    operation: &'static str,
2825) -> Result<(), ArtifactPathError> {
2826    let source_metadata =
2827        fs::symlink_metadata(source).map_err(|source_error| ArtifactPathError::Io {
2828            operation,
2829            path: source.to_path_buf(),
2830            source: source_error,
2831        })?;
2832    if source_metadata.file_type().is_symlink() {
2833        return Err(ArtifactPathError::SymlinkComponent {
2834            path: source.to_path_buf(),
2835        });
2836    }
2837    if !source_metadata.is_file() {
2838        return Err(ArtifactPathError::FileDestinationNotFile {
2839            path: source.to_path_buf(),
2840        });
2841    }
2842    reject_hardlinked_regular_file(&source_metadata, source)?;
2843    let mut source_file = File::open(source).map_err(|source_error| ArtifactPathError::Io {
2844        operation,
2845        path: source.to_path_buf(),
2846        source: source_error,
2847    })?;
2848    let mut destination_file = OpenOptions::new()
2849        .write(true)
2850        .create_new(true)
2851        .open(destination)
2852        .map_err(|source_error| ArtifactPathError::Io {
2853            operation,
2854            path: destination.to_path_buf(),
2855            source: source_error,
2856        })?;
2857    std::io::copy(&mut source_file, &mut destination_file).map_err(|source_error| {
2858        ArtifactPathError::Io {
2859            operation,
2860            path: destination.to_path_buf(),
2861            source: source_error,
2862        }
2863    })?;
2864    destination_file
2865        .sync_all()
2866        .map_err(|source_error| ArtifactPathError::Io {
2867            operation: "sync copied artifact",
2868            path: destination.to_path_buf(),
2869            source: source_error,
2870        })?;
2871    record_durability_event("sync_file");
2872    Ok(())
2873}
2874
2875fn sync_artifact_tree(root: &Path) -> Result<(), ArtifactPathError> {
2876    sync_artifact_tree_recursive(root)
2877}
2878
2879#[cfg(unix)]
2880fn reject_hardlinked_regular_file(
2881    metadata: &fs::Metadata,
2882    path: &Path,
2883) -> Result<(), ArtifactPathError> {
2884    use std::os::unix::fs::MetadataExt;
2885
2886    if metadata.nlink() > 1 {
2887        return Err(ArtifactPathError::HardLinkedArtifact {
2888            path: path.to_path_buf(),
2889            links: metadata.nlink(),
2890        });
2891    }
2892    Ok(())
2893}
2894
2895#[cfg(not(unix))]
2896fn reject_hardlinked_regular_file(
2897    _metadata: &fs::Metadata,
2898    _path: &Path,
2899) -> Result<(), ArtifactPathError> {
2900    Ok(())
2901}
2902
2903fn sync_artifact_tree_recursive(directory: &Path) -> Result<(), ArtifactPathError> {
2904    let entries = fs::read_dir(directory)
2905        .map_err(|source| ArtifactPathError::Io {
2906            operation: "list artifact tree for sync",
2907            path: directory.to_path_buf(),
2908            source,
2909        })?
2910        .collect::<Result<Vec<_>, _>>()
2911        .map_err(|source| ArtifactPathError::Io {
2912            operation: "inspect artifact tree entry for sync",
2913            path: directory.to_path_buf(),
2914            source,
2915        })?;
2916    for entry in entries {
2917        let path = entry.path();
2918        let metadata = fs::symlink_metadata(&path).map_err(|source| ArtifactPathError::Io {
2919            operation: "inspect artifact tree entry for sync",
2920            path: path.clone(),
2921            source,
2922        })?;
2923        if metadata.file_type().is_symlink() {
2924            return Err(ArtifactPathError::SymlinkComponent { path });
2925        }
2926        if metadata.is_dir() {
2927            sync_artifact_tree_recursive(&path)?;
2928        } else if metadata.is_file() {
2929            File::open(&path)
2930                .and_then(|file| file.sync_all())
2931                .map_err(|source| ArtifactPathError::Io {
2932                    operation: "sync artifact file",
2933                    path: path.clone(),
2934                    source,
2935                })?;
2936            record_durability_event("sync_file");
2937        } else {
2938            return Err(ArtifactPathError::FileDestinationNotFile { path });
2939        }
2940    }
2941    sync_directory(directory)
2942}
2943
2944#[cfg(unix)]
2945fn sync_directory(path: &Path) -> Result<(), ArtifactPathError> {
2946    maybe_fail_directory_sync(path)?;
2947    File::open(path)
2948        .and_then(|directory| directory.sync_all())
2949        .map_err(|source| ArtifactPathError::Io {
2950            operation: "sync artifact directory",
2951            path: path.to_path_buf(),
2952            source,
2953        })?;
2954    record_durability_event("sync_directory");
2955    Ok(())
2956}
2957
2958#[cfg(not(unix))]
2959fn sync_directory(_path: &Path) -> Result<(), ArtifactPathError> {
2960    maybe_fail_directory_sync(_path)?;
2961    // Rust's portable filesystem API does not expose a directory handle that
2962    // can be flushed on every non-Unix target. File contents are still synced;
2963    // power-loss durability of directory entries is therefore weaker there.
2964    record_durability_event("directory_sync_unavailable");
2965    Ok(())
2966}
2967
2968#[cfg(test)]
2969thread_local! {
2970    static DURABILITY_EVENTS: std::cell::RefCell<Vec<&'static str>> = const {
2971        std::cell::RefCell::new(Vec::new())
2972    };
2973    static FAIL_ALIAS_INSTALL_AFTER: std::cell::Cell<Option<usize>> = const {
2974        std::cell::Cell::new(None)
2975    };
2976    static MUTATE_LATEST_SOURCE_AFTER_VALIDATION: std::cell::RefCell<Option<(PathBuf, Vec<u8>)>> =
2977        const { std::cell::RefCell::new(None) };
2978    static FAIL_DIRECTORY_SYNC_AFTER_EVENT: std::cell::Cell<Option<&'static str>> =
2979        const { std::cell::Cell::new(None) };
2980    static FAIL_NEXT_DIRECTORY_SYNC: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
2981}
2982
2983#[cfg(test)]
2984fn record_durability_event(event: &'static str) {
2985    DURABILITY_EVENTS.with(|events| events.borrow_mut().push(event));
2986    FAIL_DIRECTORY_SYNC_AFTER_EVENT.with(|target| {
2987        if target.get() == Some(event) {
2988            target.set(None);
2989            FAIL_NEXT_DIRECTORY_SYNC.with(|fail| fail.set(true));
2990        }
2991    });
2992}
2993
2994#[cfg(not(test))]
2995fn record_durability_event(_event: &'static str) {}
2996
2997#[cfg(test)]
2998fn maybe_fail_directory_sync(path: &Path) -> Result<(), ArtifactPathError> {
2999    let fail = FAIL_NEXT_DIRECTORY_SYNC.with(|fail| fail.replace(false));
3000    if fail {
3001        return Err(ArtifactPathError::Io {
3002            operation: "sync artifact directory",
3003            path: path.to_path_buf(),
3004            source: std::io::Error::other("injected post-commit directory sync failure"),
3005        });
3006    }
3007    Ok(())
3008}
3009
3010#[cfg(not(test))]
3011fn maybe_fail_directory_sync(_path: &Path) -> Result<(), ArtifactPathError> {
3012    Ok(())
3013}
3014
3015#[cfg(test)]
3016fn maybe_fail_alias_install(installed: usize, path: &Path) -> Result<(), ArtifactPathError> {
3017    let should_fail = FAIL_ALIAS_INSTALL_AFTER.with(|fail_after| {
3018        let should_fail = fail_after.get() == Some(installed);
3019        if should_fail {
3020            fail_after.set(None);
3021        }
3022        should_fail
3023    });
3024    if should_fail {
3025        return Err(ArtifactPathError::Io {
3026            operation: "install stable latest alias",
3027            path: path.to_path_buf(),
3028            source: std::io::Error::other("injected stable-alias failure"),
3029        });
3030    }
3031    Ok(())
3032}
3033
3034#[cfg(test)]
3035fn maybe_mutate_latest_source_after_validation() {
3036    MUTATE_LATEST_SOURCE_AFTER_VALIDATION.with(|mutation| {
3037        if let Some((path, contents)) = mutation.borrow_mut().take() {
3038            fs::write(path, contents).expect("apply injected post-validation source mutation");
3039        }
3040    });
3041}
3042
3043#[cfg(not(test))]
3044fn maybe_mutate_latest_source_after_validation() {}
3045
3046#[cfg(not(test))]
3047fn maybe_fail_alias_install(_installed: usize, _path: &Path) -> Result<(), ArtifactPathError> {
3048    Ok(())
3049}
3050
3051fn active_workspaces() -> &'static Mutex<BTreeSet<PathBuf>> {
3052    ACTIVE_WORKSPACES.get_or_init(|| Mutex::new(BTreeSet::new()))
3053}
3054
3055fn register_active_workspace(path: &Path) {
3056    active_workspaces()
3057        .lock()
3058        .unwrap_or_else(std::sync::PoisonError::into_inner)
3059        .insert(path.to_path_buf());
3060}
3061
3062fn unregister_active_workspace(path: &Path) {
3063    active_workspaces()
3064        .lock()
3065        .unwrap_or_else(std::sync::PoisonError::into_inner)
3066        .remove(path);
3067}
3068
3069fn is_active_workspace(path: &Path) -> bool {
3070    active_workspaces()
3071        .lock()
3072        .unwrap_or_else(std::sync::PoisonError::into_inner)
3073        .contains(path)
3074}
3075
3076fn workspace_nonce() -> String {
3077    let timestamp = SystemTime::now()
3078        .duration_since(UNIX_EPOCH)
3079        .unwrap_or_default()
3080        .as_nanos();
3081    let sequence = WORKSPACE_SEQUENCE.fetch_add(1, Ordering::Relaxed);
3082    format!("{timestamp:x}-{:x}-{sequence:x}", std::process::id())
3083}
3084
3085/// An existing directory approved as the boundary for artifact writes.
3086#[derive(Debug, Clone, PartialEq, Eq)]
3087pub struct ApprovedRoot {
3088    path: PathBuf,
3089}
3090
3091impl ApprovedRoot {
3092    /// Approve an existing, non-symlink directory and retain its canonical path.
3093    pub fn existing(path: impl AsRef<Path>) -> Result<Self, ArtifactPathError> {
3094        let path = path.as_ref();
3095        let metadata = fs::symlink_metadata(path).map_err(|source| ArtifactPathError::Io {
3096            operation: "inspect approved root",
3097            path: path.to_path_buf(),
3098            source,
3099        })?;
3100        if metadata.file_type().is_symlink() {
3101            return Err(ArtifactPathError::SymlinkComponent {
3102                path: path.to_path_buf(),
3103            });
3104        }
3105        if !metadata.is_dir() {
3106            return Err(ArtifactPathError::RootNotDirectory(path.to_path_buf()));
3107        }
3108
3109        let path = fs::canonicalize(path).map_err(|source| ArtifactPathError::Io {
3110            operation: "canonicalize approved root",
3111            path: path.to_path_buf(),
3112            source,
3113        })?;
3114        Ok(Self { path })
3115    }
3116
3117    /// Return the canonical approved root.
3118    pub fn path(&self) -> &Path {
3119        &self.path
3120    }
3121
3122    /// Resolve a relative directory beneath this root without creating it.
3123    ///
3124    /// Existing components must be real directories, never symbolic links.
3125    /// Once a missing component is reached there cannot be a deeper existing
3126    /// component, so the remaining validated path can be projected safely.
3127    pub fn project_dir(&self, relative: impl AsRef<Path>) -> Result<PathBuf, ArtifactPathError> {
3128        let relative = relative.as_ref();
3129        validate_relative_path(relative)?;
3130
3131        let mut current = self.path.clone();
3132        for component in relative.components() {
3133            let std::path::Component::Normal(name) = component else {
3134                unreachable!("relative path was validated before directory projection")
3135            };
3136            current.push(name);
3137
3138            match fs::symlink_metadata(&current) {
3139                Ok(metadata) if metadata.file_type().is_symlink() => {
3140                    return Err(ArtifactPathError::SymlinkComponent {
3141                        path: current.clone(),
3142                    });
3143                }
3144                Ok(metadata) if !metadata.is_dir() => {
3145                    return Err(ArtifactPathError::DirectoryComponentNotDirectory {
3146                        path: current.clone(),
3147                    });
3148                }
3149                Ok(_) => {}
3150                Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
3151                    return Ok(self.path.join(relative));
3152                }
3153                Err(source) => {
3154                    return Err(ArtifactPathError::Io {
3155                        operation: "inspect projected artifact directory",
3156                        path: current,
3157                        source,
3158                    });
3159                }
3160            }
3161        }
3162
3163        Ok(current)
3164    }
3165
3166    /// Prepare a relative directory, rejecting traversal and symlink components.
3167    pub fn prepare_dir(&self, relative: impl AsRef<Path>) -> Result<PathBuf, ArtifactPathError> {
3168        let relative = relative.as_ref();
3169        validate_relative_path(relative)?;
3170        prepare_directory_components(&self.path, relative)
3171    }
3172
3173    /// Prepare a relative file's parent directories and validate its destination.
3174    ///
3175    /// This does not create or truncate the file. An existing destination must
3176    /// be a regular file and must not be a symlink.
3177    pub fn prepare_file(&self, relative: impl AsRef<Path>) -> Result<PathBuf, ArtifactPathError> {
3178        let relative = relative.as_ref();
3179        validate_relative_path(relative)?;
3180        if let Some(parent) = relative
3181            .parent()
3182            .filter(|parent| !parent.as_os_str().is_empty())
3183        {
3184            prepare_directory_components(&self.path, parent)?;
3185        }
3186        let path = self.path.join(relative);
3187        match fs::symlink_metadata(&path) {
3188            Ok(metadata) if metadata.file_type().is_symlink() => {
3189                Err(ArtifactPathError::SymlinkComponent { path })
3190            }
3191            Ok(metadata) if !metadata.is_file() => {
3192                Err(ArtifactPathError::FileDestinationNotFile { path })
3193            }
3194            Ok(_) => Ok(path),
3195            Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(path),
3196            Err(source) => Err(ArtifactPathError::Io {
3197                operation: "inspect artifact file",
3198                path,
3199                source,
3200            }),
3201        }
3202    }
3203}
3204
3205/// Failure to approve a root or safely prepare one of its descendants.
3206#[derive(Debug, Error)]
3207pub enum ArtifactPathError {
3208    #[error("artifact identifier must be one non-empty relative path component: {value:?}")]
3209    InvalidArtifactId { value: String },
3210    #[error("could not allocate a unique workspace for logical run ID {logical_id}")]
3211    WorkspaceAllocationExhausted { logical_id: ArtifactId },
3212    #[error("artifact publication destination already exists: {path}")]
3213    PublicationDestinationExists { path: PathBuf },
3214    #[error("staged run already contains the reserved manifest path: {path}")]
3215    ReservedManifestPath { path: PathBuf },
3216    #[error("required staged artifact is missing: {path}")]
3217    RequiredArtifactMissing { path: PathBuf },
3218    #[error("latest artifact destination was requested more than once: {id}")]
3219    DuplicateLatestDestination { id: ArtifactId },
3220    #[error("could not allocate a unique latest-artifact transaction")]
3221    LatestTransactionAllocationExhausted,
3222    #[error("timed out waiting for the latest-artifact update lock at {path}")]
3223    LatestLockTimeout { path: PathBuf },
3224    #[error("no committed latest snapshot is available at {path}")]
3225    LatestSnapshotUnavailable { path: PathBuf },
3226    #[error("latest pointer is not one valid generation identity: {path}")]
3227    InvalidLatestPointer { path: PathBuf },
3228    #[error("latest snapshot does not contain alias {id}")]
3229    LatestArtifactNotInSnapshot { id: ArtifactId },
3230    #[error("latest generation is corrupt at {path}: {detail}")]
3231    CorruptGeneration { path: PathBuf, detail: String },
3232    #[error("latest-generation chain is broken at {generation}: missing {predecessor}")]
3233    BrokenGenerationChain {
3234        generation: String,
3235        predecessor: String,
3236    },
3237    #[error("latest-generation chain contains a cycle at {generation}")]
3238    GenerationChainCycle { generation: String },
3239    #[error("latest-generation chain forks after {predecessor} into {children:?}")]
3240    AmbiguousGenerationFork {
3241        predecessor: String,
3242        children: Vec<String>,
3243    },
3244    #[error("latest-generation recovery has ambiguous tips: {tips:?}")]
3245    AmbiguousGenerationTips { tips: Vec<String> },
3246    #[error(
3247        "published run has stale latest-generation token: expected {expected:?}, observed {observed:?}"
3248    )]
3249    StaleLatestGeneration {
3250        expected: Option<String>,
3251        observed: Option<String>,
3252    },
3253    #[error("refusing latest-generation producer downgrade from {existing} to {attempted}")]
3254    ProducerDowngrade { existing: String, attempted: String },
3255    #[error("invalid artifact producer version {version:?}: {source}")]
3256    InvalidProducerVersion {
3257        version: String,
3258        #[source]
3259        source: semver::Error,
3260    },
3261    #[error("unsupported manifest version {found} at {path}; supported version is {supported}")]
3262    UnknownManifestVersion {
3263        path: PathBuf,
3264        found: u32,
3265        supported: u32,
3266    },
3267    #[error("artifact manifest integrity check failed at {path}: {detail}")]
3268    ManifestIntegrity { path: PathBuf, detail: String },
3269    #[error("could not encode or decode artifact manifest at {path}: {source}")]
3270    ManifestJson {
3271        path: PathBuf,
3272        #[source]
3273        source: serde_json::Error,
3274    },
3275    #[error(
3276        "latest-artifact update failed ({original}) and rollback was incomplete ({rollback}); recovery material retained at {recovery_path}"
3277    )]
3278    LatestRollback {
3279        original: Box<ArtifactPathError>,
3280        rollback: Box<ArtifactPathError>,
3281        /// Transaction directory retained with any unrestored backups.
3282        recovery_path: PathBuf,
3283    },
3284    #[error(
3285        "latest-generation preparation failed ({original}) and staging cleanup failed ({cleanup}) at {staging_path}"
3286    )]
3287    LatestStagingCleanup {
3288        original: Box<ArtifactPathError>,
3289        cleanup: Box<ArtifactPathError>,
3290        staging_path: PathBuf,
3291    },
3292    #[error("run is committed at {path}, but directory durability is uncertain: {source}")]
3293    PublicationDurabilityUncertain {
3294        path: PathBuf,
3295        #[source]
3296        source: Box<ArtifactPathError>,
3297    },
3298    #[error("run is committed at {path}, but post-commit retention maintenance failed: {source}")]
3299    PublicationPostCommitMaintenance {
3300        path: PathBuf,
3301        #[source]
3302        source: Box<ArtifactPathError>,
3303    },
3304    #[error(
3305        "latest generation {generation} is committed, but pointer durability is uncertain: {source}"
3306    )]
3307    LatestCommitDurabilityUncertain {
3308        generation: String,
3309        #[source]
3310        source: Box<ArtifactPathError>,
3311    },
3312    #[error("latest generation {generation} is committed, but retention failed: {source}")]
3313    LatestRetentionAfterCommit {
3314        generation: String,
3315        #[source]
3316        source: Box<ArtifactPathError>,
3317    },
3318    #[error("approved artifact root is not a directory: {0}")]
3319    RootNotDirectory(PathBuf),
3320    #[error(
3321        "artifact child path must be a non-empty relative path without parent traversal: {path}"
3322    )]
3323    InvalidRelativePath { path: PathBuf },
3324    #[error("artifact path component is a symbolic link: {path}")]
3325    SymlinkComponent { path: PathBuf },
3326    #[error("artifact directory component is not a directory: {path}")]
3327    DirectoryComponentNotDirectory { path: PathBuf },
3328    #[error("artifact file destination is not a regular file: {path}")]
3329    FileDestinationNotFile { path: PathBuf },
3330    #[error("artifact file has {links} hard links and is not publication-private: {path}")]
3331    HardLinkedArtifact { path: PathBuf, links: u64 },
3332    #[error("artifact path is not valid UTF-8 and cannot be recorded in a manifest: {path:?}")]
3333    NonUtf8ArtifactPath { path: PathBuf },
3334    #[error("failed to {operation} at {path}: {source}")]
3335    Io {
3336        operation: &'static str,
3337        path: PathBuf,
3338        #[source]
3339        source: std::io::Error,
3340    },
3341}
3342
3343fn validate_relative_path(path: &Path) -> Result<(), ArtifactPathError> {
3344    if path.as_os_str().is_empty()
3345        || path
3346            .components()
3347            .any(|component| !matches!(component, std::path::Component::Normal(_)))
3348    {
3349        return Err(ArtifactPathError::InvalidRelativePath {
3350            path: path.to_path_buf(),
3351        });
3352    }
3353    Ok(())
3354}
3355
3356fn prepare_directory_components(
3357    root: &Path,
3358    relative: &Path,
3359) -> Result<PathBuf, ArtifactPathError> {
3360    let mut current = root.to_path_buf();
3361
3362    for component in relative.components() {
3363        let std::path::Component::Normal(name) = component else {
3364            unreachable!("relative path was validated before directory preparation")
3365        };
3366        current.push(name);
3367
3368        match fs::symlink_metadata(&current) {
3369            Ok(metadata) if metadata.file_type().is_symlink() => {
3370                return Err(ArtifactPathError::SymlinkComponent {
3371                    path: current.clone(),
3372                });
3373            }
3374            Ok(metadata) if !metadata.is_dir() => {
3375                return Err(ArtifactPathError::DirectoryComponentNotDirectory {
3376                    path: current.clone(),
3377                });
3378            }
3379            Ok(_) => {}
3380            Err(source) if source.kind() == std::io::ErrorKind::NotFound => {
3381                match fs::create_dir(&current) {
3382                    Ok(()) => {}
3383                    Err(source) if source.kind() == std::io::ErrorKind::AlreadyExists => {
3384                        match fs::symlink_metadata(&current) {
3385                            Ok(metadata) if metadata.file_type().is_symlink() => {
3386                                return Err(ArtifactPathError::SymlinkComponent {
3387                                    path: current.clone(),
3388                                });
3389                            }
3390                            Ok(metadata) if !metadata.is_dir() => {
3391                                return Err(ArtifactPathError::DirectoryComponentNotDirectory {
3392                                    path: current.clone(),
3393                                });
3394                            }
3395                            Ok(_) => {}
3396                            Err(source) => {
3397                                return Err(ArtifactPathError::Io {
3398                                    operation: "inspect raced artifact directory",
3399                                    path: current.clone(),
3400                                    source,
3401                                });
3402                            }
3403                        }
3404                    }
3405                    Err(source) => {
3406                        return Err(ArtifactPathError::Io {
3407                            operation: "create artifact directory",
3408                            path: current.clone(),
3409                            source,
3410                        });
3411                    }
3412                }
3413            }
3414            Err(source) => {
3415                return Err(ArtifactPathError::Io {
3416                    operation: "inspect artifact directory",
3417                    path: current.clone(),
3418                    source,
3419                });
3420            }
3421        }
3422    }
3423
3424    Ok(current)
3425}
3426
3427#[cfg(test)]
3428mod tests {
3429    use super::*;
3430
3431    #[test]
3432    fn artifact_id_rejects_an_empty_component() {
3433        assert!(matches!(
3434            ArtifactId::new(""),
3435            Err(ArtifactPathError::InvalidArtifactId { .. })
3436        ));
3437    }
3438
3439    #[test]
3440    fn artifact_id_rejects_values_that_are_not_one_normal_path_component() {
3441        for invalid in [
3442            ".",
3443            "..",
3444            "/absolute",
3445            "nested/child",
3446            r"nested\child",
3447            r"C:\absolute",
3448        ] {
3449            assert!(
3450                matches!(
3451                    ArtifactId::new(invalid),
3452                    Err(ArtifactPathError::InvalidArtifactId { .. })
3453                ),
3454                "accepted invalid artifact ID {invalid:?}"
3455            );
3456        }
3457    }
3458
3459    #[test]
3460    fn repeated_run_workspace_allocations_are_isolated() {
3461        let temp = tempfile::tempdir().expect("tempdir");
3462        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3463
3464        let first = RunWorkspace::allocate(temp.path(), &logical_id).expect("first workspace");
3465        let second = RunWorkspace::allocate(temp.path(), &logical_id).expect("second workspace");
3466
3467        assert_ne!(first.staging_path(), second.staging_path());
3468        assert_ne!(first.published_path(), second.published_path());
3469        assert!(first.staging_path().is_dir());
3470        assert!(second.staging_path().is_dir());
3471        assert!(!first.published_path().exists());
3472        assert!(!second.published_path().exists());
3473    }
3474
3475    #[test]
3476    fn concurrent_run_workspace_allocations_are_isolated() {
3477        use std::collections::BTreeSet;
3478        use std::sync::{Arc, Barrier};
3479
3480        let temp = tempfile::tempdir().expect("tempdir");
3481        let root = Arc::new(temp.path().to_path_buf());
3482        let barrier = Arc::new(Barrier::new(8));
3483        let handles: Vec<_> = (0..8)
3484            .map(|_| {
3485                let root = Arc::clone(&root);
3486                let barrier = Arc::clone(&barrier);
3487                std::thread::spawn(move || {
3488                    let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3489                    barrier.wait();
3490                    let workspace =
3491                        RunWorkspace::allocate(root.as_ref(), &logical_id).expect("workspace");
3492                    (
3493                        workspace.staging_path().to_path_buf(),
3494                        workspace.published_path().to_path_buf(),
3495                    )
3496                })
3497            })
3498            .collect();
3499
3500        let paths: Vec<_> = handles
3501            .into_iter()
3502            .map(|handle| handle.join().expect("allocation thread"))
3503            .collect();
3504        let staging_paths: BTreeSet<_> = paths.iter().map(|(path, _)| path).collect();
3505        let published_paths: BTreeSet<_> = paths.iter().map(|(_, path)| path).collect();
3506
3507        assert_eq!(staging_paths.len(), paths.len());
3508        assert_eq!(published_paths.len(), paths.len());
3509    }
3510
3511    #[test]
3512    fn repeated_published_runs_do_not_inherit_stale_artifacts() {
3513        let temp = tempfile::tempdir().expect("tempdir");
3514        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3515        let required = [
3516            ArtifactId::new("profile.json").expect("valid profile ID"),
3517            ArtifactId::new("summary.md").expect("valid summary ID"),
3518        ];
3519        let first = RunWorkspace::allocate(temp.path(), &logical_id).expect("first workspace");
3520        fs::write(first.staging_path().join("profile.json"), "first profile")
3521            .expect("write first profile");
3522        fs::write(first.staging_path().join("summary.md"), "first summary")
3523            .expect("write first summary");
3524        fs::write(first.staging_path().join("stale.data"), "stale").expect("write stale artifact");
3525        let first = first.publish(&required).expect("publish first run");
3526
3527        let second = RunWorkspace::allocate(temp.path(), &logical_id).expect("second workspace");
3528        fs::write(second.staging_path().join("profile.json"), "second profile")
3529            .expect("write second profile");
3530        fs::write(second.staging_path().join("summary.md"), "second summary")
3531            .expect("write second summary");
3532        let second = second.publish(&required).expect("publish second run");
3533
3534        assert_ne!(first.path(), second.path());
3535        assert!(first.path().join("stale.data").is_file());
3536        assert!(!second.path().join("stale.data").exists());
3537    }
3538
3539    #[cfg(unix)]
3540    #[test]
3541    fn run_workspace_rejects_a_symlink_publication_collision() {
3542        use std::os::unix::fs::symlink;
3543
3544        let temp = tempfile::tempdir().expect("tempdir");
3545        let outside = tempfile::tempdir().expect("outside tempdir");
3546        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3547        let workspace =
3548            RunWorkspace::allocate(temp.path(), &logical_id).expect("allocate workspace");
3549        fs::write(workspace.staging_path().join("profile.json"), "complete")
3550            .expect("write staged manifest");
3551        symlink(outside.path(), workspace.published_path()).expect("create collision symlink");
3552
3553        let required = [ArtifactId::new("profile.json").expect("valid required file")];
3554        let error = workspace
3555            .publish(&required)
3556            .expect_err("reject symlink collision");
3557
3558        assert!(matches!(error, ArtifactPathError::SymlinkComponent { .. }));
3559        assert!(outside.path().is_dir());
3560        assert!(!outside.path().join("profile.json").exists());
3561    }
3562
3563    #[test]
3564    fn incomplete_run_workspace_does_not_replace_latest_artifacts() {
3565        let temp = tempfile::tempdir().expect("tempdir");
3566        fs::write(temp.path().join("profile.json"), "old profile").expect("seed latest profile");
3567        fs::write(temp.path().join("summary.md"), "old summary").expect("seed latest summary");
3568        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3569        let workspace =
3570            RunWorkspace::allocate(temp.path(), &logical_id).expect("allocate workspace");
3571        let published_path = workspace.published_path().to_path_buf();
3572        fs::write(workspace.staging_path().join("profile.json"), "new profile")
3573            .expect("write staged profile");
3574        let required = [
3575            ArtifactId::new("profile.json").expect("valid profile ID"),
3576            ArtifactId::new("summary.md").expect("valid summary ID"),
3577        ];
3578
3579        let error = workspace
3580            .publish(&required)
3581            .expect_err("reject incomplete run");
3582
3583        assert!(matches!(
3584            error,
3585            ArtifactPathError::RequiredArtifactMissing { .. }
3586        ));
3587        assert!(!published_path.exists());
3588        assert_eq!(
3589            fs::read_to_string(temp.path().join("profile.json")).expect("read latest profile"),
3590            "old profile"
3591        );
3592        assert_eq!(
3593            fs::read_to_string(temp.path().join("summary.md")).expect("read latest summary"),
3594            "old summary"
3595        );
3596    }
3597
3598    #[test]
3599    fn published_run_refreshes_stable_latest_artifacts() {
3600        let temp = tempfile::tempdir().expect("tempdir");
3601        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3602        let workspace =
3603            RunWorkspace::allocate(temp.path(), &logical_id).expect("allocate workspace");
3604        fs::write(workspace.staging_path().join("profile.json"), "new profile")
3605            .expect("write staged profile");
3606        fs::write(workspace.staging_path().join("summary.md"), "new summary")
3607            .expect("write staged summary");
3608        let required = [
3609            ArtifactId::new("profile.json").expect("valid profile ID"),
3610            ArtifactId::new("summary.md").expect("valid summary ID"),
3611        ];
3612
3613        let published = workspace.publish(&required).expect("publish complete run");
3614        published
3615            .refresh_latest(&[
3616                LatestArtifact::same(required[0].clone()),
3617                LatestArtifact::same(required[1].clone()),
3618            ])
3619            .expect("refresh latest artifacts");
3620
3621        assert!(published.path().join("profile.json").is_file());
3622        assert!(published.path().join("summary.md").is_file());
3623        assert_eq!(
3624            fs::read_to_string(temp.path().join("profile.json")).expect("read latest profile"),
3625            "new profile"
3626        );
3627        assert_eq!(
3628            fs::read_to_string(temp.path().join("summary.md")).expect("read latest summary"),
3629            "new summary"
3630        );
3631    }
3632
3633    #[test]
3634    fn latest_refresh_removes_aliases_omitted_by_the_next_generation() {
3635        let temp = tempfile::tempdir().expect("tempdir");
3636        let first = publish_pair(temp.path(), "first-run", "first");
3637        refresh_pair(&first).expect("refresh initial aliases");
3638        assert!(temp.path().join("summary.md").is_file());
3639
3640        let logical_id = ArtifactId::new("profile-only-run").expect("logical ID");
3641        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
3642        fs::write(workspace.staging_path().join("profile.json"), "second").expect("write profile");
3643        let profile = ArtifactId::new("profile.json").expect("profile ID");
3644        let published = workspace
3645            .publish(std::slice::from_ref(&profile))
3646            .expect("publish profile-only run");
3647        published
3648            .refresh_latest(&[LatestArtifact::same(profile)])
3649            .expect("refresh profile-only aliases");
3650
3651        assert_eq!(
3652            fs::read_to_string(temp.path().join("profile.json")).expect("profile alias"),
3653            "second"
3654        );
3655        assert!(
3656            !temp.path().join("summary.md").exists(),
3657            "obsolete alias survived the generation change"
3658        );
3659    }
3660
3661    #[test]
3662    fn latest_generation_retention_keeps_a_bounded_valid_chain() {
3663        let temp = tempfile::tempdir().expect("tempdir");
3664        for index in 0..(RETAIN_LATEST_GENERATIONS + 4) {
3665            let published = publish_pair(
3666                temp.path(),
3667                &format!("retained-run-{index}"),
3668                &format!("receipt-{index}"),
3669            );
3670            refresh_pair(&published).expect("refresh retained generation");
3671        }
3672
3673        let generations = temp
3674            .path()
3675            .join(LATEST_DIRECTORY)
3676            .join(LATEST_GENERATIONS_DIRECTORY);
3677        assert_eq!(
3678            fs::read_dir(&generations)
3679                .expect("list retained generations")
3680                .count(),
3681            RETAIN_LATEST_GENERATIONS
3682        );
3683        let snapshot = LatestSnapshot::open(temp.path()).expect("open retained latest snapshot");
3684        assert_eq!(
3685            snapshot
3686                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
3687                .expect("read retained latest artifact"),
3688            format!("receipt-{}", RETAIN_LATEST_GENERATIONS + 3).as_bytes()
3689        );
3690        validate_predecessor_chain(
3691            &ApprovedRoot::existing(temp.path().join(LATEST_DIRECTORY)).expect("open latest root"),
3692            snapshot.manifest(),
3693        )
3694        .expect("retained predecessor chain remains valid");
3695    }
3696
3697    #[test]
3698    fn latest_generation_retention_defers_while_an_old_reader_is_leased() {
3699        let temp = tempfile::tempdir().expect("tempdir");
3700        let first = publish_pair(temp.path(), "leased-run-0", "receipt-0");
3701        refresh_pair(&first).expect("refresh first generation");
3702        let old_reader = LatestSnapshot::open(temp.path()).expect("lease first generation");
3703
3704        for index in 1..=RETAIN_LATEST_GENERATIONS {
3705            let published = publish_pair(
3706                temp.path(),
3707                &format!("leased-run-{index}"),
3708                &format!("receipt-{index}"),
3709            );
3710            refresh_pair(&published).expect("refresh while old reader is leased");
3711        }
3712        let generations = temp
3713            .path()
3714            .join(LATEST_DIRECTORY)
3715            .join(LATEST_GENERATIONS_DIRECTORY);
3716        assert_eq!(
3717            fs::read_dir(&generations)
3718                .expect("list deferred generations")
3719                .count(),
3720            RETAIN_LATEST_GENERATIONS + 1
3721        );
3722        assert_eq!(
3723            old_reader
3724                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
3725                .expect("read leased generation"),
3726            b"receipt-0"
3727        );
3728        drop(old_reader);
3729
3730        let final_run = publish_pair(temp.path(), "leased-run-final", "receipt-final");
3731        refresh_pair(&final_run).expect("refresh after releasing old reader");
3732        assert_eq!(
3733            fs::read_dir(&generations)
3734                .expect("list pruned generations")
3735                .count(),
3736            RETAIN_LATEST_GENERATIONS
3737        );
3738    }
3739
3740    #[test]
3741    fn published_run_retention_is_bounded_and_defers_leased_runs() {
3742        let temp = tempfile::tempdir().expect("tempdir");
3743        let leased = publish_pair(temp.path(), "published-run-0", "receipt-0");
3744        for index in 1..=RETAIN_PUBLISHED_RUNS {
3745            drop(publish_pair(
3746                temp.path(),
3747                &format!("published-run-{index}"),
3748                &format!("receipt-{index}"),
3749            ));
3750        }
3751
3752        let count_runs = || {
3753            fs::read_dir(temp.path())
3754                .expect("list artifact root")
3755                .filter_map(Result::ok)
3756                .filter(|entry| entry.path().join(RUN_MANIFEST_FILE).is_file())
3757                .count()
3758        };
3759        assert_eq!(count_runs(), RETAIN_PUBLISHED_RUNS + 1);
3760        assert_eq!(
3761            fs::read_to_string(leased.path().join("profile.json")).expect("read leased run"),
3762            "receipt-0"
3763        );
3764        drop(leased);
3765
3766        drop(publish_pair(
3767            temp.path(),
3768            "published-run-final",
3769            "receipt-final",
3770        ));
3771        assert_eq!(count_runs(), RETAIN_PUBLISHED_RUNS);
3772    }
3773
3774    #[test]
3775    fn abandoned_workspace_quarantine_is_bounded() {
3776        let temp = tempfile::tempdir().expect("tempdir");
3777        let staging = temp.path().join(STAGING_DIRECTORY);
3778        fs::create_dir_all(&staging).expect("create staging root");
3779        for index in 0..(RETAIN_QUARANTINE_ENTRIES + 5) {
3780            let crashed = staging.join(format!("crashed-{index}"));
3781            fs::create_dir(&crashed).expect("create crashed workspace");
3782            fs::write(crashed.join(WORKSPACE_LOCK_FILE), "").expect("write released lock");
3783        }
3784
3785        let workspace = RunWorkspace::allocate(
3786            temp.path(),
3787            &ArtifactId::new("quarantine-writer").expect("logical ID"),
3788        )
3789        .expect("recover abandoned workspaces");
3790        assert_eq!(
3791            fs::read_dir(temp.path().join(STAGING_QUARANTINE_DIRECTORY))
3792                .expect("list bounded quarantine")
3793                .count(),
3794            RETAIN_QUARANTINE_ENTRIES
3795        );
3796        drop(workspace);
3797    }
3798
3799    #[test]
3800    fn concurrent_latest_refreshes_leave_one_complete_alias_pair() {
3801        use std::sync::{Arc, Barrier};
3802
3803        const WRITERS: usize = 8;
3804        let temp = tempfile::tempdir().expect("tempdir");
3805        let required = [
3806            ArtifactId::new("profile.json").expect("valid profile ID"),
3807            ArtifactId::new("summary.md").expect("valid summary ID"),
3808        ];
3809        let mut published = Vec::new();
3810        for index in 0..WRITERS {
3811            let logical_id = ArtifactId::new(format!("run-{index}")).expect("valid logical ID");
3812            let workspace =
3813                RunWorkspace::allocate(temp.path(), &logical_id).expect("allocate workspace");
3814            let receipt = format!("run-{index}");
3815            fs::write(workspace.staging_path().join("profile.json"), &receipt)
3816                .expect("write staged profile");
3817            fs::write(workspace.staging_path().join("summary.md"), &receipt)
3818                .expect("write staged summary");
3819            published.push(workspace.publish(&required).expect("publish complete run"));
3820        }
3821
3822        let barrier = Arc::new(Barrier::new(WRITERS));
3823        let handles: Vec<_> = published
3824            .into_iter()
3825            .map(|published| {
3826                let barrier = Arc::clone(&barrier);
3827                let required = required.clone();
3828                std::thread::spawn(move || {
3829                    barrier.wait();
3830                    published.refresh_latest(&[
3831                        LatestArtifact::same(required[0].clone()),
3832                        LatestArtifact::same(required[1].clone()),
3833                    ])
3834                })
3835            })
3836            .collect();
3837
3838        let mut committed = 0;
3839        let mut rejected_as_stale = 0;
3840        for handle in handles {
3841            match handle.join().expect("latest refresh thread") {
3842                Ok(()) => committed += 1,
3843                Err(ArtifactPathError::StaleLatestGeneration { .. }) => rejected_as_stale += 1,
3844                Err(error) => panic!("unexpected latest refresh error: {error}"),
3845            }
3846        }
3847        assert_eq!(committed, 1, "exactly one shared CAS token may commit");
3848        assert_eq!(rejected_as_stale, WRITERS - 1);
3849
3850        let profile =
3851            fs::read_to_string(temp.path().join("profile.json")).expect("read latest profile");
3852        let summary =
3853            fs::read_to_string(temp.path().join("summary.md")).expect("read latest summary");
3854        assert_eq!(profile, summary, "latest aliases came from different runs");
3855        assert!(profile.starts_with("run-"));
3856    }
3857
3858    #[test]
3859    fn rollback_restores_other_backups_after_installed_file_removal_fails() {
3860        let temp = tempfile::tempdir().expect("tempdir");
3861        let blocked_destination = temp.path().join("blocked");
3862        let restorable_destination = temp.path().join("restorable");
3863        let blocked_backup = temp.path().join("blocked.backup");
3864        let restorable_backup = temp.path().join("restorable.backup");
3865        fs::create_dir(&blocked_destination).expect("create removal failure directory");
3866        fs::write(&restorable_destination, "new").expect("write installed file");
3867        fs::write(&blocked_backup, "old blocked").expect("write blocked backup");
3868        fs::write(&restorable_backup, "old restored").expect("write restorable backup");
3869        let original = ArtifactPathError::Io {
3870            operation: "install latest artifact",
3871            path: restorable_destination.clone(),
3872            source: std::io::Error::other("simulated install failure"),
3873        };
3874
3875        let error = latest_update_failure(
3876            original,
3877            &[restorable_destination.clone(), blocked_destination.clone()],
3878            &[
3879                (blocked_destination.clone(), blocked_backup.clone()),
3880                (restorable_destination.clone(), restorable_backup),
3881            ],
3882            temp.path().join("retained-rollback"),
3883        );
3884
3885        assert!(matches!(error, ArtifactPathError::LatestRollback { .. }));
3886        assert_eq!(
3887            fs::read_to_string(&restorable_destination).expect("read restored backup"),
3888            "old restored"
3889        );
3890        assert!(blocked_backup.is_file());
3891    }
3892
3893    #[test]
3894    fn incomplete_transaction_rollback_retains_backup_directory() {
3895        let temp = tempfile::tempdir().expect("tempdir");
3896        let transaction_path = temp.path().join("alias-transaction");
3897        let backup_root = transaction_path.join("backup");
3898        fs::create_dir_all(&backup_root).expect("create transaction backup root");
3899        let destination = temp.path().join("profile.json");
3900        fs::create_dir(&destination).expect("create blocking destination directory");
3901        let backup = backup_root.join("profile.json");
3902        fs::write(&backup, "previous").expect("write backup");
3903        let transaction = StableAliasTransaction {
3904            path: transaction_path.clone(),
3905            root_path: temp.path().to_path_buf(),
3906            installed: vec![destination.clone()],
3907            backups: vec![(destination, backup.clone())],
3908        };
3909        let original = ArtifactPathError::Io {
3910            operation: "install stable latest alias",
3911            path: temp.path().join("profile.json"),
3912            source: std::io::Error::other("injected install failure"),
3913        };
3914
3915        let error = transaction.rollback(original);
3916        assert!(matches!(
3917            error,
3918            ArtifactPathError::LatestRollback {
3919                ref recovery_path,
3920                ..
3921            } if recovery_path == &transaction_path
3922        ));
3923        assert!(transaction_path.is_dir());
3924        assert_eq!(
3925            fs::read_to_string(backup).expect("read retained backup"),
3926            "previous"
3927        );
3928    }
3929
3930    #[test]
3931    fn latest_refresh_prepares_every_source_before_replacing_destinations() {
3932        let temp = tempfile::tempdir().expect("tempdir");
3933        fs::write(temp.path().join("profile.json"), "old profile").expect("seed latest profile");
3934        fs::write(temp.path().join("summary.md"), "old summary").expect("seed latest summary");
3935        let logical_id = ArtifactId::new("android-sample").expect("valid logical ID");
3936        let workspace =
3937            RunWorkspace::allocate(temp.path(), &logical_id).expect("allocate workspace");
3938        fs::write(workspace.staging_path().join("profile.json"), "new profile")
3939            .expect("write staged profile");
3940        fs::write(workspace.staging_path().join("summary.md"), "new summary")
3941            .expect("write staged summary");
3942        let required = [
3943            ArtifactId::new("profile.json").expect("valid profile ID"),
3944            ArtifactId::new("summary.md").expect("valid summary ID"),
3945        ];
3946        let published = workspace.publish(&required).expect("publish complete run");
3947        fs::remove_file(published.path().join("summary.md"))
3948            .expect("remove second published source");
3949
3950        let error = published
3951            .refresh_latest(&[
3952                LatestArtifact::same(required[0].clone()),
3953                LatestArtifact::same(required[1].clone()),
3954            ])
3955            .expect_err("reject incomplete latest source set");
3956
3957        assert!(matches!(
3958            error,
3959            ArtifactPathError::RequiredArtifactMissing { .. }
3960        ));
3961        assert_eq!(
3962            fs::read_to_string(temp.path().join("profile.json")).expect("read latest profile"),
3963            "old profile"
3964        );
3965        assert_eq!(
3966            fs::read_to_string(temp.path().join("summary.md")).expect("read latest summary"),
3967            "old summary"
3968        );
3969    }
3970
3971    #[test]
3972    fn approved_root_prepares_nested_directory_beneath_root() {
3973        let temp = tempfile::tempdir().expect("tempdir");
3974        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
3975
3976        let prepared = root
3977            .prepare_dir("plots/nested")
3978            .expect("prepare nested directory");
3979
3980        assert_eq!(prepared, root.path().join("plots/nested"));
3981        assert!(prepared.is_dir());
3982    }
3983
3984    #[test]
3985    fn approved_root_projects_a_missing_directory_without_creating_it() {
3986        let temp = tempfile::tempdir().expect("tempdir");
3987        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
3988
3989        let projected = root
3990            .project_dir("target/mobench")
3991            .expect("project output directory");
3992
3993        assert_eq!(projected, root.path().join("target/mobench"));
3994        assert!(!projected.exists());
3995        assert!(!root.path().join("target").exists());
3996    }
3997
3998    #[test]
3999    fn approved_root_rejects_paths_that_are_not_relative_descendants() {
4000        let temp = tempfile::tempdir().expect("tempdir");
4001        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4002
4003        for invalid in [Path::new("../escape"), Path::new("/absolute/escape")] {
4004            assert!(matches!(
4005                root.prepare_dir(invalid),
4006                Err(ArtifactPathError::InvalidRelativePath { .. })
4007            ));
4008        }
4009
4010        assert!(!root.path().join("../escape").exists());
4011    }
4012
4013    #[test]
4014    fn approved_root_projection_rejects_paths_that_are_not_relative_descendants() {
4015        let temp = tempfile::tempdir().expect("tempdir");
4016        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4017
4018        for invalid in [Path::new("../escape"), Path::new("/absolute/escape")] {
4019            assert!(matches!(
4020                root.project_dir(invalid),
4021                Err(ArtifactPathError::InvalidRelativePath { .. })
4022            ));
4023        }
4024    }
4025
4026    #[cfg(unix)]
4027    #[test]
4028    fn approved_root_rejects_symlink_directory_components_without_following_them() {
4029        use std::os::unix::fs::symlink;
4030
4031        let temp = tempfile::tempdir().expect("tempdir");
4032        let outside = tempfile::tempdir().expect("outside tempdir");
4033        symlink(outside.path(), temp.path().join("linked")).expect("create symlink");
4034        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4035
4036        assert!(matches!(
4037            root.prepare_dir("linked/nested"),
4038            Err(ArtifactPathError::SymlinkComponent { .. })
4039        ));
4040        assert!(!outside.path().join("nested").exists());
4041    }
4042
4043    #[cfg(unix)]
4044    #[test]
4045    fn approved_root_projection_rejects_symlink_directory_components_without_writing() {
4046        use std::os::unix::fs::symlink;
4047
4048        let temp = tempfile::tempdir().expect("tempdir");
4049        let outside = tempfile::tempdir().expect("outside tempdir");
4050        symlink(outside.path(), temp.path().join("linked")).expect("create symlink");
4051        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4052
4053        assert!(matches!(
4054            root.project_dir("linked/nested"),
4055            Err(ArtifactPathError::SymlinkComponent { .. })
4056        ));
4057        assert!(!outside.path().join("nested").exists());
4058    }
4059
4060    #[test]
4061    fn approved_root_prepares_file_parent_without_creating_the_file() {
4062        let temp = tempfile::tempdir().expect("tempdir");
4063        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4064
4065        let prepared = root
4066            .prepare_file("plots/nested/chart.svg")
4067            .expect("prepare artifact file");
4068
4069        assert_eq!(prepared, root.path().join("plots/nested/chart.svg"));
4070        assert!(prepared.parent().expect("parent").is_dir());
4071        assert!(!prepared.exists());
4072    }
4073
4074    #[cfg(unix)]
4075    #[test]
4076    fn approved_root_rejects_preexisting_symlink_file_without_following_it() {
4077        use std::os::unix::fs::symlink;
4078
4079        let temp = tempfile::tempdir().expect("tempdir");
4080        let root = ApprovedRoot::existing(temp.path()).expect("approve root");
4081        let plots = root.prepare_dir("plots").expect("prepare plots");
4082        let outside = tempfile::NamedTempFile::new().expect("outside file");
4083        fs::write(outside.path(), "unchanged").expect("seed outside file");
4084        symlink(outside.path(), plots.join("chart.svg")).expect("create symlink");
4085
4086        assert!(matches!(
4087            root.prepare_file("plots/chart.svg"),
4088            Err(ArtifactPathError::SymlinkComponent { .. })
4089        ));
4090        assert_eq!(
4091            fs::read_to_string(outside.path()).expect("read outside file"),
4092            "unchanged"
4093        );
4094    }
4095
4096    #[cfg(unix)]
4097    #[test]
4098    fn approved_root_rejects_a_symlink_as_the_root() {
4099        use std::os::unix::fs::symlink;
4100
4101        let parent = tempfile::tempdir().expect("parent tempdir");
4102        let actual_root = tempfile::tempdir().expect("actual root");
4103        let linked_root = parent.path().join("linked-root");
4104        symlink(actual_root.path(), &linked_root).expect("create root symlink");
4105
4106        assert!(matches!(
4107            ApprovedRoot::existing(&linked_root),
4108            Err(ArtifactPathError::SymlinkComponent { path }) if path == linked_root
4109        ));
4110    }
4111
4112    fn publish_pair(root: &Path, logical: &str, receipt: &str) -> PublishedRun {
4113        let logical_id = ArtifactId::new(logical).expect("valid logical ID");
4114        let workspace = RunWorkspace::allocate(root, &logical_id).expect("allocate workspace");
4115        fs::write(workspace.staging_path().join("profile.json"), receipt).expect("write profile");
4116        fs::write(workspace.staging_path().join("summary.md"), receipt).expect("write summary");
4117        workspace
4118            .publish(&[
4119                ArtifactId::new("profile.json").expect("profile ID"),
4120                ArtifactId::new("summary.md").expect("summary ID"),
4121            ])
4122            .expect("publish pair")
4123    }
4124
4125    fn refresh_pair(published: &PublishedRun) -> Result<(), ArtifactPathError> {
4126        published.refresh_latest(&[
4127            LatestArtifact::same(ArtifactId::new("profile.json").expect("profile ID")),
4128            LatestArtifact::same(ArtifactId::new("summary.md").expect("summary ID")),
4129        ])
4130    }
4131
4132    fn durability_events() -> Vec<&'static str> {
4133        DURABILITY_EVENTS.with(|events| std::mem::take(&mut *events.borrow_mut()))
4134    }
4135
4136    fn rewrite_latest_manifest(generation_path: &Path, update: impl FnOnce(&mut LatestManifest)) {
4137        let manifest_path = generation_path.join(LATEST_MANIFEST_FILE);
4138        let mut manifest: LatestManifest =
4139            read_json_file(&manifest_path, "read test manifest").expect("read manifest");
4140        update(&mut manifest);
4141        fs::remove_file(&manifest_path).expect("remove old manifest");
4142        write_json_file(&manifest_path, &manifest, "rewrite test manifest")
4143            .expect("rewrite manifest");
4144    }
4145
4146    fn clone_generation(
4147        source: &LatestSnapshot,
4148        generation: &str,
4149        predecessor: Option<&str>,
4150    ) -> PathBuf {
4151        let generations = source
4152            .path()
4153            .parent()
4154            .expect("generation parent")
4155            .to_path_buf();
4156        let destination = generations.join(generation);
4157        fs::create_dir(&destination).expect("create cloned generation");
4158        for artifact in &source.manifest().artifacts {
4159            fs::copy(
4160                source.path().join(&artifact.destination_relative_path),
4161                destination.join(&artifact.destination_relative_path),
4162            )
4163            .expect("copy generation artifact");
4164        }
4165        create_reader_lease_file(&destination).expect("create cloned generation lease");
4166        let mut manifest = source.manifest().clone();
4167        manifest.generation = generation.to_owned();
4168        manifest.predecessor_generation = predecessor.map(str::to_owned);
4169        write_json_file(
4170            &destination.join(LATEST_MANIFEST_FILE),
4171            &manifest,
4172            "write cloned generation manifest",
4173        )
4174        .expect("write cloned manifest");
4175        destination
4176    }
4177
4178    #[test]
4179    fn writer_startup_quarantines_an_unlocked_crash_workspace() {
4180        let temp = tempfile::tempdir().expect("tempdir");
4181        let staging = temp.path().join(STAGING_DIRECTORY);
4182        let crashed = staging.join("crashed-workspace");
4183        fs::create_dir_all(&crashed).expect("create crashed workspace");
4184        fs::write(crashed.join(WORKSPACE_LOCK_FILE), "").expect("create released lock");
4185        fs::write(crashed.join("partial.data"), "partial").expect("write partial artifact");
4186
4187        let workspace = RunWorkspace::allocate(
4188            temp.path(),
4189            &ArtifactId::new("recovery-writer").expect("logical ID"),
4190        )
4191        .expect("allocate after crash");
4192        assert!(!crashed.exists());
4193        let quarantine = temp.path().join(STAGING_QUARANTINE_DIRECTORY);
4194        assert!(
4195            fs::read_dir(quarantine)
4196                .expect("read run quarantine")
4197                .any(|entry| entry
4198                    .expect("quarantine entry")
4199                    .file_name()
4200                    .to_string_lossy()
4201                    .starts_with("crashed-workspace--"))
4202        );
4203        drop(workspace);
4204    }
4205
4206    #[test]
4207    fn writer_startup_never_quarantines_an_active_workspace() {
4208        let temp = tempfile::tempdir().expect("tempdir");
4209        let first = RunWorkspace::allocate(
4210            temp.path(),
4211            &ArtifactId::new("active-first").expect("logical ID"),
4212        )
4213        .expect("first workspace");
4214        let first_path = first.staging_path().to_path_buf();
4215
4216        let second = RunWorkspace::allocate(
4217            temp.path(),
4218            &ArtifactId::new("active-second").expect("logical ID"),
4219        )
4220        .expect("second workspace");
4221        assert!(first_path.is_dir(), "active workspace was moved");
4222        assert!(first_path.join(WORKSPACE_LOCK_FILE).is_file());
4223        drop(second);
4224        drop(first);
4225    }
4226
4227    #[test]
4228    fn run_manifest_records_deterministic_recursive_digests_and_identity() {
4229        let temp = tempfile::tempdir().expect("tempdir");
4230        let logical_id = ArtifactId::new("android-sample").expect("logical ID");
4231        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
4232        fs::write(workspace.staging_path().join("profile.json"), "hello").expect("write profile");
4233        fs::create_dir(workspace.staging_path().join("nested")).expect("create nested");
4234        fs::write(workspace.staging_path().join("nested/notes.txt"), "notes").expect("write notes");
4235        let published = workspace
4236            .publish(&[ArtifactId::new("profile.json").expect("profile ID")])
4237            .expect("publish");
4238
4239        let manifest = published.manifest().expect("verified manifest");
4240        assert_eq!(manifest.format_version, RUN_MANIFEST_VERSION);
4241        assert_eq!(manifest.producer_version, env!("CARGO_PKG_VERSION"));
4242        assert_eq!(manifest.logical_id, "android-sample");
4243        assert_eq!(
4244            manifest.publication_id,
4245            published
4246                .path()
4247                .file_name()
4248                .expect("publication name")
4249                .to_string_lossy()
4250        );
4251        assert_eq!(
4252            manifest
4253                .artifacts
4254                .iter()
4255                .map(|artifact| artifact.relative_path.as_str())
4256                .collect::<Vec<_>>(),
4257            ["nested/notes.txt", "profile.json"]
4258        );
4259        let profile = manifest
4260            .artifacts
4261            .iter()
4262            .find(|artifact| artifact.relative_path == "profile.json")
4263            .expect("profile manifest entry");
4264        assert_eq!(profile.size, 5);
4265        assert_eq!(
4266            profile.sha256,
4267            "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
4268        );
4269    }
4270
4271    #[test]
4272    fn publication_syncs_files_and_directories_before_the_atomic_rename() {
4273        let temp = tempfile::tempdir().expect("tempdir");
4274        let logical_id = ArtifactId::new("ordered-run").expect("logical ID");
4275        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
4276        fs::write(workspace.staging_path().join("profile.json"), "complete")
4277            .expect("write profile");
4278        durability_events();
4279
4280        workspace
4281            .publish(&[ArtifactId::new("profile.json").expect("profile ID")])
4282            .expect("publish");
4283        let events = durability_events();
4284        let publication = events
4285            .iter()
4286            .position(|event| *event == "publish_run")
4287            .expect("publication event");
4288        assert!(events[..publication].contains(&"sync_manifest"));
4289        assert!(events[..publication].contains(&"sync_file"));
4290        assert!(events[..publication].contains(&"sync_directory"));
4291        assert!(events[publication + 1..].contains(&"sync_directory"));
4292    }
4293
4294    #[test]
4295    fn post_rename_sync_failure_reports_committed_run_explicitly() {
4296        let temp = tempfile::tempdir().expect("tempdir");
4297        let logical_id = ArtifactId::new("uncertain-run").expect("logical ID");
4298        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
4299        fs::write(workspace.staging_path().join("profile.json"), "complete")
4300            .expect("write profile");
4301        let expected_path = workspace.published_path().to_path_buf();
4302        FAIL_DIRECTORY_SYNC_AFTER_EVENT.with(|target| target.set(Some("publish_run")));
4303
4304        let error = workspace
4305            .publish(&[ArtifactId::new("profile.json").expect("profile ID")])
4306            .expect_err("surface post-commit durability uncertainty");
4307        assert!(matches!(
4308            error,
4309            ArtifactPathError::PublicationDurabilityUncertain { ref path, .. }
4310                if path == &expected_path
4311        ));
4312        assert!(expected_path.join(RUN_MANIFEST_FILE).is_file());
4313        validate_run_manifest(&expected_path).expect("committed run remains valid");
4314    }
4315
4316    #[test]
4317    fn latest_commit_orders_generation_sync_before_the_single_pointer_swap() {
4318        let temp = tempfile::tempdir().expect("tempdir");
4319        let published = publish_pair(temp.path(), "ordered-latest", "receipt");
4320        durability_events();
4321
4322        refresh_pair(&published).expect("refresh latest");
4323        let events = durability_events();
4324        let generation = events
4325            .iter()
4326            .position(|event| *event == "publish_generation")
4327            .expect("generation publication event");
4328        let pointer_sync = events
4329            .iter()
4330            .position(|event| *event == "sync_pointer_file")
4331            .expect("pointer sync event");
4332        let pointer_commit = events
4333            .iter()
4334            .position(|event| *event == "commit_pointer")
4335            .expect("pointer commit event");
4336        assert!(events[..generation].contains(&"sync_manifest"));
4337        assert!(events[generation + 1..pointer_sync].contains(&"sync_directory"));
4338        let aliases = events
4339            .iter()
4340            .position(|event| *event == "install_aliases")
4341            .expect("alias install event");
4342        assert!(generation < pointer_sync);
4343        assert!(
4344            aliases < pointer_sync,
4345            "aliases must precede pointer commit"
4346        );
4347        assert!(pointer_sync < pointer_commit);
4348        assert!(events[pointer_commit + 1..].contains(&"sync_directory"));
4349    }
4350
4351    #[test]
4352    fn post_pointer_sync_failure_reports_committed_generation_explicitly() {
4353        let temp = tempfile::tempdir().expect("tempdir");
4354        let published = publish_pair(temp.path(), "uncertain-latest", "receipt");
4355        FAIL_DIRECTORY_SYNC_AFTER_EVENT.with(|target| target.set(Some("commit_pointer")));
4356
4357        let error = refresh_pair(&published).expect_err("surface pointer durability uncertainty");
4358        let generation = match error {
4359            ArtifactPathError::LatestCommitDurabilityUncertain { generation, .. } => generation,
4360            other => panic!("unexpected error: {other}"),
4361        };
4362        let snapshot = LatestSnapshot::open(temp.path()).expect("committed snapshot remains open");
4363        assert_eq!(snapshot.generation(), generation);
4364        assert_eq!(
4365            snapshot
4366                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4367                .expect("read committed artifact"),
4368            b"receipt"
4369        );
4370    }
4371
4372    #[test]
4373    fn latest_snapshot_pins_one_generation_across_concurrent_refreshes() {
4374        let temp = tempfile::tempdir().expect("tempdir");
4375        let first = publish_pair(temp.path(), "first-run", "first");
4376        refresh_pair(&first).expect("first refresh");
4377        let pinned = LatestSnapshot::open(temp.path()).expect("first snapshot");
4378
4379        let second = publish_pair(temp.path(), "second-run", "second");
4380        refresh_pair(&second).expect("second refresh");
4381        let current = LatestSnapshot::open(temp.path()).expect("second snapshot");
4382
4383        assert_ne!(pinned.generation(), current.generation());
4384        for id in ["profile.json", "summary.md"] {
4385            let id = ArtifactId::new(id).expect("artifact ID");
4386            assert_eq!(pinned.read_artifact(&id).expect("pinned read"), b"first");
4387            assert_eq!(current.read_artifact(&id).expect("current read"), b"second");
4388        }
4389    }
4390
4391    #[cfg(unix)]
4392    #[test]
4393    fn latest_snapshot_open_succeeds_with_a_read_only_artifact_root() {
4394        use std::os::unix::fs::PermissionsExt;
4395
4396        fn set_tree_mode(path: &Path, directory_mode: u32, file_mode: u32) {
4397            let metadata = fs::symlink_metadata(path).expect("tree metadata");
4398            if metadata.is_dir() {
4399                for entry in fs::read_dir(path).expect("read tree") {
4400                    set_tree_mode(
4401                        &entry.expect("tree entry").path(),
4402                        directory_mode,
4403                        file_mode,
4404                    );
4405                }
4406                fs::set_permissions(path, fs::Permissions::from_mode(directory_mode))
4407                    .expect("set directory mode");
4408            } else {
4409                fs::set_permissions(path, fs::Permissions::from_mode(file_mode))
4410                    .expect("set file mode");
4411            }
4412        }
4413
4414        let temp = tempfile::tempdir().expect("tempdir");
4415        let published = publish_pair(temp.path(), "read-only", "receipt");
4416        refresh_pair(&published).expect("refresh latest");
4417        set_tree_mode(temp.path(), 0o555, 0o444);
4418
4419        let opened = LatestSnapshot::open(temp.path()).expect("read-only snapshot open");
4420        assert_eq!(
4421            opened
4422                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4423                .expect("read snapshot artifact"),
4424            b"receipt"
4425        );
4426
4427        set_tree_mode(temp.path(), 0o755, 0o644);
4428    }
4429
4430    #[test]
4431    fn writer_allocation_recovers_mixed_legacy_aliases() {
4432        let temp = tempfile::tempdir().expect("tempdir");
4433        let published = publish_pair(temp.path(), "recover-run", "committed");
4434        refresh_pair(&published).expect("refresh latest");
4435        fs::write(temp.path().join("profile.json"), "interrupted-new")
4436            .expect("simulate interrupted alias install");
4437        fs::write(temp.path().join("summary.md"), "committed").expect("retain old alias");
4438
4439        LatestSnapshot::open(temp.path()).expect("read-only open");
4440        assert_eq!(
4441            fs::read_to_string(temp.path().join("profile.json")).expect("unrecovered profile"),
4442            "interrupted-new",
4443            "reader open must not mutate compatibility aliases"
4444        );
4445
4446        let recovery_probe = RunWorkspace::allocate(
4447            temp.path(),
4448            &ArtifactId::new("recovery-probe").expect("probe ID"),
4449        )
4450        .expect("writer startup recovers aliases");
4451        drop(recovery_probe);
4452        let recovered = LatestSnapshot::open(temp.path()).expect("open recovered snapshot");
4453        assert_eq!(
4454            fs::read(temp.path().join("profile.json")).expect("profile alias"),
4455            b"committed"
4456        );
4457        assert_eq!(
4458            fs::read(temp.path().join("summary.md")).expect("summary alias"),
4459            b"committed"
4460        );
4461        assert_eq!(
4462            recovered
4463                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4464                .expect("read recovered snapshot"),
4465            b"committed"
4466        );
4467    }
4468
4469    #[test]
4470    fn missing_current_recovers_the_unique_highest_chain_tip() {
4471        let temp = tempfile::tempdir().expect("tempdir");
4472        let first = publish_pair(temp.path(), "first-chain-run", "first");
4473        refresh_pair(&first).expect("first refresh");
4474        let second = publish_pair(temp.path(), "second-chain-run", "second");
4475        refresh_pair(&second).expect("second refresh");
4476        let expected = LatestSnapshot::open(temp.path()).expect("tip snapshot");
4477        let expected_generation = expected.generation().to_owned();
4478        let current = temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE);
4479        fs::remove_file(&current).expect("remove current pointer");
4480        fs::write(temp.path().join("profile.json"), "interrupted")
4481            .expect("tamper compatibility alias");
4482
4483        assert!(matches!(
4484            LatestSnapshot::open(temp.path()),
4485            Err(ArtifactPathError::LatestSnapshotUnavailable { .. })
4486        ));
4487        let recovered =
4488            LatestSnapshot::recover_stable_aliases(temp.path()).expect("recover missing pointer");
4489        assert_eq!(recovered.generation(), expected_generation);
4490        assert_eq!(
4491            fs::read_to_string(temp.path().join("profile.json")).expect("recovered profile"),
4492            "second"
4493        );
4494    }
4495
4496    #[test]
4497    fn corrupt_current_recovers_the_unique_valid_chain_tip() {
4498        let temp = tempfile::tempdir().expect("tempdir");
4499        let published = publish_pair(temp.path(), "corrupt-pointer-run", "receipt");
4500        refresh_pair(&published).expect("refresh latest");
4501        let expected = LatestSnapshot::open(temp.path())
4502            .expect("snapshot")
4503            .generation()
4504            .to_owned();
4505        let current = temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE);
4506        fs::write(&current, "../../invalid\n").expect("corrupt pointer");
4507
4508        assert!(matches!(
4509            LatestSnapshot::open(temp.path()),
4510            Err(ArtifactPathError::InvalidLatestPointer { .. })
4511        ));
4512        let recovered =
4513            LatestSnapshot::recover_stable_aliases(temp.path()).expect("recover corrupt pointer");
4514        assert_eq!(recovered.generation(), expected);
4515        assert_eq!(
4516            fs::read_to_string(current)
4517                .expect("read repaired pointer")
4518                .trim(),
4519            expected
4520        );
4521    }
4522
4523    #[test]
4524    fn missing_current_fails_closed_on_an_ambiguous_generation_fork() {
4525        let temp = tempfile::tempdir().expect("tempdir");
4526        let published = publish_pair(temp.path(), "fork-root-run", "root");
4527        refresh_pair(&published).expect("refresh root");
4528        let root_snapshot = LatestSnapshot::open(temp.path()).expect("root snapshot");
4529        let root_generation = root_snapshot.generation().to_owned();
4530        clone_generation(&root_snapshot, "fork-child-a", Some(&root_generation));
4531        clone_generation(&root_snapshot, "fork-child-b", Some(&root_generation));
4532        fs::remove_file(temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE))
4533            .expect("remove current pointer");
4534
4535        assert!(matches!(
4536            LatestSnapshot::recover_stable_aliases(temp.path()),
4537            Err(ArtifactPathError::AmbiguousGenerationFork { .. })
4538        ));
4539    }
4540
4541    #[test]
4542    fn predecessor_cycle_is_rejected_by_open_and_recovery() {
4543        let temp = tempfile::tempdir().expect("tempdir");
4544        let first = publish_pair(temp.path(), "cycle-first", "first");
4545        refresh_pair(&first).expect("first refresh");
4546        let first_snapshot = LatestSnapshot::open(temp.path()).expect("first snapshot");
4547        let first_generation = first_snapshot.generation().to_owned();
4548        let second = publish_pair(temp.path(), "cycle-second", "second");
4549        refresh_pair(&second).expect("second refresh");
4550        let second_snapshot = LatestSnapshot::open(temp.path()).expect("second snapshot");
4551        let second_generation = second_snapshot.generation().to_owned();
4552        rewrite_latest_manifest(&first_snapshot.generation_path, |manifest| {
4553            manifest.predecessor_generation = Some(second_generation.clone());
4554        });
4555
4556        assert!(matches!(
4557            LatestSnapshot::open(temp.path()),
4558            Err(ArtifactPathError::GenerationChainCycle { .. })
4559        ));
4560        fs::remove_file(temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE))
4561            .expect("remove current pointer");
4562        assert!(matches!(
4563            LatestSnapshot::recover_stable_aliases(temp.path()),
4564            Err(ArtifactPathError::GenerationChainCycle { .. })
4565        ));
4566        assert_ne!(first_generation, second_generation);
4567    }
4568
4569    #[test]
4570    fn broken_predecessor_chain_fails_closed_during_recovery() {
4571        let temp = tempfile::tempdir().expect("tempdir");
4572        let published = publish_pair(temp.path(), "broken-chain", "receipt");
4573        refresh_pair(&published).expect("refresh latest");
4574        let snapshot = LatestSnapshot::open(temp.path()).expect("snapshot");
4575        rewrite_latest_manifest(snapshot.path(), |manifest| {
4576            manifest.predecessor_generation = Some("missing-predecessor".to_owned());
4577        });
4578        fs::remove_file(temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE))
4579            .expect("remove current pointer");
4580
4581        assert!(matches!(
4582            LatestSnapshot::recover_stable_aliases(temp.path()),
4583            Err(ArtifactPathError::BrokenGenerationChain { .. })
4584        ));
4585    }
4586
4587    #[test]
4588    fn alias_install_failure_rolls_back_before_pointer_commit() {
4589        let temp = tempfile::tempdir().expect("tempdir");
4590        let first = publish_pair(temp.path(), "first-run", "first");
4591        refresh_pair(&first).expect("first refresh");
4592        let before = LatestSnapshot::open(temp.path()).expect("snapshot before failure");
4593        let before_generation = before.generation().to_owned();
4594        let second = publish_pair(temp.path(), "second-run", "second");
4595
4596        FAIL_ALIAS_INSTALL_AFTER.with(|fail_after| fail_after.set(Some(1)));
4597        let error = refresh_pair(&second).expect_err("injected alias failure");
4598        assert!(matches!(error, ArtifactPathError::Io { .. }));
4599
4600        let after = LatestSnapshot::open(temp.path()).expect("snapshot after rollback");
4601        assert_eq!(after.generation(), before_generation);
4602        assert_eq!(
4603            fs::read_to_string(temp.path().join("profile.json")).expect("profile alias"),
4604            "first"
4605        );
4606        assert_eq!(
4607            fs::read_to_string(temp.path().join("summary.md")).expect("summary alias"),
4608            "first"
4609        );
4610    }
4611
4612    #[test]
4613    fn stale_allocated_run_cannot_replace_a_newer_generation() {
4614        let temp = tempfile::tempdir().expect("tempdir");
4615        let stale = publish_pair(temp.path(), "stale-run", "stale");
4616        let winner = publish_pair(temp.path(), "winner-run", "winner");
4617        refresh_pair(&winner).expect("winner refresh");
4618        let winner_generation = LatestSnapshot::open(temp.path())
4619            .expect("winner snapshot")
4620            .generation()
4621            .to_owned();
4622
4623        let error = refresh_pair(&stale).expect_err("reject stale CAS token");
4624        assert!(matches!(
4625            error,
4626            ArtifactPathError::StaleLatestGeneration {
4627                expected: None,
4628                observed: Some(_)
4629            }
4630        ));
4631        let current = LatestSnapshot::open(temp.path()).expect("current snapshot");
4632        assert_eq!(current.generation(), winner_generation);
4633        assert_eq!(
4634            current
4635                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4636                .expect("read winner"),
4637            b"winner"
4638        );
4639    }
4640
4641    #[test]
4642    fn abandoned_latest_staging_is_quarantined_before_the_next_commit() {
4643        let temp = tempfile::tempdir().expect("tempdir");
4644        let first = publish_pair(temp.path(), "first-run", "first");
4645        refresh_pair(&first).expect("first refresh");
4646        let staging = temp
4647            .path()
4648            .join(LATEST_DIRECTORY)
4649            .join(LATEST_STAGING_DIRECTORY);
4650        let abandoned = staging.join("abandoned");
4651        fs::create_dir(&abandoned).expect("create abandoned staging");
4652        fs::write(abandoned.join("partial"), "partial").expect("write partial staging");
4653
4654        let second = publish_pair(temp.path(), "second-run", "second");
4655        refresh_pair(&second).expect("second refresh");
4656        assert!(!abandoned.exists());
4657        let quarantine = temp
4658            .path()
4659            .join(LATEST_DIRECTORY)
4660            .join(LATEST_QUARANTINE_DIRECTORY);
4661        assert!(
4662            fs::read_dir(quarantine)
4663                .expect("read quarantine")
4664                .any(|entry| entry
4665                    .expect("quarantine entry")
4666                    .file_name()
4667                    .to_string_lossy()
4668                    .starts_with("abandoned--"))
4669        );
4670    }
4671
4672    #[test]
4673    fn active_latest_staging_lease_defers_recovery_quarantine() {
4674        let temp = tempfile::tempdir().expect("tempdir");
4675        let root = ApprovedRoot::existing(temp.path()).expect("approved root");
4676        let latest_path = root
4677            .prepare_dir(LATEST_DIRECTORY)
4678            .expect("latest directory");
4679        let latest_root = ApprovedRoot::existing(latest_path).expect("approved latest root");
4680        latest_root
4681            .prepare_dir(LATEST_QUARANTINE_DIRECTORY)
4682            .expect("quarantine directory");
4683        let (_, staging_path) = allocate_latest_generation(&latest_root).expect("staging path");
4684        create_reader_lease_file(&staging_path).expect("staging lease file");
4685        let lease_path = staging_path.join(LATEST_READER_LOCK_FILE);
4686        let lease = OpenOptions::new()
4687            .read(true)
4688            .write(true)
4689            .open(&lease_path)
4690            .expect("open staging lease");
4691        FileExt::lock_exclusive(&lease).expect("lock active staging lease");
4692
4693        quarantine_abandoned_latest_staging(&latest_root).expect("active staging recovery pass");
4694        assert!(staging_path.exists(), "active staging was quarantined");
4695
4696        FileExt::unlock(&lease).expect("unlock staging lease");
4697        drop(lease);
4698        quarantine_abandoned_latest_staging(&latest_root).expect("abandoned staging recovery pass");
4699        assert!(!staging_path.exists(), "abandoned staging was retained");
4700    }
4701
4702    #[test]
4703    fn newer_latest_producer_version_rejects_a_downgrade() {
4704        let temp = tempfile::tempdir().expect("tempdir");
4705        let first = publish_pair(temp.path(), "first-run", "first");
4706        refresh_pair(&first).expect("first refresh");
4707        let second = publish_pair(temp.path(), "second-run", "second");
4708        let snapshot = LatestSnapshot::open(temp.path()).expect("snapshot");
4709        let manifest_path = snapshot.path().join(LATEST_MANIFEST_FILE);
4710        let mut manifest: LatestManifest =
4711            read_json_file(&manifest_path, "read test manifest").expect("manifest");
4712        manifest.producer_version = "999.0.0".to_owned();
4713        fs::remove_file(&manifest_path).expect("remove old manifest");
4714        write_json_file(&manifest_path, &manifest, "rewrite test manifest")
4715            .expect("write newer producer manifest");
4716
4717        let error = refresh_pair(&second).expect_err("reject producer downgrade");
4718        assert!(matches!(error, ArtifactPathError::ProducerDowngrade { .. }));
4719        assert_eq!(
4720            fs::read_to_string(temp.path().join("profile.json")).expect("stable profile"),
4721            "first"
4722        );
4723    }
4724
4725    #[test]
4726    fn unknown_latest_manifest_version_fails_closed() {
4727        let temp = tempfile::tempdir().expect("tempdir");
4728        let published = publish_pair(temp.path(), "unknown-version", "receipt");
4729        refresh_pair(&published).expect("refresh latest");
4730        let snapshot = LatestSnapshot::open(temp.path()).expect("snapshot");
4731        let manifest_path = snapshot.path().join(LATEST_MANIFEST_FILE);
4732        let mut manifest: LatestManifest =
4733            read_json_file(&manifest_path, "read test manifest").expect("manifest");
4734        manifest.format_version = LATEST_MANIFEST_VERSION + 1;
4735        fs::remove_file(&manifest_path).expect("remove old manifest");
4736        write_json_file(&manifest_path, &manifest, "rewrite test manifest")
4737            .expect("write unknown manifest");
4738
4739        assert!(matches!(
4740            LatestSnapshot::open(temp.path()),
4741            Err(ArtifactPathError::UnknownManifestVersion { .. })
4742        ));
4743    }
4744
4745    #[test]
4746    fn corrupt_latest_generation_fails_digest_validation() {
4747        let temp = tempfile::tempdir().expect("tempdir");
4748        let published = publish_pair(temp.path(), "corrupt-generation", "receipt");
4749        refresh_pair(&published).expect("refresh latest");
4750        let snapshot = LatestSnapshot::open(temp.path()).expect("snapshot");
4751        fs::write(snapshot.path().join("profile.json"), "tampered").expect("tamper artifact");
4752
4753        assert!(matches!(
4754            LatestSnapshot::open(temp.path()),
4755            Err(ArtifactPathError::CorruptGeneration { .. })
4756        ));
4757    }
4758
4759    #[test]
4760    fn pinned_snapshot_rejects_bytes_changed_after_open() {
4761        let temp = tempfile::tempdir().expect("tempdir");
4762        let published = publish_pair(temp.path(), "mutated-after-open", "original");
4763        refresh_pair(&published).expect("refresh latest");
4764        let snapshot = LatestSnapshot::open(temp.path()).expect("open pinned snapshot");
4765        fs::write(snapshot.path().join("profile.json"), "changed").expect("mutate artifact");
4766
4767        assert!(matches!(
4768            snapshot.read_artifact(&ArtifactId::new("profile.json").expect("profile ID")),
4769            Err(ArtifactPathError::CorruptGeneration { .. })
4770        ));
4771    }
4772
4773    #[test]
4774    fn latest_refresh_rejects_source_changed_after_run_validation() {
4775        let temp = tempfile::tempdir().expect("tempdir");
4776        let published = publish_pair(temp.path(), "mutated-source", "original");
4777        MUTATE_LATEST_SOURCE_AFTER_VALIDATION.with(|mutation| {
4778            mutation.replace(Some((
4779                published.path().join("profile.json"),
4780                b"changed after validation".to_vec(),
4781            )));
4782        });
4783
4784        let error = refresh_pair(&published).expect_err("reject changed source");
4785        assert!(matches!(error, ArtifactPathError::ManifestIntegrity { .. }));
4786        assert!(matches!(
4787            LatestSnapshot::open(temp.path()),
4788            Err(ArtifactPathError::LatestSnapshotUnavailable { .. })
4789        ));
4790    }
4791
4792    #[test]
4793    fn invalid_current_directory_does_not_block_first_commit() {
4794        let temp = tempfile::tempdir().expect("tempdir");
4795        let latest = temp.path().join(LATEST_DIRECTORY);
4796        fs::create_dir_all(latest.join(LATEST_CURRENT_FILE)).expect("seed invalid pointer dir");
4797        let published = publish_pair(temp.path(), "first-after-invalid-pointer", "receipt");
4798
4799        refresh_pair(&published).expect("commit after quarantining invalid pointer");
4800        let snapshot = LatestSnapshot::open(temp.path()).expect("open committed snapshot");
4801        assert_eq!(
4802            snapshot
4803                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4804                .expect("read committed artifact"),
4805            b"receipt"
4806        );
4807        assert!(
4808            fs::read_dir(latest.join(LATEST_QUARANTINE_DIRECTORY))
4809                .expect("read quarantine")
4810                .any(|entry| entry
4811                    .expect("quarantine entry")
4812                    .file_name()
4813                    .to_string_lossy()
4814                    .starts_with("current-corrupt--"))
4815        );
4816    }
4817
4818    #[test]
4819    fn corrupt_off_chain_generation_is_quarantined_without_blocking_writer() {
4820        let temp = tempfile::tempdir().expect("tempdir");
4821        let first = publish_pair(temp.path(), "anchored", "first");
4822        refresh_pair(&first).expect("anchor current generation");
4823        let orphan = temp
4824            .path()
4825            .join(LATEST_DIRECTORY)
4826            .join(LATEST_GENERATIONS_DIRECTORY)
4827            .join("generation-corrupt-orphan");
4828        fs::create_dir(&orphan).expect("create corrupt orphan");
4829        fs::write(orphan.join(LATEST_MANIFEST_FILE), b"not-json").expect("write corrupt orphan");
4830
4831        let second = publish_pair(temp.path(), "after-orphan", "second");
4832        refresh_pair(&second).expect("refresh despite corrupt orphan");
4833        assert!(!orphan.exists());
4834        let snapshot = LatestSnapshot::open(temp.path()).expect("open latest snapshot");
4835        assert_eq!(
4836            snapshot
4837                .read_artifact(&ArtifactId::new("profile.json").expect("profile ID"))
4838                .expect("read latest artifact"),
4839            b"second"
4840        );
4841    }
4842
4843    #[cfg(unix)]
4844    #[test]
4845    fn recursive_publication_validation_rejects_nested_symlinks() {
4846        use std::os::unix::fs::symlink;
4847
4848        let temp = tempfile::tempdir().expect("tempdir");
4849        let outside = tempfile::NamedTempFile::new().expect("outside file");
4850        let logical_id = ArtifactId::new("nested-symlink").expect("logical ID");
4851        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
4852        fs::write(workspace.staging_path().join("profile.json"), "complete")
4853            .expect("write profile");
4854        fs::create_dir(workspace.staging_path().join("nested")).expect("nested directory");
4855        symlink(
4856            outside.path(),
4857            workspace.staging_path().join("nested/linked"),
4858        )
4859        .expect("nested symlink");
4860
4861        assert!(matches!(
4862            workspace.publish(&[ArtifactId::new("profile.json").expect("profile ID")]),
4863            Err(ArtifactPathError::SymlinkComponent { .. })
4864        ));
4865    }
4866
4867    #[cfg(unix)]
4868    #[test]
4869    fn publication_rejects_artifacts_hardlinked_outside_the_workspace() {
4870        let temp = tempfile::tempdir().expect("tempdir");
4871        let outside = tempfile::tempdir().expect("outside tempdir");
4872        let logical_id = ArtifactId::new("hardlinked-artifact").expect("logical ID");
4873        let workspace = RunWorkspace::allocate(temp.path(), &logical_id).expect("workspace");
4874        let profile = workspace.staging_path().join("profile.json");
4875        fs::write(&profile, "complete").expect("write profile");
4876        fs::hard_link(&profile, outside.path().join("alias.json")).expect("create hard link");
4877
4878        assert!(matches!(
4879            workspace.publish(&[ArtifactId::new("profile.json").expect("profile ID")]),
4880            Err(ArtifactPathError::HardLinkedArtifact { .. })
4881        ));
4882    }
4883
4884    #[cfg(any(
4885        target_os = "macos",
4886        target_os = "ios",
4887        target_os = "linux",
4888        target_os = "android"
4889    ))]
4890    #[test]
4891    fn platform_publication_rename_never_replaces_a_raced_directory() {
4892        let temp = tempfile::tempdir().expect("tempdir");
4893        let source = temp.path().join("source");
4894        let destination = temp.path().join("destination");
4895        fs::create_dir(&source).expect("create source");
4896        fs::create_dir(&destination).expect("create raced destination");
4897        fs::write(source.join("source.txt"), "source").expect("write source");
4898        fs::write(destination.join("destination.txt"), "destination").expect("write destination");
4899
4900        let error = rename_directory_noreplace(&source, &destination)
4901            .expect_err("no-replace rename must reject collision");
4902        assert!(matches!(
4903            error.kind(),
4904            std::io::ErrorKind::AlreadyExists | std::io::ErrorKind::Other
4905        ));
4906        assert!(source.join("source.txt").is_file());
4907        assert_eq!(
4908            fs::read_to_string(destination.join("destination.txt"))
4909                .expect("read original destination"),
4910            "destination"
4911        );
4912    }
4913
4914    #[cfg(unix)]
4915    #[test]
4916    fn latest_reader_rejects_a_symlink_pointer() {
4917        use std::os::unix::fs::symlink;
4918
4919        let temp = tempfile::tempdir().expect("tempdir");
4920        let outside = tempfile::NamedTempFile::new().expect("outside pointer");
4921        fs::write(outside.path(), "generation-attacker").expect("write outside pointer");
4922        let published = publish_pair(temp.path(), "symlink-pointer", "receipt");
4923        refresh_pair(&published).expect("refresh latest");
4924        let current = temp.path().join(LATEST_DIRECTORY).join(LATEST_CURRENT_FILE);
4925        fs::remove_file(&current).expect("remove current pointer");
4926        symlink(outside.path(), &current).expect("install pointer symlink");
4927
4928        assert!(matches!(
4929            LatestSnapshot::open(temp.path()),
4930            Err(ArtifactPathError::SymlinkComponent { .. })
4931        ));
4932    }
4933}