diff --git a/core/archipelago/Cargo.toml b/core/archipelago/Cargo.toml index eb1ad18c..4a2e974a 100644 --- a/core/archipelago/Cargo.toml +++ b/core/archipelago/Cargo.toml @@ -148,6 +148,7 @@ iroh-blobs = { version = "0.103", optional = true } lofty = "0.24.0" cashu = { version = "0.17.5", default-features = false, features = ["wallet"] } +tempfile = "3.10" + [dev-dependencies] tokio-test = "0.4" -tempfile = "3.10" diff --git a/core/archipelago/src/api/handler/content.rs b/core/archipelago/src/api/handler/content.rs index b41fc286..d9453bec 100644 --- a/core/archipelago/src/api/handler/content.rs +++ b/core/archipelago/src/api/handler/content.rs @@ -139,10 +139,23 @@ impl ApiHandler { } // Parse Range header for streaming support - let range = headers - .get("range") - .and_then(|v| v.to_str().ok()) - .and_then(content_server::parse_range_header); + let range = match headers.get("range") { + None => None, + Some(value) => match value + .to_str() + .ok() + .and_then(content_server::parse_range_header) + { + Some(range) => Some(range), + None => { + return Ok(build_response( + StatusCode::BAD_REQUEST, + "text/plain", + hyper::Body::from("Invalid byte range"), + )) + } + }, + }; match content_server::serve_content( &config.data_dir, @@ -155,6 +168,7 @@ impl ApiHandler { ) .await { + Ok(content_server::ServeResult::Stream(body)) => body.into_response(), Ok(content_server::ServeResult::Ok(bytes, mime_type)) => { let len = bytes.len(); Ok(Response::builder() @@ -496,6 +510,7 @@ impl ApiHandler { } match content_server::serve_content_preview(&config.data_dir, content_id).await { + Ok(content_server::PreviewResult::Stream(body)) => body.into_response(), Ok(content_server::PreviewResult::FullContent(bytes, mime_type)) => { let len = bytes.len(); Ok(Response::builder() diff --git a/core/archipelago/src/content_server.rs b/core/archipelago/src/content_server.rs index f9e30196..5f6e59f1 100644 --- a/core/archipelago/src/content_server.rs +++ b/core/archipelago/src/content_server.rs @@ -202,26 +202,37 @@ pub async fn set_availability(data_dir: &Path, id: &str, availability: Availabil } /// A byte range request (start, optional end). -pub struct ByteRange { - pub start: u64, - pub end: Option, +pub enum ByteRange { + From { start: u64, end: Option }, + Suffix(u64), } /// Parse an HTTP Range header value like "bytes=0-1023". pub fn parse_range_header(header: &str) -> Option { - let s = header.strip_prefix("bytes=")?; - let mut parts = s.splitn(2, '-'); - let start_str = parts.next()?.trim(); - let end_str = parts.next().map(|s| s.trim()); - let start = start_str.parse::().ok()?; - let end = end_str - .filter(|s| !s.is_empty()) - .and_then(|s| s.parse::().ok()); - Some(ByteRange { start, end }) + let (start, end) = header.strip_prefix("bytes=")?.split_once('-')?; + let number = |s: &str| { + (!s.is_empty() && s.bytes().all(|b| b.is_ascii_digit())) + .then(|| s.parse::().ok()) + .flatten() + }; + if start.is_empty() { + let count = number(end)?; + return (count > 0).then_some(ByteRange::Suffix(count)); + } + Some(ByteRange::From { + start: number(start)?, + end: if end.is_empty() { + None + } else { + Some(number(end)?) + }, + }) } /// Result of attempting to serve content. pub enum ServeResult { + /// Bounded file-backed response; payment is checked before returning it. + Stream(crate::prepared_media::PreparedMedia), /// Content served successfully (full body). Ok(Vec, String), /// Partial content served (range response). @@ -265,7 +276,15 @@ pub async fn serve_content( peer_did, range, owner_session, - |path, range, mime| prepare_content(data_dir, path, range, mime), + |path, range, mime| { + prepare_content_mode( + data_dir, + path, + range, + mime, + payment_token.is_some() && !owner_session, + ) + }, |token, amount| async move { verify_payment_token(data_dir, &token, amount).await }, ) .await @@ -364,10 +383,16 @@ where } } + // Validate response metadata before touching a bearer payment too. + if hyper::header::HeaderValue::from_str(&item.mime_type).is_err() { + return Ok(ServeResult::Unavailable); + } // Finish all file I/O before consuming bearer payment. Merely opening then // reopening after charging still lost payments on read errors or deletion. let prepared = match read(file_path, range, item.mime_type.clone()).await { - Ok(result @ (ServeResult::Ok(..) | ServeResult::Partial { .. })) => result, + Ok( + result @ (ServeResult::Ok(..) | ServeResult::Partial { .. } | ServeResult::Stream(..)), + ) => result, Ok(other) => return Ok(other), Err(error) => { warn!(content_id = %id, "Cannot prepare shared content: {error:#}"); @@ -422,13 +447,26 @@ where Ok(prepared) } +// Preserve a fully readable snapshot before redeeming a bearer token. Free, +// owner and durable invoice downloads can stream their open file directly. +#[cfg(test)] async fn prepare_content( data_dir: &Path, path: PathBuf, range: Option, mime: String, ) -> Result { - use tokio::io::{AsyncReadExt, AsyncSeekExt}; + prepare_content_mode(data_dir, path, range, mime, true).await +} + +async fn prepare_content_mode( + data_dir: &Path, + path: PathBuf, + range: Option, + mime: String, + snapshot: bool, +) -> Result { + use tokio::io::AsyncSeekExt; let mut file = match fs::OpenOptions::new() .read(true) .custom_flags(libc::O_NONBLOCK) @@ -437,45 +475,79 @@ async fn prepare_content( { Ok(file) => file, Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => { - let bytes = read_filebrowser_via_userns(data_dir, &path).await?; - return slice_prepared_content(bytes, range, mime); + return prepare_filebrowser_via_userns(data_dir, &path, range, mime).await; } Err(error) => return Err(error).context("Opening shared content"), }; let metadata = file.metadata().await?; anyhow::ensure!(metadata.is_file(), "Shared content is not a regular file"); let total = metadata.len(); - if let Some(range) = range { - let Some((start, end)) = checked_range(&range, total) else { - return Ok(ServeResult::RangeNotSatisfiable(total)); - }; - file.seek(std::io::SeekFrom::Start(start)).await?; - let len = usize::try_from(end - start + 1).context("Content range is too large")?; - let mut bytes = vec![0; len]; - file.read_exact(&mut bytes) - .await - .context("Reading shared content range")?; - return Ok(ServeResult::Partial { + let selected = match range { + Some(range) => match checked_range(&range, total) { + Some((start, end)) => Some((start, end, total)), + None => return Ok(ServeResult::RangeNotSatisfiable(total)), + }, + None => None, + }; + let (start, length) = selected + .map(|(start, end, _)| (start, end - start + 1)) + .unwrap_or((0, total)); + file.seek(std::io::SeekFrom::Start(start)).await?; + if snapshot || length <= 1024 * 1024 { + prepare_reader(data_dir, file, length, mime, selected).await + } else { + Ok(ServeResult::Stream( + crate::prepared_media::PreparedMedia::direct(file, start, length, mime, selected) + .await?, + )) + } +} + +async fn prepare_reader( + data_dir: &Path, + mut source: R, + length: u64, + mime: String, + selected: Option<(u64, u64, u64)>, +) -> Result { + use tokio::io::AsyncReadExt; + if length > 1024 * 1024 { + return Ok(ServeResult::Stream( + crate::prepared_media::PreparedMedia::snapshot( + data_dir, source, length, mime, selected, + ) + .await?, + )); + } + let mut bytes = vec![0; length as usize]; + source + .read_exact(&mut bytes) + .await + .context("Preparing shared content")?; + Ok(match selected { + Some((start, end, total)) => ServeResult::Partial { bytes, mime_type: mime, start, end, total, - }); - } - let mut bytes = Vec::new(); - file.read_to_end(&mut bytes) - .await - .context("Reading shared content")?; - Ok(ServeResult::Ok(bytes, mime)) + }, + None => ServeResult::Ok(bytes, mime), + }) } fn checked_range(range: &ByteRange, total: u64) -> Option<(u64, u64)> { let last = total.checked_sub(1)?; - let end = range.end.unwrap_or(last).min(last); - (range.start <= end && range.start < total).then_some((range.start, end)) + match range { + ByteRange::Suffix(count) => (*count > 0).then_some((total.saturating_sub(*count), last)), + ByteRange::From { start, end } => { + let end = end.unwrap_or(last).min(last); + (*start <= end && *start < total).then_some((*start, end)) + } + } } +#[cfg(test)] fn slice_prepared_content( bytes: Vec, range: Option, @@ -513,37 +585,69 @@ async fn filebrowser_read_path(data_dir: &Path, path: &Path) -> Result Ok(target) } -async fn read_filebrowser_via_userns(data_dir: &Path, path: &Path) -> Result> { +async fn prepare_filebrowser_via_userns( + data_dir: &Path, + path: &Path, + range: Option, + mime: String, +) -> Result { let path = filebrowser_read_path(data_dir, path).await?; - // Tests exercise the boundary explicitly; they never launch the host Podman. #[cfg(test)] { - let _ = path; + let _ = (path, range, mime); anyhow::bail!("Files namespace read disabled in unit tests") } #[cfg(not(test))] { - let output = tokio::time::timeout( - std::time::Duration::from_secs(900), - tokio::process::Command::new("podman") - .args(["unshare", "cat", "--"]) - .arg(path) + let total = fs::metadata(&path).await?.len(); + let selected = match range { + Some(range) => match checked_range(&range, total) { + Some((start, end)) => Some((start, end, total)), + None => return Ok(ServeResult::RangeNotSatisfiable(total)), + }, + None => None, + }; + let (start, length) = selected + .map(|(start, end, _)| (start, end - start + 1)) + .unwrap_or((0, total)); + tokio::time::timeout(std::time::Duration::from_secs(900), async { + // No shell, no full stdout buffering, and only the selected bytes. + let mut input = std::ffi::OsString::from("if="); + input.push(&path); + let mut child = tokio::process::Command::new("podman") + .args([ + "unshare", + "dd", + "iflag=skip_bytes,count_bytes,nonblock,nofollow", + "status=none", + ]) + .arg(input) + .arg(format!("skip={start}")) + .arg(format!("count={length}")) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::null()) .kill_on_drop(true) - .output(), - ) + .spawn() + .context("Starting Files namespace read")?; + let stdout = child + .stdout + .take() + .context("Missing Files namespace output")?; + let result = prepare_reader(data_dir, stdout, length, mime, selected).await?; + anyhow::ensure!( + child.wait().await?.success(), + "Files namespace read failed; no payment was redeemed" + ); + Ok(result) + }) .await - .context("Files namespace read timed out")??; - anyhow::ensure!( - output.status.success(), - "Files namespace read failed: {}", - output.status - ); - Ok(output.stdout) + .context("Files namespace read timed out")? } } /// Result of attempting to serve a preview. pub enum PreviewResult { + Stream(crate::prepared_media::PreparedMedia), /// Full publicly shared free content. FullContent(Vec, String), /// Small, server-blurred JPEG for a paid image; never the original bytes. @@ -748,11 +852,14 @@ pub async fn serve_content_preview(data_dir: &Path, id: &str) -> Result { - // Free or peers-only — serve full content as preview - let bytes = fs::read(&file_path) - .await - .context("Failed to read content file")?; - Ok(PreviewResult::FullContent(bytes, item.mime_type.clone())) + // Only publicly available free content reaches this branch. + match prepare_content_mode(data_dir, file_path, None, item.mime_type.clone(), false) + .await? + { + ServeResult::Ok(bytes, mime) => Ok(PreviewResult::FullContent(bytes, mime)), + ServeResult::Stream(body) => Ok(PreviewResult::Stream(body)), + _ => Ok(PreviewResult::PreviewUnavailable), + } } } } @@ -979,6 +1086,139 @@ mod paid_read_order_tests { } } + #[test] + fn peer_ranges_reject_malformed_headers_and_support_suffixes() { + for invalid in [ + "bytes=0-invalid", + "bytes=0", + "bytes=+0-2", + "bytes=0-1,3-4", + "bytes=-0", + "bytes=0--1", + "bytes=18446744073709551616-", + "bytes=-", + "nope", + ] { + assert!(parse_range_header(invalid).is_none(), "{invalid}"); + } + for (value, expected) in [ + ("bytes=-4", Some((6, 9))), + ("bytes=-100", Some((0, 9))), + ("bytes=2-", Some((2, 9))), + ("bytes=2-100", Some((2, 9))), + ("bytes=8-3", None), + ("bytes=10-", None), + ] { + assert_eq!( + checked_range(&parse_range_header(value).unwrap(), 10), + expected, + "{value}" + ); + } + assert_eq!( + checked_range(&parse_range_header("bytes=0-").unwrap(), 0), + None + ); + } + + #[tokio::test] + async fn large_paid_stream_still_requires_redemption_and_survives_source_deletion() { + use hyper::body::HttpBody; + let bytes = vec![91; 2 * 1024 * 1024]; + let dir = fixture(&bytes).await; + let charged = std::sync::atomic::AtomicUsize::new(0); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, amount| async { + assert_eq!(amount, 10); + charged.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + fs::remove_file(dir.path().join("content/files/test.bin")) + .await + .unwrap(); + true + }, + ) + .await + .unwrap(); + assert_eq!(charged.load(std::sync::atomic::Ordering::SeqCst), 1); + let ServeResult::Stream(body) = result else { + panic!("expected bounded stream") + }; + let mut response = body.into_response().unwrap(); + let mut received = Vec::new(); + while let Some(chunk) = response.body_mut().data().await { + let chunk = chunk.unwrap(); + assert!(chunk.len() <= 65536); + received.extend_from_slice(&chunk); + } + assert_eq!(received, bytes); + } + + #[tokio::test] + async fn rejected_payment_never_returns_the_prepared_large_stream() { + let bytes = vec![91; 2 * 1024 * 1024]; + let dir = fixture(&bytes).await; + let checked = std::sync::atomic::AtomicUsize::new(0); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBinvalid"), + None, + None, + None, + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { + checked.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + false + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::PaymentRequired(10))); + assert_eq!(checked.load(std::sync::atomic::Ordering::SeqCst), 1); + assert_eq!( + std::fs::read_dir(dir.path().join("content-staging")) + .unwrap() + .count(), + 0 + ); + } + + #[tokio::test] + async fn free_large_preview_uses_bounded_stream_without_private_snapshot() { + use hyper::body::HttpBody; + let dir = fixture(&[0; 1]).await; + let mut catalog = load_catalog(dir.path()).await.unwrap(); + catalog.items[0].access = AccessControl::Free; + save_catalog(dir.path(), &catalog).await.unwrap(); + let file = fs::OpenOptions::new() + .write(true) + .open(dir.path().join("content/files/test.bin")) + .await + .unwrap(); + file.set_len(4 * 1024 * 1024 * 1024).await.unwrap(); + let PreviewResult::Stream(body) = serve_content_preview(dir.path(), "paid").await.unwrap() + else { + panic!("expected bounded preview") + }; + let mut response = body.into_response().unwrap(); + assert_eq!(response.headers()["content-length"], "4294967296"); + assert_eq!( + response.body_mut().data().await.unwrap().unwrap().len(), + 65536 + ); + drop(response); + assert!(!dir.path().join("content-staging").exists()); + } + #[tokio::test] async fn deletion_during_payment_cannot_lose_prepared_bytes() { let dir = fixture(b"original").await; @@ -1020,7 +1260,7 @@ mod paid_read_order_tests { Some("cashuBtest"), None, None, - Some(ByteRange { start, end }), + Some(ByteRange::From { start, end }), false, |path, range, mime| prepare_content(dir.path(), path, range, mime), |_, _| async { panic!("invalid range reached payment") }, @@ -1042,7 +1282,7 @@ mod paid_read_order_tests { Some("cashuBtest"), None, None, - Some(ByteRange { + Some(ByteRange::From { start: 2, end: Some(999), }), @@ -1194,7 +1434,7 @@ mod paid_read_order_tests { assert!(matches!( slice_prepared_content( vec![], - Some(ByteRange { + Some(ByteRange::From { start: 0, end: None }), @@ -1204,7 +1444,7 @@ mod paid_read_order_tests { ServeResult::RangeNotSatisfiable(0) )); assert!( - matches!(slice_prepared_content(b"abc".to_vec(), Some(ByteRange { start: 1, end: None }), "x".into()).unwrap(), ServeResult::Partial { bytes, start: 1, end: 2, total: 3, .. } if bytes == b"bc") + matches!(slice_prepared_content(b"abc".to_vec(), Some(ByteRange::From { start: 1, end: None }), "x".into()).unwrap(), ServeResult::Partial { bytes, start: 1, end: 2, total: 3, .. } if bytes == b"bc") ); } } diff --git a/core/archipelago/src/main.rs b/core/archipelago/src/main.rs index 9fb3ac57..b8a18bc5 100644 --- a/core/archipelago/src/main.rs +++ b/core/archipelago/src/main.rs @@ -45,6 +45,7 @@ mod content_indeehub; mod content_invoice; mod content_owned; mod media_stream; +mod prepared_media; mod content_server; mod crash_recovery; mod credentials; diff --git a/core/archipelago/src/prepared_media.rs b/core/archipelago/src/prepared_media.rs new file mode 100644 index 00000000..7b9f42e3 --- /dev/null +++ b/core/archipelago/src/prepared_media.rs @@ -0,0 +1,263 @@ +//! File-backed peer responses with bounded buffers and private payment snapshots. +//! Snapshot construction finishes before bearer ecash is redeemed. Anonymous +//! temporary files are removed automatically when the response is dropped. +use anyhow::{Context, Result}; +use hyper::{Body, Response, StatusCode}; +use std::path::Path; +use std::sync::{Arc, Mutex, OnceLock}; +use tokio::fs::File; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeekExt, AsyncWriteExt}; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; + +const CHUNK: usize = 64 * 1024; +const DISK_RESERVE: u64 = 256 * 1024 * 1024; +static RESERVED: Mutex = Mutex::new(0); +static SLOTS: OnceLock> = OnceLock::new(); + +struct Reservation { + bytes: u64, + _slot: OwnedSemaphorePermit, +} +impl Drop for Reservation { + fn drop(&mut self) { + let mut reserved = RESERVED.lock().unwrap_or_else(|e| e.into_inner()); + *reserved = reserved.saturating_sub(self.bytes); + } +} + +fn reserve(file: &File, length: u64) -> Result { + use std::os::fd::AsRawFd; + let slot = SLOTS + .get_or_init(|| Arc::new(Semaphore::new(4))) + .clone() + .try_acquire_owned() + .context("Content preparation is busy")?; + let mut stat = std::mem::MaybeUninit::::uninit(); + // The live descriptor supplies the filesystem; no attacker-controlled C path. + if unsafe { libc::fstatvfs(file.as_raw_fd(), stat.as_mut_ptr()) } != 0 { + return Err(std::io::Error::last_os_error()).context("Checking content staging space"); + } + let stat = unsafe { stat.assume_init() }; + let available = (stat.f_bavail as u64).saturating_mul(stat.f_frsize as u64); + let mut reserved = RESERVED.lock().unwrap_or_else(|e| e.into_inner()); + let next = reserved + .checked_add(length) + .context("Content size overflow")?; + anyhow::ensure!( + next.checked_add(DISK_RESERVE.max(available / 20)) + .is_some_and(|n| n <= available), + "Insufficient private staging space; no payment was redeemed" + ); + *reserved = next; + Ok(Reservation { + bytes: length, + _slot: slot, + }) +} + +pub struct PreparedMedia { + file: File, + length: u64, + mime: String, + range: Option<(u64, u64, u64)>, + reservation: Option, +} + +impl PreparedMedia { + pub async fn direct( + mut file: File, + start: u64, + length: u64, + mime: String, + range: Option<(u64, u64, u64)>, + ) -> Result { + anyhow::ensure!( + file.metadata().await?.is_file(), + "Content is not a regular file" + ); + file.seek(std::io::SeekFrom::Start(start)).await?; + Ok(Self { + file, + length, + mime, + range, + reservation: None, + }) + } + + pub async fn snapshot( + data_dir: &Path, + mut source: R, + length: u64, + mime: String, + range: Option<(u64, u64, u64)>, + ) -> Result { + let dir = data_dir.join("content-staging"); + tokio::fs::create_dir_all(&dir).await?; + let temporary = tokio::task::spawn_blocking(move || tempfile::tempfile_in(dir)).await??; + let mut file = File::from_std(temporary); + let reservation = reserve(&file, length)?; + let mut left = length; + let mut buffer = vec![0; CHUNK]; + while left > 0 { + let limit = left.min(CHUNK as u64) as usize; + let count = source.read(&mut buffer[..limit]).await?; + anyhow::ensure!( + count > 0, + "Content changed while preparing payment response" + ); + file.write_all(&buffer[..count]).await?; + left -= count as u64; + } + // Detect writeback errors before the caller attempts bearer redemption. + file.flush().await?; + file.sync_data().await?; + file.seek(std::io::SeekFrom::Start(0)).await?; + Ok(Self { + file, + length, + mime, + range, + reservation: Some(reservation), + }) + } + + pub fn into_response(self) -> Result> { + let Self { + file, + length, + mime, + range, + reservation, + } = self; + let chunks = futures_util::stream::try_unfold( + (file, length, reservation), + |(mut file, left, reservation)| async move { + if left == 0 { + return Ok::<_, std::io::Error>(None); + } + let mut buffer = vec![0; left.min(CHUNK as u64) as usize]; + let count = file.read(&mut buffer).await?; + if count == 0 { + return Err(std::io::Error::new( + std::io::ErrorKind::UnexpectedEof, + "Content changed during transfer", + )); + } + buffer.truncate(count); + Ok(Some((buffer, (file, left - count as u64, reservation)))) + }, + ); + let mut response = Response::builder() + .status(if range.is_some() { + StatusCode::PARTIAL_CONTENT + } else { + StatusCode::OK + }) + .header("Content-Type", mime) + .header("Content-Length", length) + .header("Accept-Ranges", "bytes") + .header("X-Content-Type-Options", "nosniff") + .header("Cache-Control", "private, no-store"); + if let Some((start, end, total)) = range { + response = response.header("Content-Range", format!("bytes {start}-{end}/{total}")); + } + Ok(response.body(Body::wrap_stream(chunks))?) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use hyper::body::HttpBody; + + #[tokio::test] + async fn direct_large_sparse_file_does_not_read_ahead_and_short_reads_fail() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("film"); + let write = File::create(&path).await.unwrap(); + write.set_len(4 * 1024 * 1024 * 1024).await.unwrap(); + let mut response = PreparedMedia::direct( + File::open(&path).await.unwrap(), + 0, + 4 * 1024 * 1024 * 1024, + "video/mp4".into(), + None, + ) + .await + .unwrap() + .into_response() + .unwrap(); + assert_eq!(response.headers()["content-length"], "4294967296"); + assert_eq!( + response.body_mut().data().await.unwrap().unwrap().len(), + CHUNK + ); + write.set_len(0).await.unwrap(); + assert!(response.body_mut().data().await.unwrap().is_err()); + drop(response); + } + + #[tokio::test] + async fn snapshot_survives_original_removal_and_has_no_named_temporary_file() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("original"); + let bytes = vec![73; 3 * 1024 * 1024]; + tokio::fs::write(&path, &bytes).await.unwrap(); + let prepared = PreparedMedia::snapshot( + dir.path(), + File::open(&path).await.unwrap(), + bytes.len() as u64, + "application/octet-stream".into(), + None, + ) + .await + .unwrap(); + tokio::fs::remove_file(path).await.unwrap(); + assert_eq!( + std::fs::read_dir(dir.path().join("content-staging")) + .unwrap() + .count(), + 0 + ); + let mut response = prepared.into_response().unwrap(); + let mut received = Vec::new(); + while let Some(chunk) = response.body_mut().data().await { + let chunk = chunk.unwrap(); + assert!(chunk.len() <= CHUNK); + received.extend_from_slice(&chunk); + } + assert_eq!(received, bytes); + } + + #[tokio::test] + async fn incomplete_snapshot_fails_before_a_payment_can_be_attempted() { + let dir = tempfile::tempdir().unwrap(); + assert!( + PreparedMedia::snapshot(dir.path(), &b"short"[..], 100, "x".into(), None) + .await + .is_err() + ); + assert_eq!( + std::fs::read_dir(dir.path().join("content-staging")) + .unwrap() + .count(), + 0 + ); + // A subsequent snapshot still works; the failed preparation releases its slot. + let body = PreparedMedia::snapshot( + dir.path(), + &b"ok"[..], + 2, + "text/plain".into(), + Some((4, 5, 10)), + ) + .await + .unwrap() + .into_response() + .unwrap(); + assert_eq!(body.status(), StatusCode::PARTIAL_CONTENT); + assert_eq!(body.headers()["content-range"], "bytes 4-5/10"); + assert_eq!(hyper::body::to_bytes(body.into_body()).await.unwrap(), "ok"); + } +}