diff --git a/core/archipelago/src/api/rpc/content.rs b/core/archipelago/src/api/rpc/content.rs index bf42e200..89ff3303 100644 --- a/core/archipelago/src/api/rpc/content.rs +++ b/core/archipelago/src/api/rpc/content.rs @@ -43,6 +43,26 @@ async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: & } } +// Inline RPC responses are for previews/small legacy downloads. Films use the +// Range-capable HTTP path; never let a peer force whole-film base64 allocation. +async fn bounded_content_bytes(mut response: reqwest::Response, limit: usize) -> Result> { + anyhow::ensure!( + response + .content_length() + .is_none_or(|size| size <= limit as u64), + "Content exceeds the inline limit; open it through the streaming viewer" + ); + let mut bytes = Vec::new(); + while let Some(chunk) = response.chunk().await? { + anyhow::ensure!( + chunk.len() <= limit.saturating_sub(bytes.len()), + "Content exceeds the inline limit; use streaming" + ); + bytes.extend_from_slice(&chunk); + } + Ok(bytes) +} + async fn bounded_seller_error(mut response: reqwest::Response) -> String { let mut bytes = Vec::new(); let _ = tokio::time::timeout(std::time::Duration::from_secs(5), async { @@ -450,9 +470,7 @@ impl RpcHandler { .header("X-Federation-DID", local_did) .timeout(std::time::Duration::from_secs(120)) .fips_timeout(std::time::Duration::from_secs(8)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await .context("Failed to connect to peer")?; // Record which transport actually reached the peer (B14) so the UI @@ -494,10 +512,9 @@ impl RpcHandler { return Err(anyhow::anyhow!("Peer returned: {}", response.status())); } - let bytes = response - .bytes() + let bytes = bounded_content_bytes(response, 16 * 1024 * 1024) .await - .context("Failed to read response body")?; + .context("Failed to read bounded content")?; use base64::Engine; let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes); @@ -542,9 +559,7 @@ impl RpcHandler { // against the UI's 30s deadline — users saw errors, not // fallback. .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await .context("Failed to connect to peer")?; // Record which transport actually reached the peer (B14). @@ -646,6 +661,9 @@ impl RpcHandler { // one system can still pay (#3). let method = params.get("method").and_then(|v| v.as_str()); + let (data, _) = self.state_manager.get_snapshot().await; + let local_did = crate::identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?; + let mint_cashu = || ecash::send_token(&self.config.data_dir, price_sats); let mint_fedimint = || crate::wallet::fedimint_client::spend_from_any(&self.config.data_dir, price_sats); @@ -699,9 +717,6 @@ impl RpcHandler { "paid download: paying {price_sats} sats to {onion} via {used_backend} ecash" ); - let (data, _) = self.state_manager.get_snapshot().await; - let local_did = crate::identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?; - let path = format!("/content/{}", content_id); // Surface a real reason instead of the generic sanitized error (#30): // A bearer token must not be replayed after an ambiguous delivery. @@ -714,9 +729,7 @@ impl RpcHandler { .header("X-Payment-Token", token_str.clone()) .single_delivery() .timeout(std::time::Duration::from_secs(900)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -835,9 +848,7 @@ impl RpcHandler { .header("X-Federation-DID", local_did) .timeout(std::time::Duration::from_secs(25)) .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -901,9 +912,7 @@ impl RpcHandler { .service(crate::settings::transport::PeerService::PeerFiles) .timeout(std::time::Duration::from_secs(15)) .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -995,9 +1004,7 @@ impl RpcHandler { .header("X-Federation-DID", local_did) .header("X-Invoice-Hash", payment_hash.to_string()) .timeout(std::time::Duration::from_secs(900)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -1095,9 +1102,7 @@ impl RpcHandler { .header("X-Federation-DID", local_did) .timeout(std::time::Duration::from_secs(25)) .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -1156,9 +1161,7 @@ impl RpcHandler { .service(crate::settings::transport::PeerService::PeerFiles) .timeout(std::time::Duration::from_secs(15)) .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -1215,9 +1218,7 @@ impl RpcHandler { .header("X-Federation-DID", local_did) .header("X-Onchain-Address", address.to_string()) .timeout(std::time::Duration::from_secs(900)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await { Ok(v) => v, @@ -1297,9 +1298,7 @@ impl RpcHandler { .require_fips() .timeout(std::time::Duration::from_secs(30)) .fips_timeout(std::time::Duration::from_secs(6)) - .authenticate_content(&self.config.data_dir) - .await? - .send_get() + .send_content_get(&self.config.data_dir) .await .context("Failed to connect to peer for preview")?; // Record which transport actually reached the peer (B14). @@ -1335,10 +1334,9 @@ impl RpcHandler { .unwrap_or("application/octet-stream") .to_string(); - let bytes = response - .bytes() + let bytes = bounded_content_bytes(response, 8 * 1024 * 1024) .await - .context("Failed to read preview response")?; + .context("Failed to read bounded preview")?; use base64::Engine; let encoded = base64::engine::general_purpose::STANDARD.encode(&bytes); diff --git a/core/archipelago/src/api/rpc/content_tests.rs b/core/archipelago/src/api/rpc/content_tests.rs index b53fd3d1..0697246a 100644 --- a/core/archipelago/src/api/rpc/content_tests.rs +++ b/core/archipelago/src/api/rpc/content_tests.rs @@ -168,3 +168,26 @@ async fn known_purchase_never_becomes_a_new_spend_when_cache_or_index_is_unavail b"damaged" ); } + +#[tokio::test] +async fn inline_download_limits_known_and_chunked_bodies_without_draining_them() { + use futures_util::StreamExt; + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }; + let response: reqwest::Response = hyper::Response::new("small").into(); + assert!(bounded_content_bytes(response, 4).await.is_err()); + let response: reqwest::Response = hyper::Response::new("small").into(); + assert_eq!(bounded_content_bytes(response, 5).await.unwrap(), b"small"); + let consumed = Arc::new(AtomicUsize::new(0)); + let counter = consumed.clone(); + let chunks = futures_util::stream::iter(0..1000).map(move |_| { + counter.fetch_add(1, Ordering::SeqCst); + Ok::<_, std::io::Error>(bytes::Bytes::from_static(b"abc")) + }); + let body = reqwest::Body::wrap_stream(chunks); + let response: reqwest::Response = hyper::Response::new(body).into(); + assert!(bounded_content_bytes(response, 4).await.is_err()); + assert_eq!(consumed.load(Ordering::SeqCst), 2); +} diff --git a/core/archipelago/src/content_auth.rs b/core/archipelago/src/content_auth.rs index 111b8ef6..15382034 100644 --- a/core/archipelago/src/content_auth.rs +++ b/core/archipelago/src/content_auth.rs @@ -219,6 +219,26 @@ mod tests { assert!(!visible_to(&item, Some("verified-peer"), true, false)); } + #[tokio::test] + async fn authentication_failure_is_a_delivery_error_before_any_network_attempt() { + let dir = tempfile::tempdir().unwrap(); + tokio::fs::create_dir(dir.path().join("federation")) + .await + .unwrap(); + tokio::fs::write(dir.path().join("federation/nodes.json"), b"invalid") + .await + .unwrap(); + let error = crate::fips::dial::PeerRequest::new(None, "peer.onion", "/content/file") + .require_fips() + .send_content_get(dir.path()) + .await + .unwrap_err(); + // The corrupt identity store fails before the separate missing-FIPS + // route error. It reaches the payment caller's existing refund branch. + assert!(!error.to_string().contains("FIPS")); + assert!(!dir.path().join("identity").exists()); + } + #[test] fn plain_claimed_did_never_becomes_an_authenticated_peer() { let mut headers = hyper::HeaderMap::new(); diff --git a/core/archipelago/src/fips/dial.rs b/core/archipelago/src/fips/dial.rs index 85ff1247..d2966e54 100644 --- a/core/archipelago/src/fips/dial.rs +++ b/core/archipelago/src/fips/dial.rs @@ -502,6 +502,15 @@ impl<'a> PeerRequest<'a> { Ok(self) } + /// Return authentication and delivery failures through one result. Payment + /// callers must be able to reclaim an unsent token if proof preparation fails. + pub async fn send_content_get( + self, + data_dir: &std::path::Path, + ) -> Result<(reqwest::Response, crate::transport::TransportKind)> { + self.authenticate_content(data_dir).await?.send_get().await + } + pub fn header(mut self, name: &'a str, value: impl Into) -> Self { self.headers.push((name, value.into())); self diff --git a/docs/post-1.9.0-progress-20261006.md b/docs/post-1.9.0-progress-20261006.md index 20d4021e..e33fe96e 100644 --- a/docs/post-1.9.0-progress-20261006.md +++ b/docs/post-1.9.0-progress-20261006.md @@ -304,3 +304,31 @@ Seller streaming's first full compile found a lifetime error in one new test; `df7677d2` corrects it. The repeated isolated full suite is compiling from that frozen source (`/tmp/archy-bounded-media-backend-tests-2.log`). The failed run is retained and is not counted as a pass. + + +### Latest peer-content review + +Seller streaming at `df7677d2` passed the full isolated suite: 1,729 passed, +zero failures, five explicit skips. Buyer streaming required updating existing +regression fixtures to the new streamed Files-copy boundary; two failed compile +runs are retained. The final buyer UI has 21 focused tests passing and its +production build passes (`/tmp/archy-buyer-cache-ui-tests-final.log`, +`/tmp/archy-buyer-stream-ui-build.log`). No deployment has occurred. + +A separate source review found restricted requests trusting an unsigned peer DID. +The candidate now verifies recipient/path/range/time-bound node signatures, and +uses the same visibility gate for metadata, invoice issuance and bytes. The +isolated combined suite is compiling from the frozen authentication worktree; +see `peer-content-authentication.md` for compatibility and remaining gates. + +Further review caught an authentication-preparation failure escaping the payment +request's refund branch. Authentication and transport now return through one +result, and local identity validation happens before ecash creation. Inline +preview/legacy RPC bodies are also bounded, including unknown-length responses: +large free videos must not be base64-loaded merely to populate a card preview. +Those final changes still require backend qualification. + +Read-only checks at09:49UTC confirm dev and Yaya Monitoring/federation CPU, memory +and disk measurements match, with each tested RPC below0.5seconds. Framework +still returns401 for the saved dashboard session. These are current live-source +checks, not acceptance of the new undeployed file-streaming candidate.