From b001b3c7b1d2c0b9a624021302a8f7f3c58791ad Mon Sep 17 00:00:00 2001 From: Avi Date: Mon, 24 Aug 2026 00:03:12 -0500 Subject: [PATCH] Trim note publish latency: skip global connection wait, parallel sends, 6s cap --- src/publish.rs | 102 ++++++++++++++++++++++++++++++++----------------- 1 file changed, 67 insertions(+), 35 deletions(-) diff --git a/src/publish.rs b/src/publish.rs index 5e4f476..23187e6 100644 --- a/src/publish.rs +++ b/src/publish.rs @@ -11,10 +11,9 @@ use crate::relays; use crate::settings::Settings; use crate::vault::Vault; -/// How long to wait for relays to accept a connection attempt. -const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); -/// How long to wait for a single relay to accept an event. -const RELAY_SEND_TIMEOUT: Duration = Duration::from_secs(15); +/// How long to wait for a single relay to accept an event. Relays are sent +/// to in parallel, so this caps the whole publish, not each relay. +const RELAY_SEND_TIMEOUT: Duration = Duration::from_secs(6); /// A relay that rejected a published note. #[derive(Debug, Clone, Serialize)] @@ -185,40 +184,73 @@ async fn publish_with_keys( .map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?; } client.connect().await; - client.wait_for_connection(CONNECT_TIMEOUT).await; + // No explicit wait for connections here: `send_event` waits for each relay + // to become writable itself, and the per-send timeout below already bounds + // the whole publish. Waiting for *all* relays first would burn the full + // timeout whenever a single relay is unreachable. - // Send to each relay individually so partial failures are fully reported. - let mut succeeded: Vec = Vec::new(); - let mut failed: Vec = Vec::new(); - for url in &relay_urls { - let relay = match client.relay(url.as_str()).await { - Ok(relay) => relay, - Err(err) => { - let (message, details) = relay_error_message(&err); - failed.push(RelayFailure { - url: url.clone(), - error: message, - details: Some(details), - }); - continue; + // Send to every relay in parallel so one slow or dead relay cannot drag + // the whole publish out to (relays x timeout). + let mut outcomes: Vec>> = + (0..relay_urls.len()).map(|_| None).collect(); + let mut sends = tokio::task::JoinSet::new(); + for (index, url) in relay_urls.iter().cloned().enumerate() { + let client = client.clone(); + let event = event.clone(); + sends.spawn(async move { + match client.relay(url.as_str()).await { + Ok(relay) => { + match tokio::time::timeout(RELAY_SEND_TIMEOUT, relay.send_event(&event)).await { + Ok(Ok(_)) => (index, Ok(url)), + Ok(Err(err)) => { + let (message, details) = relay_error_message(&err); + ( + index, + Err(RelayFailure { + url, + error: message, + details: Some(details), + }), + ) + } + Err(_) => ( + index, + Err(RelayFailure { + url, + error: "The relay did not respond in time.".to_string(), + details: Some( + "Timed out while waiting for the relay to accept the note." + .to_string(), + ), + }), + ), + } + } + Err(err) => { + let (message, details) = relay_error_message(&err); + ( + index, + Err(RelayFailure { + url, + error: message, + details: Some(details), + }), + ) + } } - }; + }); + } + while let Some(joined) = sends.join_next().await { + let (index, outcome) = joined.expect("publish send task panicked"); + outcomes[index] = Some(outcome); + } - match tokio::time::timeout(RELAY_SEND_TIMEOUT, relay.send_event(&event)).await { - Ok(Ok(_)) => succeeded.push(url.clone()), - Ok(Err(err)) => { - let (message, details) = relay_error_message(&err); - failed.push(RelayFailure { - url: url.clone(), - error: message, - details: Some(details), - }); - } - Err(_) => failed.push(RelayFailure { - url: url.clone(), - error: "The relay did not respond in time.".to_string(), - details: Some("Timed out while waiting for the relay to accept the note.".into()), - }), + let mut succeeded = Vec::new(); + let mut failed = Vec::new(); + for outcome in outcomes.into_iter().flatten() { + match outcome { + Ok(url) => succeeded.push(url), + Err(failure) => failed.push(failure), } }