diff --git a/core/archipelago/src/api/handler/content.rs b/core/archipelago/src/api/handler/content.rs index e1fb9ade..24100269 100644 --- a/core/archipelago/src/api/handler/content.rs +++ b/core/archipelago/src/api/handler/content.rs @@ -162,11 +162,33 @@ impl ApiHandler { r#"{"error":"This file is shared with the host's federation peers only. Federate with that node (exchange invites) so it recognizes you, then try again."}"#, ), )), - Ok(content_server::ServeResult::NotFound) | Err(_) => Ok(build_response( + Ok(content_server::ServeResult::Unavailable) => Ok(build_response( + StatusCode::SERVICE_UNAVAILABLE, + "application/json", + hyper::Body::from( + r#"{"error":"The seller's node can't read this file right now. This request did not redeem an ecash payment."}"#, + ), + )), + Ok(content_server::ServeResult::RangeNotSatisfiable(total)) => Ok(Response::builder() + .status(StatusCode::RANGE_NOT_SATISFIABLE) + .header("Content-Range", format!("bytes */{total}")) + .body(hyper::Body::empty()) + .unwrap()), + Ok(content_server::ServeResult::NotFound) => Ok(build_response( StatusCode::NOT_FOUND, "text/plain", hyper::Body::from("Content not found"), )), + // Not a 404: a paid request may already have been charged by the + // time this fails, and "not found" hid the real error entirely. + Err(e) => { + tracing::error!("Serving content {content_id} failed: {e:#}"); + Ok(build_response( + StatusCode::INTERNAL_SERVER_ERROR, + "text/plain", + hyper::Body::from("Failed to serve content"), + )) + } } } diff --git a/core/archipelago/src/api/rpc/content.rs b/core/archipelago/src/api/rpc/content.rs index dc6cde36..77558529 100644 --- a/core/archipelago/src/api/rpc/content.rs +++ b/core/archipelago/src/api/rpc/content.rs @@ -43,6 +43,25 @@ async fn reclaim_spent_ecash(data_dir: &std::path::Path, token: &str, backend: & } } +/// Only pass through the peer's bounded, printable explanation; refund status +/// is always determined locally and must never come from the peer's wording. +fn seller_error_message(status: reqwest::StatusCode, body: &str) -> String { + let reason = serde_json::from_str::(body) + .ok() + .and_then(|v| v.get("error").and_then(|e| e.as_str()).map(str::to_owned)); + match reason { + Some(reason) if !reason.trim().is_empty() => { + let clean: String = reason + .chars() + .filter(|c| !c.is_control()) + .take(240) + .collect(); + format!("Seller response ({status}): {clean}") + } + _ => format!("Peer returned an error ({status})."), + } +} + /// Keep first purchases and cached repeats compatible with both existing clients. fn paid_content_response(bytes: &[u8], mime: &str, paid_sats: u64) -> serde_json::Value { use base64::Engine; @@ -54,13 +73,9 @@ fn paid_content_response(bytes: &[u8], mime: &str, paid_sats: u64) -> serde_json }) } -/// FileBrowser owns its files through a rootless UID mapping. Use its authenticated -/// API rather than writing host paths with the backend's unrelated UID. Its -/// override=false upload atomically refuses existing names, including races. +/// File purchases through an atomic no-clobber write in Files' own namespace. async fn file_purchase_in_files( - client: &reqwest::Client, - base_url: &str, - token: &str, + data_dir: &std::path::Path, filename: &str, mime: &str, bytes: &[u8], @@ -72,59 +87,24 @@ async fn file_purchase_in_files( } else { "Documents" }; - let mut folder_url = reqwest::Url::parse(base_url)?; - folder_url - .path_segments_mut() - .map_err(|_| anyhow::anyhow!("Invalid Files URL"))? - .extend(["api", "resources", folder, ""]); - let response = client - .get(folder_url.clone()) - .header("X-Auth", token) - .send() - .await?; - if response.status() == reqwest::StatusCode::NOT_FOUND { - let response = client - .post(folder_url.clone()) - .header("X-Auth", token) - .send() - .await?; - if response.status() != reqwest::StatusCode::CONFLICT { - response.error_for_status()?; - } - } else { - response.error_for_status()?; - } - let base = std::path::Path::new(filename) + let root = data_dir.join("filebrowser"); + anyhow::ensure!( + tokio::fs::metadata(&root).await?.is_dir(), + "Files storage is unavailable" + ); + let name = std::path::Path::new(filename) .file_name() .and_then(|n| n.to_str()) .filter(|n| !n.is_empty()) .unwrap_or("download"); - let (stem, extension) = match base.rsplit_once('.') { - Some((stem, ext)) if !stem.is_empty() => (stem, format!(".{ext}")), - _ => (base, String::new()), - }; - for attempt in 1..=100 { - let name = if attempt == 1 { - base.to_string() - } else { - format!("{stem} ({attempt}){extension}") - }; - let mut url = folder_url.clone(); - url.path_segments_mut().unwrap().pop_if_empty().push(&name); - url.query_pairs_mut().append_pair("override", "false"); - let response = client - .post(url) - .header("X-Auth", token) - .body(bytes.to_vec()) - .send() - .await?; - if response.status() == reqwest::StatusCode::CONFLICT { - continue; - } - response.error_for_status()?; - return Ok(format!("{folder}/{name}")); - } - anyhow::bail!("Too many existing copies; purchased file remains in the purchase cache") + let path = + crate::container::filebrowser::save_new_file(&root.join(folder), name, bytes).await?; + Ok(format!( + "{folder}/{}", + path.file_name() + .and_then(|n| n.to_str()) + .context("Invalid Files name")? + )) } impl RpcHandler { @@ -623,13 +603,14 @@ impl RpcHandler { let path = format!("/content/{}", content_id); // Surface a real reason instead of the generic sanitized error (#30): - // the dial already tries FIPS/mesh then falls back to Tor, so a failure - // here means the peer is genuinely unreachable on both transports. + // A bearer token must not be replayed after an ambiguous delivery. + // A transport error can mean the seller received it without replying. let (response, transport) = match crate::fips::dial::PeerRequest::new(fips_npub.as_deref(), onion, &path) .service(crate::settings::transport::PeerService::PeerFiles) .header("X-Federation-DID", local_did) .header("X-Payment-Token", token_str.clone()) + .single_delivery() .timeout(std::time::Duration::from_secs(900)) .send_get() .await @@ -642,7 +623,7 @@ impl RpcHandler { let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ - "error": format!("Could not reach the peer over mesh or Tor. {refund}") + "error": format!("The purchase could not be completed. {refund}") })); } }; @@ -679,7 +660,7 @@ impl RpcHandler { tracing::warn!("paid download: seller {onion} returned {status}: {body}"); let refund = reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; return Ok(serde_json::json!({ - "error": format!("Peer returned an error ({status}). {refund}") + "error": format!("{} {refund}", seller_error_message(status, &body)) })); } @@ -693,10 +674,17 @@ impl RpcHandler { .filter(|s| !s.is_empty()) .unwrap_or_else(|| "application/octet-stream".to_string()); - let bytes = response - .bytes() - .await - .context("Failed to read response body")?; + let bytes = match response.bytes().await { + Ok(bytes) => bytes, + Err(error) => { + tracing::warn!("paid download: response body failed: {error}"); + let refund = + reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await; + return Ok(serde_json::json!({ + "error": format!("The file transfer was interrupted after payment was sent. {refund}") + })); + } + }; // Persist the purchase so it "stays unlocked" for this buyer: cache the // bytes + metadata keyed by (onion, content_id). The gallery then renders @@ -728,28 +716,8 @@ impl RpcHandler { // The durable purchased-content cache above is primary. A Files copy // remains optional: a stopped FileBrowser must not undo a paid download. - let filed = async { - let auth = self.handle_filebrowser_token().await?; - let token = auth - .get("token") - .and_then(|v| v.as_str()) - .context("FileBrowser omitted its authentication token")?; - let client = reqwest::Client::builder() - .no_proxy() - .redirect(reqwest::redirect::Policy::none()) - .timeout(std::time::Duration::from_secs(30)) - .build()?; - file_purchase_in_files( - &client, - "http://127.0.0.1:8083", - token, - &filename, - &mime_type, - &bytes, - ) - .await - } - .await; + let filed = + file_purchase_in_files(&self.config.data_dir, &filename, &mime_type, &bytes).await; match filed { Ok(path) => tracing::info!("paid download: filed into Files/{path}"), Err(error) => tracing::warn!( diff --git a/core/archipelago/src/api/rpc/content_tests.rs b/core/archipelago/src/api/rpc/content_tests.rs index 78fa9d51..7df1742b 100644 --- a/core/archipelago/src/api/rpc/content_tests.rs +++ b/core/archipelago/src/api/rpc/content_tests.rs @@ -1,69 +1,4 @@ use super::*; -use hyper::{ - service::{make_service_fn, service_fn}, - Body, Response, Server, -}; -use std::{ - collections::VecDeque, - convert::Infallible, - sync::{Arc, Mutex}, -}; - -struct FilesApi { - url: String, - seen: Arc)>>>, - task: tokio::task::JoinHandle<()>, -} -impl Drop for FilesApi { - fn drop(&mut self) { - self.task.abort(); - } -} -fn files_api(statuses: Vec) -> FilesApi { - let statuses = Arc::new(Mutex::new(VecDeque::from(statuses))); - let seen = Arc::new(Mutex::new(Vec::new())); - let history = seen.clone(); - let server = Server::bind(&([127, 0, 0, 1], 0).into()); - let address = server.local_addr(); - let service = make_service_fn(move |_| { - let statuses = statuses.clone(); - let seen = history.clone(); - async move { - Ok::<_, Infallible>(service_fn(move |request: hyper::Request| { - let statuses = statuses.clone(); - let seen = seen.clone(); - async move { - assert_eq!(request.headers().get("X-Auth").unwrap(), "test-session"); - let method = request.method().to_string(); - let uri = request.uri().to_string(); - let body = hyper::body::to_bytes(request.into_body()) - .await - .unwrap() - .to_vec(); - seen.lock().unwrap().push((method, uri, body)); - let status = statuses - .lock() - .unwrap() - .pop_front() - .expect("unexpected extra Files request"); - Ok::<_, Infallible>( - Response::builder() - .status(status) - .body(Body::empty()) - .unwrap(), - ) - } - })) - } - }); - FilesApi { - url: format!("http://{address}"), - seen, - task: tokio::spawn(async move { - server.serve(service).await.unwrap(); - }), - } -} #[test] fn first_and_cached_paid_downloads_have_the_same_client_payload_contract() { @@ -85,80 +20,54 @@ fn first_and_cached_paid_downloads_have_the_same_client_payload_contract() { } #[tokio::test] -async fn files_copy_uses_authenticated_api_and_preserves_existing_names() { - let api = files_api(vec![200, 409, 200]); - let client = reqwest::Client::new(); - let path = file_purchase_in_files( - &client, - &api.url, - "test-session", - "../my #file?.txt", - "text/plain", - b"paid bytes", - ) - .await - .unwrap(); - assert_eq!(path, "Documents/my #file? (2).txt"); - let seen = api.seen.lock().unwrap(); - assert_eq!(seen[0].0, "GET"); - assert_eq!(seen[0].1, "/api/resources/Documents/"); - assert_eq!(seen.len(), 3); - for (_, uri, body) in &seen[1..] { - assert!(uri.contains("override=false")); - assert!(uri.contains("%23file%3F")); - assert!(!uri.contains("../")); - assert_eq!(body, b"paid bytes"); - } -} - -#[tokio::test] -async fn files_copy_creates_missing_media_folder() { +async fn files_copy_routes_media_and_sanitizes_the_filename() { + let dir = tempfile::tempdir().unwrap(); + tokio::fs::create_dir(dir.path().join("filebrowser")) + .await + .unwrap(); for (mime, folder) in [ ("image/png", "Photos"), ("video/mp4", "Photos"), - ("audio/ogg", "Music"), + ("audio/mpeg", "Music"), + ("text/plain", "Documents"), ] { - let api = files_api(vec![404, 200, 200]); - let path = file_purchase_in_files( - &reqwest::Client::new(), - &api.url, - "test-session", - "file", - mime, - b"bytes", - ) - .await - .unwrap(); - assert_eq!(path, format!("{folder}/file")); - let seen = api.seen.lock().unwrap(); - assert_eq!(seen[1].0, "POST"); - assert!(seen[1].1.ends_with('/')); - assert!(seen[1].2.is_empty()); - assert_eq!(seen[2].2, b"bytes"); + let relative = file_purchase_in_files(dir.path(), "../name #?.bin", mime, b"paid") + .await + .unwrap(); + assert!(relative.starts_with(&format!("{folder}/name #?"))); + assert_eq!( + tokio::fs::read(dir.path().join("filebrowser").join(relative)) + .await + .unwrap(), + b"paid" + ); } } #[tokio::test] -async fn files_copy_fails_without_overwriting_or_claiming_success_on_errors() { - for statuses in [ - vec![401], - vec![503], - vec![404, 500], - vec![200, 507], - vec![200, 403], - ] { - let expected = statuses.len(); - let api = files_api(statuses); - assert!(file_purchase_in_files( - &reqwest::Client::new(), - &api.url, - "test-session", - "file.txt", - "text/plain", - b"bytes" - ) - .await - .is_err()); - assert_eq!(api.seen.lock().unwrap().len(), expected); +async fn unavailable_files_storage_is_reported_without_creating_a_fake_installation() { + let dir = tempfile::tempdir().unwrap(); + assert!( + file_purchase_in_files(dir.path(), "name", "text/plain", b"bytes") + .await + .is_err() + ); + assert!(!dir.path().join("filebrowser").exists()); +} + +#[test] +fn seller_errors_are_bounded_printable_and_identified_as_peer_text() { + let status = reqwest::StatusCode::SERVICE_UNAVAILABLE; + let message = seller_error_message(status, r#"{"error":"Cannot read file\n\u0000"}"#); + assert!(message.starts_with("Seller response (503")); + assert!(message.ends_with("Cannot read file")); + assert!(!message.contains('\n') && !message.contains('\0')); + let body = serde_json::json!({"error": "é".repeat(1000)}).to_string(); + assert!(seller_error_message(status, &body).chars().count() < 300); + for body in ["not JSON", r#"{"error": 7}"#, r#"{"error":" "}"#] { + assert_eq!( + seller_error_message(status, body), + "Peer returned an error (503 Service Unavailable)." + ); } } diff --git a/core/archipelago/src/container/filebrowser.rs b/core/archipelago/src/container/filebrowser.rs index e51b11fe..1e85fa11 100644 --- a/core/archipelago/src/container/filebrowser.rs +++ b/core/archipelago/src/container/filebrowser.rs @@ -5,7 +5,7 @@ //! starting the container with `--config /data/.filebrowser.json`. use anyhow::{Context, Result}; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use tokio::fs; use crate::update::host_sudo; @@ -117,6 +117,197 @@ fn shell_quote(s: &str) -> String { s.replace('\'', "'\\''") } +/// Save a complete purchase without overwriting any existing directory entry. +/// Both host and rootless-namespace paths publish with a no-clobber hard link. +pub async fn save_new_file(dir: &Path, name: &str, bytes: &[u8]) -> Result { + save_new_file_with(dir, name, bytes, write_via_userns).await +} + +fn validate_filename(name: &str) -> Result<()> { + anyhow::ensure!( + !name.is_empty() + && name != "." + && name != ".." + && !name.contains(['/', '\\', '\0']) + && name.len() <= 255, + "Invalid purchased filename" + ); + Ok(()) +} + +async fn save_new_file_with( + dir: &Path, + name: &str, + bytes: &[u8], + fallback: F, +) -> Result +where + F: FnOnce(PathBuf, String, Vec) -> Fut, + Fut: std::future::Future>, +{ + validate_filename(name)?; + // Never follow a user-created destination directory symlink. + match fs::symlink_metadata(dir).await { + Ok(meta) => anyhow::ensure!(meta.is_dir(), "Files destination is not a directory"), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => {} + Err(error) => return Err(error.into()), + } + save_after_direct_result( + write_direct(dir, name, bytes).await, + dir, + name, + bytes, + fallback, + ) + .await +} + +async fn save_after_direct_result( + result: std::io::Result, + dir: &Path, + name: &str, + bytes: &[u8], + fallback: F, +) -> Result +where + F: FnOnce(PathBuf, String, Vec) -> Fut, + Fut: std::future::Future>, +{ + match result { + Ok(path) => Ok(path), + Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => { + fallback(dir.to_owned(), name.to_owned(), bytes.to_vec()) + .await + .context("Saving purchase in Files user namespace") + } + Err(error) => Err(error).context("Saving purchase in Files"), + } +} + +fn numbered_name(name: &str, attempt: usize) -> String { + if attempt == 1 { + return name.to_owned(); + } + match name.rsplit_once('.') { + Some((stem, extension)) if !stem.is_empty() => format!("{stem} ({attempt}).{extension}"), + _ => format!("{name} ({attempt})"), + } +} + +struct PendingFile(PathBuf); +impl Drop for PendingFile { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.0); + } +} + +async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result { + use std::os::unix::fs::PermissionsExt; + use tokio::io::AsyncWriteExt; + fs::create_dir_all(dir).await?; + let temp_path = dir.join(format!(".archy-saving-{}", uuid::Uuid::new_v4())); + let mut file = fs::OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&temp_path) + .await?; + let temp = PendingFile(temp_path); + file.write_all(bytes).await?; + file.set_permissions(std::fs::Permissions::from_mode(0o644)) + .await?; + file.sync_all().await?; + for attempt in 1..=100 { + let target = dir.join(numbered_name(name, attempt)); + match fs::hard_link(&temp.0, &target).await { + Ok(()) => return Ok(target), + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue, + Err(error) => return Err(error), + } + } + Err(std::io::Error::new( + std::io::ErrorKind::AlreadyExists, + "Too many existing copies; purchase cache retained", + )) +} + +// Positional arguments carry all user-controlled text. mktemp prevents temp-name +// collisions; ln -T refuses files, symlinks and directories, including races. +const WRITE_VIA_USERNS: &str = r#"set -eu +dir=$1 +name=$2 +expected=$3 +[ ! -L "$dir" ] || exit 1 +if [ ! -d "$dir" ]; then + mkdir -p -- "$dir" + chown --reference="$(dirname -- "$dir")" -- "$dir" +fi +tmp=$(mktemp "$dir/.archy-saving.XXXXXXXXXX") +trap 'rm -f -- "$tmp"' EXIT HUP INT TERM +cat > "$tmp" +[ "$(wc -c < "$tmp")" -eq "$expected" ] || exit 1 +chown --reference="$dir" -- "$tmp" +chmod 0644 -- "$tmp" +sync -f -- "$tmp" +stem=$name +ext= +case "$name" in + *.*) prefix=${name%.*}; if [ -n "$prefix" ]; then stem=$prefix; ext=.${name##*.}; fi ;; +esac +n=1 +while [ "$n" -le 100 ]; do + candidate=$name + if [ "$n" -gt 1 ]; then candidate="$stem ($n)$ext"; fi + dst="$dir/$candidate" + if ln -T -- "$tmp" "$dst" 2>/dev/null; then + printf '%s' "$candidate" + exit 0 + fi + # A conflict may be a dangling symlink; never follow it or overwrite it. + if [ ! -e "$dst" ] && [ ! -L "$dst" ]; then exit 1; fi + n=$((n + 1)) +done +exit 1 +"#; + +async fn write_via_userns(dir: PathBuf, name: String, bytes: Vec) -> Result { + use tokio::io::AsyncWriteExt; + let mut child = tokio::process::Command::new("podman") + .args(["unshare", "sh", "-c", WRITE_VIA_USERNS, "sh"]) + .arg(&dir) + .arg(&name) + .arg(bytes.len().to_string()) + .kill_on_drop(true) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn() + .context("Starting Files namespace writer")?; + let mut stdin = child.stdin.take().context("Files writer stdin missing")?; + let operation = async { + let fed = stdin.write_all(&bytes).await; + drop(stdin); + let output = child.wait_with_output().await?; + anyhow::ensure!( + output.status.success(), + "Files namespace writer failed: {}", + output.status + ); + fed.context("Sending purchase bytes to Files")?; + let chosen = + String::from_utf8(output.stdout).context("Files writer returned an invalid name")?; + validate_filename(&chosen)?; + anyhow::ensure!( + (1..=100).any(|n| numbered_name(&name, n) == chosen), + "Files writer returned an unexpected name" + ); + Ok(dir.join(chosen)) + }; + tokio::time::timeout(std::time::Duration::from_secs(120), operation) + .await + .context("Files namespace writer timed out")? +} + #[cfg(test)] mod tests { use super::*; @@ -152,3 +343,231 @@ mod tests { assert_eq!(second, EnsureOutcome::Unchanged); } } + +#[cfg(test)] +mod purchase_write_tests { + use super::*; + use std::{ + collections::HashSet, + os::unix::fs::{symlink, PermissionsExt}, + }; + + fn no_temps(dir: &Path) { + assert!(std::fs::read_dir(dir).unwrap().all(|e| !e + .unwrap() + .file_name() + .to_string_lossy() + .starts_with(".archy-saving"))); + } + + #[tokio::test] + async fn direct_write_uses_complete_bytes_and_preserves_originals() { + let dir = tempfile::tempdir().unwrap(); + fs::write(dir.path().join("song.mp3"), b"original") + .await + .unwrap(); + let target = save_new_file(dir.path(), "song.mp3", b"new").await.unwrap(); + assert_eq!(target.file_name().unwrap(), "song (2).mp3"); + assert_eq!(fs::read(target).await.unwrap(), b"new"); + assert_eq!( + fs::read(dir.path().join("song.mp3")).await.unwrap(), + b"original" + ); + no_temps(dir.path()); + } + + #[tokio::test] + async fn simultaneous_saves_publish_unique_complete_files() { + let dir = tempfile::tempdir().unwrap(); + let mut tasks = Vec::new(); + for n in 0..24u8 { + let dir = dir.path().to_owned(); + tasks.push(tokio::spawn(async move { + let bytes = vec![n; 32768]; + let path = save_new_file(&dir, "same.bin", &bytes).await.unwrap(); + assert_eq!(fs::read(&path).await.unwrap(), bytes); + path + })); + } + let mut paths = HashSet::new(); + for task in tasks { + assert!(paths.insert(task.await.unwrap())); + } + assert_eq!(paths.len(), 24); + no_temps(dir.path()); + } + + #[tokio::test] + async fn existing_directories_and_dangling_symlinks_are_conflicts() { + let dir = tempfile::tempdir().unwrap(); + fs::create_dir(dir.path().join("name")).await.unwrap(); + symlink("missing", dir.path().join("name (2)")).unwrap(); + let path = save_new_file(dir.path(), "name", b"new").await.unwrap(); + assert_eq!(path.file_name().unwrap(), "name (3)"); + assert!(dir.path().join("name").is_dir()); + assert!(fs::symlink_metadata(dir.path().join("name (2)")) + .await + .unwrap() + .is_symlink()); + no_temps(dir.path()); + } + + #[tokio::test] + async fn invalid_names_and_symlink_destination_are_refused() { + let dir = tempfile::tempdir().unwrap(); + for name in [ + "", + ".", + "..", + "../escape", + "/absolute", + "a/b", + "a\\b", + "a\0b", + ] { + assert!(save_new_file(dir.path(), name, b"bytes").await.is_err()); + } + let outside = tempfile::tempdir().unwrap(); + symlink(outside.path(), dir.path().join("Music")).unwrap(); + assert!(save_new_file(&dir.path().join("Music"), "song", b"bytes") + .await + .is_err()); + assert_eq!(std::fs::read_dir(outside.path()).unwrap().count(), 0); + } + + #[tokio::test] + async fn collision_limit_preserves_all_files_and_cleans_temporary_data() { + let dir = tempfile::tempdir().unwrap(); + for n in 1..=100 { + fs::write(dir.path().join(numbered_name("a.txt", n)), b"keep") + .await + .unwrap(); + } + assert!(save_new_file(dir.path(), "a.txt", b"new").await.is_err()); + for n in 1..=100 { + assert_eq!( + fs::read(dir.path().join(numbered_name("a.txt", n))) + .await + .unwrap(), + b"keep" + ); + } + no_temps(dir.path()); + } + + #[tokio::test] + async fn permission_fallback_is_exercised_without_skipping_as_root() { + let dir = tempfile::tempdir().unwrap(); + let result = save_after_direct_result( + Err(std::io::ErrorKind::PermissionDenied.into()), + dir.path(), + "a", + b"abc", + |dir, name, bytes| async move { + assert_eq!(bytes, b"abc"); + Ok(dir.join(name)) + }, + ) + .await + .unwrap(); + assert_eq!(result, dir.path().join("a")); + assert!(save_after_direct_result( + Err(std::io::ErrorKind::PermissionDenied.into()), + dir.path(), + "a", + b"abc", + |_, _, _| async { anyhow::bail!("namespace unavailable") } + ) + .await + .unwrap_err() + .to_string() + .contains("namespace")); + assert!(save_after_direct_result( + Err(std::io::ErrorKind::StorageFull.into()), + dir.path(), + "a", + b"abc", + |_, _, _| async { panic!("disk full must not trigger permission fallback") } + ) + .await + .is_err()); + } + + async fn run_script( + dir: &Path, + name: &str, + bytes: &[u8], + expected: usize, + ) -> std::process::Output { + use tokio::io::AsyncWriteExt; + let mut child = tokio::process::Command::new("sh") + .args(["-c", WRITE_VIA_USERNS, "sh"]) + .arg(dir) + .arg(name) + .arg(expected.to_string()) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn() + .unwrap(); + let mut input = child.stdin.take().unwrap(); + input.write_all(bytes).await.unwrap(); + drop(input); + child.wait_with_output().await.unwrap() + } + + #[tokio::test] + async fn namespace_script_preserves_names_bytes_modes_and_existing_entries() { + let dir = tempfile::tempdir().unwrap(); + let folder = dir.path().join("Music"); + let name = "song ' $() ; #.mp3"; + for n in 1..=2 { + let output = run_script(&folder, name, b"abc", 3).await; + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + let chosen = String::from_utf8(output.stdout).unwrap(); + assert_eq!(chosen, numbered_name(name, n)); + let path = folder.join(chosen); + assert_eq!(fs::read(&path).await.unwrap(), b"abc"); + assert_eq!( + fs::metadata(path).await.unwrap().permissions().mode() & 0o777, + 0o644 + ); + } + no_temps(&folder); + } + + #[tokio::test] + async fn namespace_script_refuses_truncated_input_and_cleans_up() { + let dir = tempfile::tempdir().unwrap(); + let output = run_script(dir.path(), "never.bin", b"partial", 100).await; + assert!(!output.status.success()); + assert!(!dir.path().join("never.bin").exists()); + no_temps(dir.path()); + } + + #[tokio::test] + async fn namespace_script_does_not_link_inside_existing_directory() { + let dir = tempfile::tempdir().unwrap(); + fs::create_dir(dir.path().join("name")).await.unwrap(); + symlink("missing", dir.path().join("name (2)")).unwrap(); + let output = run_script(dir.path(), "name", b"abc", 3).await; + assert!(output.status.success()); + assert_eq!(output.stdout, b"name (3)"); + assert_eq!( + std::fs::read_dir(dir.path().join("name")).unwrap().count(), + 0 + ); + no_temps(dir.path()); + } + + #[test] + fn names_keep_extensions_and_dotfiles() { + assert_eq!(numbered_name("a.tar.gz", 2), "a.tar (2).gz"); + assert_eq!(numbered_name(".hidden", 2), ".hidden (2)"); + assert_eq!(numbered_name("README", 2), "README (2)"); + } +} diff --git a/core/archipelago/src/content_server.rs b/core/archipelago/src/content_server.rs index f12d2f24..0858fda0 100644 --- a/core/archipelago/src/content_server.rs +++ b/core/archipelago/src/content_server.rs @@ -238,6 +238,11 @@ pub enum ServeResult { Forbidden, /// Content not found. NotFound, + /// The catalog entry and file exist but this node can't read the file. + /// Returned before any payment is taken. + Unavailable, + /// Requested byte range cannot be served; no payment was taken. + RangeNotSatisfiable(u64), } /// Serve a content item by ID with access control and optional range request. @@ -252,6 +257,39 @@ pub async fn serve_content( range: Option, owner_session: bool, ) -> Result { + serve_content_with( + data_dir, + id, + payment_token, + invoice_hash, + peer_did, + range, + owner_session, + |path, range, mime| prepare_content(data_dir, path, range, mime), + |token, amount| async move { verify_payment_token(data_dir, &token, amount).await }, + ) + .await +} + +// Inject only the read and payment boundaries, so tests can prove ordering +// without mint access, file-permission assumptions or privileged commands. +async fn serve_content_with( + data_dir: &Path, + id: &str, + payment_token: Option<&str>, + invoice_hash: Option<&str>, + peer_did: Option<&str>, + range: Option, + owner_session: bool, + read: R, + verify: V, +) -> Result +where + R: FnOnce(PathBuf, Option, String) -> RF, + RF: std::future::Future>, + V: FnOnce(String, u64) -> VF, + VF: std::future::Future, +{ let catalog = load_catalog(data_dir).await?; let item = match catalog.items.iter().find(|i| i.id == id) { Some(i) => i, @@ -314,6 +352,29 @@ pub async fn serve_content( return Ok(ServeResult::NotFound); } + // Refuse unauthorized viewers before opening or reading any bytes. + if !owner_session && matches!(item.access, AccessControl::PeersOnly) && !is_known_peer { + return Ok(ServeResult::Forbidden); + } + if !owner_session { + if let AccessControl::Paid { price_sats, .. } = &item.access { + if payment_token.is_none() && invoice_hash.is_none() { + return Ok(ServeResult::PaymentRequired(*price_sats)); + } + } + } + + // 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(other) => return Ok(other), + Err(error) => { + warn!(content_id = %id, "Cannot prepare shared content: {error:#}"); + return Ok(ServeResult::Unavailable); + } + }; + // Check access control if !owner_session { match &item.access { @@ -331,7 +392,7 @@ pub async fn serve_content( "fedimint" }; if method_accepted(&item.access, method) - && verify_payment_token(data_dir, token, *price_sats).await + && verify(token.to_owned(), *price_sats).await { authorized = true; } @@ -358,55 +419,127 @@ pub async fn serve_content( } } - let metadata = fs::metadata(&file_path) + Ok(prepared) +} + +async fn prepare_content( + data_dir: &Path, + path: PathBuf, + range: Option, + mime: String, +) -> Result { + use tokio::io::{AsyncReadExt, AsyncSeekExt}; + let mut file = match fs::OpenOptions::new() + .read(true) + .custom_flags(libc::O_NONBLOCK) + .open(&path) .await - .context("Failed to read file metadata")?; - let total_size = metadata.len(); - - // Handle range request for streaming - if let Some(range) = range { - let start = range.start.min(total_size.saturating_sub(1)); - let end = range - .end - .map(|e| e.min(total_size - 1)) - .unwrap_or(total_size - 1); - - if start > end || start >= total_size { - return Ok(ServeResult::NotFound); + { + 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); } - - let len = (end - start + 1) as usize; - use tokio::io::{AsyncReadExt, AsyncSeekExt}; - let mut file = tokio::fs::File::open(&file_path) + 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("Failed to open content file")?; - file.seek(std::io::SeekFrom::Start(start)) - .await - .context("Failed to seek")?; - let mut buf = vec![0u8; len]; - file.read_exact(&mut buf) - .await - .context("Failed to read range")?; - - debug!( - "Serving content '{}' range {}-{}/{} ({} bytes)", - id, start, end, total_size, len - ); + .context("Reading shared content range")?; return Ok(ServeResult::Partial { - bytes: buf, - mime_type: item.mime_type.clone(), + bytes, + mime_type: mime, start, end, - total: total_size, + total, }); } - - let bytes = fs::read(&file_path) + let mut bytes = Vec::new(); + file.read_to_end(&mut bytes) .await - .context("Failed to read content file")?; + .context("Reading shared content")?; + Ok(ServeResult::Ok(bytes, mime)) +} - debug!("Serving content '{}' ({} bytes)", id, bytes.len()); - Ok(ServeResult::Ok(bytes, item.mime_type.clone())) +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)) +} + +fn slice_prepared_content( + bytes: Vec, + range: Option, + mime: String, +) -> Result { + let total = bytes.len() as u64; + match range { + None => Ok(ServeResult::Ok(bytes, mime)), + Some(range) => match checked_range(&range, total) { + Some((start, end)) => Ok(ServeResult::Partial { + bytes: bytes[start as usize..=end as usize].to_vec(), + mime_type: mime, + start, + end, + total, + }), + None => Ok(ServeResult::RangeNotSatisfiable(total)), + }, + } +} + +/// Read only an explicitly shared, regular file within FileBrowser storage. +/// Do not change its mode or grant world-readable access to paid/private data. +async fn filebrowser_read_path(data_dir: &Path, path: &Path) -> Result { + let root = fs::canonicalize(data_dir.join("filebrowser")).await?; + let target = fs::canonicalize(path).await?; + anyhow::ensure!( + target.starts_with(&root) && target != root, + "Shared file is outside Files storage" + ); + anyhow::ensure!( + fs::metadata(&target).await?.is_file(), + "Shared content is not a regular file" + ); + Ok(target) +} + +async fn read_filebrowser_via_userns(data_dir: &Path, path: &Path) -> 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; + 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) + .kill_on_drop(true) + .output(), + ) + .await + .context("Files namespace read timed out")??; + anyhow::ensure!( + output.status.success(), + "Files namespace read failed: {}", + output.status + ); + Ok(output.stdout) + } } /// Result of attempting to serve a preview. @@ -729,3 +862,301 @@ mod prune_missing_content_tests { assert_eq!(reloaded.items[0].id, "present-item"); } } + +#[cfg(test)] +mod paid_read_order_tests { + use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + + async fn fixture(bytes: &[u8]) -> tempfile::TempDir { + let dir = tempfile::tempdir().unwrap(); + fs::create_dir_all(dir.path().join("content/files")) + .await + .unwrap(); + fs::write(dir.path().join("content/files/test.bin"), bytes) + .await + .unwrap(); + save_catalog( + dir.path(), + &ContentCatalog { + items: vec![ContentItem { + id: "paid".into(), + filename: "test.bin".into(), + mime_type: "application/octet-stream".into(), + size_bytes: bytes.len() as u64, + description: String::new(), + access: AccessControl::Paid { + price_sats: 10, + accepted: vec!["ecash".into()], + }, + availability: Availability::AllPeers, + added_at: "2026-09-30".into(), + }], + }, + ) + .await + .unwrap(); + dir + } + + #[tokio::test] + async fn all_read_failures_precede_redemption_even_as_root() { + for kind in [ + std::io::ErrorKind::PermissionDenied, + std::io::ErrorKind::UnexpectedEof, + std::io::ErrorKind::NotFound, + std::io::ErrorKind::Other, + ] { + let dir = fixture(b"abc").await; + let charged = AtomicUsize::new(0); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |_, _, _| async move { Err(std::io::Error::from(kind).into()) }, + |_, _| async { + charged.fetch_add(1, Ordering::SeqCst); + true + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Unavailable)); + assert_eq!(charged.load(Ordering::SeqCst), 0); + assert_eq!(load_catalog(dir.path()).await.unwrap().items.len(), 1); + } + } + + #[tokio::test] + async fn deletion_during_payment_cannot_lose_prepared_bytes() { + let dir = fixture(b"original").await; + 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| { + assert_eq!(amount, 10); + async { + fs::remove_file(dir.path().join("content/files/test.bin")) + .await + .unwrap(); + true + } + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Ok(bytes, _) if bytes == b"original")); + } + + #[tokio::test] + async fn empty_out_of_bounds_and_reversed_ranges_never_charge() { + for (bytes, start, end) in [ + (b"".as_slice(), 0, None), + (b"abc".as_slice(), 3, None), + (b"abc".as_slice(), 2, Some(1)), + ] { + let dir = fixture(bytes).await; + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + Some(ByteRange { start, end }), + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("invalid range reached payment") }, + ) + .await + .unwrap(); + assert!( + matches!(result, ServeResult::RangeNotSatisfiable(n) if n == bytes.len() as u64) + ); + } + } + + #[tokio::test] + async fn prepared_range_survives_file_change_while_payment_is_verified() { + let dir = fixture(b"abcdef").await; + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + Some(ByteRange { + start: 2, + end: Some(999), + }), + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { + fs::write(dir.path().join("content/files/test.bin"), b"x") + .await + .unwrap(); + true + }, + ) + .await + .unwrap(); + assert!( + matches!(result, ServeResult::Partial { bytes, start: 2, end: 5, total: 6, .. } if bytes == b"cdef") + ); + } + + #[tokio::test] + async fn payment_denial_never_returns_prepared_content() { + let dir = fixture(b"secret").await; + let charged = 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), + |_, _| async { + charged.fetch_add(1, Ordering::SeqCst); + false + }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::PaymentRequired(10))); + assert_eq!(charged.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn missing_payment_and_peer_restrictions_precede_file_reads() { + let dir = fixture(b"secret").await; + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + false, + |_, _, _| async { panic!("unauthorized file read") }, + |_, _| async { panic!("unexpected payment") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::PaymentRequired(10))); + let mut catalog = load_catalog(dir.path()).await.unwrap(); + catalog.items[0].access = AccessControl::PeersOnly; + save_catalog(dir.path(), &catalog).await.unwrap(); + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + false, + |_, _, _| async { panic!("unauthorized file read") }, + |_, _| async { panic!("unexpected payment") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Forbidden)); + } + + #[tokio::test] + async fn owner_reads_paid_content_without_redemption() { + let dir = fixture(b"own file").await; + let result = serve_content_with( + dir.path(), + "paid", + None, + None, + None, + None, + true, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("owner charged") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Ok(bytes, _) if bytes == b"own file")); + } + + #[tokio::test] + async fn directory_in_place_of_file_does_not_charge() { + let dir = fixture(b"abc").await; + let path = dir.path().join("content/files/test.bin"); + fs::remove_file(&path).await.unwrap(); + fs::create_dir(&path).await.unwrap(); + let result = serve_content_with( + dir.path(), + "paid", + Some("cashuBtest"), + None, + None, + None, + false, + |path, range, mime| prepare_content(dir.path(), path, range, mime), + |_, _| async { panic!("directory charged") }, + ) + .await + .unwrap(); + assert!(matches!(result, ServeResult::Unavailable)); + } + + #[tokio::test] + async fn files_namespace_read_is_scoped_to_regular_files_and_keeps_mode() { + use std::os::unix::fs::{symlink, PermissionsExt}; + let dir = fixture(b"outside").await; + let root = dir.path().join("filebrowser"); + fs::create_dir(&root).await.unwrap(); + let inside = root.join("song"); + fs::write(&inside, b"song").await.unwrap(); + fs::set_permissions(&inside, std::fs::Permissions::from_mode(0o640)) + .await + .unwrap(); + assert_eq!( + filebrowser_read_path(dir.path(), &inside).await.unwrap(), + inside + ); + assert_eq!( + fs::metadata(&inside).await.unwrap().permissions().mode() & 0o777, + 0o640 + ); + let outside = dir.path().join("content/files/test.bin"); + symlink(&outside, root.join("escape")).unwrap(); + for path in [outside, root.join("escape"), root.clone()] { + assert!(filebrowser_read_path(dir.path(), &path).await.is_err()); + } + } + + #[test] + fn user_namespace_bytes_use_the_same_range_rules() { + assert!(matches!( + slice_prepared_content( + vec![], + Some(ByteRange { + start: 0, + end: None + }), + "x".into() + ) + .unwrap(), + 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") + ); + } +} diff --git a/core/archipelago/src/fips/dial.rs b/core/archipelago/src/fips/dial.rs index 451a0a54..02838efc 100644 --- a/core/archipelago/src/fips/dial.rs +++ b/core/archipelago/src/fips/dial.rs @@ -46,6 +46,25 @@ fn fips_should_fall_back(status: reqwest::StatusCode) -> bool { status == reqwest::StatusCode::NOT_FOUND || status.is_server_error() } +/// Is this FIPS answer the final one, or should the request go again over +/// Tor? A single-delivery request already reached the peer, so any answer +/// is final: a Tor replay would carry the same (possibly spent) payload. +fn fips_answer_is_final( + pref: crate::settings::transport::TransportPref, + single_delivery: bool, + status: reqwest::StatusCode, +) -> bool { + pref == crate::settings::transport::TransportPref::Fips + || single_delivery + || !fips_should_fall_back(status) +} + +/// May a failed FIPS attempt be sent again? Only a failed connect proves the +/// peer never saw it; a timeout can land after the request was delivered. +fn fips_retryable(single_delivery: bool, e: &reqwest::Error) -> bool { + e.is_connect() || (!single_delivery && e.is_timeout()) +} + /// DNS suffix appended to a peer's bech32 npub. pub const FIPS_DNS_SUFFIX: &str = "fips"; @@ -113,7 +132,21 @@ pub fn client() -> reqwest::Client { /// before the Tor fallback ever gets a chance. The generous `connect_timeout` /// is preserved so a cold hole-punched path still gets time to establish. pub fn client_with_timeout(timeout: Duration) -> reqwest::Client { + client_with_delivery_policy(timeout, false) +} + +fn delivery_redirect_policy(single: bool) -> reqwest::redirect::Policy { + if single { + reqwest::redirect::Policy::none() + } else { + reqwest::redirect::Policy::default() + } +} + +fn client_with_delivery_policy(timeout: Duration, single: bool) -> reqwest::Client { reqwest::Client::builder() + .no_proxy() + .redirect(delivery_redirect_policy(single)) .timeout(timeout) .connect_timeout(Duration::from_secs(8)) .user_agent("archipelago-fips/1") @@ -130,10 +163,18 @@ pub fn client_with_timeout(timeout: Duration) -> reqwest::Client { /// robust". Only connect/timeout errors are retried (a real HTTP response, /// including 4xx/5xx, is returned as-is for the caller to interpret). async fn send_with_retry(rb: reqwest::RequestBuilder) -> Result { + send_with_retry_if(rb, |e| e.is_connect() || e.is_timeout()).await +} + +/// [`send_with_retry`], retrying only on errors `retryable` accepts. +async fn send_with_retry_if( + rb: reqwest::RequestBuilder, + retryable: impl Fn(&reqwest::Error) -> bool, +) -> Result { let retry = rb.try_clone(); match rb.send().await { Ok(resp) => Ok(resp), - Err(e) if (e.is_connect() || e.is_timeout()) && retry.is_some() => { + Err(e) if retryable(&e) && retry.is_some() => { // Brief pause so the hole-punch packets from the first attempt can // traverse before we re-dial onto the warmed path. tokio::time::sleep(Duration::from_millis(600)).await; @@ -350,6 +391,9 @@ pub struct PeerRequest<'a> { /// the per-peer FIPS/Tor badge reflects reality. Opt-in because not /// every caller has a data dir in scope. pub record_data_dir: Option, + /// The request carries something that must reach the peer at most once + /// (a bearer ecash token). See [`PeerRequest::single_delivery`]. + pub single_delivery: bool, } impl<'a> PeerRequest<'a> { @@ -363,9 +407,25 @@ impl<'a> PeerRequest<'a> { fips_timeout: None, service: None, record_data_dir: None, + single_delivery: false, } } + /// 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 + /// buyer is charged and gets a 402 instead of the file (2026-09-29: FIPS + /// answered 404 after the seller redeemed, the Tor retry got 402). + /// + /// With this set, whatever FIPS answers is final, the FIPS retry fires + /// only when the first attempt never connected, and Tor is used only when + /// FIPS could not have delivered the request. An attempt that may have + /// been delivered but timed out is an error, not a fallback. + pub fn single_delivery(mut self) -> Self { + self.single_delivery = true; + self + } + /// Record the transport that serves this request into federation storage /// (matched by this request's onion host). Best-effort, off the hot path. pub fn record_transport(mut self, data_dir: impl Into) -> Self { @@ -442,7 +502,7 @@ impl<'a> PeerRequest<'a> { // Use the FIPS reply unless it's one a Tor retry could // fix (404 path-not-served / 5xx) and we're allowed to // fall back. FIPS-only never falls back. - if pref == TransportPref::Fips || !fips_should_fall_back(resp.status()) { + if fips_answer_is_final(pref, self.single_delivery, resp.status()) { telemetry::record_fips_ok(); self.spawn_record(crate::transport::TransportKind::Fips); return Ok((resp, crate::transport::TransportKind::Fips)); @@ -481,7 +541,7 @@ impl<'a> PeerRequest<'a> { if matches!(pref, TransportPref::Auto | TransportPref::Fips) { match self.try_fips_get().await? { Some(resp) => { - if pref == TransportPref::Fips || !fips_should_fall_back(resp.status()) { + if fips_answer_is_final(pref, self.single_delivery, resp.status()) { telemetry::record_fips_ok(); self.spawn_record(crate::transport::TransportKind::Fips); return Ok((resp, crate::transport::TransportKind::Fips)); @@ -551,13 +611,21 @@ impl<'a> PeerRequest<'a> { } else { budget }; - let c = client_with_timeout(per_attempt); + let c = client_with_delivery_policy(per_attempt, self.single_delivery); let mut rb = c.post(&url).json(body); for (k, v) in &self.headers { rb = rb.header(*k, v); } - match tokio::time::timeout(budget, send_with_retry(rb)).await { + let single = self.single_delivery; + let attempt = send_with_retry_if(rb, |e| fips_retryable(single, e)); + match tokio::time::timeout(budget, attempt).await { Ok(Ok(r)) => Ok(Some(r)), + Ok(Err(e)) if single && !e.is_connect() => Err(anyhow::anyhow!( + "FIPS POST failed after possible delivery; not replaying: {e}" + )), + Err(_) if single => Err(anyhow::anyhow!( + "FIPS POST exceeded its budget after possible delivery; not replaying" + )), Ok(Err(e)) => { telemetry::record_fallback(FallbackReason::ConnectFail); tracing::info!( @@ -612,13 +680,28 @@ impl<'a> PeerRequest<'a> { } else { budget }; - let c = client_with_timeout(per_attempt); + let c = client_with_delivery_policy(per_attempt, self.single_delivery); let mut rb = c.get(&url); for (k, v) in &self.headers { rb = rb.header(*k, v); } - match tokio::time::timeout(budget, send_with_retry(rb)).await { + let single = self.single_delivery; + let attempt = send_with_retry_if(rb, |e| fips_retryable(single, e)); + match tokio::time::timeout(budget, attempt).await { Ok(Ok(r)) => Ok(Some(r)), + // Anything but a failed connect may have reached the peer. + Ok(Err(e)) if single && !e.is_connect() => Err(anyhow::anyhow!( + "FIPS GET {} failed after the request may have been delivered \ + (not retrying over Tor): {}", + self.path, + e + )), + Err(_) if single => Err(anyhow::anyhow!( + "FIPS GET {} exceeded its {:?} budget after the request may have \ + been delivered (not retrying over Tor)", + self.path, + budget + )), Ok(Err(e)) => { telemetry::record_fallback(FallbackReason::ConnectFail); tracing::info!( @@ -676,6 +759,7 @@ impl<'a> PeerRequest<'a> { .context("Invalid Tor SOCKS proxy URL")?; reqwest::Client::builder() .proxy(proxy) + .redirect(delivery_redirect_policy(self.single_delivery)) .timeout(self.timeout) .build() .context("Build Tor HTTP client") @@ -759,4 +843,181 @@ mod tests { let err = decode_response(0xAABB, &r, "x").unwrap_err(); assert!(err.to_string().contains("no AAAA")); } + + #[test] + fn a_single_delivery_answer_is_final_whatever_its_status() { + use crate::settings::transport::TransportPref; + use reqwest::StatusCode; + // Regression (2026-09-29): the seller redeemed a paid download's + // token, answered 404, and the Tor fallback replayed the spent token. + for status in [ + StatusCode::NOT_FOUND, + StatusCode::INTERNAL_SERVER_ERROR, + StatusCode::SERVICE_UNAVAILABLE, + StatusCode::OK, + ] { + assert!(fips_answer_is_final(TransportPref::Auto, true, status)); + } + // Everything else keeps the existing fallback rules. + assert!(!fips_answer_is_final( + TransportPref::Auto, + false, + StatusCode::NOT_FOUND + )); + assert!(!fips_answer_is_final( + TransportPref::Auto, + false, + StatusCode::BAD_GATEWAY + )); + assert!(fips_answer_is_final( + TransportPref::Auto, + false, + StatusCode::PAYMENT_REQUIRED + )); + assert!(fips_answer_is_final( + TransportPref::Fips, + false, + StatusCode::NOT_FOUND + )); + } + + /// A listener that accepts connections and never answers, counting them. + async fn silent_peer() -> (String, std::sync::Arc) { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let seen = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let counter = seen.clone(); + tokio::spawn(async move { + let mut held = Vec::new(); + while let Ok((stream, _)) = listener.accept().await { + counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + held.push(stream); // keep it open, never reply + } + }); + (format!("http://{addr}/content/x"), seen) + } + + #[tokio::test] + async fn a_single_delivery_request_is_not_resent_after_a_timeout() { + let (url, seen) = silent_peer().await; + let c = client_with_timeout(Duration::from_millis(300)); + let err = send_with_retry_if(c.get(&url), |e| fips_retryable(true, e)) + .await + .expect_err("peer never answers"); + assert!(err.is_timeout()); + assert_eq!(seen.load(std::sync::atomic::Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn an_ordinary_request_is_still_retried_once_after_a_timeout() { + let (url, seen) = silent_peer().await; + let c = client_with_timeout(Duration::from_millis(300)); + let _ = send_with_retry_if(c.get(&url), |e| fips_retryable(false, e)).await; + assert_eq!(seen.load(std::sync::atomic::Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn a_single_delivery_request_still_retries_a_refused_connect() { + // Nothing listening: the peer provably never saw the request. + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + drop(listener); + let c = client_with_timeout(Duration::from_millis(500)); + let err = send_with_retry_if(c.get(format!("http://{addr}/")), |e| { + fips_retryable(true, e) + }) + .await + .expect_err("nothing listening"); + assert!(err.is_connect()); + assert!(fips_retryable(true, &err)); + } +} + +#[cfg(test)] +mod delivery_redirect_tests { + use super::*; + use hyper::{ + service::{make_service_fn, service_fn}, + Body, Response, Server, + }; + use std::{ + convert::Infallible, + sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }, + }; + + #[tokio::test] + async fn paid_bearer_request_does_not_follow_redirects_but_normal_get_does() { + let seen = Arc::new(AtomicUsize::new(0)); + let counter = seen.clone(); + let server = Server::bind(&([127, 0, 0, 1], 0).into()); + let address = server.local_addr(); + let service = make_service_fn(move |_| { + let counter = counter.clone(); + async move { + Ok::<_, Infallible>(service_fn(move |request: hyper::Request| { + let counter = counter.clone(); + async move { + counter.fetch_add(1, Ordering::SeqCst); + let response = if request.uri().path() == "/first" { + Response::builder() + .status(302) + .header("Location", "/replay") + .body(Body::empty()) + .unwrap() + } else { + Response::new(Body::from("replayed")) + }; + Ok::<_, Infallible>(response) + } + })) + } + }); + let task = tokio::spawn(server.serve(service)); + let url = format!("http://{address}/first"); + let response = client_with_delivery_policy(Duration::from_secs(2), true) + .get(&url) + .header("X-Payment-Token", "dummy-test-token") + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::FOUND); + assert_eq!(seen.load(Ordering::SeqCst), 1); + let response = client_with_delivery_policy(Duration::from_secs(2), false) + .get(url) + .send() + .await + .unwrap(); + assert_eq!(response.status(), reqwest::StatusCode::OK); + assert_eq!(seen.load(Ordering::SeqCst), 3); + task.abort(); + } + + #[tokio::test] + async fn paid_request_is_not_resent_when_peer_disconnects_after_reading_it() { + use tokio::io::AsyncReadExt; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let seen = Arc::new(AtomicUsize::new(0)); + let counter = seen.clone(); + let task = tokio::spawn(async move { + while let Ok((mut stream, _)) = listener.accept().await { + let mut buf = [0; 4096]; + let _ = stream.read(&mut buf).await; + counter.fetch_add(1, Ordering::SeqCst); + drop(stream); + } + }); + let c = client_with_delivery_policy(Duration::from_secs(2), true); + let error = send_with_retry_if(c.get(format!("http://{address}/")), |e| { + fips_retryable(true, e) + }) + .await + .unwrap_err(); + assert!(!error.is_connect()); + assert_eq!(seen.load(Ordering::SeqCst), 1); + task.abort(); + } }