Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
03e38d1ca3 | ||
|
|
e5fc99d66c | ||
|
|
8b74803290 |
@@ -162,11 +162,28 @@ 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."}"#,
|
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. No payment was taken."}"#,
|
||||||
|
),
|
||||||
|
)),
|
||||||
|
Ok(content_server::ServeResult::NotFound) => Ok(build_response(
|
||||||
StatusCode::NOT_FOUND,
|
StatusCode::NOT_FOUND,
|
||||||
"text/plain",
|
"text/plain",
|
||||||
hyper::Body::from("Content not found"),
|
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"),
|
||||||
|
))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -555,6 +555,9 @@ impl RpcHandler {
|
|||||||
.service(crate::settings::transport::PeerService::PeerFiles)
|
.service(crate::settings::transport::PeerService::PeerFiles)
|
||||||
.header("X-Federation-DID", local_did)
|
.header("X-Federation-DID", local_did)
|
||||||
.header("X-Payment-Token", token_str.clone())
|
.header("X-Payment-Token", token_str.clone())
|
||||||
|
// The token is a bearer instrument the seller redeems on first sight:
|
||||||
|
// a Tor replay after FIPS delivered it can only arrive spent.
|
||||||
|
.single_delivery()
|
||||||
.timeout(std::time::Duration::from_secs(900))
|
.timeout(std::time::Duration::from_secs(900))
|
||||||
.send_get()
|
.send_get()
|
||||||
.await
|
.await
|
||||||
@@ -610,8 +613,12 @@ impl RpcHandler {
|
|||||||
let body = response.text().await.unwrap_or_default();
|
let body = response.text().await.unwrap_or_default();
|
||||||
tracing::warn!("paid download: seller {onion} returned {status}: {body}");
|
tracing::warn!("paid download: seller {onion} returned {status}: {body}");
|
||||||
reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await;
|
reclaim_spent_ecash(&self.config.data_dir, &token_str, used_backend).await;
|
||||||
|
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_string))
|
||||||
|
.unwrap_or_else(|| format!("Peer returned an error ({status})."));
|
||||||
return Ok(serde_json::json!({
|
return Ok(serde_json::json!({
|
||||||
"error": format!("Peer returned an error ({status}). Your ecash was refunded to your wallet.")
|
"error": format!("{reason} Your ecash was refunded to your wallet.")
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -7,7 +7,7 @@ use anyhow::{Context, Result};
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
use tokio::fs;
|
use tokio::fs;
|
||||||
use tracing::{debug, warn};
|
use tracing::{debug, info, warn};
|
||||||
|
|
||||||
const CATALOG_FILE: &str = "content/catalog.json";
|
const CATALOG_FILE: &str = "content/catalog.json";
|
||||||
const CONTENT_DIR: &str = "content/files";
|
const CONTENT_DIR: &str = "content/files";
|
||||||
@@ -238,6 +238,9 @@ pub enum ServeResult {
|
|||||||
Forbidden,
|
Forbidden,
|
||||||
/// Content not found.
|
/// Content not found.
|
||||||
NotFound,
|
NotFound,
|
||||||
|
/// The catalog entry and file exist but this node can't read the file.
|
||||||
|
/// Returned before any payment is taken.
|
||||||
|
Unavailable,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Serve a content item by ID with access control and optional range request.
|
/// Serve a content item by ID with access control and optional range request.
|
||||||
@@ -296,6 +299,37 @@ pub async fn serve_content(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let file_path = content_file_path(data_dir, item);
|
||||||
|
if !file_path.exists() {
|
||||||
|
// The catalog entry survived (it's a separate JSON file) but its
|
||||||
|
// backing file is gone — most likely lost in an unrelated data-dir
|
||||||
|
// reset (a shared filebrowser file, 2026-07-01: two catalog entries
|
||||||
|
// outlived a filebrowser reinstall that wiped the files themselves).
|
||||||
|
// Leaving the entry in place would keep advertising it as available
|
||||||
|
// to every peer forever, each hitting the exact same dead end this
|
||||||
|
// one just did. Prune it so it stops being offered.
|
||||||
|
warn!(
|
||||||
|
content_id = %id,
|
||||||
|
filename = %item.filename,
|
||||||
|
"content catalog entry's file is missing on disk — pruning the stale entry"
|
||||||
|
);
|
||||||
|
prune_missing_content_entry(data_dir, id).await;
|
||||||
|
return Ok(ServeResult::NotFound);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Confirm the file is readable BEFORE the paid gate below redeems the
|
||||||
|
// buyer's token. Reading it only afterwards meant a permission error
|
||||||
|
// surfaced after the sale: the buyer was charged and got an error
|
||||||
|
// instead of the file (2026-09-29, a FileBrowser upload left 0640).
|
||||||
|
if let Err(e) = ensure_readable(&file_path).await {
|
||||||
|
warn!(
|
||||||
|
content_id = %id,
|
||||||
|
path = %file_path.display(),
|
||||||
|
"shared content file is not readable by this node: {e:#}"
|
||||||
|
);
|
||||||
|
return Ok(ServeResult::Unavailable);
|
||||||
|
}
|
||||||
|
|
||||||
// Check access control
|
// Check access control
|
||||||
if !owner_session {
|
if !owner_session {
|
||||||
match &item.access {
|
match &item.access {
|
||||||
@@ -336,23 +370,6 @@ pub async fn serve_content(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
let file_path = content_file_path(data_dir, item);
|
|
||||||
if !file_path.exists() {
|
|
||||||
// The catalog entry survived (it's a separate JSON file) but its
|
|
||||||
// backing file is gone — most likely lost in an unrelated data-dir
|
|
||||||
// reset (a shared filebrowser file, 2026-07-01: two catalog entries
|
|
||||||
// outlived a filebrowser reinstall that wiped the files themselves).
|
|
||||||
// Leaving the entry in place would keep advertising it as available
|
|
||||||
// to every peer forever, each hitting the exact same dead end this
|
|
||||||
// one just did. Prune it so it stops being offered.
|
|
||||||
warn!(
|
|
||||||
content_id = %id,
|
|
||||||
filename = %item.filename,
|
|
||||||
"content catalog entry's file is missing on disk — pruning the stale entry"
|
|
||||||
);
|
|
||||||
prune_missing_content_entry(data_dir, id).await;
|
|
||||||
return Ok(ServeResult::NotFound);
|
|
||||||
}
|
|
||||||
|
|
||||||
let metadata = fs::metadata(&file_path)
|
let metadata = fs::metadata(&file_path)
|
||||||
.await
|
.await
|
||||||
@@ -572,6 +589,62 @@ pub async fn serve_content_preview(data_dir: &Path, id: &str) -> Result<PreviewR
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Make sure this service can open `path`, granting read access if it can't.
|
||||||
|
///
|
||||||
|
/// FileBrowser writes uploads as its container user (a rootless subuid such
|
||||||
|
/// as 100999), and some arrive 0640 — unreadable by this service, although
|
||||||
|
/// most shared files are already 0644. Inside the rootless user namespace
|
||||||
|
/// that subuid is ours, so `podman unshare chmod a+r` grants the same read
|
||||||
|
/// access the other shared files have, without sudo.
|
||||||
|
async fn ensure_readable(path: &Path) -> Result<()> {
|
||||||
|
ensure_readable_with(path, grant_read_access).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn ensure_readable_with<F, Fut>(path: &Path, grant: F) -> Result<()>
|
||||||
|
where
|
||||||
|
F: FnOnce(PathBuf) -> Fut,
|
||||||
|
Fut: std::future::Future<Output = Result<()>>,
|
||||||
|
{
|
||||||
|
match fs::File::open(path).await {
|
||||||
|
Ok(_) => return Ok(()),
|
||||||
|
Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {}
|
||||||
|
Err(e) => return Err(e).context("Failed to open content file"),
|
||||||
|
}
|
||||||
|
grant(path.to_path_buf()).await?;
|
||||||
|
info!("Granted read access to shared content file {}", path.display());
|
||||||
|
fs::File::open(path)
|
||||||
|
.await
|
||||||
|
.context("Content file still unreadable after chmod")?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
// Tests must not shell out to podman: whether it exists (and can chmod a
|
||||||
|
// file the test user owns) would decide the outcome.
|
||||||
|
#[cfg(not(test))]
|
||||||
|
use grant_read_via_podman as grant_read_access;
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
async fn grant_read_access(_path: PathBuf) -> Result<()> {
|
||||||
|
anyhow::bail!("granting read access is disabled in tests")
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg_attr(test, allow(dead_code))]
|
||||||
|
async fn grant_read_via_podman(path: PathBuf) -> Result<()> {
|
||||||
|
let out = tokio::process::Command::new("podman")
|
||||||
|
.args(["unshare", "chmod", "a+r"])
|
||||||
|
.arg(&path)
|
||||||
|
.output()
|
||||||
|
.await
|
||||||
|
.context("Failed to run podman unshare chmod")?;
|
||||||
|
if !out.status.success() {
|
||||||
|
anyhow::bail!(
|
||||||
|
"podman unshare chmod a+r failed: {}",
|
||||||
|
String::from_utf8_lossy(&out.stderr).trim()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// Verify a payment token covers the required amount.
|
/// Verify a payment token covers the required amount.
|
||||||
/// Accepts both cashuA tokens (real Cashu) and legacy cashuSend_ format.
|
/// Accepts both cashuA tokens (real Cashu) and legacy cashuSend_ format.
|
||||||
/// Swaps proofs at the mint to verify they're unspent before accepting.
|
/// Swaps proofs at the mint to verify they're unspent before accepting.
|
||||||
@@ -725,3 +798,137 @@ mod prune_missing_content_tests {
|
|||||||
assert_eq!(reloaded.items[0].id, "present-item");
|
assert_eq!(reloaded.items[0].id, "present-item");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod unreadable_content_tests {
|
||||||
|
use super::*;
|
||||||
|
use std::os::unix::fs::PermissionsExt;
|
||||||
|
|
||||||
|
/// Writes `bytes` to the FileBrowser area and makes it unreadable, the
|
||||||
|
/// way a 0640 upload owned by a container subuid looks to this service.
|
||||||
|
/// `None` when the test runs as root, where mode bits don't stop reads.
|
||||||
|
fn unreadable_file(data_dir: &Path, name: &str) -> Option<PathBuf> {
|
||||||
|
let dir = data_dir.join("filebrowser").join("Music");
|
||||||
|
std::fs::create_dir_all(&dir).unwrap();
|
||||||
|
let path = dir.join(name);
|
||||||
|
std::fs::write(&path, b"audio").unwrap();
|
||||||
|
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o000)).unwrap();
|
||||||
|
std::fs::File::open(&path).is_err().then_some(path)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn paid_item(filename: &str) -> ContentItem {
|
||||||
|
ContentItem {
|
||||||
|
id: "paid-item".to_string(),
|
||||||
|
filename: filename.to_string(),
|
||||||
|
mime_type: "audio/mpeg".to_string(),
|
||||||
|
size_bytes: 5,
|
||||||
|
description: String::new(),
|
||||||
|
access: AccessControl::Paid {
|
||||||
|
price_sats: 10,
|
||||||
|
accepted: vec!["ecash".to_string()],
|
||||||
|
},
|
||||||
|
availability: Availability::AllPeers,
|
||||||
|
added_at: "2026-01-01T00:00:00Z".to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Regression (2026-09-29): the seller redeemed the buyer's token and
|
||||||
|
/// only then failed to read the file, so the buyer paid for nothing.
|
||||||
|
/// An unreadable file must be refused before the payment gate runs,
|
||||||
|
/// which is why a token that would never verify still gets Unavailable
|
||||||
|
/// rather than PaymentRequired.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn an_unreadable_paid_file_is_refused_before_any_payment_is_taken() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let data_dir = dir.path();
|
||||||
|
let Some(_path) = unreadable_file(data_dir, "song.mp3") else {
|
||||||
|
return; // running as root
|
||||||
|
};
|
||||||
|
save_catalog(
|
||||||
|
data_dir,
|
||||||
|
&ContentCatalog {
|
||||||
|
items: vec![paid_item("Music/song.mp3")],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let result = serve_content(
|
||||||
|
data_dir,
|
||||||
|
"paid-item",
|
||||||
|
Some("cashuBnot-a-real-token"),
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(matches!(result, ServeResult::Unavailable));
|
||||||
|
// An unreadable file is not a missing one: keep the catalog entry.
|
||||||
|
assert_eq!(load_catalog(data_dir).await.unwrap().items.len(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_readable_paid_file_still_demands_payment() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let data_dir = dir.path();
|
||||||
|
let music = data_dir.join("filebrowser").join("Music");
|
||||||
|
std::fs::create_dir_all(&music).unwrap();
|
||||||
|
std::fs::write(music.join("song.mp3"), b"audio").unwrap();
|
||||||
|
save_catalog(
|
||||||
|
data_dir,
|
||||||
|
&ContentCatalog {
|
||||||
|
items: vec![paid_item("Music/song.mp3")],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let result = serve_content(data_dir, "paid-item", None, None, None, None, false)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(matches!(result, ServeResult::PaymentRequired(10)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn ensure_readable_grants_access_once_then_reopens() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let Some(path) = unreadable_file(dir.path(), "a.mp3") else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let calls = std::sync::atomic::AtomicUsize::new(0);
|
||||||
|
ensure_readable_with(&path, |p| {
|
||||||
|
calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||||
|
async move {
|
||||||
|
std::fs::set_permissions(&p, std::fs::Permissions::from_mode(0o644))?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn ensure_readable_leaves_a_readable_file_alone() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let path = dir.path().join("ok.mp3");
|
||||||
|
std::fs::write(&path, b"x").unwrap();
|
||||||
|
ensure_readable_with(&path, |_| async { anyhow::bail!("must not grant") })
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn ensure_readable_reports_a_failed_grant() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let Some(path) = unreadable_file(dir.path(), "b.mp3") else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let err = ensure_readable_with(&path, |_| async { anyhow::bail!("no podman") })
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
|
assert!(err.to_string().contains("no podman"));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -46,6 +46,25 @@ fn fips_should_fall_back(status: reqwest::StatusCode) -> bool {
|
|||||||
status == reqwest::StatusCode::NOT_FOUND || status.is_server_error()
|
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.
|
/// DNS suffix appended to a peer's bech32 npub.
|
||||||
pub const FIPS_DNS_SUFFIX: &str = "fips";
|
pub const FIPS_DNS_SUFFIX: &str = "fips";
|
||||||
|
|
||||||
@@ -130,10 +149,18 @@ pub fn client_with_timeout(timeout: Duration) -> reqwest::Client {
|
|||||||
/// robust". Only connect/timeout errors are retried (a real HTTP response,
|
/// robust". Only connect/timeout errors are retried (a real HTTP response,
|
||||||
/// including 4xx/5xx, is returned as-is for the caller to interpret).
|
/// 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> {
|
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();
|
let retry = rb.try_clone();
|
||||||
match rb.send().await {
|
match rb.send().await {
|
||||||
Ok(resp) => Ok(resp),
|
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
|
// Brief pause so the hole-punch packets from the first attempt can
|
||||||
// traverse before we re-dial onto the warmed path.
|
// traverse before we re-dial onto the warmed path.
|
||||||
tokio::time::sleep(Duration::from_millis(600)).await;
|
tokio::time::sleep(Duration::from_millis(600)).await;
|
||||||
@@ -350,6 +377,9 @@ pub struct PeerRequest<'a> {
|
|||||||
/// the per-peer FIPS/Tor badge reflects reality. Opt-in because not
|
/// the per-peer FIPS/Tor badge reflects reality. Opt-in because not
|
||||||
/// every caller has a data dir in scope.
|
/// every caller has a data dir in scope.
|
||||||
pub record_data_dir: Option<std::path::PathBuf>,
|
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> {
|
impl<'a> PeerRequest<'a> {
|
||||||
@@ -363,9 +393,25 @@ impl<'a> PeerRequest<'a> {
|
|||||||
fips_timeout: None,
|
fips_timeout: None,
|
||||||
service: None,
|
service: None,
|
||||||
record_data_dir: 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
|
/// Record the transport that serves this request into federation storage
|
||||||
/// (matched by this request's onion host). Best-effort, off the hot path.
|
/// (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 {
|
pub fn record_transport(mut self, data_dir: impl Into<std::path::PathBuf>) -> Self {
|
||||||
@@ -481,7 +527,7 @@ impl<'a> PeerRequest<'a> {
|
|||||||
if matches!(pref, TransportPref::Auto | TransportPref::Fips) {
|
if matches!(pref, TransportPref::Auto | TransportPref::Fips) {
|
||||||
match self.try_fips_get().await? {
|
match self.try_fips_get().await? {
|
||||||
Some(resp) => {
|
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();
|
telemetry::record_fips_ok();
|
||||||
self.spawn_record(crate::transport::TransportKind::Fips);
|
self.spawn_record(crate::transport::TransportKind::Fips);
|
||||||
return Ok((resp, crate::transport::TransportKind::Fips));
|
return Ok((resp, crate::transport::TransportKind::Fips));
|
||||||
@@ -617,8 +663,23 @@ impl<'a> PeerRequest<'a> {
|
|||||||
for (k, v) in &self.headers {
|
for (k, v) in &self.headers {
|
||||||
rb = rb.header(*k, v);
|
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(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)) => {
|
Ok(Err(e)) => {
|
||||||
telemetry::record_fallback(FallbackReason::ConnectFail);
|
telemetry::record_fallback(FallbackReason::ConnectFail);
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
@@ -759,4 +820,92 @@ mod tests {
|
|||||||
let err = decode_response(0xAABB, &r, "x").unwrap_err();
|
let err = decode_response(0xAABB, &r, "x").unwrap_err();
|
||||||
assert!(err.to_string().contains("no AAAA"));
|
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));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -514,6 +514,13 @@ impl MintClient {
|
|||||||
pub async fn swap(&self, inputs: &[Proof], target_amounts: &[u64]) -> Result<SwapResult> {
|
pub async fn swap(&self, inputs: &[Proof], target_amounts: &[u64]) -> Result<SwapResult> {
|
||||||
let keyset = self.get_active_sat_keyset().await?;
|
let keyset = self.get_active_sat_keyset().await?;
|
||||||
|
|
||||||
|
// cashuB tokens carry NUT-02 v2 keyset ids in their 8-byte short form,
|
||||||
|
// which the mint rejects (bare 422). Repair here, not at each caller:
|
||||||
|
// the paid-download seller path swapped directly and every Minibits
|
||||||
|
// payment failed once the mint rotated to a v2 keyset.
|
||||||
|
let repaired = self.resolve_truncated_keyset_ids(inputs).await;
|
||||||
|
let inputs: &[Proof] = &repaired;
|
||||||
|
|
||||||
// NUT-02: a mint may charge a per-input fee, and it rejects the swap
|
// NUT-02: a mint may charge a per-input fee, and it rejects the swap
|
||||||
// outright unless outputs == inputs - fee (`11005 Transaction inputs
|
// outright unless outputs == inputs - fee (`11005 Transaction inputs
|
||||||
// should equal outputs less fee`). Applied here rather than at each
|
// should equal outputs less fee`). Applied here rather than at each
|
||||||
@@ -813,8 +820,7 @@ impl MintClient {
|
|||||||
let total: u64 = entry.proofs.iter().map(|p| p.amount).sum();
|
let total: u64 = entry.proofs.iter().map(|p| p.amount).sum();
|
||||||
let target_amounts = amount_to_denominations(total);
|
let target_amounts = amount_to_denominations(total);
|
||||||
|
|
||||||
let proofs = self.resolve_truncated_keyset_ids(&entry.proofs).await;
|
let result = self.swap(&entry.proofs, &target_amounts).await?;
|
||||||
let result = self.swap(&proofs, &target_amounts).await?;
|
|
||||||
all_new_proofs.extend(result.new_proofs);
|
all_new_proofs.extend(result.new_proofs);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -863,4 +869,136 @@ mod tests {
|
|||||||
let client = MintClient::new("http://mint.example.com").unwrap();
|
let client = MintClient::new("http://mint.example.com").unwrap();
|
||||||
assert_eq!(client.url(), "http://mint.example.com");
|
assert_eq!(client.url(), "http://mint.example.com");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A minimal mint on 127.0.0.1 answering `/v1/keys` and `/v1/keysets`
|
||||||
|
/// with one v2 keyset, and rejecting every `/v1/swap` as already spent.
|
||||||
|
/// Each swap request body is sent back on the returned channel.
|
||||||
|
async fn stub_mint(
|
||||||
|
full_id: &'static str,
|
||||||
|
) -> (String, tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>) {
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let addr = listener.local_addr().unwrap();
|
||||||
|
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
loop {
|
||||||
|
let Ok((mut stream, _)) = listener.accept().await else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let mut buf = Vec::new();
|
||||||
|
let mut chunk = [0u8; 4096];
|
||||||
|
// Read headers, then as much body as Content-Length says.
|
||||||
|
let (head, body) = loop {
|
||||||
|
let n = stream.read(&mut chunk).await.unwrap_or(0);
|
||||||
|
if n == 0 {
|
||||||
|
break (String::new(), Vec::new());
|
||||||
|
}
|
||||||
|
buf.extend_from_slice(&chunk[..n]);
|
||||||
|
let Some(pos) = buf.windows(4).position(|w| w == b"\r\n\r\n") else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
let head = String::from_utf8_lossy(&buf[..pos]).to_string();
|
||||||
|
let len = head
|
||||||
|
.lines()
|
||||||
|
.find_map(|l| {
|
||||||
|
let (k, v) = l.split_once(':')?;
|
||||||
|
k.eq_ignore_ascii_case("content-length")
|
||||||
|
.then(|| v.trim().parse::<usize>().ok())?
|
||||||
|
})
|
||||||
|
.unwrap_or(0);
|
||||||
|
while buf.len() < pos + 4 + len {
|
||||||
|
let n = stream.read(&mut chunk).await.unwrap_or(0);
|
||||||
|
if n == 0 {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
buf.extend_from_slice(&chunk[..n]);
|
||||||
|
}
|
||||||
|
break (head, buf[pos + 4..].to_vec());
|
||||||
|
};
|
||||||
|
let request_line = head.lines().next().unwrap_or_default().to_string();
|
||||||
|
let (status, reply) = if request_line.starts_with("GET /v1/keys ") {
|
||||||
|
let key = "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798";
|
||||||
|
let keys: serde_json::Map<String, serde_json::Value> = (0..16)
|
||||||
|
.map(|i| ((1u64 << i).to_string(), serde_json::json!(key)))
|
||||||
|
.collect();
|
||||||
|
(
|
||||||
|
"200 OK",
|
||||||
|
serde_json::json!({"keysets": [
|
||||||
|
{"id": full_id, "unit": "sat", "active": true, "keys": keys}
|
||||||
|
]}),
|
||||||
|
)
|
||||||
|
} else if request_line.starts_with("GET /v1/keysets ") {
|
||||||
|
(
|
||||||
|
"200 OK",
|
||||||
|
serde_json::json!({"keysets": [
|
||||||
|
{"id": full_id, "unit": "sat", "active": true, "input_fee_ppk": 0}
|
||||||
|
]}),
|
||||||
|
)
|
||||||
|
} else if request_line.starts_with("POST /v1/swap ") {
|
||||||
|
let _ = tx.send(serde_json::from_slice(&body).unwrap_or_default());
|
||||||
|
(
|
||||||
|
"400 Bad Request",
|
||||||
|
serde_json::json!({"code": 11001, "detail": "Token Already Spent"}),
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
("404 Not Found", serde_json::json!({}))
|
||||||
|
};
|
||||||
|
let reply = reply.to_string();
|
||||||
|
let _ = stream
|
||||||
|
.write_all(
|
||||||
|
format!(
|
||||||
|
"HTTP/1.1 {status}\r\nContent-Type: application/json\r\n\
|
||||||
|
Content-Length: {}\r\nConnection: close\r\n\r\n{reply}",
|
||||||
|
reply.len()
|
||||||
|
)
|
||||||
|
.as_bytes(),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
(format!("http://{addr}"), rx)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn proof_with_id(id: &str) -> Proof {
|
||||||
|
Proof {
|
||||||
|
amount: 8,
|
||||||
|
id: id.to_string(),
|
||||||
|
secret: "test-secret".to_string(),
|
||||||
|
c: "02".to_string() + &"11".repeat(32),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Regression (2026-09-29): the paid-download seller called `swap`
|
||||||
|
/// directly with a cashuB token's short v2 keyset id and the mint
|
||||||
|
/// answered 422. `swap` itself must send the full id.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn swap_expands_a_short_v2_keyset_id_before_calling_the_mint() {
|
||||||
|
const FULL: &str = "01fc0ec0e59cd6fa01b7a88f8cd77fce81fd1e64bca67d752e984992b7a3c3a821";
|
||||||
|
let (url, mut swaps) = stub_mint(FULL).await;
|
||||||
|
let client = MintClient::new(&url).unwrap();
|
||||||
|
|
||||||
|
let Err(err) = client
|
||||||
|
.swap(&[proof_with_id("01fc0ec0e59cd6fa")], &[8])
|
||||||
|
.await
|
||||||
|
else {
|
||||||
|
panic!("stub mint rejects every swap");
|
||||||
|
};
|
||||||
|
assert!(err.is::<AlreadyRedeemed>(), "unexpected error: {err:#}");
|
||||||
|
|
||||||
|
let body = swaps.recv().await.expect("swap reached the mint");
|
||||||
|
assert_eq!(body["inputs"][0]["id"], FULL);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn swap_passes_complete_keyset_ids_through_unchanged() {
|
||||||
|
const FULL: &str = "01fc0ec0e59cd6fa01b7a88f8cd77fce81fd1e64bca67d752e984992b7a3c3a821";
|
||||||
|
let (url, mut swaps) = stub_mint(FULL).await;
|
||||||
|
let client = MintClient::new(&url).unwrap();
|
||||||
|
|
||||||
|
for id in [FULL, "009a1f293253e41e"] {
|
||||||
|
let _ = client.swap(&[proof_with_id(id)], &[8]).await;
|
||||||
|
let body = swaps.recv().await.expect("swap reached the mint");
|
||||||
|
assert_eq!(body["inputs"][0]["id"], id);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user