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(¬_directory, b"file").expect("plain file"); 381 assert_eq!( 382 acquire(¬_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 }