From cebf5253faec60bdf11081a573b6c2fae36ddb7d Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Wed, 17 Jun 2026 13:59:22 -0300 Subject: [PATCH 1/7] Add immediate read to mailroom database --- payjoin-mailroom/src/db/files.rs | 41 ++++++++++++++++++++++++++++++++ payjoin-mailroom/src/db/mod.rs | 33 +++++++++++++++++++++++++ 2 files changed, 74 insertions(+) diff --git a/payjoin-mailroom/src/db/files.rs b/payjoin-mailroom/src/db/files.rs index 6ab27098f..063cf1129 100644 --- a/payjoin-mailroom/src/db/files.rs +++ b/payjoin-mailroom/src/db/files.rs @@ -238,6 +238,14 @@ impl DbTrait for FilesDb { Ok(guard.post_v2(id, payload).await?) } + async fn peek_v2_payload( + &self, + id: &ShortId, + ) -> Result>>, DbError> { + let mut guard = self.mailboxes.lock().await; + Ok(guard.read(id).await?) + } + async fn wait_for_v2_payload( &self, id: &ShortId, @@ -1070,4 +1078,37 @@ mod tests { Ok(()) } + + #[tokio::test(start_paused = true)] + async fn peek_returns_immediately_on_empty_mailbox() { + let dir = tempfile::tempdir().unwrap(); + let db = FilesDb::init( + Duration::from_secs(30), + dir.path().to_owned(), + Duration::from_secs(60 * 60 * 24 * 7), + ) + .await + .unwrap(); + let id = ShortId([0u8; 8]); + let start = tokio::time::Instant::now(); + let got = db.peek_v2_payload(&id).await.expect("peek"); + assert!(got.is_none()); + assert_eq!(start.elapsed(), Duration::ZERO, "peek must not block"); + } + + #[tokio::test] + async fn peek_returns_present_payload() { + let dir = tempfile::tempdir().unwrap(); + let db = FilesDb::init( + Duration::from_millis(10), + dir.path().to_owned(), + Duration::from_secs(60 * 60 * 24 * 7), + ) + .await + .unwrap(); + let id = ShortId([0u8; 8]); + db.post_v2_payload(&id, b"hi".to_vec()).await.unwrap().unwrap(); + let got = db.peek_v2_payload(&id).await.expect("peek").expect("present"); + assert_eq!(&got[..], b"hi"); + } } diff --git a/payjoin-mailroom/src/db/mod.rs b/payjoin-mailroom/src/db/mod.rs index 210bac4e2..73a878536 100644 --- a/payjoin-mailroom/src/db/mod.rs +++ b/payjoin-mailroom/src/db/mod.rs @@ -72,6 +72,12 @@ pub trait Db: Clone + Send + Sync + 'static { mailbox_id: &ShortId, ) -> impl Future>, Error>> + Send; + /// Read a stored v2 payload if present, without waiting. + fn peek_v2_payload( + &self, + mailbox_id: &ShortId, + ) -> impl Future>>, Error>> + Send; + /// Write a v1 response payload. fn post_v1_response( &self, @@ -91,6 +97,7 @@ pub trait Db: Clone + Send + Sync + 'static { pub enum DbRequest { PostV2Payload { mailbox_id: ShortId, payload: Vec }, WaitForV2Payload { mailbox_id: ShortId }, + PeekV2Payload { mailbox_id: ShortId }, PostV1Response { mailbox_id: ShortId, payload: Vec }, PostV1RequestAndWaitForResponse { mailbox_id: ShortId, payload: Vec }, } @@ -99,6 +106,7 @@ pub enum DbRequest { pub enum DbResponse { PostV2Payload(Option<()>), WaitForV2Payload(Arc>), + PeekV2Payload(Option>>), PostV1Response(()), PostV1RequestAndWaitForResponse(Arc>), } @@ -134,6 +142,8 @@ impl Service for FilesDbService { Ok(DbResponse::PostV2Payload(db.post_v2_payload(&mailbox_id, payload).await?)), DbRequest::WaitForV2Payload { mailbox_id } => Ok(DbResponse::WaitForV2Payload(db.wait_for_v2_payload(&mailbox_id).await?)), + DbRequest::PeekV2Payload { mailbox_id } => + Ok(DbResponse::PeekV2Payload(db.peek_v2_payload(&mailbox_id).await?)), DbRequest::PostV1Response { mailbox_id, payload } => { db.post_v1_response(&mailbox_id, payload).await?; Ok(DbResponse::PostV1Response(())) @@ -199,6 +209,21 @@ impl Db for DbServiceAdapter { } } + async fn peek_v2_payload( + &self, + mailbox_id: &ShortId, + ) -> Result>>, Error> { + let response = self + .inner + .clone() + .oneshot(DbRequest::PeekV2Payload { mailbox_id: *mailbox_id }) + .await?; + match response { + DbResponse::PeekV2Payload(result) => Ok(result), + _ => Err(Self::invalid_response("peek_v2_payload")), + } + } + async fn post_v1_response( &self, mailbox_id: &ShortId, @@ -272,6 +297,14 @@ impl Db for MetricsDb { self.inner.wait_for_v2_payload(mailbox_id).await } + async fn peek_v2_payload( + &self, + mailbox_id: &ShortId, + ) -> Result>>, Error> { + self.metrics.record_short_id(mailbox_id); + self.inner.peek_v2_payload(mailbox_id).await + } + async fn post_v1_response( &self, mailbox_id: &ShortId, From 44e9bca0587685ceb939314e2ff99e308be7aa79 Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Fri, 19 Jun 2026 20:36:22 -0300 Subject: [PATCH 2/7] Add Poisson poll schedule to payjoin Shared so all clients converge on one rate; a divergent rate is itself a fingerprint. --- payjoin/src/core/mod.rs | 2 + payjoin/src/core/schedule.rs | 80 ++++++++++++++++++++++++++++++++++++ 2 files changed, 82 insertions(+) create mode 100644 payjoin/src/core/schedule.rs diff --git a/payjoin/src/core/mod.rs b/payjoin/src/core/mod.rs index ec64e9963..04cdd6306 100644 --- a/payjoin/src/core/mod.rs +++ b/payjoin/src/core/mod.rs @@ -19,6 +19,8 @@ pub use into_url::{Error as IntoUrlError, IntoUrl}; pub(crate) mod url; pub use url::{ParseError as UrlParseError, Url}; #[cfg(feature = "v2")] +pub mod schedule; +#[cfg(feature = "v2")] pub mod time; pub mod uri; pub use uri::{PjParam, PjParseError, PjUri, Uri, UriExt}; diff --git a/payjoin/src/core/schedule.rs b/payjoin/src/core/schedule.rs new file mode 100644 index 000000000..cd3f9e7d9 --- /dev/null +++ b/payjoin/src/core/schedule.rs @@ -0,0 +1,80 @@ +use std::collections::hash_map::RandomState; +use std::hash::{BuildHasher, Hasher}; +use std::time::Duration; + +/// Mean gap of the Poisson poll schedule. The directory learns only this rate, +/// so it must stay uniform across clients, not a per-user knob. +pub const POLL_MEAN: Duration = Duration::from_secs(5); + +/// Samples inter-poll gaps from an Exp(1/mean) distribution. Emit polls on +/// this clock independently of responses (reset on fire, before awaiting the +/// poll) so the observed interval is the gap, not gap + round-trip. +#[derive(Debug)] +pub struct PollSchedule { + state: u64, + mean: Duration, +} + +impl PollSchedule { + /// Create a schedule seeded from OS entropy at the standard [`POLL_MEAN`] rate. + pub fn new() -> Self { + // Seed from OS entropy via the standard library's randomly-keyed hasher: + // hashing a fixed value mixes those random keys into a u64. + let mut h = RandomState::new().build_hasher(); + h.write_u64(0); + Self { state: h.finish(), mean: POLL_MEAN } + } + + /// Sample the next inter-poll gap. + pub fn next_gap(&mut self) -> Duration { + Duration::from_secs_f64(-self.mean.as_secs_f64() * self.next_uniform().ln()) + } + + /// Draw the next uniform sample in (0, 1) from the SplitMix64 state. + fn next_uniform(&mut self) -> f64 { + // SplitMix64 (reference constants: golden-ratio increment + two mixers) + self.state = self.state.wrapping_add(0x9E37_79B9_7F4A_7C15); + let mut z = self.state; + z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); + z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); + z ^= z >> 31; + ((z >> 11) as f64 + 1.0) / (1u64 << 53) as f64 + } +} + +impl Default for PollSchedule { + fn default() -> Self { Self::new() } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn from_seed(seed: u64, mean: Duration) -> PollSchedule { PollSchedule { state: seed, mean } } + + #[test] + fn uniform_sequence_is_pinned() { + let mut s = from_seed(1, Duration::from_secs(5)); + assert_eq!(s.next_uniform(), 0.566561575172281); + assert_eq!(s.next_uniform(), 0.7457817572627012); + assert_eq!(s.next_uniform(), 0.9710027535867963); + } + + #[test] + fn same_seed_is_deterministic() { + let mut a = from_seed(1, Duration::from_secs(5)); + let mut b = from_seed(1, Duration::from_secs(5)); + for _ in 0..100 { + assert_eq!(a.next_gap(), b.next_gap()); + } + } + + #[test] + fn mean_is_near_lambda_inverse() { + let mut s = from_seed(440, Duration::from_secs(5)); + let n = 20_000; + let total: f64 = (0..n).map(|_| s.next_gap().as_secs_f64()).sum(); + let mean = total / n as f64; + assert!((mean - 5.0).abs() < 0.2, "mean gap {mean} not ~5s"); + } +} From 368498d9df3c5c196ef70a4b985f6e2476bfcd8b Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Fri, 19 Jun 2026 20:36:23 -0300 Subject: [PATCH 3/7] Return mailbox GET without long-polling GET now peeks instead of blocking. The v2 waitmap it no longer uses is still woven into the Db trait, post_v2, and tests, so removal is a follow-up. --- payjoin-mailroom/src/db/files.rs | 1 + payjoin-mailroom/src/directory.rs | 31 +++++++++++++++++++++++++++++-- 2 files changed, 30 insertions(+), 2 deletions(-) diff --git a/payjoin-mailroom/src/db/files.rs b/payjoin-mailroom/src/db/files.rs index 063cf1129..a9d97f34a 100644 --- a/payjoin-mailroom/src/db/files.rs +++ b/payjoin-mailroom/src/db/files.rs @@ -246,6 +246,7 @@ impl DbTrait for FilesDb { Ok(guard.read(id).await?) } + // Unused by GET after the non-blocking switch; v2 waitmap removal is a follow-up. async fn wait_for_v2_payload( &self, id: &ShortId, diff --git a/payjoin-mailroom/src/directory.rs b/payjoin-mailroom/src/directory.rs index 038cc4070..9a4750209 100644 --- a/payjoin-mailroom/src/directory.rs +++ b/payjoin-mailroom/src/directory.rs @@ -284,8 +284,16 @@ impl Service { async fn get_mailbox(&self, id: &str) -> Result, HandlerError> { let id = ShortId::from_str(id)?; - let timeout_response = Response::builder().status(StatusCode::ACCEPTED).body(empty())?; - handle_peek(self.db.wait_for_v2_payload(&id).await, timeout_response) + let empty_response = Response::builder().status(StatusCode::ACCEPTED).body(empty())?; + match self.db.peek_v2_payload(&id).await { + Ok(Some(payload)) => Ok(Response::new(full((*payload).clone()))), + Ok(None) => Ok(empty_response), + Err(DbError::Operational(err)) => { + error!("Storage error: {err}"); + Err(HandlerError::InternalServerError(anyhow::Error::msg("Internal server error"))) + } + Err(_) => Ok(empty_response), + } } /// Screen a V1 PSBT body against the address blocklist. @@ -872,6 +880,25 @@ mod tests { } } + #[tokio::test(start_paused = true)] + async fn get_mailbox_returns_immediately_when_empty() { + let svc = test_service(None).await; + let id = valid_short_id_path(); + let start = tokio::time::Instant::now(); + let res = svc.get_mailbox(&id).await.expect("get_mailbox"); + assert_eq!(res.status(), StatusCode::ACCEPTED); + assert_eq!(start.elapsed(), Duration::ZERO, "GET must not block"); + } + + #[tokio::test] + async fn get_mailbox_returns_payload_when_present() { + let svc = test_service(None).await; + let id = valid_short_id_path(); + svc.post_mailbox(&id, Body::from(b"hi".to_vec())).await.expect("post"); + let res = svc.get_mailbox(&id).await.expect("get_mailbox"); + assert_eq!(res.status(), StatusCode::OK); + } + #[tokio::test] async fn post_mailbox_records_short_id_cardinality() { use opentelemetry_sdk::metrics::{ From ad63f4c18492dcd0d9f269325ec90f4935df256a Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Fri, 19 Jun 2026 20:36:23 -0300 Subject: [PATCH 4/7] Poll the directory on a Poisson schedule --- payjoin-cli/src/app/v2/mod.rs | 100 ++++++++++++++++++++++------------ payjoin-cli/tests/e2e.rs | 6 +- 2 files changed, 69 insertions(+), 37 deletions(-) diff --git a/payjoin-cli/src/app/v2/mod.rs b/payjoin-cli/src/app/v2/mod.rs index 42601e205..a7640d4a9 100644 --- a/payjoin-cli/src/app/v2/mod.rs +++ b/payjoin-cli/src/app/v2/mod.rs @@ -12,6 +12,7 @@ use payjoin::receive::v2::{ ReceiverBuilder, SessionOutcome as ReceiverSessionOutcome, UncheckedOriginalPayload, WantsFeeRange, WantsInputs, WantsOutputs, }; +use payjoin::schedule::PollSchedule; use payjoin::send::v2::{ replay_event_log as replay_sender_event_log, PendingFallback as SenderPendingFallback, PollingForProposal, SendSession, Sender, SenderBuilder, SessionOutcome as SenderSessionOutcome, @@ -35,6 +36,7 @@ const W_ID: usize = 12; const W_ROLE: usize = 25; const W_DONE: usize = 15; const W_STATUS: usize = 15; +const POLL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); #[derive(Clone)] pub(crate) struct App { @@ -704,27 +706,44 @@ impl App { sender: Sender, persister: &SenderPersister, ) -> Result<()> { - let mut session = sender.clone(); - // Long poll until we get a response + let session = sender; + let mut schedule = PollSchedule::new(); + let mut polls = tokio::task::JoinSet::new(); + let next = tokio::time::sleep(schedule.next_gap()); + tokio::pin!(next); loop { - let (response, ctx) = - self.post_via_relay(|relay| session.create_poll_request(relay)).await?; - let res = session.process_response(&response.bytes().await?, ctx).save(persister); - match res { - Ok(OptionalTransitionOutcome::Progress(psbt)) => { - println!("Proposal received. Processing..."); - self.process_pj_response(psbt)?; - return Ok(()); - } - Ok(OptionalTransitionOutcome::Stasis(current_state)) => { - println!("No response yet."); - session = current_state; - continue; + tokio::select! { + Some(joined) = polls.join_next(), if !polls.is_empty() => { + let (body, ctx): (Vec, _) = match joined { + Ok(Ok(v)) => v, + _ => continue, + }; + match session.clone().process_response(&body, ctx).save(persister) { + Ok(OptionalTransitionOutcome::Progress(psbt)) => { + println!("Proposal received. Processing..."); + self.process_pj_response(psbt)?; + return Ok(()); + } + Ok(OptionalTransitionOutcome::Stasis(_)) => { + println!("No response yet."); + } + Err(re) => { + println!("{re}"); + tracing::debug!("{re:?}"); + return Err(anyhow!("Response error").context(re)); + } + } } - Err(re) => { - println!("{re}"); - tracing::debug!("{re:?}"); - return Err(anyhow!("Response error").context(re)); + () = &mut next => { + next.as_mut().reset(tokio::time::Instant::now() + schedule.next_gap()); + let relay = self.relay_manager.choose_relay()?; + let (req, ctx) = session.create_poll_request(relay.as_str())?; + let app = self.clone(); + polls.spawn(async move { + let resp = + tokio::time::timeout(POLL_TIMEOUT, app.post_request(req)).await??; + Ok::<_, anyhow::Error>((resp.bytes().await?.to_vec(), ctx)) + }); } } } @@ -735,24 +754,37 @@ impl App { session: Receiver, persister: &ReceiverPersister, ) -> Result> { - let mut session = session; + let mut schedule = PollSchedule::new(); + let mut polls = tokio::task::JoinSet::new(); + let next = tokio::time::sleep(schedule.next_gap()); + tokio::pin!(next); loop { - println!("Polling receive request..."); - let (ohttp_response, context) = - self.post_via_relay(|relay| session.create_poll_request(relay)).await?; - let state_transition = session - .process_response(ohttp_response.bytes().await?.to_vec().as_slice(), context) - .save(persister); - match state_transition { - Ok(OptionalTransitionOutcome::Progress(next_state)) => { - println!("Got a request from the sender. Responding with a Payjoin proposal."); - return Ok(next_state); + tokio::select! { + Some(joined) = polls.join_next(), if !polls.is_empty() => { + let (body, ctx): (Vec, _) = match joined { + Ok(Ok(v)) => v, + _ => continue, + }; + match session.clone().process_response(&body, ctx).save(persister) { + Ok(OptionalTransitionOutcome::Progress(next_state)) => { + println!("Got a request from the sender. Responding with a Payjoin proposal."); + return Ok(next_state); + } + Ok(OptionalTransitionOutcome::Stasis(_)) => {} + Err(e) => return Err(e.into()), + } } - Ok(OptionalTransitionOutcome::Stasis(current_state)) => { - session = current_state; - continue; + () = &mut next => { + next.as_mut().reset(tokio::time::Instant::now() + schedule.next_gap()); + let relay = self.relay_manager.choose_relay()?; + let (req, ctx) = session.create_poll_request(relay.as_str())?; + let app = self.clone(); + polls.spawn(async move { + let resp = + tokio::time::timeout(POLL_TIMEOUT, app.post_request(req)).await??; + Ok::<_, anyhow::Error>((resp.bytes().await?.to_vec(), ctx)) + }); } - Err(e) => return Err(e.into()), } } } diff --git a/payjoin-cli/tests/e2e.rs b/payjoin-cli/tests/e2e.rs index a97e23b7c..9e454078d 100644 --- a/payjoin-cli/tests/e2e.rs +++ b/payjoin-cli/tests/e2e.rs @@ -421,7 +421,7 @@ mod e2e { async fn respond_with_payjoin(mut cli_receive_resumer: Child) -> Result<()> { let mut stdout = cli_receive_resumer.stdout.take().expect("Failed to take stdout of child process"); - let timeout = tokio::time::Duration::from_secs(10); + let timeout = tokio::time::Duration::from_secs(45); let res = tokio::time::timeout( timeout, wait_for_stdout_match(&mut stdout, |line| line.contains("Response successful")), @@ -436,7 +436,7 @@ mod e2e { async fn check_payjoin_sent(mut cli_send_resumer: Child) -> Result<()> { let mut stdout = cli_send_resumer.stdout.take().expect("Failed to take stdout of child process"); - let timeout = tokio::time::Duration::from_secs(10); + let timeout = tokio::time::Duration::from_secs(45); let res = tokio::time::timeout( timeout, wait_for_stdout_match(&mut stdout, |line| line.contains("Payjoin sent")), @@ -466,7 +466,7 @@ mod e2e { async fn check_resume_completed(mut cli_resumer: Child) -> Result<()> { let mut stdout = cli_resumer.stdout.take().expect("Failed to take stdout of child process"); - let timeout = tokio::time::Duration::from_secs(10); + let timeout = tokio::time::Duration::from_secs(45); let res = tokio::time::timeout( timeout, wait_for_stdout_match(&mut stdout, |line| { From 0b59b9bf147cadfaa8bbab441dd4996d07fa1ff6 Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Mon, 22 Jun 2026 19:00:42 -0300 Subject: [PATCH 5/7] Expose transient errors on PersistedError --- payjoin/src/core/persist.rs | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/payjoin/src/core/persist.rs b/payjoin/src/core/persist.rs index 969513120..85a464ac4 100644 --- a/payjoin/src/core/persist.rs +++ b/payjoin/src/core/persist.rs @@ -745,6 +745,10 @@ where } } + pub fn is_transient(&self) -> bool { + matches!(&self.0, InternalPersistedError::Api(ApiError::Transient(_))) + } + pub fn error_state(self) -> Option { match self.0 { InternalPersistedError::Api(ApiError::FatalWithState(_, state)) => Some(state), @@ -1597,4 +1601,27 @@ mod tests { assert!(transient_error.storage_error_ref().is_none()); assert!(transient_error.api_error_ref().is_some()); } + + #[test] + fn is_transient_only_for_transient_api_error() { + let transient = PersistedError::( + InternalPersistedError::Api(ApiError::Transient(InMemoryTestError {})), + ); + assert!(transient.is_transient()); + + let fatal = PersistedError::( + InternalPersistedError::Api(ApiError::Fatal(InMemoryTestError {})), + ); + assert!(!fatal.is_transient()); + + let storage = PersistedError::( + InternalPersistedError::Storage(InMemoryTestError {}), + ); + assert!(!storage.is_transient()); + + let fatal_with_state = PersistedError::( + InternalPersistedError::Api(ApiError::FatalWithState(InMemoryTestError {}, ())), + ); + assert!(!fatal_with_state.is_transient()); + } } From bd5b6b37d45fa60d89ff83ed2f9b9f89c6f96899 Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Mon, 22 Jun 2026 19:42:45 -0300 Subject: [PATCH 6/7] Test poll-path transient keeps session open --- payjoin/src/core/receive/v2/mod.rs | 17 +++++++++++ payjoin/src/core/send/v2/mod.rs | 49 ++++++++++++++++++++++++++++++ 2 files changed, 66 insertions(+) diff --git a/payjoin/src/core/receive/v2/mod.rs b/payjoin/src/core/receive/v2/mod.rs index a57d28256..a9d1f8a8d 100644 --- a/payjoin/src/core/receive/v2/mod.rs +++ b/payjoin/src/core/receive/v2/mod.rs @@ -2524,4 +2524,21 @@ pub mod test { other => panic!("Expected HasReplyableError, got {other:?}"), } } + + #[test] + fn poll_transient_directory_error_leaves_session_open() -> Result<(), BoxError> { + let receiver = receiver(Initialized {}); + let (req, ctx) = receiver.create_poll_request(EXAMPLE_URL)?; + let response = ohttp_response_for(&req.body, http::StatusCode::INTERNAL_SERVER_ERROR); + let persister = InMemoryPersister::::default(); + + let err = receiver + .process_response(&response, ctx) + .save(&persister) + .expect_err("transient response should error"); + + assert!(err.is_transient()); + assert_events(&persister, &[], false); + Ok(()) + } } diff --git a/payjoin/src/core/send/v2/mod.rs b/payjoin/src/core/send/v2/mod.rs index e2dda4a11..f7e3f6cba 100644 --- a/payjoin/src/core/send/v2/mod.rs +++ b/payjoin/src/core/send/v2/mod.rs @@ -817,4 +817,53 @@ mod test { do_cancel_test!(PollingForProposal); Ok(()) } + + fn ohttp_response_for( + ohttp_keys: &OhttpKeys, + req_body: &[u8], + status: http::StatusCode, + ) -> Vec { + let server = + ohttp::Server::new(ohttp_keys.0.clone()).expect("test OHTTP server should be valid"); + let (_, probe_response) = server.decapsulate(req_body).expect("request should decapsulate"); + let response_overhead = + probe_response.encapsulate(&[]).expect("probe should encrypt").len(); + + let (_, server_response) = + server.decapsulate(req_body).expect("request should decapsulate again"); + let mut bhttp_response = + vec![0u8; crate::directory::ENCAPSULATED_MESSAGE_BYTES - response_overhead]; + bhttp::Message::response( + bhttp::StatusCode::try_from(status.as_u16()).expect("status should be valid"), + ) + .write_bhttp(bhttp::Mode::KnownLength, &mut bhttp_response.as_mut_slice()) + .expect("BHTTP response should encode"); + let encrypted = + server_response.encapsulate(&bhttp_response).expect("response should encrypt"); + assert_eq!(encrypted.len(), crate::directory::ENCAPSULATED_MESSAGE_BYTES); + encrypted + } + + #[test] + fn poll_transient_directory_error_is_transient() -> Result<(), BoxError> { + let expiration = + Time::from_now(Duration::from_secs(60)).expect("expiration should be valid"); + let sender = create_sender_context(expiration)?; + let ohttp_keys = sender.session_context.pj_param.ohttp_keys().clone(); + let sender = Sender { state: PollingForProposal, session_context: sender.session_context }; + let (req, ctx) = sender.create_poll_request(EXAMPLE_URL)?; + let response = + ohttp_response_for(&ohttp_keys, &req.body, http::StatusCode::INTERNAL_SERVER_ERROR); + let persister = InMemoryPersister::::default(); + + let err = sender + .process_response(&response, ctx) + .save(&persister) + .expect_err("transient response should error"); + + assert!(err.is_transient()); + assert!(!persister.inner.lock().expect("Shouldn't be poisoned").is_closed); + assert_eq!(persister.inner.lock().expect("Shouldn't be poisoned").events.len(), 0); + Ok(()) + } } From 362a44cb9c5a2fd358e4bb6ea37e77b573b64d98 Mon Sep 17 00:00:00 2001 From: bc1cindy Date: Mon, 22 Jun 2026 20:06:00 -0300 Subject: [PATCH 7/7] Retry transient directory errors while polling --- payjoin-cli/src/app/v2/mod.rs | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/payjoin-cli/src/app/v2/mod.rs b/payjoin-cli/src/app/v2/mod.rs index a7640d4a9..4ad1c1ce0 100644 --- a/payjoin-cli/src/app/v2/mod.rs +++ b/payjoin-cli/src/app/v2/mod.rs @@ -727,6 +727,9 @@ impl App { Ok(OptionalTransitionOutcome::Stasis(_)) => { println!("No response yet."); } + Err(re) if re.is_transient() => { + tracing::debug!("Transient directory error, retrying poll: {re:?}"); + } Err(re) => { println!("{re}"); tracing::debug!("{re:?}"); @@ -771,6 +774,9 @@ impl App { return Ok(next_state); } Ok(OptionalTransitionOutcome::Stasis(_)) => {} + Err(e) if e.is_transient() => { + tracing::debug!("Transient directory error, retrying poll: {e:?}"); + } Err(e) => return Err(e.into()), } }