lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }