From bf7fb425ebd86f12875c67a36e9c95059f69f0c1 Mon Sep 17 00:00:00 2001 From: archipelago Date: Mon, 5 Oct 2026 23:51:29 -0400 Subject: [PATCH] Require FIPS for peer playback and stream owned media with bounded reads --- core/archipelago/src/api/handler/proxy.rs | 71 +++------- core/archipelago/src/content_owned.rs | 81 +++++++++++ core/archipelago/src/fips/dial.rs | 32 ++++- core/archipelago/src/main.rs | 1 + core/archipelago/src/media_stream.rs | 161 ++++++++++++++++++++++ 5 files changed, 291 insertions(+), 55 deletions(-) create mode 100644 core/archipelago/src/media_stream.rs diff --git a/core/archipelago/src/api/handler/proxy.rs b/core/archipelago/src/api/handler/proxy.rs index cda46c47..3e92c5a4 100644 --- a/core/archipelago/src/api/handler/proxy.rs +++ b/core/archipelago/src/api/handler/proxy.rs @@ -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); diff --git a/core/archipelago/src/content_owned.rs b/core/archipelago/src/content_owned.rs index 40ce9c14..be7c4846 100644 --- a/core/archipelago/src/content_owned.rs +++ b/core/archipelago/src/content_owned.rs @@ -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> { + 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(); diff --git a/core/archipelago/src/fips/dial.rs b/core/archipelago/src/fips/dial.rs index 02838efc..144ff8cd 100644 --- a/core/archipelago/src/fips/dial.rs +++ b/core/archipelago/src/fips/dial.rs @@ -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(); diff --git a/core/archipelago/src/main.rs b/core/archipelago/src/main.rs index 1e8b99dc..9fb3ac57 100644 --- a/core/archipelago/src/main.rs +++ b/core/archipelago/src/main.rs @@ -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; diff --git a/core/archipelago/src/media_stream.rs b/core/archipelago/src/media_stream.rs new file mode 100644 index 00000000..b5c3a875 --- /dev/null +++ b/core/archipelago/src/media_stream.rs @@ -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::().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> { + 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 + ); + } +}