Merge remote-tracking branch 'origin/main'
Demo images / Build & push demo images (push) Failing after 1m10s

This commit is contained in:
archipelago
2026-09-30 09:32:20 -04:00
6 changed files with 1276 additions and 266 deletions
+23 -1
View File
@@ -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"),
))
}
}
}
+53 -85
View File
@@ -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::<serde_json::Value>(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!(
+41 -132
View File
@@ -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<Mutex<Vec<(String, String, Vec<u8>)>>>,
task: tokio::task::JoinHandle<()>,
}
impl Drop for FilesApi {
fn drop(&mut self) {
self.task.abort();
}
}
fn files_api(statuses: Vec<u16>) -> 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<Body>| {
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)."
);
}
}
+420 -1
View File
@@ -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<PathBuf> {
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<F, Fut>(
dir: &Path,
name: &str,
bytes: &[u8],
fallback: F,
) -> Result<PathBuf>
where
F: FnOnce(PathBuf, String, Vec<u8>) -> Fut,
Fut: std::future::Future<Output = Result<PathBuf>>,
{
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<F, Fut>(
result: std::io::Result<PathBuf>,
dir: &Path,
name: &str,
bytes: &[u8],
fallback: F,
) -> Result<PathBuf>
where
F: FnOnce(PathBuf, String, Vec<u8>) -> Fut,
Fut: std::future::Future<Output = Result<PathBuf>>,
{
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<PathBuf> {
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<u8>) -> Result<PathBuf> {
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)");
}
}
+471 -40
View File
@@ -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<ByteRange>,
owner_session: bool,
) -> Result<ServeResult> {
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<R, RF, V, VF>(
data_dir: &Path,
id: &str,
payment_token: Option<&str>,
invoice_hash: Option<&str>,
peer_did: Option<&str>,
range: Option<ByteRange>,
owner_session: bool,
read: R,
verify: V,
) -> Result<ServeResult>
where
R: FnOnce(PathBuf, Option<ByteRange>, String) -> RF,
RF: std::future::Future<Output = Result<ServeResult>>,
V: FnOnce(String, u64) -> VF,
VF: std::future::Future<Output = bool>,
{
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<ByteRange>,
mime: String,
) -> Result<ServeResult> {
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<u8>,
range: Option<ByteRange>,
mime: String,
) -> Result<ServeResult> {
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<PathBuf> {
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<Vec<u8>> {
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")
);
}
}
+268 -7
View File
@@ -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<reqwest::Response, reqwest::Error> {
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<reqwest::Response, reqwest::Error> {
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<std::path::PathBuf>,
/// 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<std::path::PathBuf>) -> 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<std::sync::atomic::AtomicUsize>) {
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<Body>| {
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();
}
}