lib

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

download.rs (14458B)


      1 //! Explicit, caller-driven asset acquisition.
      2 
      3 use std::fmt;
      4 use std::fs::{File, OpenOptions};
      5 use std::io::{self, Write};
      6 use std::path::{Path, PathBuf};
      7 
      8 use fs2::FileExt;
      9 use tempfile::NamedTempFile;
     10 
     11 use crate::asset::{inspect, io_error, verify_file};
     12 use crate::{AssetSpec, AssetStatus, Error};
     13 
     14 /// Stable phase attached to an injected fetch failure.
     15 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     16 #[non_exhaustive]
     17 pub enum FetchFailurePhase {
     18     /// Establishing the source connection or opening the source.
     19     Connect,
     20     /// Receiving source metadata or an initial response.
     21     Response,
     22     /// Streaming source bytes.
     23     Read,
     24     /// Caller-requested cancellation.
     25     Cancelled,
     26 }
     27 
     28 impl fmt::Display for FetchFailurePhase {
     29     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     30         formatter.write_str(match self {
     31             Self::Connect => "connect",
     32             Self::Response => "response",
     33             Self::Read => "read",
     34             Self::Cancelled => "cancellation",
     35         })
     36     }
     37 }
     38 
     39 /// Explicit byte source supplied by the host.
     40 ///
     41 /// Implementations own transport execution and cancellation. Failure details
     42 /// cannot enter this API, preventing source URLs or credentials from leaking.
     43 pub trait Fetcher {
     44     /// Streams the requested source into the bounded destination.
     45     fn fetch(&self, source: &str, destination: &mut dyn Write) -> Result<(), Error>;
     46 }
     47 
     48 /// Installs a verified asset into a caller-selected existing directory.
     49 ///
     50 /// The fetcher is invoked only when the final asset is absent or invalid. The
     51 /// existing destination remains untouched until a fully written staging file
     52 /// passes exact length and SHA-256 validation.
     53 pub fn acquire(
     54     directory: impl AsRef<Path>,
     55     spec: &AssetSpec,
     56     fetcher: &dyn Fetcher,
     57 ) -> Result<AssetStatus, Error> {
     58     let directory = safe_directory(directory.as_ref())?;
     59     let destination = directory.join(spec.file_name());
     60     match inspect(&destination, spec)? {
     61         AssetStatus::Available => return Ok(AssetStatus::Available),
     62         AssetStatus::Missing | AssetStatus::Invalid => {}
     63     }
     64 
     65     reject_symlink(&destination)?;
     66     let lock_path = directory.join(format!(".{}.lock", spec.file_name()));
     67     reject_symlink(&lock_path)?;
     68     let lock = open_lock(&lock_path)?;
     69     lock.try_lock_exclusive()
     70         .map_err(|_| Error::AssetDestinationBusy)?;
     71 
     72     if inspect(&destination, spec)? == AssetStatus::Available {
     73         return Ok(AssetStatus::Available);
     74     }
     75 
     76     let mut staging = NamedTempFile::new_in(&directory)
     77         .map_err(|error| io_error("create asset staging file", error))?;
     78     let (fetch_result, observed, overflowed) = {
     79         let mut writer = BoundedWriter::new(staging.as_file_mut(), spec.byte_size());
     80         let fetch_result = fetcher.fetch(spec.source(), &mut writer);
     81         (fetch_result, writer.observed, writer.overflowed)
     82     };
     83     if overflowed {
     84         return Err(Error::AssetSizeMismatch {
     85             expected: spec.byte_size(),
     86             actual: spec.byte_size().saturating_add(1),
     87         });
     88     }
     89     fetch_result?;
     90     if observed != spec.byte_size() {
     91         return Err(Error::AssetSizeMismatch {
     92             expected: spec.byte_size(),
     93             actual: observed,
     94         });
     95     }
     96     staging
     97         .as_file_mut()
     98         .sync_all()
     99         .map_err(|error| io_error("sync asset staging file", error))?;
    100     verify_file(staging.path(), spec)?;
    101     reject_symlink(&destination)?;
    102     staging
    103         .persist(&destination)
    104         .map_err(|error| io_error("finalize asset", error.error))?;
    105     sync_directory(&directory)?;
    106     verify_file(&destination, spec)?;
    107     Ok(AssetStatus::Available)
    108 }
    109 
    110 fn safe_directory(directory: &Path) -> Result<PathBuf, Error> {
    111     let metadata = directory
    112         .symlink_metadata()
    113         .map_err(|error| io_error("inspect asset directory", error))?;
    114     if metadata.file_type().is_symlink() || !metadata.is_dir() {
    115         return Err(Error::UnsafeAssetDestination);
    116     }
    117     directory
    118         .canonicalize()
    119         .map_err(|error| io_error("resolve asset directory", error))
    120 }
    121 
    122 fn reject_symlink(path: &Path) -> Result<(), Error> {
    123     match path.symlink_metadata() {
    124         Ok(metadata) if metadata.file_type().is_symlink() => Err(Error::UnsafeAssetDestination),
    125         Ok(_) => Ok(()),
    126         Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
    127         Err(error) => Err(io_error("inspect asset destination", error)),
    128     }
    129 }
    130 
    131 fn open_lock(path: &Path) -> Result<File, Error> {
    132     OpenOptions::new()
    133         .read(true)
    134         .write(true)
    135         .create(true)
    136         .truncate(false)
    137         .open(path)
    138         .map_err(|error| io_error("open asset lock", error))
    139 }
    140 
    141 #[cfg(unix)]
    142 fn sync_directory(directory: &Path) -> Result<(), Error> {
    143     File::open(directory)
    144         .and_then(|file| file.sync_all())
    145         .map_err(|error| io_error("sync asset directory", error))
    146 }
    147 
    148 #[cfg(not(unix))]
    149 fn sync_directory(_directory: &Path) -> Result<(), Error> {
    150     Ok(())
    151 }
    152 
    153 struct BoundedWriter<'a> {
    154     destination: &'a mut File,
    155     maximum: u64,
    156     observed: u64,
    157     overflowed: bool,
    158 }
    159 
    160 impl<'a> BoundedWriter<'a> {
    161     fn new(destination: &'a mut File, maximum: u64) -> Self {
    162         Self {
    163             destination,
    164             maximum,
    165             observed: 0,
    166             overflowed: false,
    167         }
    168     }
    169 }
    170 
    171 impl Write for BoundedWriter<'_> {
    172     fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
    173         let incoming = u64::try_from(buffer.len()).unwrap_or(u64::MAX);
    174         if self.observed.saturating_add(incoming) > self.maximum {
    175             self.overflowed = true;
    176             return Err(io::Error::new(
    177                 io::ErrorKind::FileTooLarge,
    178                 "asset exceeds declared size",
    179             ));
    180         }
    181         let written = self.destination.write(buffer)?;
    182         self.observed = self
    183             .observed
    184             .saturating_add(u64::try_from(written).unwrap_or(u64::MAX));
    185         Ok(written)
    186     }
    187 
    188     fn flush(&mut self) -> io::Result<()> {
    189         self.destination.flush()
    190     }
    191 }
    192 
    193 #[cfg(test)]
    194 mod tests {
    195     use std::fs::{self, OpenOptions};
    196     use std::io::Write;
    197 
    198     use fs2::FileExt;
    199     use sha2::{Digest, Sha256};
    200     use tempfile::tempdir;
    201 
    202     use super::{BoundedWriter, FetchFailurePhase, Fetcher, acquire};
    203     use crate::asset::inspect;
    204     use crate::{AssetSpec, AssetStatus, Error};
    205 
    206     struct BytesFetcher(Vec<u8>);
    207 
    208     impl Fetcher for BytesFetcher {
    209         fn fetch(&self, _source: &str, destination: &mut dyn Write) -> Result<(), Error> {
    210             destination.write_all(&self.0).map_err(|_| Error::Fetch {
    211                 phase: FetchFailurePhase::Read,
    212             })
    213         }
    214     }
    215 
    216     struct InterruptedFetcher;
    217 
    218     impl Fetcher for InterruptedFetcher {
    219         fn fetch(&self, _source: &str, destination: &mut dyn Write) -> Result<(), Error> {
    220             destination
    221                 .write_all(b"partial")
    222                 .map_err(|_| Error::Fetch {
    223                     phase: FetchFailurePhase::Read,
    224                 })?;
    225             Err(Error::Fetch {
    226                 phase: FetchFailurePhase::Cancelled,
    227             })
    228         }
    229     }
    230 
    231     struct PanicFetcher;
    232 
    233     impl Fetcher for PanicFetcher {
    234         fn fetch(&self, _source: &str, _destination: &mut dyn Write) -> Result<(), Error> {
    235             panic!("available assets must not invoke the fetcher")
    236         }
    237     }
    238 
    239     fn spec(bytes: &[u8]) -> AssetSpec {
    240         AssetSpec::new(
    241             "test-v1",
    242             "geonames-test.db",
    243             "https://assets.example/geonames-test.db",
    244             "assets.example",
    245             u64::try_from(bytes.len()).expect("fixture length"),
    246             Sha256::digest(bytes).into(),
    247         )
    248         .expect("asset spec")
    249     }
    250 
    251     #[test]
    252     fn missing_and_successful_acquisition_are_explicit() {
    253         let directory = tempdir().expect("tempdir");
    254         let bytes = b"verified geonames fixture";
    255         let spec = spec(bytes);
    256         let path = directory.path().join(spec.file_name());
    257         assert_eq!(inspect(&path, &spec), Ok(AssetStatus::Missing));
    258         assert_eq!(
    259             acquire(directory.path(), &spec, &BytesFetcher(bytes.to_vec())),
    260             Ok(AssetStatus::Available)
    261         );
    262         assert_eq!(inspect(&path, &spec), Ok(AssetStatus::Available));
    263         assert_eq!(fs::read(path).expect("read installed asset"), bytes);
    264     }
    265 
    266     #[test]
    267     fn interrupted_acquisition_leaves_no_destination_or_staging_file() {
    268         let directory = tempdir().expect("tempdir");
    269         let spec = spec(b"complete bytes");
    270         let error = acquire(directory.path(), &spec, &InterruptedFetcher)
    271             .expect_err("interrupted fetch must fail");
    272         assert!(matches!(
    273             error,
    274             Error::Fetch {
    275                 phase: FetchFailurePhase::Cancelled
    276             }
    277         ));
    278         assert!(!directory.path().join(spec.file_name()).exists());
    279         let entries = fs::read_dir(directory.path())
    280             .expect("read directory")
    281             .filter_map(Result::ok)
    282             .map(|entry| entry.file_name())
    283             .collect::<Vec<_>>();
    284         assert_eq!(
    285             entries,
    286             vec![std::ffi::OsString::from(format!(
    287                 ".{}.lock",
    288                 spec.file_name()
    289             ))]
    290         );
    291     }
    292 
    293     #[test]
    294     fn oversized_and_hash_mismatched_streams_never_replace_existing_asset() {
    295         let directory = tempdir().expect("tempdir");
    296         let expected = b"expected";
    297         let spec = spec(expected);
    298         let path = directory.path().join(spec.file_name());
    299         fs::write(&path, b"old").expect("old destination");
    300 
    301         assert!(matches!(
    302             acquire(
    303                 directory.path(),
    304                 &spec,
    305                 &BytesFetcher(b"expected-extra".to_vec())
    306             ),
    307             Err(Error::AssetSizeMismatch { .. })
    308         ));
    309         assert_eq!(fs::read(&path).expect("preserved old bytes"), b"old");
    310 
    311         assert_eq!(
    312             acquire(directory.path(), &spec, &BytesFetcher(b"notright".to_vec())),
    313             Err(Error::AssetHashMismatch)
    314         );
    315         assert_eq!(fs::read(path).expect("preserved old bytes"), b"old");
    316     }
    317 
    318     #[cfg(unix)]
    319     #[test]
    320     fn symlink_directories_and_destinations_are_rejected() {
    321         use std::os::unix::fs::symlink;
    322 
    323         let root = tempdir().expect("tempdir");
    324         let real = root.path().join("real");
    325         fs::create_dir(&real).expect("real directory");
    326         let linked = root.path().join("linked");
    327         symlink(&real, &linked).expect("directory symlink");
    328         let spec = spec(b"asset");
    329         assert_eq!(
    330             acquire(&linked, &spec, &BytesFetcher(b"asset".to_vec())),
    331             Err(Error::UnsafeAssetDestination)
    332         );
    333 
    334         let target = real.join(spec.file_name());
    335         let outside = root.path().join("outside");
    336         fs::write(&outside, b"outside").expect("outside file");
    337         symlink(&outside, &target).expect("destination symlink");
    338         assert_eq!(
    339             acquire(&real, &spec, &BytesFetcher(b"asset".to_vec())),
    340             Err(Error::UnsafeAssetDestination)
    341         );
    342     }
    343 
    344     #[test]
    345     fn available_short_busy_and_invalid_directory_paths_are_explicit() {
    346         let directory = tempdir().expect("tempdir");
    347         let bytes = b"asset";
    348         let spec = spec(bytes);
    349         fs::write(directory.path().join(spec.file_name()), bytes).expect("available asset");
    350         let owned_directory = directory.path().to_path_buf();
    351         assert_eq!(
    352             acquire(&owned_directory, &spec, &PanicFetcher),
    353             Ok(AssetStatus::Available)
    354         );
    355 
    356         fs::write(directory.path().join(spec.file_name()), b"old").expect("invalid asset");
    357         assert!(matches!(
    358             acquire(directory.path(), &spec, &BytesFetcher(b"a".to_vec())),
    359             Err(Error::AssetSizeMismatch {
    360                 expected: 5,
    361                 actual: 1
    362             })
    363         ));
    364 
    365         let lock_path = directory.path().join(format!(".{}.lock", spec.file_name()));
    366         let lock = OpenOptions::new()
    367             .read(true)
    368             .write(true)
    369             .create(true)
    370             .truncate(false)
    371             .open(lock_path)
    372             .expect("lock file");
    373         lock.lock_exclusive().expect("exclusive lock");
    374         assert_eq!(
    375             acquire(directory.path(), &spec, &BytesFetcher(bytes.to_vec())),
    376             Err(Error::AssetDestinationBusy)
    377         );
    378 
    379         let not_directory = directory.path().join("plain-file");
    380         fs::write(&not_directory, b"file").expect("plain file");
    381         assert_eq!(
    382             acquire(&not_directory, &spec, &BytesFetcher(bytes.to_vec())),
    383             Err(Error::UnsafeAssetDestination)
    384         );
    385         assert!(matches!(
    386             acquire(
    387                 directory.path().join("missing"),
    388                 &spec,
    389                 &BytesFetcher(bytes.to_vec())
    390             ),
    391             Err(Error::Io {
    392                 operation: "inspect asset directory",
    393                 kind: std::io::ErrorKind::NotFound,
    394             })
    395         ));
    396     }
    397 
    398     #[test]
    399     fn fetch_phases_and_bounded_writer_flush_are_covered() {
    400         assert_eq!(FetchFailurePhase::Connect.to_string(), "connect");
    401         assert_eq!(FetchFailurePhase::Response.to_string(), "response");
    402         assert_eq!(FetchFailurePhase::Read.to_string(), "read");
    403         assert_eq!(FetchFailurePhase::Cancelled.to_string(), "cancellation");
    404 
    405         let mut file = tempfile::tempfile().expect("temporary file");
    406         let mut writer = BoundedWriter::new(&mut file, 4);
    407         writer.write_all(b"data").expect("bounded write");
    408         writer.flush().expect("flush");
    409         assert_eq!(writer.observed, 4);
    410         assert!(!writer.overflowed);
    411     }
    412 
    413     #[test]
    414     fn invalid_lock_entry_is_reported_as_an_io_failure() {
    415         let directory = tempdir().expect("tempdir");
    416         let bytes = b"asset";
    417         let spec = spec(bytes);
    418         let lock_path = directory.path().join(format!(".{}.lock", spec.file_name()));
    419         fs::create_dir(lock_path).expect("invalid lock directory");
    420         let error = acquire(directory.path(), &spec, &BytesFetcher(bytes.to_vec()))
    421             .expect_err("directory lock entry must fail to open");
    422         assert!(matches!(
    423             error,
    424             Error::Io {
    425                 operation: "open asset lock",
    426                 ..
    427             }
    428         ));
    429     }
    430 }