1use 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#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
25pub struct ArtifactId(String);
26
27impl ArtifactId {
28 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 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;
77pub const RUN_MANIFEST_FILE: &str = "mobench-run-manifest.json";
79pub 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#[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#[derive(Debug, Clone)]
102pub struct PublishedRun {
103 root: ApprovedRoot,
104 path: PathBuf,
105 _reader_lease: Arc<File>,
106}
107
108#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
110pub struct ManifestArtifact {
111 pub relative_path: String,
113 pub size: u64,
115 pub sha256: String,
117}
118
119#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121pub struct RunManifest {
122 pub format_version: u32,
124 pub producer_version: String,
126 pub logical_id: String,
128 pub publication_id: String,
130 pub expected_latest_generation: Option<String>,
133 pub artifacts: Vec<ManifestArtifact>,
135}
136
137#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
139pub struct LatestManifestArtifact {
140 pub source_relative_path: String,
142 pub destination_relative_path: String,
144 pub size: u64,
146 pub sha256: String,
148}
149
150#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub struct LatestManifest {
153 pub format_version: u32,
155 pub producer_version: String,
157 pub generation: String,
159 pub predecessor_generation: Option<String>,
161 pub source_logical_id: String,
163 pub source_publication_id: String,
165 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#[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#[derive(Debug, Clone, PartialEq, Eq)]
235pub struct LatestArtifact {
236 source: ArtifactId,
237 destination: ArtifactId,
238}
239
240impl LatestArtifact {
241 pub fn same(id: ArtifactId) -> Self {
243 Self {
244 source: id.clone(),
245 destination: id,
246 }
247 }
248
249 pub fn new(source: ArtifactId, destination: ArtifactId) -> Self {
251 Self {
252 source,
253 destination,
254 }
255 }
256}
257
258impl RunWorkspace {
259 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 pub fn staging_path(&self) -> &Path {
345 self.staging_path
346 .as_deref()
347 .expect("unpublished workspace retains its staging path")
348 }
349
350 pub fn published_path(&self) -> &Path {
352 &self.published_path
353 }
354
355 pub fn root(&self) -> &ApprovedRoot {
357 &self.root
358 }
359
360 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 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 pub fn path(&self) -> &Path {
571 &self.path
572 }
573
574 pub fn root(&self) -> &ApprovedRoot {
576 &self.root
577 }
578
579 pub fn manifest(&self) -> Result<RunManifest, ArtifactPathError> {
581 validate_run_manifest(&self.path)
582 }
583
584 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 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 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 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 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 pub fn generation(&self) -> &str {
863 &self.manifest.generation
864 }
865
866 pub fn path(&self) -> &Path {
868 &self.generation_path
869 }
870
871 pub fn root(&self) -> &ApprovedRoot {
873 &self.root
874 }
875
876 pub fn manifest(&self) -> &LatestManifest {
878 &self.manifest
879 }
880
881 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 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(¤t) {
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, ¤t).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 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 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(¤t) {
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(¤t, &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(¤t) {
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(¤t)
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 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 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 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 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 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 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 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 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 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 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 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#[derive(Debug, Clone, PartialEq, Eq)]
3087pub struct ApprovedRoot {
3088 path: PathBuf,
3089}
3090
3091impl ApprovedRoot {
3092 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 pub fn path(&self) -> &Path {
3119 &self.path
3120 }
3121
3122 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(¤t) {
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 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 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#[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 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(¤t) {
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(¤t) {
3382 Ok(()) => {}
3383 Err(source) if source.kind() == std::io::ErrorKind::AlreadyExists => {
3384 match fs::symlink_metadata(¤t) {
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(¤t).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(¤t, "../../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(¤t).expect("remove current pointer");
4926 symlink(outside.path(), ¤t).expect("install pointer symlink");
4927
4928 assert!(matches!(
4929 LatestSnapshot::open(temp.path()),
4930 Err(ArtifactPathError::SymlinkComponent { .. })
4931 ));
4932 }
4933}