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