artifact_bundle.rs (39149B)
1 use fs2::FileExt; 2 use serde::{Deserialize, Serialize}; 3 use sha2::{Digest, Sha256}; 4 use std::{ 5 collections::BTreeSet, 6 fs::{self, OpenOptions}, 7 io::{self, ErrorKind, Write}, 8 path::{Component, Path, PathBuf}, 9 }; 10 use tempfile::NamedTempFile; 11 12 const TRANSACTION_SCHEMA_VERSION: u32 = 1; 13 const TRANSACTION_DIRECTORY: &str = ".radroots-contract-artifact-transaction-v1"; 14 const TRANSACTION_JOURNAL: &str = "journal.json"; 15 const LOCK_DIRECTORY: &str = "radroots-xtask-contract-artifact-locks-v1"; 16 17 pub(crate) struct GeneratedArtifact { 18 pub(crate) relative: &'static str, 19 pub(crate) contents: Vec<u8>, 20 } 21 22 #[derive(Debug, Deserialize, Serialize)] 23 #[serde(deny_unknown_fields)] 24 struct ArtifactTransactionJournal { 25 schema_version: u32, 26 artifacts: Vec<ArtifactTransactionEntry>, 27 } 28 29 #[derive(Debug, Deserialize, Serialize)] 30 #[serde(deny_unknown_fields)] 31 struct ArtifactTransactionEntry { 32 relative: String, 33 original_byte_length: Option<u64>, 34 original_sha256: Option<String>, 35 } 36 37 struct OriginalArtifact { 38 path: PathBuf, 39 relative: &'static str, 40 contents: Option<Vec<u8>>, 41 permissions: Option<fs::Permissions>, 42 } 43 44 struct PendingArtifact { 45 original: OriginalArtifact, 46 contents: Vec<u8>, 47 } 48 49 struct StagedArtifact { 50 original: OriginalArtifact, 51 temporary: NamedTempFile, 52 } 53 54 #[derive(Clone, Copy, Eq, PartialEq)] 55 enum SimulatedInterruption { 56 AfterStaging, 57 AfterCommits(usize), 58 } 59 60 pub(crate) struct ArtifactBundleTransaction<'a> { 61 workspace_root: &'a Path, 62 } 63 64 impl ArtifactBundleTransaction<'_> { 65 pub(crate) fn write(&self, artifacts: Vec<GeneratedArtifact>) -> Result<(), String> { 66 write_artifact_bundle_impl(self.workspace_root, artifacts, None) 67 } 68 } 69 70 pub(crate) fn read_regular_file(workspace_root: &Path, relative: &str) -> Result<Vec<u8>, String> { 71 let path = validate_workspace_path(workspace_root, relative, false)?; 72 fs::read(path).map_err(|error| format!("read {relative}: {error}")) 73 } 74 75 pub(super) fn validate_canonical_json_artifact(relative: &str, bytes: &[u8]) -> Result<(), String> { 76 if bytes.contains(&b'\r') { 77 return Err(format!("{relative} must use LF line endings")); 78 } 79 if !bytes.ends_with(b"\n") || bytes.ends_with(b"\n\n") { 80 return Err(format!("{relative} must end with exactly one LF")); 81 } 82 Ok(()) 83 } 84 85 pub(super) fn validate_sha256_artifact(relative: &str, bytes: &[u8]) -> Result<(), String> { 86 if bytes.len() != 65 || bytes[64] != b'\n' { 87 return Err(format!( 88 "{relative} must contain 64 lowercase hexadecimal bytes and one LF" 89 )); 90 } 91 if !bytes[..64] 92 .iter() 93 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(byte)) 94 { 95 return Err(format!( 96 "{relative} must contain a lowercase SHA-256 digest" 97 )); 98 } 99 Ok(()) 100 } 101 102 pub(crate) fn with_artifact_bundle_transaction<T>( 103 workspace_root: &Path, 104 operation: impl FnOnce(&ArtifactBundleTransaction<'_>) -> Result<T, String>, 105 ) -> Result<T, String> { 106 let lock = acquire_workspace_lock(workspace_root)?; 107 let result = recover_pending_transaction(workspace_root) 108 .and_then(|()| operation(&ArtifactBundleTransaction { workspace_root })); 109 let unlock = FileExt::unlock(&lock).map_err(|error| { 110 format!( 111 "unlock generated artifact transaction for {}: {error}", 112 workspace_root.display() 113 ) 114 }); 115 match (result, unlock) { 116 (Ok(value), Ok(())) => Ok(value), 117 (Err(error), Ok(())) => Err(error), 118 (Ok(_), Err(error)) => Err(error), 119 (Err(error), Err(unlock)) => Err(format!("{error}; {unlock}")), 120 } 121 } 122 123 fn write_artifact_bundle_impl( 124 workspace_root: &Path, 125 artifacts: Vec<GeneratedArtifact>, 126 simulated_interruption: Option<SimulatedInterruption>, 127 ) -> Result<(), String> { 128 let mut relative_paths = BTreeSet::new(); 129 for artifact in &artifacts { 130 if artifact.relative == TRANSACTION_DIRECTORY 131 || artifact 132 .relative 133 .starts_with(&format!("{TRANSACTION_DIRECTORY}/")) 134 { 135 return Err(format!( 136 "generated artifact path conflicts with transaction authority: {}", 137 artifact.relative 138 )); 139 } 140 if !relative_paths.insert(artifact.relative) { 141 return Err(format!( 142 "generated artifact bundle contains duplicate path: {}", 143 artifact.relative 144 )); 145 } 146 } 147 148 let mut pending = Vec::new(); 149 for artifact in artifacts { 150 let path = validate_workspace_path(workspace_root, artifact.relative, true)?; 151 let parent = path 152 .parent() 153 .ok_or_else(|| format!("artifact path has no parent: {}", artifact.relative))?; 154 fs::create_dir_all(parent) 155 .map_err(|error| format!("create {}: {error}", parent.display()))?; 156 validate_workspace_path(workspace_root, artifact.relative, true)?; 157 158 let (contents, permissions) = match fs::read(&path) { 159 Ok(current) if current == artifact.contents => continue, 160 Ok(current) => ( 161 Some(current), 162 Some( 163 fs::metadata(&path) 164 .map_err(|error| format!("inspect {}: {error}", artifact.relative))? 165 .permissions(), 166 ), 167 ), 168 Err(error) if error.kind() == ErrorKind::NotFound => (None, default_permissions()), 169 Err(error) => return Err(format!("read {}: {error}", artifact.relative)), 170 }; 171 172 pending.push(PendingArtifact { 173 original: OriginalArtifact { 174 path, 175 relative: artifact.relative, 176 contents, 177 permissions, 178 }, 179 contents: artifact.contents, 180 }); 181 } 182 183 if pending.is_empty() { 184 return Ok(()); 185 } 186 let transaction_directory = create_transaction_directory(workspace_root)?; 187 let staged = match stage_artifacts(&transaction_directory, pending) { 188 Ok(staged) => staged, 189 Err(error) => { 190 let cleanup = discard_transaction_directory(workspace_root); 191 return Err(match cleanup { 192 Ok(()) => error, 193 Err(cleanup) => format!("{error}; discard unprepared transaction: {cleanup}"), 194 }); 195 } 196 }; 197 if simulated_interruption == Some(SimulatedInterruption::AfterStaging) { 198 retain_staged_files(staged)?; 199 return Err("simulated interruption after staging artifact transaction".to_owned()); 200 } 201 if let Err(error) = prepare_transaction(workspace_root, &staged) { 202 let cleanup = discard_transaction_directory(workspace_root); 203 return Err(match cleanup { 204 Ok(()) => error, 205 Err(cleanup) => format!("{error}; discard unprepared transaction: {cleanup}"), 206 }); 207 } 208 209 let mut committed = Vec::new(); 210 if simulated_interruption == Some(SimulatedInterruption::AfterCommits(0)) { 211 return Err("simulated interruption after preparing artifact transaction".to_owned()); 212 } 213 for staged_artifact in staged { 214 let StagedArtifact { 215 original, 216 temporary, 217 } = staged_artifact; 218 if let Err(error) = temporary.persist(&original.path) { 219 return fail_and_rollback( 220 workspace_root, 221 &committed, 222 format!("persist {}: {}", original.relative, error.error), 223 ); 224 } 225 let relative = original.relative; 226 let path = original.path.clone(); 227 committed.push(original); 228 if let Err(error) = sync_parent(&path) { 229 return fail_and_rollback( 230 workspace_root, 231 &committed, 232 format!("sync parent for {relative}: {error}"), 233 ); 234 } 235 if simulated_interruption == Some(SimulatedInterruption::AfterCommits(committed.len())) { 236 return Err(format!( 237 "simulated interruption after committing {} artifact(s)", 238 committed.len() 239 )); 240 } 241 } 242 finish_transaction(workspace_root) 243 } 244 245 fn create_transaction_directory(workspace_root: &Path) -> Result<PathBuf, String> { 246 let transaction_directory = transaction_directory(workspace_root); 247 fs::create_dir(&transaction_directory).map_err(|error| { 248 format!( 249 "create artifact transaction directory {}: {error}", 250 transaction_directory.display() 251 ) 252 })?; 253 sync_directory(workspace_root).map_err(|error| { 254 format!("sync workspace after preparing artifact transaction directory: {error}") 255 })?; 256 Ok(transaction_directory) 257 } 258 259 fn stage_artifacts( 260 transaction_directory: &Path, 261 pending: Vec<PendingArtifact>, 262 ) -> Result<Vec<StagedArtifact>, String> { 263 let mut staged = Vec::with_capacity(pending.len()); 264 for pending_artifact in pending { 265 let PendingArtifact { original, contents } = pending_artifact; 266 let mut temporary = NamedTempFile::new_in(transaction_directory) 267 .map_err(|error| format!("stage {}: {error}", original.relative))?; 268 temporary 269 .write_all(&contents) 270 .and_then(|_| temporary.flush()) 271 .map_err(|error| format!("stage {}: {error}", original.relative))?; 272 if let Some(permissions) = original.permissions.as_ref() { 273 fs::set_permissions(temporary.path(), permissions.clone()) 274 .map_err(|error| format!("set permissions for {}: {error}", original.relative))?; 275 } 276 temporary 277 .as_file() 278 .sync_all() 279 .map_err(|error| format!("sync staged {}: {error}", original.relative))?; 280 let staged_bytes = fs::read(temporary.path()) 281 .map_err(|error| format!("verify staged {}: {error}", original.relative))?; 282 if staged_bytes != contents { 283 return Err(format!( 284 "staged bytes do not match generated artifact {}", 285 original.relative 286 )); 287 } 288 staged.push(StagedArtifact { 289 original, 290 temporary, 291 }); 292 } 293 sync_directory(transaction_directory).map_err(|error| { 294 format!( 295 "sync staged artifact transaction {}: {error}", 296 transaction_directory.display() 297 ) 298 })?; 299 Ok(staged) 300 } 301 302 fn retain_staged_files(staged: Vec<StagedArtifact>) -> Result<(), String> { 303 for staged_artifact in staged { 304 let relative = staged_artifact.original.relative; 305 staged_artifact.temporary.keep().map_err(|error| { 306 format!( 307 "retain simulated staged artifact {relative}: {}", 308 error.error 309 ) 310 })?; 311 } 312 Ok(()) 313 } 314 315 fn acquire_workspace_lock(workspace_root: &Path) -> Result<fs::File, String> { 316 let lock_path = workspace_lock_path(workspace_root)?; 317 let lock = OpenOptions::new() 318 .create(true) 319 .truncate(false) 320 .read(true) 321 .write(true) 322 .open(&lock_path) 323 .map_err(|error| { 324 format!( 325 "open generated artifact lock {}: {error}", 326 lock_path.display() 327 ) 328 })?; 329 lock.lock_exclusive().map_err(|error| { 330 format!( 331 "lock generated artifact bundle {}: {error}", 332 lock_path.display() 333 ) 334 })?; 335 Ok(lock) 336 } 337 338 fn workspace_lock_path(workspace_root: &Path) -> Result<PathBuf, String> { 339 validate_workspace_root(workspace_root)?; 340 let canonical_root = fs::canonicalize(workspace_root).map_err(|error| { 341 format!( 342 "canonicalize workspace root {}: {error}", 343 workspace_root.display() 344 ) 345 })?; 346 let digest = Sha256::digest(canonical_root.as_os_str().as_encoded_bytes()); 347 let lock_directory = std::env::temp_dir().join(LOCK_DIRECTORY); 348 fs::create_dir_all(&lock_directory).map_err(|error| { 349 format!( 350 "create lock directory {}: {error}", 351 lock_directory.display() 352 ) 353 })?; 354 let metadata = lock_directory.symlink_metadata().map_err(|error| { 355 format!( 356 "inspect lock directory {}: {error}", 357 lock_directory.display() 358 ) 359 })?; 360 if metadata.file_type().is_symlink() || !metadata.is_dir() { 361 return Err(format!( 362 "generated artifact lock path must be a non-symlink directory: {}", 363 lock_directory.display() 364 )); 365 } 366 let lock_path = lock_directory.join(format!("{}.lock", hex::encode(digest))); 367 if matches!( 368 lock_path.symlink_metadata(), 369 Ok(metadata) if metadata.file_type().is_symlink() 370 ) { 371 return Err(format!( 372 "generated artifact lock must not be a symlink: {}", 373 lock_path.display() 374 )); 375 } 376 Ok(lock_path) 377 } 378 379 fn prepare_transaction(workspace_root: &Path, staged: &[StagedArtifact]) -> Result<(), String> { 380 let transaction_directory = transaction_directory(workspace_root); 381 382 let mut artifacts = Vec::with_capacity(staged.len()); 383 for (index, staged_artifact) in staged.iter().enumerate() { 384 let original = &staged_artifact.original; 385 let (original_byte_length, original_sha256) = 386 if let Some(contents) = original.contents.as_deref() { 387 let backup_path = transaction_directory.join(backup_name(index)); 388 let mut backup = OpenOptions::new() 389 .create_new(true) 390 .write(true) 391 .open(&backup_path) 392 .map_err(|error| { 393 format!( 394 "create transaction backup {}: {error}", 395 backup_path.display() 396 ) 397 })?; 398 backup 399 .write_all(contents) 400 .and_then(|()| backup.flush()) 401 .map_err(|error| { 402 format!( 403 "write transaction backup {}: {error}", 404 backup_path.display() 405 ) 406 })?; 407 if let Some(permissions) = original.permissions.as_ref() { 408 fs::set_permissions(&backup_path, permissions.clone()).map_err(|error| { 409 format!( 410 "set transaction backup permissions {}: {error}", 411 backup_path.display() 412 ) 413 })?; 414 } 415 backup.sync_all().map_err(|error| { 416 format!("sync transaction backup {}: {error}", backup_path.display()) 417 })?; 418 ( 419 Some(u64::try_from(contents.len()).map_err(|_| { 420 format!("{} byte length does not fit in u64", original.relative) 421 })?), 422 Some(sha256_hex(contents)), 423 ) 424 } else { 425 (None, None) 426 }; 427 artifacts.push(ArtifactTransactionEntry { 428 relative: original.relative.to_owned(), 429 original_byte_length, 430 original_sha256, 431 }); 432 } 433 sync_directory(&transaction_directory).map_err(|error| { 434 format!( 435 "sync artifact transaction backups {}: {error}", 436 transaction_directory.display() 437 ) 438 })?; 439 440 let journal = ArtifactTransactionJournal { 441 schema_version: TRANSACTION_SCHEMA_VERSION, 442 artifacts, 443 }; 444 let mut journal_bytes = serde_json::to_vec_pretty(&journal) 445 .map_err(|error| format!("serialize artifact transaction journal: {error}"))?; 446 journal_bytes.push(b'\n'); 447 let journal_path = transaction_directory.join(TRANSACTION_JOURNAL); 448 let mut temporary = NamedTempFile::new_in(&transaction_directory).map_err(|error| { 449 format!( 450 "stage artifact transaction journal {}: {error}", 451 journal_path.display() 452 ) 453 })?; 454 temporary 455 .write_all(&journal_bytes) 456 .and_then(|()| temporary.flush()) 457 .and_then(|()| temporary.as_file().sync_all()) 458 .map_err(|error| { 459 format!( 460 "write artifact transaction journal {}: {error}", 461 journal_path.display() 462 ) 463 })?; 464 temporary 465 .persist_noclobber(&journal_path) 466 .map_err(|error| { 467 format!( 468 "persist artifact transaction journal {}: {}", 469 journal_path.display(), 470 error.error 471 ) 472 })?; 473 sync_directory(&transaction_directory).map_err(|error| { 474 format!( 475 "sync prepared artifact transaction {}: {error}", 476 transaction_directory.display() 477 ) 478 }) 479 } 480 481 fn recover_pending_transaction(workspace_root: &Path) -> Result<(), String> { 482 let transaction_directory = transaction_directory(workspace_root); 483 match transaction_directory.symlink_metadata() { 484 Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => { 485 return Err(format!( 486 "artifact transaction path must be a non-symlink directory: {}", 487 transaction_directory.display() 488 )); 489 } 490 Ok(_) => {} 491 Err(error) if error.kind() == ErrorKind::NotFound => return Ok(()), 492 Err(error) => { 493 return Err(format!( 494 "inspect artifact transaction path {}: {error}", 495 transaction_directory.display() 496 )); 497 } 498 } 499 500 let journal_path = transaction_directory.join(TRANSACTION_JOURNAL); 501 let journal_bytes = match fs::read(&journal_path) { 502 Ok(bytes) => bytes, 503 Err(error) if error.kind() == ErrorKind::NotFound => { 504 return discard_transaction_directory(workspace_root); 505 } 506 Err(error) => { 507 return Err(format!( 508 "read artifact transaction journal {}: {error}", 509 journal_path.display() 510 )); 511 } 512 }; 513 let journal: ArtifactTransactionJournal = serde_json::from_slice(&journal_bytes) 514 .map_err(|error| format!("parse artifact transaction journal: {error}"))?; 515 validate_transaction_journal(workspace_root, &journal)?; 516 517 let mut failures = Vec::new(); 518 for (index, artifact) in journal.artifacts.iter().enumerate().rev() { 519 let path = match validate_workspace_path(workspace_root, &artifact.relative, true) { 520 Ok(path) => path, 521 Err(error) => { 522 failures.push(error); 523 continue; 524 } 525 }; 526 let result = match ( 527 artifact.original_byte_length, 528 artifact.original_sha256.as_deref(), 529 ) { 530 (Some(expected_length), Some(expected_sha256)) => { 531 let backup_path = transaction_directory.join(backup_name(index)); 532 restore_transaction_backup(&path, &backup_path, expected_length, expected_sha256) 533 } 534 (None, None) => match fs::remove_file(&path) { 535 Ok(()) => sync_parent(&path), 536 Err(error) if error.kind() == ErrorKind::NotFound => Ok(()), 537 Err(error) => Err(error), 538 }, 539 _ => unreachable!("validated transaction journal backup shape"), 540 }; 541 if let Err(error) = result { 542 failures.push(format!("{}: {error}", artifact.relative)); 543 } 544 } 545 if !failures.is_empty() { 546 return Err(format!( 547 "recover prepared artifact transaction: {}", 548 failures.join(", ") 549 )); 550 } 551 finish_transaction(workspace_root) 552 } 553 554 fn validate_transaction_journal( 555 workspace_root: &Path, 556 journal: &ArtifactTransactionJournal, 557 ) -> Result<(), String> { 558 if journal.schema_version != TRANSACTION_SCHEMA_VERSION { 559 return Err(format!( 560 "artifact transaction schema_version must be {TRANSACTION_SCHEMA_VERSION}" 561 )); 562 } 563 if journal.artifacts.is_empty() { 564 return Err("artifact transaction must contain at least one artifact".to_owned()); 565 } 566 let mut paths = BTreeSet::new(); 567 for artifact in &journal.artifacts { 568 validate_workspace_path(workspace_root, &artifact.relative, true)?; 569 if artifact.relative == TRANSACTION_DIRECTORY 570 || artifact 571 .relative 572 .starts_with(&format!("{TRANSACTION_DIRECTORY}/")) 573 { 574 return Err(format!( 575 "artifact transaction path conflicts with transaction authority: {}", 576 artifact.relative 577 )); 578 } 579 if !paths.insert(artifact.relative.as_str()) { 580 return Err(format!( 581 "artifact transaction contains duplicate path: {}", 582 artifact.relative 583 )); 584 } 585 match ( 586 artifact.original_byte_length, 587 artifact.original_sha256.as_deref(), 588 ) { 589 (Some(_), Some(digest)) => validate_sha256(digest)?, 590 (None, None) => {} 591 _ => { 592 return Err(format!( 593 "artifact transaction backup metadata is incomplete for {}", 594 artifact.relative 595 )); 596 } 597 } 598 } 599 Ok(()) 600 } 601 602 fn restore_transaction_backup( 603 artifact_path: &Path, 604 backup_path: &Path, 605 expected_length: u64, 606 expected_sha256: &str, 607 ) -> io::Result<()> { 608 let metadata = backup_path.symlink_metadata()?; 609 if metadata.file_type().is_symlink() || !metadata.is_file() { 610 return Err(io::Error::other(format!( 611 "transaction backup must be a regular non-symlink file: {}", 612 backup_path.display() 613 ))); 614 } 615 let contents = fs::read(backup_path)?; 616 let actual_length = u64::try_from(contents.len()) 617 .map_err(|_| io::Error::other("transaction backup length does not fit in u64"))?; 618 if actual_length != expected_length || sha256_hex(&contents) != expected_sha256 { 619 return Err(io::Error::other(format!( 620 "transaction backup integrity mismatch: {}", 621 backup_path.display() 622 ))); 623 } 624 let permissions = metadata.permissions(); 625 let staging_directory = backup_path 626 .parent() 627 .ok_or_else(|| io::Error::other("transaction backup path has no parent"))?; 628 restore_bytes( 629 artifact_path, 630 &contents, 631 Some(&permissions), 632 staging_directory, 633 ) 634 } 635 636 fn fail_and_rollback( 637 workspace_root: &Path, 638 committed: &[OriginalArtifact], 639 failure: String, 640 ) -> Result<(), String> { 641 match rollback_artifacts(workspace_root, committed) { 642 Ok(()) => match finish_transaction(workspace_root) { 643 Ok(()) => Err(format!("{failure}; restored committed artifacts")), 644 Err(cleanup) => Err(format!( 645 "{failure}; restored committed artifacts; transaction cleanup failed: {cleanup}" 646 )), 647 }, 648 Err(rollback) => Err(format!( 649 "{failure}; rollback also failed: {rollback}; prepared transaction retained for recovery" 650 )), 651 } 652 } 653 654 fn finish_transaction(workspace_root: &Path) -> Result<(), String> { 655 let transaction_directory = transaction_directory(workspace_root); 656 let journal_path = transaction_directory.join(TRANSACTION_JOURNAL); 657 fs::remove_file(&journal_path).map_err(|error| { 658 format!( 659 "remove artifact transaction journal {}: {error}", 660 journal_path.display() 661 ) 662 })?; 663 sync_directory(&transaction_directory).map_err(|error| { 664 format!( 665 "sync completed artifact transaction {}: {error}", 666 transaction_directory.display() 667 ) 668 })?; 669 discard_transaction_directory(workspace_root) 670 } 671 672 fn discard_transaction_directory(workspace_root: &Path) -> Result<(), String> { 673 let transaction_directory = transaction_directory(workspace_root); 674 match fs::remove_dir_all(&transaction_directory) { 675 Ok(()) => sync_directory(workspace_root).map_err(|error| { 676 format!( 677 "sync workspace after removing artifact transaction {}: {error}", 678 transaction_directory.display() 679 ) 680 }), 681 Err(error) if error.kind() == ErrorKind::NotFound => Ok(()), 682 Err(error) => Err(format!( 683 "remove artifact transaction directory {}: {error}", 684 transaction_directory.display() 685 )), 686 } 687 } 688 689 fn transaction_directory(workspace_root: &Path) -> PathBuf { 690 workspace_root.join(TRANSACTION_DIRECTORY) 691 } 692 693 fn backup_name(index: usize) -> String { 694 format!("{index:08}.backup") 695 } 696 697 fn sha256_hex(bytes: &[u8]) -> String { 698 hex::encode(Sha256::digest(bytes)) 699 } 700 701 fn validate_sha256(digest: &str) -> Result<(), String> { 702 if digest.len() != 64 703 || !digest 704 .as_bytes() 705 .iter() 706 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(byte)) 707 { 708 return Err( 709 "artifact transaction digest must be 64 lowercase hexadecimal bytes".to_owned(), 710 ); 711 } 712 Ok(()) 713 } 714 715 pub(super) fn validate_workspace_path( 716 workspace_root: &Path, 717 relative: &str, 718 allow_missing: bool, 719 ) -> Result<PathBuf, String> { 720 let relative_path = Path::new(relative); 721 if relative.is_empty() 722 || relative.contains('\\') 723 || relative_path.is_absolute() 724 || !relative_path 725 .components() 726 .all(|component| matches!(component, Component::Normal(_))) 727 { 728 return Err(format!( 729 "artifact path must be normalized and workspace-relative: {relative}" 730 )); 731 } 732 733 validate_workspace_root(workspace_root)?; 734 735 let component_count = relative_path.components().count(); 736 let mut current = workspace_root.to_path_buf(); 737 for (index, component) in relative_path.components().enumerate() { 738 let Component::Normal(segment) = component else { 739 return Err(format!("artifact path is not normalized: {relative}")); 740 }; 741 current.push(segment); 742 match current.symlink_metadata() { 743 Ok(metadata) if metadata.file_type().is_symlink() => { 744 return Err(format!( 745 "artifact path contains a symlink component: {}", 746 current.display() 747 )); 748 } 749 Ok(metadata) if index + 1 == component_count && !metadata.is_file() => { 750 return Err(format!("artifact must be a regular file: {relative}")); 751 } 752 Ok(metadata) if index + 1 < component_count && !metadata.is_dir() => { 753 return Err(format!( 754 "artifact parent must be a directory: {}", 755 current.display() 756 )); 757 } 758 Ok(_) => {} 759 Err(error) if error.kind() == ErrorKind::NotFound && allow_missing => break, 760 Err(error) => { 761 return Err(format!( 762 "inspect artifact path {}: {error}", 763 current.display() 764 )); 765 } 766 } 767 } 768 Ok(workspace_root.join(relative_path)) 769 } 770 771 fn validate_workspace_root(workspace_root: &Path) -> Result<(), String> { 772 let root_metadata = workspace_root.symlink_metadata().map_err(|error| { 773 format!( 774 "inspect workspace root {}: {error}", 775 workspace_root.display() 776 ) 777 })?; 778 if root_metadata.file_type().is_symlink() || !root_metadata.is_dir() { 779 return Err(format!( 780 "workspace root must be a non-symlink directory: {}", 781 workspace_root.display() 782 )); 783 } 784 Ok(()) 785 } 786 787 fn rollback_artifacts(workspace_root: &Path, committed: &[OriginalArtifact]) -> Result<(), String> { 788 let mut failures = Vec::new(); 789 let staging_directory = transaction_directory(workspace_root); 790 for original in committed.iter().rev() { 791 let result = if let Some(contents) = original.contents.as_ref() { 792 restore_file(original, contents, &staging_directory) 793 } else { 794 match fs::remove_file(&original.path) { 795 Ok(()) => sync_parent(&original.path), 796 Err(error) if error.kind() == ErrorKind::NotFound => Ok(()), 797 Err(error) => Err(error), 798 } 799 }; 800 if let Err(error) = result { 801 failures.push(format!("{}: {error}", original.relative)); 802 } 803 } 804 if failures.is_empty() { 805 Ok(()) 806 } else { 807 Err(failures.join(", ")) 808 } 809 } 810 811 fn restore_file( 812 original: &OriginalArtifact, 813 contents: &[u8], 814 staging_directory: &Path, 815 ) -> io::Result<()> { 816 restore_bytes( 817 &original.path, 818 contents, 819 original.permissions.as_ref(), 820 staging_directory, 821 ) 822 } 823 824 fn restore_bytes( 825 path: &Path, 826 contents: &[u8], 827 permissions: Option<&fs::Permissions>, 828 staging_directory: &Path, 829 ) -> io::Result<()> { 830 let parent = path 831 .parent() 832 .ok_or_else(|| io::Error::other("artifact path has no parent"))?; 833 fs::create_dir_all(parent)?; 834 let mut temporary = NamedTempFile::new_in(staging_directory)?; 835 temporary.write_all(contents)?; 836 temporary.flush()?; 837 if let Some(permissions) = permissions { 838 fs::set_permissions(temporary.path(), permissions.clone())?; 839 } 840 temporary.as_file().sync_all()?; 841 temporary.persist(path).map_err(|error| error.error)?; 842 sync_parent(path) 843 } 844 845 fn sync_parent(path: &Path) -> io::Result<()> { 846 let parent = path 847 .parent() 848 .ok_or_else(|| io::Error::other("artifact path has no parent"))?; 849 sync_directory(parent) 850 } 851 852 fn sync_directory(path: &Path) -> io::Result<()> { 853 fs::File::open(path)?.sync_all() 854 } 855 856 #[cfg(unix)] 857 fn default_permissions() -> Option<fs::Permissions> { 858 use std::os::unix::fs::PermissionsExt; 859 Some(fs::Permissions::from_mode(0o644)) 860 } 861 862 #[cfg(not(unix))] 863 fn default_permissions() -> Option<fs::Permissions> { 864 None 865 } 866 867 #[cfg(test)] 868 mod tests { 869 use super::*; 870 871 #[test] 872 fn bundle_rejects_duplicate_paths_before_writing() { 873 let workspace = tempfile::TempDir::new().expect("workspace"); 874 let error = with_artifact_bundle_transaction(workspace.path(), |transaction| { 875 transaction.write(vec![ 876 GeneratedArtifact { 877 relative: "generated/value.txt", 878 contents: b"first\n".to_vec(), 879 }, 880 GeneratedArtifact { 881 relative: "generated/value.txt", 882 contents: b"second\n".to_vec(), 883 }, 884 ]) 885 }) 886 .expect_err("duplicate artifact paths must fail"); 887 assert!(error.contains("duplicate path")); 888 assert!(!workspace.path().join("generated/value.txt").exists()); 889 } 890 891 #[test] 892 fn workspace_lock_excludes_a_second_file_descriptor() { 893 let workspace = tempfile::TempDir::new().expect("workspace"); 894 let first = acquire_workspace_lock(workspace.path()).expect("first workspace lock"); 895 let lock_path = workspace_lock_path(workspace.path()).expect("workspace lock path"); 896 let second = OpenOptions::new() 897 .read(true) 898 .write(true) 899 .open(lock_path) 900 .expect("second lock descriptor"); 901 902 assert!( 903 second.try_lock_exclusive().is_err(), 904 "the second descriptor must observe the held advisory lock" 905 ); 906 FileExt::unlock(&first).expect("release first lock"); 907 second 908 .try_lock_exclusive() 909 .expect("second descriptor acquires released lock"); 910 FileExt::unlock(&second).expect("release second lock"); 911 } 912 913 #[test] 914 fn next_transaction_recovers_an_interrupted_multi_file_commit() { 915 let workspace = tempfile::TempDir::new().expect("workspace"); 916 fs::create_dir_all(workspace.path().join("generated")).expect("artifact directory"); 917 fs::write(workspace.path().join("generated/first.txt"), b"first-old\n") 918 .expect("first original"); 919 fs::write( 920 workspace.path().join("generated/second.txt"), 921 b"second-old\n", 922 ) 923 .expect("second original"); 924 925 let interrupted = with_artifact_bundle_transaction(workspace.path(), |_| { 926 write_artifact_bundle_impl( 927 workspace.path(), 928 vec![ 929 GeneratedArtifact { 930 relative: "generated/first.txt", 931 contents: b"first-interrupted\n".to_vec(), 932 }, 933 GeneratedArtifact { 934 relative: "generated/second.txt", 935 contents: b"second-interrupted\n".to_vec(), 936 }, 937 ], 938 Some(SimulatedInterruption::AfterCommits(1)), 939 ) 940 }) 941 .expect_err("simulated process interruption"); 942 assert!(interrupted.contains("simulated interruption")); 943 assert_eq!( 944 fs::read(workspace.path().join("generated/first.txt")).expect("mixed first"), 945 b"first-interrupted\n" 946 ); 947 assert_eq!( 948 fs::read(workspace.path().join("generated/second.txt")).expect("mixed second"), 949 b"second-old\n" 950 ); 951 assert!( 952 transaction_directory(workspace.path()) 953 .join(TRANSACTION_JOURNAL) 954 .is_file(), 955 "prepared journal must survive an interrupted commit" 956 ); 957 958 with_artifact_bundle_transaction(workspace.path(), |transaction| { 959 assert_eq!( 960 fs::read(workspace.path().join("generated/first.txt")).expect("recovered first"), 961 b"first-old\n" 962 ); 963 assert_eq!( 964 fs::read(workspace.path().join("generated/second.txt")).expect("recovered second"), 965 b"second-old\n" 966 ); 967 transaction.write(vec![ 968 GeneratedArtifact { 969 relative: "generated/first.txt", 970 contents: b"first-final\n".to_vec(), 971 }, 972 GeneratedArtifact { 973 relative: "generated/second.txt", 974 contents: b"second-final\n".to_vec(), 975 }, 976 ]) 977 }) 978 .expect("recover and replace bundle"); 979 980 assert_eq!( 981 fs::read(workspace.path().join("generated/first.txt")).expect("final first"), 982 b"first-final\n" 983 ); 984 assert_eq!( 985 fs::read(workspace.path().join("generated/second.txt")).expect("final second"), 986 b"second-final\n" 987 ); 988 assert!(!transaction_directory(workspace.path()).exists()); 989 } 990 991 #[test] 992 fn unjournaled_stages_are_recovered_without_target_directory_residue() { 993 let workspace = tempfile::TempDir::new().expect("workspace"); 994 let generated = workspace.path().join("generated"); 995 fs::create_dir_all(&generated).expect("artifact directory"); 996 let artifact = generated.join("value.txt"); 997 fs::write(&artifact, b"old\n").expect("original artifact"); 998 999 with_artifact_bundle_transaction(workspace.path(), |_| { 1000 write_artifact_bundle_impl( 1001 workspace.path(), 1002 vec![GeneratedArtifact { 1003 relative: "generated/value.txt", 1004 contents: b"interrupted\n".to_vec(), 1005 }], 1006 Some(SimulatedInterruption::AfterStaging), 1007 ) 1008 }) 1009 .expect_err("simulated pre-journal interruption"); 1010 1011 assert_eq!(fs::read(&artifact).expect("unchanged target"), b"old\n"); 1012 assert!( 1013 transaction_directory(workspace.path()).is_dir(), 1014 "simulated abrupt death must retain a recoverable transaction authority" 1015 ); 1016 assert!( 1017 !transaction_directory(workspace.path()) 1018 .join(TRANSACTION_JOURNAL) 1019 .exists(), 1020 "the simulated interruption must precede the durable journal" 1021 ); 1022 let target_entries = fs::read_dir(&generated) 1023 .expect("target directory") 1024 .map(|entry| entry.expect("target entry").file_name()) 1025 .collect::<Vec<_>>(); 1026 assert_eq!( 1027 target_entries, 1028 vec![std::ffi::OsString::from("value.txt")], 1029 "staging must not leave temporary files beside governed targets" 1030 ); 1031 assert!( 1032 fs::read_dir(transaction_directory(workspace.path())) 1033 .expect("transaction stages") 1034 .next() 1035 .is_some(), 1036 "the test interruption must retain at least one staged file" 1037 ); 1038 1039 with_artifact_bundle_transaction(workspace.path(), |_| Ok(())) 1040 .expect("next locked operation recovers pre-journal stages"); 1041 assert_eq!(fs::read(&artifact).expect("recovered target"), b"old\n"); 1042 assert!( 1043 !transaction_directory(workspace.path()).exists(), 1044 "recovery must remove every transaction-owned stage" 1045 ); 1046 } 1047 1048 #[test] 1049 fn recovery_fails_closed_and_retains_a_corrupt_prepared_transaction() { 1050 let workspace = tempfile::TempDir::new().expect("workspace"); 1051 fs::create_dir_all(workspace.path().join("generated")).expect("artifact directory"); 1052 fs::write(workspace.path().join("generated/first.txt"), b"first-old\n") 1053 .expect("first original"); 1054 fs::write( 1055 workspace.path().join("generated/second.txt"), 1056 b"second-old\n", 1057 ) 1058 .expect("second original"); 1059 1060 with_artifact_bundle_transaction(workspace.path(), |_| { 1061 write_artifact_bundle_impl( 1062 workspace.path(), 1063 vec![ 1064 GeneratedArtifact { 1065 relative: "generated/first.txt", 1066 contents: b"first-new\n".to_vec(), 1067 }, 1068 GeneratedArtifact { 1069 relative: "generated/second.txt", 1070 contents: b"second-new\n".to_vec(), 1071 }, 1072 ], 1073 Some(SimulatedInterruption::AfterCommits(1)), 1074 ) 1075 }) 1076 .expect_err("simulated process interruption"); 1077 fs::write( 1078 transaction_directory(workspace.path()).join(backup_name(0)), 1079 b"corrupt\n", 1080 ) 1081 .expect("corrupt transaction backup"); 1082 1083 let error = with_artifact_bundle_transaction(workspace.path(), |_| Ok(())) 1084 .expect_err("corrupt backup must prevent validation or writing"); 1085 assert!(error.contains("integrity mismatch")); 1086 assert!( 1087 transaction_directory(workspace.path()) 1088 .join(TRANSACTION_JOURNAL) 1089 .is_file(), 1090 "failed recovery must retain its durable journal" 1091 ); 1092 } 1093 }