diff --git a/crates/buzz-relay/src/operator_listener.rs b/crates/buzz-relay/src/operator_listener.rs index a21292d2aec..8e788b0d0e2 100644 --- a/crates/buzz-relay/src/operator_listener.rs +++ b/crates/buzz-relay/src/operator_listener.rs @@ -2,11 +2,10 @@ use std::{collections::HashMap, sync::Arc, time::Duration}; -use chrono::{DateTime, TimeDelta, Utc}; +use chrono::{TimeDelta, Utc}; use futures_util::future::join_all; use serde::Serialize; use tracing::{error, warn}; -use uuid::Uuid; use crate::{nip98::nip98_header, state::AppState}; @@ -89,93 +88,22 @@ fn build_delivery_http_client(timeout: Duration) -> Result WorkerIteration { - let transport = ReqwestDeliveryTransport { http }; - run_delivery_once_with_transport(state, &transport).await + run_delivery_once_with_db( + &state.db, + &state.config.operator_listener_delivery_urls, + &state.relay_keypair, + http, + ) + .await } -#[async_trait::async_trait] -trait DeliveryTransport: Sync { - async fn post( - &self, - url: &url::Url, - authorization: &str, - body: Vec, - ) -> Result; -} - -#[async_trait::async_trait] -trait DeliveryStore: Sync { - async fn release(&self, id: Uuid, claim_id: Uuid, next: DateTime) -> Result; - async fn complete(&self, id: Uuid, claim_id: Uuid) -> Result; - async fn retry(&self, id: Uuid, claim_id: Uuid, next: DateTime) -> Result; - async fn fail(&self, id: Uuid, claim_id: Uuid) -> Result; -} - -struct DbDeliveryStore<'a> { - db: &'a buzz_db::Db, -} - -#[async_trait::async_trait] -impl DeliveryStore for DbDeliveryStore<'_> { - async fn release(&self, id: Uuid, claim_id: Uuid, next: DateTime) -> Result { - self.db - .release_operator_listener_delivery(id, claim_id, next) - .await - .map_err(|error| error.to_string()) - } - - async fn complete(&self, id: Uuid, claim_id: Uuid) -> Result { - self.db - .complete_operator_listener_delivery(id, claim_id) - .await - .map_err(|error| error.to_string()) - } - - async fn retry(&self, id: Uuid, claim_id: Uuid, next: DateTime) -> Result { - self.db - .retry_operator_listener_delivery(id, claim_id, next) - .await - .map_err(|error| error.to_string()) - } - - async fn fail(&self, id: Uuid, claim_id: Uuid) -> Result { - self.db - .fail_operator_listener_delivery(id, claim_id) - .await - .map_err(|error| error.to_string()) - } -} - -struct ReqwestDeliveryTransport<'a> { - http: &'a reqwest::Client, -} - -#[async_trait::async_trait] -impl DeliveryTransport for ReqwestDeliveryTransport<'_> { - async fn post( - &self, - url: &url::Url, - authorization: &str, - body: Vec, - ) -> Result { - self.http - .post(url.clone()) - .header("Authorization", authorization) - .header("Content-Type", "application/json") - .body(body) - .send() - .await - .map(|response| response.status()) - .map_err(|error| error.to_string()) - } -} - -async fn run_delivery_once_with_transport( - state: &AppState, - transport: &T, +async fn run_delivery_once_with_db( + db: &buzz_db::Db, + routes: &HashMap, + relay_keypair: &nostr::Keys, + http: &reqwest::Client, ) -> WorkerIteration { - let claimed = match state - .db + let claimed = match db .claim_operator_listener_deliveries( DELIVERY_BATCH_LIMIT, Utc::now() + TimeDelta::seconds(CLAIM_SECS), @@ -191,25 +119,20 @@ async fn run_delivery_once_with_transport( if claimed.is_empty() { return WorkerIteration::Idle; } - let store = DbDeliveryStore { db: &state.db }; - join_all(claimed.into_iter().map(|delivery| { - deliver_one( - &state.config.operator_listener_delivery_urls, - &state.relay_keypair, - &store, - transport, - delivery, - ) - })) + join_all( + claimed + .into_iter() + .map(|delivery| deliver_one(db, routes, relay_keypair, http, delivery)), + ) .await; WorkerIteration::Worked } -async fn deliver_one( +async fn deliver_one( + db: &buzz_db::Db, routes: &HashMap, relay_keypair: &nostr::Keys, - store: &S, - transport: &T, + http: &reqwest::Client, delivery: buzz_db::operator_listener::ClaimedDelivery, ) { let listener_hex = hex::encode(&delivery.listener_pubkey); @@ -219,8 +142,8 @@ async fn deliver_one( listener=%listener_hex, "operator-listener delivery has no route on this pod; releasing claim" ); - if let Err(error) = store - .release( + if let Err(error) = db + .release_operator_listener_delivery( delivery.id, delivery.claim_id, Utc::now() + TimeDelta::seconds(CLAIM_SECS), @@ -244,7 +167,7 @@ async fn deliver_one( Ok(body) => body, Err(error) => { fail_permanently( - store, + db, &delivery, &format!("notification encoding failed: {error}"), ) @@ -255,19 +178,24 @@ async fn deliver_one( let auth = match nip98_header(relay_keypair, url.as_str(), &body) { Ok(auth) => auth, Err(error) => { - fail_permanently( - store, - &delivery, - &format!("notification auth failed: {error}"), - ) - .await; + fail_permanently(db, &delivery, &format!("notification auth failed: {error}")).await; return; } }; - let response = transport.post(url, &auth, body).await; + let response = http + .post(url.clone()) + .header("Authorization", auth) + .header("Content-Type", "application/json") + .body(body) + .send() + .await + .map(|response| response.status()); match response { Ok(status) if status.is_success() => { - match store.complete(delivery.id, delivery.claim_id).await { + match db + .complete_operator_listener_delivery(delivery.id, delivery.claim_id) + .await + { Ok(true) => {} Ok(false) => warn!( delivery=%delivery.id, @@ -287,35 +215,38 @@ async fn deliver_one( || status == StatusCode::TOO_MANY_REQUESTS || status.is_server_error() => { - retry_or_fail(store, &delivery, format!("HTTP {status}")).await + retry_or_fail(db, &delivery, format!("HTTP {status}")).await } Ok(status) => { - fail_permanently(store, &delivery, &format!("HTTP {status}")).await; + fail_permanently(db, &delivery, &format!("HTTP {status}")).await; } - Err(error) => retry_or_fail(store, &delivery, error.to_string()).await, + Err(error) => retry_or_fail(db, &delivery, error.to_string()).await, } } -async fn fail_permanently( - store: &S, +async fn fail_permanently( + db: &buzz_db::Db, delivery: &buzz_db::operator_listener::ClaimedDelivery, reason: &str, ) { error!(delivery=%delivery.id, %reason, "operator-listener delivery failed permanently"); - if let Err(error) = store.fail(delivery.id, delivery.claim_id).await { + if let Err(error) = db + .fail_operator_listener_delivery(delivery.id, delivery.claim_id) + .await + { error!(delivery=%delivery.id, %error, "failed to delete terminal operator-listener delivery"); } metrics::counter!("buzz_operator_listener_deliveries_total", "outcome" => "failed") .increment(1); } -async fn retry_or_fail( - store: &S, +async fn retry_or_fail( + db: &buzz_db::Db, delivery: &buzz_db::operator_listener::ClaimedDelivery, reason: String, ) { if delivery.attempt >= buzz_db::operator_listener::MAX_DELIVERY_ATTEMPTS { - fail_permanently(store, delivery, &format!("retries exhausted: {reason}")).await; + fail_permanently(db, delivery, &format!("retries exhausted: {reason}")).await; return; } let delay = 2_i64.pow((delivery.attempt - 1).clamp(0, 7) as u32); @@ -326,8 +257,8 @@ async fn retry_or_fail( %reason, "operator-listener delivery failed; retrying" ); - if let Err(error) = store - .retry( + if let Err(error) = db + .retry_operator_listener_delivery( delivery.id, delivery.claim_id, Utc::now() + TimeDelta::seconds(delay), @@ -341,16 +272,23 @@ async fn retry_or_fail( #[cfg(test)] mod tests { - use std::sync::Mutex; - - use tokio::{ - io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader}, - net::{TcpListener, TcpStream}, + use axum::{ + body::Bytes, + extract::{OriginalUri, State}, + http::{HeaderMap, StatusCode as AxumStatusCode}, + response::IntoResponse, + routing::post, + Router, }; - - use super::*; use buzz_core::tenant::CommunityId; + use chrono::DateTime; use nostr::Keys; + use serde_json::Value; + use std::sync::Arc; + use tokio::sync::Mutex; + use uuid::Uuid; + + use super::*; #[test] fn reapers_use_their_expected_intervals() { @@ -361,416 +299,297 @@ mod tests { ); } - #[derive(Debug, Clone, PartialEq, Eq)] - enum StoreCall { - Release(Uuid, Uuid, DateTime), - Complete(Uuid, Uuid), - Retry(Uuid, Uuid, DateTime), - Fail(Uuid, Uuid), - } + // The worker claims from the deployment-global outbox, so these tests + // must not claim each other's fixture rows in parallel. + static DB_TEST_LOCK: Mutex<()> = Mutex::const_new(()); - #[derive(Default)] - struct MockStore { - calls: Mutex>, - } + mod postgres_tests { + use super::*; - impl MockStore { - fn calls(&self) -> Vec { - self.calls.lock().expect("store calls lock").clone() + struct Listener { + status: AxumStatusCode, + redirect: Option, + received: Mutex>, } - } - #[async_trait::async_trait] - impl DeliveryStore for MockStore { - async fn release( - &self, - id: Uuid, - claim_id: Uuid, - next: DateTime, - ) -> Result { - self.calls + async fn receive( + State(listener): State>, + OriginalUri(uri): OriginalUri, + headers: HeaderMap, + body: Bytes, + ) -> impl IntoResponse { + let url: url::Url = format!("http://{}{}", headers["host"].to_str().unwrap(), uri) + .parse() + .unwrap(); + nostr::nips::nip98::verify_auth_header( + headers["authorization"].to_str().unwrap(), + &url, + nostr::nips::nip98::HttpMethod::POST, + nostr::Timestamp::now(), + Some(&body), + ) + .unwrap(); + assert_eq!(headers["content-type"], "application/json"); + listener + .received .lock() - .expect("store calls lock") - .push(StoreCall::Release(id, claim_id, next)); - Ok(true) + .await + .push(serde_json::from_slice(&body).unwrap()); + let mut response = listener.status.into_response(); + if let Some(location) = &listener.redirect { + response + .headers_mut() + .insert(axum::http::header::LOCATION, location.parse().unwrap()); + } + response } - async fn complete(&self, id: Uuid, claim_id: Uuid) -> Result { - self.calls - .lock() - .expect("store calls lock") - .push(StoreCall::Complete(id, claim_id)); - Ok(true) + async fn start_listener( + status: AxumStatusCode, + redirect: Option, + ) -> (url::Url, Arc, tokio::task::JoinHandle<()>) { + let listener = Arc::new(Listener { + status, + redirect, + received: Mutex::new(Vec::new()), + }); + let socket = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind listener"); + let url = format!("http://{}/mentions", socket.local_addr().unwrap()) + .parse() + .unwrap(); + let router = Router::new() + .route("/mentions", post(receive)) + .with_state(listener.clone()); + let server = tokio::spawn(async move { + axum::serve(socket, router).await.expect("serve listener"); + }); + (url, listener, server) } - async fn retry( - &self, + struct Fixture { + pool: sqlx::PgPool, + db: buzz_db::Db, + community: CommunityId, id: Uuid, - claim_id: Uuid, - next: DateTime, - ) -> Result { - self.calls - .lock() - .expect("store calls lock") - .push(StoreCall::Retry(id, claim_id, next)); - Ok(true) + listener: Keys, + relay: Keys, + event_id: Vec, + target_pubkey: Vec, } - async fn fail(&self, id: Uuid, claim_id: Uuid) -> Result { - self.calls - .lock() - .expect("store calls lock") - .push(StoreCall::Fail(id, claim_id)); - Ok(true) - } - } - - struct SentRequest { - url: url::Url, - authorization: String, - body: Vec, - } + impl Fixture { + async fn new(attempt: i32) -> Self { + let pool = sqlx::PgPool::connect(&crate::test_support::database_url()) + .await + .expect("connect to test DB"); + let db = buzz_db::Db::from_pool(pool.clone()); + let community = CommunityId::from_uuid(Uuid::new_v4()); + let host = format!("operator-delivery-{}.example", community.as_uuid().simple()); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") + .bind(community.as_uuid()) + .bind(host) + .execute(&pool) + .await + .expect("insert community"); + let listener = Keys::generate(); + let target_pubkey = Keys::generate().public_key().to_bytes().to_vec(); + let event_id = vec![0x22; 32]; + let id: Uuid = sqlx::query_scalar( + "INSERT INTO operator_listener_outbox \ + (listener_pubkey, target_pubkey, community_id, event_id, event_kind, event_created_at, attempts) \ + VALUES ($1, $2, $3, $4, 9, to_timestamp(1700000000), $5) RETURNING id", + ) + .bind(listener.public_key().as_bytes().as_slice()) + .bind(&target_pubkey) + .bind(community.as_uuid()) + .bind(&event_id) + .bind(attempt - 1) + .fetch_one(&pool) + .await + .expect("insert delivery"); + Self { + pool, + db, + community, + id, + listener, + relay: Keys::generate(), + event_id, + target_pubkey, + } + } - struct MockTransport { - response: Mutex>>, - request: Mutex>, - request_count: Mutex, - } + fn routes(&self, url: url::Url) -> HashMap { + HashMap::from([(self.listener.public_key().to_hex(), url)]) + } - impl MockTransport { - fn new(response: Result) -> Self { - Self { - response: Mutex::new(Some(response)), - request: Mutex::new(None), - request_count: Mutex::new(0), + async fn deliver(&self, routes: &HashMap) { + let http = build_delivery_http_client(Duration::from_secs(2)).unwrap(); + assert_eq!( + run_delivery_once_with_db(&self.db, routes, &self.relay, &http).await, + WorkerIteration::Worked + ); } - } - fn take_request(&self) -> Option { - self.request.lock().expect("transport request lock").take() - } + async fn row( + &self, + ) -> Option<( + String, + i32, + Option, + Option>, + DateTime, + )> { + use sqlx::Row; + sqlx::query( + "SELECT state, attempts, claim_id, lease_until, next_attempt_at \ + FROM operator_listener_outbox WHERE id = $1", + ) + .bind(self.id) + .fetch_optional(&self.pool) + .await + .expect("read outbox") + .map(|row| { + ( + row.get("state"), + row.get("attempts"), + row.get("claim_id"), + row.get("lease_until"), + row.get("next_attempt_at"), + ) + }) + } - fn request_count(&self) -> usize { - *self - .request_count - .lock() - .expect("transport request count lock") + async fn cleanup(self) { + sqlx::query("DELETE FROM operator_listener_outbox WHERE id = $1") + .bind(self.id) + .execute(&self.pool) + .await + .expect("delete delivery"); + sqlx::query("DELETE FROM communities WHERE id = $1") + .bind(self.community.as_uuid()) + .execute(&self.pool) + .await + .expect("delete community"); + } } - } - #[async_trait::async_trait] - impl DeliveryTransport for MockTransport { - async fn post( - &self, - url: &url::Url, - authorization: &str, - body: Vec, - ) -> Result { - *self - .request_count - .lock() - .expect("transport request count lock") += 1; - *self.request.lock().expect("transport request lock") = Some(SentRequest { - url: url.clone(), - authorization: authorization.to_owned(), - body, - }); - self.response - .lock() - .expect("transport response lock") - .take() - .expect("delivery should make one request") + #[tokio::test] + #[ignore = "requires Postgres"] + async fn successful_delivery_posts_signed_payload_and_completes_outbox() { + let _guard = DB_TEST_LOCK.lock().await; + let fixture = Fixture::new(1).await; + let (url, listener, server) = start_listener(AxumStatusCode::NO_CONTENT, None).await; + fixture.deliver(&fixture.routes(url)).await; + let received = listener.received.lock().await; + assert_eq!(received.len(), 1); + let body = &received[0]; + assert_eq!(body["v"], 1); + assert_eq!(body["pubkey"], hex::encode(&fixture.target_pubkey)); + assert_eq!( + body["community_host"], + format!( + "operator-delivery-{}.example", + fixture.community.as_uuid().simple() + ) + ); + assert_eq!(body["event_id"], hex::encode(&fixture.event_id)); + assert_eq!(body["event_kind"], 9); + assert_eq!(body["event_created_at"], 1_700_000_000); + drop(received); + assert!( + fixture.row().await.is_none(), + "success must remove the outbox row" + ); + server.abort(); + fixture.cleanup().await; } - } - - fn delivery_fixture( - attempt: i32, - ) -> ( - HashMap, - Keys, - buzz_db::operator_listener::ClaimedDelivery, - ) { - let listener = Keys::generate(); - let listener_pubkey = listener.public_key().to_bytes().to_vec(); - let route = url::Url::parse("https://listener.example/mentions").expect("test URL"); - let routes = HashMap::from([(hex::encode(&listener_pubkey), route)]); - ( - routes, - Keys::generate(), - buzz_db::operator_listener::ClaimedDelivery { - id: Uuid::new_v4(), - claim_id: Uuid::new_v4(), - listener_pubkey, - target_pubkey: vec![0x11; 32], - community: CommunityId::from_uuid(Uuid::new_v4()), - community_host: "community.example".to_owned(), - event_id: vec![0x22; 32], - event_kind: 9, - event_created_at: DateTime::from_timestamp(1_700_000_000, 0) - .expect("test timestamp"), - attempt, - }, - ) - } - #[tokio::test] - async fn unroutable_delivery_releases_claim_without_http_request() { - let (_, relay_keypair, delivery) = delivery_fixture(1); - let store = MockStore::default(); - let transport = MockTransport::new(Ok(StatusCode::NO_CONTENT)); - let before = Utc::now() + TimeDelta::seconds(CLAIM_SECS); - - deliver_one( - &HashMap::new(), - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; - - let after = Utc::now() + TimeDelta::seconds(CLAIM_SECS); - assert!(transport.take_request().is_none()); - let calls = store.calls(); - assert_eq!(calls.len(), 1); - match &calls[0] { - StoreCall::Release(id, claim_id, next) => { - assert_eq!(*id, delivery.id); - assert_eq!(*claim_id, delivery.claim_id); - assert!((before..=after).contains(next)); + #[tokio::test] + #[ignore = "requires Postgres"] + async fn transient_responses_retry_with_persisted_delay() { + let _guard = DB_TEST_LOCK.lock().await; + for status in [ + AxumStatusCode::TOO_MANY_REQUESTS, + AxumStatusCode::SERVICE_UNAVAILABLE, + ] { + let fixture = Fixture::new(3).await; + let (url, listener, server) = start_listener(status, None).await; + let before = Utc::now() + TimeDelta::seconds(4); + fixture.deliver(&fixture.routes(url)).await; + let after = Utc::now() + TimeDelta::seconds(4); + assert_eq!(listener.received.lock().await.len(), 1, "HTTP {status}"); + let (state, attempts, claim, lease, next) = fixture.row().await.expect("retry row"); + assert_eq!(state, "pending"); + assert_eq!(attempts, 3); + assert!(claim.is_none() && lease.is_none()); + assert!((before..=after).contains(&next)); + server.abort(); + fixture.cleanup().await; } - call => panic!("expected claim release, got {call:?}"), } - } - #[tokio::test] - async fn successful_delivery_posts_payload_and_completes_claim() { - let (routes, relay_keypair, delivery) = delivery_fixture(1); - let store = MockStore::default(); - let transport = MockTransport::new(Ok(StatusCode::NO_CONTENT)); - - deliver_one( - &routes, - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; - - let request = transport.take_request().expect("delivery request"); - assert_eq!(request.url, routes[&hex::encode(&delivery.listener_pubkey)]); - assert!(request.authorization.starts_with("Nostr ")); - let body: serde_json::Value = serde_json::from_slice(&request.body).expect("JSON body"); - assert_eq!(body["v"], 1); - assert_eq!(body["pubkey"], hex::encode(&delivery.target_pubkey)); - assert_eq!(body["community_host"], delivery.community_host); - assert_eq!(body["event_id"], hex::encode(&delivery.event_id)); - assert_eq!(body["event_kind"], delivery.event_kind); - assert_eq!( - body["event_created_at"], - delivery.event_created_at.timestamp() - ); - assert_eq!( - store.calls(), - [StoreCall::Complete(delivery.id, delivery.claim_id)] - ); - } - - #[tokio::test] - async fn rate_limit_and_server_errors_retry_with_exponential_delay() { - for status in [ - StatusCode::TOO_MANY_REQUESTS, - StatusCode::SERVICE_UNAVAILABLE, - ] { - let (routes, relay_keypair, delivery) = delivery_fixture(3); - let store = MockStore::default(); - let transport = MockTransport::new(Ok(status)); - let before = Utc::now() + TimeDelta::seconds(4); - - deliver_one( - &routes, - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; - - let after = Utc::now() + TimeDelta::seconds(4); - assert_eq!(transport.request_count(), 1, "HTTP {status}"); - assert!(transport.take_request().is_some()); - let calls = store.calls(); - assert_eq!(calls.len(), 1, "HTTP {status}"); - match &calls[0] { - StoreCall::Retry(id, claim_id, next) => { - assert_eq!(*id, delivery.id); - assert_eq!(*claim_id, delivery.claim_id); - assert!((before..=after).contains(next)); - } - call => panic!("expected HTTP {status} retry, got {call:?}"), - } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn client_error_fails_and_removes_outbox_row() { + let _guard = DB_TEST_LOCK.lock().await; + let fixture = Fixture::new(1).await; + let (url, listener, server) = start_listener(AxumStatusCode::BAD_REQUEST, None).await; + fixture.deliver(&fixture.routes(url)).await; + assert_eq!(listener.received.lock().await.len(), 1); + assert!(fixture.row().await.is_none()); + server.abort(); + fixture.cleanup().await; } - } - - #[tokio::test] - async fn deterministic_client_error_fails_without_retrying() { - let (routes, relay_keypair, delivery) = delivery_fixture(1); - let store = MockStore::default(); - let transport = MockTransport::new(Ok(StatusCode::BAD_REQUEST)); - - deliver_one( - &routes, - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; - assert_eq!(transport.request_count(), 1); - assert!(transport.take_request().is_some()); - assert_eq!( - store.calls(), - [StoreCall::Fail(delivery.id, delivery.claim_id)] - ); - } - - async fn read_http_request(stream: &mut TcpStream) -> (String, Vec) { - let mut reader = BufReader::new(stream); - let mut request_line = String::new(); - reader - .read_line(&mut request_line) - .await - .expect("read request line"); - - let mut content_length = 0; - loop { - let mut line = String::new(); - reader.read_line(&mut line).await.expect("read header"); - if line == "\r\n" || line.is_empty() { - break; - } - if let Some((name, value)) = line.split_once(':') { - if name.eq_ignore_ascii_case("content-length") { - content_length = value.trim().parse().expect("valid content length"); - } - } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn redirect_is_terminal_and_target_is_not_contacted() { + let _guard = DB_TEST_LOCK.lock().await; + let (target_url, target, target_server) = + start_listener(AxumStatusCode::OK, None).await; + let (source_url, source, source_server) = + start_listener(AxumStatusCode::FOUND, Some(target_url.to_string())).await; + let fixture = Fixture::new(1).await; + fixture.deliver(&fixture.routes(source_url)).await; + assert_eq!(source.received.lock().await.len(), 1); + assert!(target.received.lock().await.is_empty()); + assert!(fixture.row().await.is_none()); + source_server.abort(); + target_server.abort(); + fixture.cleanup().await; } - let mut body = vec![0; content_length]; - reader - .read_exact(&mut body) - .await - .expect("read request body"); - (request_line, body) - } - - #[tokio::test] - async fn redirect_response_fails_delivery_without_retrying() { - let source = TcpListener::bind("127.0.0.1:0") - .await - .expect("source listener"); - let source_addr = source.local_addr().expect("source address"); - let target = TcpListener::bind("127.0.0.1:0") - .await - .expect("target listener"); - let target_addr = target.local_addr().expect("target address"); - let target_url = format!("http://{target_addr}/redirected"); - let target_task = tokio::spawn(async move { - match tokio::time::timeout(Duration::from_secs(1), target.accept()).await { - Ok(Ok((mut stream, _))) => { - let request = read_http_request(&mut stream).await; - stream - .write_all( - b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", - ) - .await - .expect("write target success response"); - Some(request) - } - Ok(Err(error)) => panic!("accept target request: {error}"), - Err(_) => None, - } - }); - let (_, relay_keypair, delivery) = delivery_fixture(1); - let source_url = - url::Url::parse(&format!("http://{source_addr}/mentions")).expect("source URL"); - let routes = HashMap::from([(hex::encode(&delivery.listener_pubkey), source_url)]); - let store = MockStore::default(); - let http = build_delivery_http_client(Duration::from_secs(2)).expect("delivery client"); - let transport = ReqwestDeliveryTransport { http: &http }; - - let source_task = tokio::spawn(async move { - let (mut stream, _) = source.accept().await.expect("accept source request"); - let request = read_http_request(&mut stream).await; - let response = format!( - "HTTP/1.1 302 Found\r\nLocation: {target_url}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n" - ); - stream - .write_all(response.as_bytes()) - .await - .expect("write redirect response"); - let second_request = tokio::time::timeout(Duration::from_millis(250), source.accept()) - .await - .is_ok(); - (request, second_request) - }); - - deliver_one( - &routes, - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; - - let (request, second_source_request) = - tokio::time::timeout(Duration::from_secs(2), source_task) - .await - .expect("source server completed") - .expect("source server task"); - assert!(request.0.starts_with("POST /mentions HTTP/1.1")); - assert!(!request.1.is_empty()); - assert!( - !second_source_request, - "source received more than one request" - ); - let target_request = tokio::time::timeout(Duration::from_secs(2), target_task) - .await - .expect("target server completed") - .expect("target server task"); - assert_eq!( - store.calls(), - [StoreCall::Fail(delivery.id, delivery.claim_id)] - ); - assert!( - target_request.is_none(), - "redirect target unexpectedly received a request" - ); - } - - #[tokio::test] - async fn exhausted_delivery_attempts_fail_permanently() { - let (routes, relay_keypair, delivery) = - delivery_fixture(buzz_db::operator_listener::MAX_DELIVERY_ATTEMPTS); - let store = MockStore::default(); - let transport = MockTransport::new(Err("connection reset".to_owned())); - - deliver_one( - &routes, - &relay_keypair, - &store, - &transport, - delivery.clone(), - ) - .await; + #[tokio::test] + #[ignore = "requires Postgres"] + async fn exhausted_attempt_is_terminal() { + let _guard = DB_TEST_LOCK.lock().await; + let fixture = Fixture::new(buzz_db::operator_listener::MAX_DELIVERY_ATTEMPTS).await; + let (url, listener, server) = + start_listener(AxumStatusCode::SERVICE_UNAVAILABLE, None).await; + fixture.deliver(&fixture.routes(url)).await; + assert_eq!(listener.received.lock().await.len(), 1); + assert!(fixture.row().await.is_none()); + server.abort(); + fixture.cleanup().await; + } - assert!(transport.take_request().is_some()); - assert_eq!( - store.calls(), - [StoreCall::Fail(delivery.id, delivery.claim_id)] - ); + #[tokio::test] + #[ignore = "requires Postgres"] + async fn unroutable_delivery_releases_claim_without_http() { + let _guard = DB_TEST_LOCK.lock().await; + let fixture = Fixture::new(1).await; + let before = Utc::now() + TimeDelta::seconds(CLAIM_SECS); + fixture.deliver(&HashMap::new()).await; + let after = Utc::now() + TimeDelta::seconds(CLAIM_SECS); + let (state, attempts, claim, lease, next) = fixture.row().await.expect("released row"); + assert_eq!(state, "pending"); + assert_eq!(attempts, 0); + assert!(claim.is_none() && lease.is_none()); + assert!((before..=after).contains(&next)); + fixture.cleanup().await; + } } }