Merge current main and harden paid-download delivery

This commit is contained in:
archipelago
2026-09-30 07:25:47 -04:00
56 changed files with 3486 additions and 710 deletions
+116 -4
View File
@@ -132,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")
@@ -488,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));
@@ -597,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!(
@@ -658,7 +680,7 @@ 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);
@@ -737,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")
@@ -909,3 +932,92 @@ mod tests {
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();
}
}