Prepare large paid files on disk and stream peer responses in bounded chunks

This commit is contained in:
archipelago
2026-10-06 05:31:19 -04:00
parent 051dc7e3df
commit 2fee0339cb
5 changed files with 588 additions and 68 deletions
+2 -1
View File
@@ -148,6 +148,7 @@ iroh-blobs = { version = "0.103", optional = true }
lofty = "0.24.0" lofty = "0.24.0"
cashu = { version = "0.17.5", default-features = false, features = ["wallet"] } cashu = { version = "0.17.5", default-features = false, features = ["wallet"] }
tempfile = "3.10"
[dev-dependencies] [dev-dependencies]
tokio-test = "0.4" tokio-test = "0.4"
tempfile = "3.10"
+19 -4
View File
@@ -139,10 +139,23 @@ impl ApiHandler {
} }
// Parse Range header for streaming support // Parse Range header for streaming support
let range = headers let range = match headers.get("range") {
.get("range") None => None,
.and_then(|v| v.to_str().ok()) Some(value) => match value
.and_then(content_server::parse_range_header); .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( match content_server::serve_content(
&config.data_dir, &config.data_dir,
@@ -155,6 +168,7 @@ impl ApiHandler {
) )
.await .await
{ {
Ok(content_server::ServeResult::Stream(body)) => body.into_response(),
Ok(content_server::ServeResult::Ok(bytes, mime_type)) => { Ok(content_server::ServeResult::Ok(bytes, mime_type)) => {
let len = bytes.len(); let len = bytes.len();
Ok(Response::builder() Ok(Response::builder()
@@ -496,6 +510,7 @@ impl ApiHandler {
} }
match content_server::serve_content_preview(&config.data_dir, content_id).await { 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)) => { Ok(content_server::PreviewResult::FullContent(bytes, mime_type)) => {
let len = bytes.len(); let len = bytes.len();
Ok(Response::builder() Ok(Response::builder()
+299 -59
View File
@@ -202,26 +202,37 @@ pub async fn set_availability(data_dir: &Path, id: &str, availability: Availabil
} }
/// A byte range request (start, optional end). /// A byte range request (start, optional end).
pub struct ByteRange { pub enum ByteRange {
pub start: u64, From { start: u64, end: Option<u64> },
pub end: Option<u64>, Suffix(u64),
} }
/// Parse an HTTP Range header value like "bytes=0-1023". /// Parse an HTTP Range header value like "bytes=0-1023".
pub fn parse_range_header(header: &str) -> Option<ByteRange> { pub fn parse_range_header(header: &str) -> Option<ByteRange> {
let s = header.strip_prefix("bytes=")?; let (start, end) = header.strip_prefix("bytes=")?.split_once('-')?;
let mut parts = s.splitn(2, '-'); let number = |s: &str| {
let start_str = parts.next()?.trim(); (!s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()))
let end_str = parts.next().map(|s| s.trim()); .then(|| s.parse::<u64>().ok())
let start = start_str.parse::<u64>().ok()?; .flatten()
let end = end_str };
.filter(|s| !s.is_empty()) if start.is_empty() {
.and_then(|s| s.parse::<u64>().ok()); let count = number(end)?;
Some(ByteRange { start, 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. /// Result of attempting to serve content.
pub enum ServeResult { pub enum ServeResult {
/// Bounded file-backed response; payment is checked before returning it.
Stream(crate::prepared_media::PreparedMedia),
/// Content served successfully (full body). /// Content served successfully (full body).
Ok(Vec<u8>, String), Ok(Vec<u8>, String),
/// Partial content served (range response). /// Partial content served (range response).
@@ -265,7 +276,15 @@ pub async fn serve_content(
peer_did, peer_did,
range, range,
owner_session, 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 }, |token, amount| async move { verify_payment_token(data_dir, &token, amount).await },
) )
.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 // Finish all file I/O before consuming bearer payment. Merely opening then
// reopening after charging still lost payments on read errors or deletion. // reopening after charging still lost payments on read errors or deletion.
let prepared = match read(file_path, range, item.mime_type.clone()).await { 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), Ok(other) => return Ok(other),
Err(error) => { Err(error) => {
warn!(content_id = %id, "Cannot prepare shared content: {error:#}"); warn!(content_id = %id, "Cannot prepare shared content: {error:#}");
@@ -422,13 +447,26 @@ where
Ok(prepared) 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( async fn prepare_content(
data_dir: &Path, data_dir: &Path,
path: PathBuf, path: PathBuf,
range: Option<ByteRange>, range: Option<ByteRange>,
mime: String, mime: String,
) -> Result<ServeResult> { ) -> Result<ServeResult> {
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<ByteRange>,
mime: String,
snapshot: bool,
) -> Result<ServeResult> {
use tokio::io::AsyncSeekExt;
let mut file = match fs::OpenOptions::new() let mut file = match fs::OpenOptions::new()
.read(true) .read(true)
.custom_flags(libc::O_NONBLOCK) .custom_flags(libc::O_NONBLOCK)
@@ -437,45 +475,79 @@ async fn prepare_content(
{ {
Ok(file) => file, Ok(file) => file,
Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => { Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
let bytes = read_filebrowser_via_userns(data_dir, &path).await?; return prepare_filebrowser_via_userns(data_dir, &path, range, mime).await;
return slice_prepared_content(bytes, range, mime);
} }
Err(error) => return Err(error).context("Opening shared content"), Err(error) => return Err(error).context("Opening shared content"),
}; };
let metadata = file.metadata().await?; let metadata = file.metadata().await?;
anyhow::ensure!(metadata.is_file(), "Shared content is not a regular file"); anyhow::ensure!(metadata.is_file(), "Shared content is not a regular file");
let total = metadata.len(); let total = metadata.len();
if let Some(range) = range { let selected = match range {
let Some((start, end)) = checked_range(&range, total) else { Some(range) => match checked_range(&range, total) {
return Ok(ServeResult::RangeNotSatisfiable(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?; file.seek(std::io::SeekFrom::Start(start)).await?;
let len = usize::try_from(end - start + 1).context("Content range is too large")?; if snapshot || length <= 1024 * 1024 {
let mut bytes = vec![0; len]; prepare_reader(data_dir, file, length, mime, selected).await
file.read_exact(&mut bytes) } else {
Ok(ServeResult::Stream(
crate::prepared_media::PreparedMedia::direct(file, start, length, mime, selected)
.await?,
))
}
}
async fn prepare_reader<R: tokio::io::AsyncRead + Unpin>(
data_dir: &Path,
mut source: R,
length: u64,
mime: String,
selected: Option<(u64, u64, u64)>,
) -> Result<ServeResult> {
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 .await
.context("Reading shared content range")?; .context("Preparing shared content")?;
return Ok(ServeResult::Partial { Ok(match selected {
Some((start, end, total)) => ServeResult::Partial {
bytes, bytes,
mime_type: mime, mime_type: mime,
start, start,
end, end,
total, total,
}); },
} None => ServeResult::Ok(bytes, mime),
let mut bytes = Vec::new(); })
file.read_to_end(&mut bytes)
.await
.context("Reading shared content")?;
Ok(ServeResult::Ok(bytes, mime))
} }
fn checked_range(range: &ByteRange, total: u64) -> Option<(u64, u64)> { fn checked_range(range: &ByteRange, total: u64) -> Option<(u64, u64)> {
let last = total.checked_sub(1)?; let last = total.checked_sub(1)?;
let end = range.end.unwrap_or(last).min(last); match range {
(range.start <= end && range.start < total).then_some((range.start, end)) 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( fn slice_prepared_content(
bytes: Vec<u8>, bytes: Vec<u8>,
range: Option<ByteRange>, range: Option<ByteRange>,
@@ -513,37 +585,69 @@ async fn filebrowser_read_path(data_dir: &Path, path: &Path) -> Result<PathBuf>
Ok(target) Ok(target)
} }
async fn read_filebrowser_via_userns(data_dir: &Path, path: &Path) -> Result<Vec<u8>> { async fn prepare_filebrowser_via_userns(
data_dir: &Path,
path: &Path,
range: Option<ByteRange>,
mime: String,
) -> Result<ServeResult> {
let path = filebrowser_read_path(data_dir, path).await?; let path = filebrowser_read_path(data_dir, path).await?;
// Tests exercise the boundary explicitly; they never launch the host Podman.
#[cfg(test)] #[cfg(test)]
{ {
let _ = path; let _ = (path, range, mime);
anyhow::bail!("Files namespace read disabled in unit tests") anyhow::bail!("Files namespace read disabled in unit tests")
} }
#[cfg(not(test))] #[cfg(not(test))]
{ {
let output = tokio::time::timeout( let total = fs::metadata(&path).await?.len();
std::time::Duration::from_secs(900), let selected = match range {
tokio::process::Command::new("podman") Some(range) => match checked_range(&range, total) {
.args(["unshare", "cat", "--"]) Some((start, end)) => Some((start, end, total)),
.arg(path) 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) .kill_on_drop(true)
.output(), .spawn()
) .context("Starting Files namespace read")?;
.await let stdout = child
.context("Files namespace read timed out")??; .stdout
.take()
.context("Missing Files namespace output")?;
let result = prepare_reader(data_dir, stdout, length, mime, selected).await?;
anyhow::ensure!( anyhow::ensure!(
output.status.success(), child.wait().await?.success(),
"Files namespace read failed: {}", "Files namespace read failed; no payment was redeemed"
output.status
); );
Ok(output.stdout) Ok(result)
})
.await
.context("Files namespace read timed out")?
} }
} }
/// Result of attempting to serve a preview. /// Result of attempting to serve a preview.
pub enum PreviewResult { pub enum PreviewResult {
Stream(crate::prepared_media::PreparedMedia),
/// Full publicly shared free content. /// Full publicly shared free content.
FullContent(Vec<u8>, String), FullContent(Vec<u8>, String),
/// Small, server-blurred JPEG for a paid image; never the original bytes. /// 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<PreviewR
} }
} }
_ => { _ => {
// Free or peers-only — serve full content as preview // Only publicly available free content reaches this branch.
let bytes = fs::read(&file_path) match prepare_content_mode(data_dir, file_path, None, item.mime_type.clone(), false)
.await .await?
.context("Failed to read content file")?; {
Ok(PreviewResult::FullContent(bytes, item.mime_type.clone())) 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] #[tokio::test]
async fn deletion_during_payment_cannot_lose_prepared_bytes() { async fn deletion_during_payment_cannot_lose_prepared_bytes() {
let dir = fixture(b"original").await; let dir = fixture(b"original").await;
@@ -1020,7 +1260,7 @@ mod paid_read_order_tests {
Some("cashuBtest"), Some("cashuBtest"),
None, None,
None, None,
Some(ByteRange { start, end }), Some(ByteRange::From { start, end }),
false, false,
|path, range, mime| prepare_content(dir.path(), path, range, mime), |path, range, mime| prepare_content(dir.path(), path, range, mime),
|_, _| async { panic!("invalid range reached payment") }, |_, _| async { panic!("invalid range reached payment") },
@@ -1042,7 +1282,7 @@ mod paid_read_order_tests {
Some("cashuBtest"), Some("cashuBtest"),
None, None,
None, None,
Some(ByteRange { Some(ByteRange::From {
start: 2, start: 2,
end: Some(999), end: Some(999),
}), }),
@@ -1194,7 +1434,7 @@ mod paid_read_order_tests {
assert!(matches!( assert!(matches!(
slice_prepared_content( slice_prepared_content(
vec![], vec![],
Some(ByteRange { Some(ByteRange::From {
start: 0, start: 0,
end: None end: None
}), }),
@@ -1204,7 +1444,7 @@ mod paid_read_order_tests {
ServeResult::RangeNotSatisfiable(0) ServeResult::RangeNotSatisfiable(0)
)); ));
assert!( 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")
); );
} }
} }
+1
View File
@@ -45,6 +45,7 @@ mod content_indeehub;
mod content_invoice; mod content_invoice;
mod content_owned; mod content_owned;
mod media_stream; mod media_stream;
mod prepared_media;
mod content_server; mod content_server;
mod crash_recovery; mod crash_recovery;
mod credentials; mod credentials;
+263
View File
@@ -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<u64> = Mutex::new(0);
static SLOTS: OnceLock<Arc<Semaphore>> = 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<Reservation> {
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::<libc::statvfs>::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<Reservation>,
}
impl PreparedMedia {
pub async fn direct(
mut file: File,
start: u64,
length: u64,
mime: String,
range: Option<(u64, u64, u64)>,
) -> Result<Self> {
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<R: AsyncRead + Unpin>(
data_dir: &Path,
mut source: R,
length: u64,
mime: String,
range: Option<(u64, u64, u64)>,
) -> Result<Self> {
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<Response<Body>> {
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");
}
}