Require FIPS for peer playback and stream owned media with bounded reads
This commit is contained in:
@@ -238,54 +238,13 @@ impl ApiHandler {
|
||||
return bad("invalid onion or content id");
|
||||
}
|
||||
|
||||
// Already purchased? Serve the local cache — no network, no
|
||||
// re-payment. The seller's node charges every fetch by design; the
|
||||
// buyer-side store (content_owned) exists precisely so an owned item
|
||||
// never has to be bought twice, and the content surface's cards were
|
||||
// hitting the seller's 402 and rendering as permanent placeholders.
|
||||
// Range is honoured by slicing, so seek/playback works from cache.
|
||||
if crate::content_owned::is_owned(&self.config.data_dir, onion, content_id).await {
|
||||
if let Some((mime_type, bytes)) =
|
||||
crate::content_owned::read_owned(&self.config.data_dir, onion, content_id).await
|
||||
{
|
||||
let total = bytes.len();
|
||||
let range = headers
|
||||
.get("range")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(crate::content_server::parse_range_header);
|
||||
if let Some(r) = range {
|
||||
let start = (r.start as usize).min(total);
|
||||
let end = r
|
||||
.end
|
||||
.map(|e| e as usize)
|
||||
.unwrap_or(total.saturating_sub(1))
|
||||
.min(total.saturating_sub(1));
|
||||
if start <= end && total > 0 {
|
||||
let slice = &bytes[start..=end];
|
||||
return Ok(Response::builder()
|
||||
.status(StatusCode::PARTIAL_CONTENT)
|
||||
.header("Content-Type", mime_type)
|
||||
.header("Content-Length", slice.len().to_string())
|
||||
.header(
|
||||
"Content-Range",
|
||||
format!("bytes {}-{}/{}", start, end, total),
|
||||
)
|
||||
.header("Accept-Ranges", "bytes")
|
||||
.body(hyper::Body::from(slice.to_vec()))
|
||||
.unwrap_or_else(|_| Response::new(hyper::Body::empty())));
|
||||
}
|
||||
}
|
||||
return Ok(Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header("Content-Type", mime_type)
|
||||
.header("Content-Length", total.to_string())
|
||||
.header("Accept-Ranges", "bytes")
|
||||
.body(hyper::Body::from(bytes))
|
||||
.unwrap_or_else(|_| Response::new(hyper::Body::empty())));
|
||||
}
|
||||
// Indexed as owned but bytes missing — fall through to the peer
|
||||
// rather than erroring: the seller can still serve it (for the
|
||||
// price already paid, the operator can re-fetch and re-cache).
|
||||
// Ownership is checked before opening a bounded file stream. Corrupt
|
||||
// records or missing purchased bytes never trigger another purchase.
|
||||
match crate::content_owned::open_owned(&self.config.data_dir, onion, content_id).await {
|
||||
Ok(Some((mime, file))) => return crate::media_stream::file_response(file, &mime, headers).await,
|
||||
Ok(None) => {},
|
||||
Err(_) => return Ok(build_response(StatusCode::CONFLICT, "application/json",
|
||||
hyper::Body::from(serde_json::json!({"error": "Purchased file unavailable locally. Recover the existing purchase without paying again."}).to_string()))),
|
||||
}
|
||||
|
||||
let fips_npub = crate::federation::fips_npub_for_onion(&self.config.data_dir, onion).await;
|
||||
@@ -293,20 +252,30 @@ impl ApiHandler {
|
||||
// Generous overall timeout: this endpoint serves both seek/Range
|
||||
// playback (small, finishes fast) and full-file downloads of large
|
||||
// media (#38). 60s was too tight for a multi-hundred-MB transfer over
|
||||
// Tor and aborted the download mid-stream.
|
||||
// slow links and aborted the download mid-stream.
|
||||
let mut req = crate::fips::dial::PeerRequest::new(fips_npub.as_deref(), onion, &peer_path)
|
||||
.service(crate::settings::transport::PeerService::PeerFiles)
|
||||
.require_fips()
|
||||
.record_transport(&self.config.data_dir)
|
||||
.timeout(std::time::Duration::from_secs(900));
|
||||
if let Some(r) = headers.get("range").and_then(|v| v.to_str().ok()) {
|
||||
req = req.header("Range", r.to_string());
|
||||
}
|
||||
match req.send_get().await {
|
||||
Ok((resp, _transport)) => {
|
||||
Ok((resp, transport)) => {
|
||||
if resp.status().is_redirection() {
|
||||
return Ok(build_response(
|
||||
StatusCode::BAD_GATEWAY,
|
||||
"application/json",
|
||||
hyper::Body::from("{\"error\":\"Peer media redirects are not allowed\"}"),
|
||||
));
|
||||
}
|
||||
let status = resp.status().as_u16();
|
||||
let rh = resp.headers().clone();
|
||||
let mut builder = Response::builder()
|
||||
.status(status)
|
||||
.header("Accept-Ranges", "bytes");
|
||||
.header("Accept-Ranges", "bytes")
|
||||
.header("X-Archipelago-Transport", transport.to_string());
|
||||
for h in ["content-type", "content-range", "content-length"] {
|
||||
if let Some(v) = rh.get(h).and_then(|v| v.to_str().ok()) {
|
||||
builder = builder.header(h, v);
|
||||
|
||||
@@ -184,6 +184,45 @@ pub async fn is_owned(data_dir: &Path, onion: &str, content_id: &str) -> bool {
|
||||
}
|
||||
|
||||
/// Read a purchased item's bytes + mime type from the local cache, if present.
|
||||
pub async fn open_owned(
|
||||
data_dir: &Path,
|
||||
onion: &str,
|
||||
content_id: &str,
|
||||
) -> Result<Option<(String, fs::File)>> {
|
||||
let index = load_index_checked(data_dir).await?;
|
||||
let Some(item) = index
|
||||
.items
|
||||
.iter()
|
||||
.find(|item| item.onion == onion && item.content_id == content_id)
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
// Reject path components even when called outside the HTTP route.
|
||||
anyhow::ensure!(
|
||||
!onion.is_empty()
|
||||
&& !content_id.is_empty()
|
||||
&& onion != "."
|
||||
&& content_id != "."
|
||||
&& !onion.contains("..")
|
||||
&& !content_id.contains("..")
|
||||
&& sanitize(onion) == onion
|
||||
&& sanitize(content_id) == content_id,
|
||||
"Invalid purchase path"
|
||||
);
|
||||
let file = fs::OpenOptions::new()
|
||||
.read(true)
|
||||
.custom_flags(libc::O_NOFOLLOW)
|
||||
.open(bytes_path(data_dir, onion, content_id))
|
||||
.await
|
||||
.context("Purchased bytes unavailable")?;
|
||||
let metadata = file.metadata().await?;
|
||||
anyhow::ensure!(
|
||||
metadata.is_file() && metadata.len() == item.size_bytes,
|
||||
"Purchased bytes incomplete"
|
||||
);
|
||||
Ok(Some((item.mime_type.clone(), file)))
|
||||
}
|
||||
|
||||
pub async fn read_owned(
|
||||
data_dir: &Path,
|
||||
onion: &str,
|
||||
@@ -205,6 +244,48 @@ pub async fn read_owned(
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[tokio::test]
|
||||
async fn owned_stream_preserves_ownership_on_missing_corrupt_or_symlinked_bytes() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
record_purchase(
|
||||
dir.path(),
|
||||
"seller.onion",
|
||||
"video",
|
||||
"video",
|
||||
"video/mp4",
|
||||
b"video",
|
||||
1,
|
||||
"cashu",
|
||||
"now",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(open_owned(dir.path(), "seller.onion", "video")
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
let path = bytes_path(dir.path(), "seller.onion", "video");
|
||||
fs::write(&path, b"bad").await.unwrap();
|
||||
assert!(open_owned(dir.path(), "seller.onion", "video")
|
||||
.await
|
||||
.is_err());
|
||||
fs::remove_file(&path).await.unwrap();
|
||||
assert!(open_owned(dir.path(), "seller.onion", "video")
|
||||
.await
|
||||
.is_err());
|
||||
let private = dir.path().join("private");
|
||||
fs::write(&private, b"other").await.unwrap();
|
||||
std::os::unix::fs::symlink(&private, &path).unwrap();
|
||||
assert!(open_owned(dir.path(), "seller.onion", "video")
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(list_owned_checked(dir.path()).await.unwrap().len(), 1);
|
||||
assert!(open_owned(dir.path(), "seller.onion", "not-bought")
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn concurrent_purchases_preserve_every_item_and_exact_bytes() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -394,6 +394,8 @@ pub struct PeerRequest<'a> {
|
||||
/// The request carries something that must reach the peer at most once
|
||||
/// (a bearer ecash token). See [`PeerRequest::single_delivery`].
|
||||
pub single_delivery: bool,
|
||||
/// Media explicitly requiring the mesh must never silently use Tor or redirects.
|
||||
pub require_fips: bool,
|
||||
}
|
||||
|
||||
impl<'a> PeerRequest<'a> {
|
||||
@@ -408,9 +410,17 @@ impl<'a> PeerRequest<'a> {
|
||||
service: None,
|
||||
record_data_dir: None,
|
||||
single_delivery: false,
|
||||
require_fips: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Enforce the media transport contract independently of general service
|
||||
/// preferences. A missing mesh route is a recoverable error, not a fallback.
|
||||
pub fn require_fips(mut self) -> Self {
|
||||
self.require_fips = true;
|
||||
self
|
||||
}
|
||||
|
||||
/// Never send this request twice. A paid download carries a bearer ecash
|
||||
/// token that the seller redeems on first sight; replaying it over Tor
|
||||
/// after FIPS already delivered it hands the seller a spent token, so the
|
||||
@@ -481,6 +491,9 @@ impl<'a> PeerRequest<'a> {
|
||||
|
||||
/// Resolved preference: user setting if `service` was set, else Auto.
|
||||
async fn preference(&self) -> crate::settings::transport::TransportPref {
|
||||
if self.require_fips {
|
||||
return crate::settings::transport::TransportPref::Fips;
|
||||
}
|
||||
match self.service {
|
||||
Some(s) => crate::settings::transport::get(s).await,
|
||||
None => crate::settings::transport::TransportPref::Auto,
|
||||
@@ -523,7 +536,7 @@ impl<'a> PeerRequest<'a> {
|
||||
None => {
|
||||
if pref == TransportPref::Fips {
|
||||
anyhow::bail!(
|
||||
"User set transport preference to FIPS only, but peer is unreachable over FIPS"
|
||||
"This request requires FIPS, but the peer is unreachable over FIPS"
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -562,7 +575,7 @@ impl<'a> PeerRequest<'a> {
|
||||
None => {
|
||||
if pref == TransportPref::Fips {
|
||||
anyhow::bail!(
|
||||
"User set transport preference to FIPS only, but peer is unreachable over FIPS"
|
||||
"This request requires FIPS, but the peer is unreachable over FIPS"
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -611,7 +624,7 @@ impl<'a> PeerRequest<'a> {
|
||||
} else {
|
||||
budget
|
||||
};
|
||||
let c = client_with_delivery_policy(per_attempt, self.single_delivery);
|
||||
let c = client_with_delivery_policy(per_attempt, self.single_delivery || self.require_fips);
|
||||
let mut rb = c.post(&url).json(body);
|
||||
for (k, v) in &self.headers {
|
||||
rb = rb.header(*k, v);
|
||||
@@ -680,7 +693,7 @@ impl<'a> PeerRequest<'a> {
|
||||
} else {
|
||||
budget
|
||||
};
|
||||
let c = client_with_delivery_policy(per_attempt, self.single_delivery);
|
||||
let c = client_with_delivery_policy(per_attempt, self.single_delivery || self.require_fips);
|
||||
let mut rb = c.get(&url);
|
||||
for (k, v) in &self.headers {
|
||||
rb = rb.header(*k, v);
|
||||
@@ -770,6 +783,17 @@ impl<'a> PeerRequest<'a> {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn required_media_never_falls_back_when_peer_has_no_fips_identity() {
|
||||
let request = PeerRequest::new(None, "unreachable.onion", "/content/video").require_fips();
|
||||
assert_eq!(
|
||||
request.preference().await,
|
||||
crate::settings::transport::TransportPref::Fips
|
||||
);
|
||||
let error = request.send_get().await.unwrap_err();
|
||||
assert!(error.to_string().contains("requires FIPS"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn encode_query_round_trip_header_is_correct() {
|
||||
let q = encode_query(0x1234, "npub1abc").unwrap();
|
||||
|
||||
@@ -44,6 +44,7 @@ mod content_hash;
|
||||
mod content_indeehub;
|
||||
mod content_invoice;
|
||||
mod content_owned;
|
||||
mod media_stream;
|
||||
mod content_server;
|
||||
mod crash_recovery;
|
||||
mod credentials;
|
||||
|
||||
@@ -0,0 +1,161 @@
|
||||
//! Bounded, seekable media responses. One file descriptor and at most 64 KiB
|
||||
//! are retained per in-flight response; dropping the body closes the file.
|
||||
use anyhow::Result;
|
||||
use hyper::{Body, HeaderMap, Response, StatusCode};
|
||||
use tokio::{
|
||||
fs::File,
|
||||
io::{AsyncReadExt, AsyncSeekExt},
|
||||
};
|
||||
|
||||
fn range(value: &str, total: u64) -> Option<(u64, u64)> {
|
||||
let (start, end) = value.strip_prefix("bytes=")?.split_once('-')?;
|
||||
if total == 0 || end.contains(',') {
|
||||
return None;
|
||||
}
|
||||
let number = |s: &str| {
|
||||
if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
|
||||
s.parse::<u64>().ok()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
if start.is_empty() {
|
||||
let length = number(end)?;
|
||||
return (length > 0).then_some((total.saturating_sub(length), total - 1));
|
||||
}
|
||||
let start = number(start)?;
|
||||
let end = if end.is_empty() {
|
||||
total - 1
|
||||
} else {
|
||||
number(end)?.min(total - 1)
|
||||
};
|
||||
(start <= end && start < total).then_some((start, end))
|
||||
}
|
||||
|
||||
pub async fn file_response(
|
||||
mut file: File,
|
||||
mime: &str,
|
||||
headers: &HeaderMap,
|
||||
) -> Result<Response<Body>> {
|
||||
let metadata = file.metadata().await?;
|
||||
anyhow::ensure!(metadata.is_file(), "Media source is not a regular file");
|
||||
let total = metadata.len();
|
||||
let selected = match headers.get("range") {
|
||||
None => None,
|
||||
Some(value) => match value.to_str().ok().and_then(|value| range(value, total)) {
|
||||
Some(range) => Some(range),
|
||||
None => {
|
||||
return Ok(Response::builder()
|
||||
.status(StatusCode::RANGE_NOT_SATISFIABLE)
|
||||
.header("Content-Range", format!("bytes */{total}"))
|
||||
.header("Accept-Ranges", "bytes")
|
||||
.header("Cache-Control", "private, no-store")
|
||||
.body(Body::empty())?)
|
||||
}
|
||||
},
|
||||
};
|
||||
let (start, length) = selected
|
||||
.map(|(start, end)| (start, end - start + 1))
|
||||
.unwrap_or((0, total));
|
||||
file.seek(std::io::SeekFrom::Start(start)).await?;
|
||||
let chunks = futures::stream::try_unfold((file, length), |(mut file, left)| async move {
|
||||
if left == 0 {
|
||||
return Ok::<_, std::io::Error>(None);
|
||||
}
|
||||
let mut chunk = vec![0; left.min(64 * 1024) as usize];
|
||||
let read = file.read(&mut chunk).await?;
|
||||
if read == 0 {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::UnexpectedEof,
|
||||
"Media changed during playback",
|
||||
));
|
||||
}
|
||||
chunk.truncate(read);
|
||||
Ok(Some((chunk, (file, left - read as u64))))
|
||||
});
|
||||
let mut response = Response::builder()
|
||||
.status(if selected.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")
|
||||
.header("X-Archipelago-Transport", "local-cache");
|
||||
if let Some((start, end)) = selected {
|
||||
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;
|
||||
#[test]
|
||||
fn ranges_cover_suffix_open_ended_clamping_and_rejection() {
|
||||
assert_eq!(range("bytes=2-5", 10), Some((2, 5)));
|
||||
assert_eq!(range("bytes=2-", 10), Some((2, 9)));
|
||||
assert_eq!(range("bytes=2-100", 10), Some((2, 9)));
|
||||
assert_eq!(range("bytes=-4", 10), Some((6, 9)));
|
||||
assert_eq!(range("bytes=-100", 10), Some((0, 9)));
|
||||
for value in [
|
||||
"bytes=-0",
|
||||
"bytes=10-",
|
||||
"bytes=8-3",
|
||||
"bytes=0-1,4-5",
|
||||
"bytes=+1-4",
|
||||
"bytes=18446744073709551616-",
|
||||
"nope",
|
||||
] {
|
||||
assert_eq!(range(value, 10), None, "{value}");
|
||||
}
|
||||
assert_eq!(range("bytes=0-", 0), None);
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn sparse_large_file_is_streamed_in_bounded_chunks_and_ranges_are_exact() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join("video");
|
||||
let file = File::create(&path).await.unwrap();
|
||||
file.set_len(4 * 1024 * 1024 * 1024).await.unwrap();
|
||||
drop(file);
|
||||
let mut full = file_response(
|
||||
File::open(&path).await.unwrap(),
|
||||
"video/mp4",
|
||||
&HeaderMap::new(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(full.headers()["content-length"], "4294967296");
|
||||
assert_eq!(full.body_mut().data().await.unwrap().unwrap().len(), 65536);
|
||||
drop(full); // Cancellation must not read the remainder.
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert("range", "bytes=-3".parse().unwrap());
|
||||
let response = file_response(File::open(&path).await.unwrap(), "video/mp4", &headers)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.status(), 206);
|
||||
assert_eq!(
|
||||
response.headers()["content-range"],
|
||||
"bytes 4294967293-4294967295/4294967296"
|
||||
);
|
||||
assert_eq!(
|
||||
hyper::body::to_bytes(response.into_body())
|
||||
.await
|
||||
.unwrap()
|
||||
.as_ref(),
|
||||
&[0, 0, 0]
|
||||
);
|
||||
headers.insert("range", "bytes=4294967296-".parse().unwrap());
|
||||
assert_eq!(
|
||||
file_response(File::open(&path).await.unwrap(), "video/mp4", &headers)
|
||||
.await
|
||||
.unwrap()
|
||||
.status(),
|
||||
416
|
||||
);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user