DRY: single parallel relay fan-out for notes and metadata
publish_with_keys and publish_metadata_async carried ~55 identical lines (JoinSet per-relay sends, 6s timeout, outcome collection, disconnect, result split) that had to be hand-synced once already. Extracted into publish::send_to_all_relays; callers keep their own relay-add policy (note path hard-fails on add errors, metadata path ignores them), so behavior and user-facing strings are unchanged. METADATA_SEND_TIMEOUT and profiles::failure_for folded away.
This commit is contained in:
parent
b21678fafb
commit
1c428a9ff2
2 changed files with 35 additions and 84 deletions
|
|
@ -1,6 +1,5 @@
|
||||||
use nostr_sdk::prelude::*;
|
use nostr_sdk::prelude::*;
|
||||||
use serde::Serialize;
|
use serde::Serialize;
|
||||||
use std::time::Duration;
|
|
||||||
use zeroize::Zeroizing;
|
use zeroize::Zeroizing;
|
||||||
|
|
||||||
use crate::crypto::VaultKey;
|
use crate::crypto::VaultKey;
|
||||||
|
|
@ -10,10 +9,6 @@ use crate::relays;
|
||||||
use crate::settings::Settings;
|
use crate::settings::Settings;
|
||||||
use crate::vault::{unix_timestamp, StoredProfile, Vault};
|
use crate::vault::{unix_timestamp, StoredProfile, Vault};
|
||||||
|
|
||||||
/// How long to wait for a single relay to accept the metadata event. Relays
|
|
||||||
/// are sent to in parallel, so this caps the whole publish, not each relay.
|
|
||||||
const METADATA_SEND_TIMEOUT: Duration = Duration::from_secs(6);
|
|
||||||
|
|
||||||
/// A safe view of a profile that contains no secret key material.
|
/// A safe view of a profile that contains no secret key material.
|
||||||
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
|
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
|
||||||
pub struct ProfileSummary {
|
pub struct ProfileSummary {
|
||||||
|
|
@ -459,72 +454,12 @@ async fn publish_metadata_async(
|
||||||
for url in &relay_urls {
|
for url in &relay_urls {
|
||||||
let _ = client.add_relay(url.as_str()).await;
|
let _ = client.add_relay(url.as_str()).await;
|
||||||
}
|
}
|
||||||
client.connect().await;
|
let (succeeded, failed) =
|
||||||
// No explicit wait for connections here: `send_event` waits for each relay
|
crate::publish::send_to_all_relays(&client, relay_urls, &event, "metadata").await;
|
||||||
// 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 every relay in parallel so one slow or dead relay cannot drag
|
|
||||||
// the whole publish out to (relays x timeout).
|
|
||||||
let mut outcomes: Vec<Option<Result<String, RelayFailure>>> =
|
|
||||||
(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(METADATA_SEND_TIMEOUT, relay.send_event(&event))
|
|
||||||
.await
|
|
||||||
{
|
|
||||||
Ok(Ok(_)) => (index, Ok(url)),
|
|
||||||
Ok(Err(err)) => (index, Err(failure_for(&url, &err))),
|
|
||||||
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 metadata."
|
|
||||||
.to_string(),
|
|
||||||
),
|
|
||||||
}),
|
|
||||||
),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Err(err) => (index, Err(failure_for(&url, &err))),
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
while let Some(joined) = sends.join_next().await {
|
|
||||||
let (index, outcome) = joined.expect("metadata send task panicked");
|
|
||||||
outcomes[index] = Some(outcome);
|
|
||||||
}
|
|
||||||
client.disconnect().await;
|
|
||||||
|
|
||||||
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),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
MetadataPublishReport { succeeded, failed }
|
MetadataPublishReport { succeeded, failed }
|
||||||
}
|
}
|
||||||
|
|
||||||
fn failure_for(url: &str, err: &impl std::fmt::Display) -> RelayFailure {
|
|
||||||
let (error, details) = crate::publish::relay_error_message(err);
|
|
||||||
RelayFailure {
|
|
||||||
url: url.to_string(),
|
|
||||||
error,
|
|
||||||
details: Some(details),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Safe summaries of every stored profile, newest last. Never includes
|
/// Safe summaries of every stored profile, newest last. Never includes
|
||||||
/// secret keys.
|
/// secret keys.
|
||||||
pub fn summaries(vault: &Vault) -> Vec<ProfileSummary> {
|
pub fn summaries(vault: &Vault) -> Vec<ProfileSummary> {
|
||||||
|
|
|
||||||
|
|
@ -183,20 +183,46 @@ async fn publish_with_keys(
|
||||||
.await
|
.await
|
||||||
.map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
|
.map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
|
||||||
}
|
}
|
||||||
|
let (succeeded, failed) = send_to_all_relays(&client, relay_urls, &event, "note").await;
|
||||||
|
|
||||||
|
if succeeded.is_empty() {
|
||||||
|
return Err(AppError::publish_failed(failed));
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(PublishReport {
|
||||||
|
event_id,
|
||||||
|
succeeded,
|
||||||
|
failed,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Send an already-signed event to every listed relay in parallel so one slow
|
||||||
|
/// or dead relay cannot drag the whole publish out to (relays x timeout).
|
||||||
|
///
|
||||||
|
/// Callers own how relays are added (and whether an add failure aborts the
|
||||||
|
/// whole publish); this starts from a configured client and handles waiting,
|
||||||
|
/// sending, timing out, and splitting results. `noun` ("note", "metadata")
|
||||||
|
/// only shapes the user-facing timeout wording.
|
||||||
|
pub(crate) async fn send_to_all_relays(
|
||||||
|
client: &Client,
|
||||||
|
relay_urls: Vec<String>,
|
||||||
|
event: &Event,
|
||||||
|
noun: &str,
|
||||||
|
) -> (Vec<String>, Vec<RelayFailure>) {
|
||||||
|
let noun = noun.to_string();
|
||||||
client.connect().await;
|
client.connect().await;
|
||||||
// No explicit wait for connections here: `send_event` waits for each relay
|
// No explicit wait for connections here: `send_event` waits for each relay
|
||||||
// to become writable itself, and the per-send timeout below already bounds
|
// to become writable itself, and the per-send timeout below already bounds
|
||||||
// the whole publish. Waiting for *all* relays first would burn the full
|
// the whole publish. Waiting for *all* relays first would burn the full
|
||||||
// timeout whenever a single relay is unreachable.
|
// timeout whenever a single relay is unreachable.
|
||||||
|
|
||||||
// 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<Option<Result<String, RelayFailure>>> =
|
let mut outcomes: Vec<Option<Result<String, RelayFailure>>> =
|
||||||
(0..relay_urls.len()).map(|_| None).collect();
|
(0..relay_urls.len()).map(|_| None).collect();
|
||||||
let mut sends = tokio::task::JoinSet::new();
|
let mut sends = tokio::task::JoinSet::new();
|
||||||
for (index, url) in relay_urls.iter().cloned().enumerate() {
|
for (index, url) in relay_urls.iter().cloned().enumerate() {
|
||||||
let client = client.clone();
|
let client = client.clone();
|
||||||
let event = event.clone();
|
let event = event.clone();
|
||||||
|
let noun = noun.clone();
|
||||||
sends.spawn(async move {
|
sends.spawn(async move {
|
||||||
match client.relay(url.as_str()).await {
|
match client.relay(url.as_str()).await {
|
||||||
Ok(relay) => {
|
Ok(relay) => {
|
||||||
|
|
@ -218,10 +244,9 @@ async fn publish_with_keys(
|
||||||
Err(RelayFailure {
|
Err(RelayFailure {
|
||||||
url,
|
url,
|
||||||
error: "The relay did not respond in time.".to_string(),
|
error: "The relay did not respond in time.".to_string(),
|
||||||
details: Some(
|
details: Some(format!(
|
||||||
"Timed out while waiting for the relay to accept the note."
|
"Timed out while waiting for the relay to accept the {noun}."
|
||||||
.to_string(),
|
)),
|
||||||
),
|
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
|
|
@ -241,7 +266,7 @@ async fn publish_with_keys(
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
while let Some(joined) = sends.join_next().await {
|
while let Some(joined) = sends.join_next().await {
|
||||||
let (index, outcome) = joined.expect("publish send task panicked");
|
let (index, outcome) = joined.expect("relay send task panicked");
|
||||||
outcomes[index] = Some(outcome);
|
outcomes[index] = Some(outcome);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -255,16 +280,7 @@ async fn publish_with_keys(
|
||||||
}
|
}
|
||||||
|
|
||||||
client.disconnect().await;
|
client.disconnect().await;
|
||||||
|
(succeeded, failed)
|
||||||
if succeeded.is_empty() {
|
|
||||||
return Err(AppError::publish_failed(failed));
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(PublishReport {
|
|
||||||
event_id,
|
|
||||||
succeeded,
|
|
||||||
failed,
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Build a concise user-facing message plus technical detail from a relay
|
/// Build a concise user-facing message plus technical detail from a relay
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue