Trim note publish latency: skip global connection wait, parallel sends, 6s cap

This commit is contained in:
Avi 2026-08-24 00:03:12 -05:00
commit b001b3c7b1

View file

@ -11,10 +11,9 @@ use crate::relays;
use crate::settings::Settings; use crate::settings::Settings;
use crate::vault::Vault; use crate::vault::Vault;
/// How long to wait for relays to accept a connection attempt. /// How long to wait for a single relay to accept an event. Relays are sent
const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); /// to in parallel, so this caps the whole publish, not each relay.
/// How long to wait for a single relay to accept an event. const RELAY_SEND_TIMEOUT: Duration = Duration::from_secs(6);
const RELAY_SEND_TIMEOUT: Duration = Duration::from_secs(15);
/// A relay that rejected a published note. /// A relay that rejected a published note.
#[derive(Debug, Clone, Serialize)] #[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}")))?; .map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
} }
client.connect().await; 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. // Send to every relay in parallel so one slow or dead relay cannot drag
let mut succeeded: Vec<String> = Vec::new(); // the whole publish out to (relays x timeout).
let mut failed: Vec<RelayFailure> = Vec::new(); let mut outcomes: Vec<Option<Result<String, RelayFailure>>> =
for url in &relay_urls { (0..relay_urls.len()).map(|_| None).collect();
let relay = match client.relay(url.as_str()).await { let mut sends = tokio::task::JoinSet::new();
Ok(relay) => relay, for (index, url) in relay_urls.iter().cloned().enumerate() {
Err(err) => { let client = client.clone();
let (message, details) = relay_error_message(&err); let event = event.clone();
failed.push(RelayFailure { sends.spawn(async move {
url: url.clone(), match client.relay(url.as_str()).await {
error: message, Ok(relay) => {
details: Some(details), match tokio::time::timeout(RELAY_SEND_TIMEOUT, relay.send_event(&event)).await {
}); Ok(Ok(_)) => (index, Ok(url)),
continue; 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 { let mut succeeded = Vec::new();
Ok(Ok(_)) => succeeded.push(url.clone()), let mut failed = Vec::new();
Ok(Err(err)) => { for outcome in outcomes.into_iter().flatten() {
let (message, details) = relay_error_message(&err); match outcome {
failed.push(RelayFailure { Ok(url) => succeeded.push(url),
url: url.clone(), Err(failure) => failed.push(failure),
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()),
}),
} }
} }