refactor: share relay pool bootstrap
This commit is contained in:
parent
40ec9a0dfa
commit
6d5063d441
4 changed files with 36 additions and 30 deletions
22
src/feed.rs
22
src/feed.rs
|
|
@ -98,15 +98,8 @@ async fn contact_pubkeys(settings: &Settings, owner_hex: &str) -> Result<Vec<Pub
|
||||||
return Ok(Vec::new());
|
return Ok(Vec::new());
|
||||||
}
|
}
|
||||||
|
|
||||||
let client = Client::new(Keys::generate());
|
let client =
|
||||||
for url in &relay_urls {
|
crate::relays::open_pool(Keys::generate(), &relay_urls, Some(CONNECT_TIMEOUT)).await?;
|
||||||
client
|
|
||||||
.add_relay(url.as_str())
|
|
||||||
.await
|
|
||||||
.map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
|
|
||||||
}
|
|
||||||
client.connect().await;
|
|
||||||
client.wait_for_connection(CONNECT_TIMEOUT).await;
|
|
||||||
|
|
||||||
let filter = Filter::new().kind(Kind::ContactList).author(owner);
|
let filter = Filter::new().kind(Kind::ContactList).author(owner);
|
||||||
client
|
client
|
||||||
|
|
@ -158,15 +151,8 @@ async fn aggregate_for(
|
||||||
let effective_limit = if limit == 0 { DEFAULT_LIMIT } else { limit };
|
let effective_limit = if limit == 0 { DEFAULT_LIMIT } else { limit };
|
||||||
|
|
||||||
// A throwaway identity keeps reading the network completely off the user's keys.
|
// A throwaway identity keeps reading the network completely off the user's keys.
|
||||||
let client = Client::new(Keys::generate());
|
let client =
|
||||||
for url in &relay_urls {
|
crate::relays::open_pool(Keys::generate(), &relay_urls, Some(CONNECT_TIMEOUT)).await?;
|
||||||
client
|
|
||||||
.add_relay(url.as_str())
|
|
||||||
.await
|
|
||||||
.map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
|
|
||||||
}
|
|
||||||
client.connect().await;
|
|
||||||
client.wait_for_connection(CONNECT_TIMEOUT).await;
|
|
||||||
|
|
||||||
let since = Timestamp::now() - LOOKBACK;
|
let since = Timestamp::now() - LOOKBACK;
|
||||||
let filter = Filter::new()
|
let filter = Filter::new()
|
||||||
|
|
|
||||||
|
|
@ -454,6 +454,7 @@ 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) =
|
let (succeeded, failed) =
|
||||||
crate::publish::send_to_all_relays(&client, relay_urls, &event, "metadata").await;
|
crate::publish::send_to_all_relays(&client, relay_urls, &event, "metadata").await;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -176,13 +176,7 @@ async fn publish_with_keys(
|
||||||
.to_bech32()
|
.to_bech32()
|
||||||
.map_err(|e| AppError::internal(format!("Could not encode the event id: {e}")))?;
|
.map_err(|e| AppError::internal(format!("Could not encode the event id: {e}")))?;
|
||||||
|
|
||||||
let client = Client::new(keys.clone());
|
let client = relays::open_pool(keys.clone(), &relay_urls, None).await?;
|
||||||
for url in &relay_urls {
|
|
||||||
client
|
|
||||||
.add_relay(url.as_str())
|
|
||||||
.await
|
|
||||||
.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;
|
let (succeeded, failed) = send_to_all_relays(&client, relay_urls, &event, "note").await;
|
||||||
|
|
||||||
if succeeded.is_empty() {
|
if succeeded.is_empty() {
|
||||||
|
|
@ -199,10 +193,12 @@ async fn publish_with_keys(
|
||||||
/// Send an already-signed event to every listed relay in parallel so one slow
|
/// 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).
|
/// 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
|
/// Callers own opening the pool (see `relays::open_pool`) — including whether
|
||||||
/// whole publish); this starts from a configured client and handles waiting,
|
/// they wait for connections, which this deliberately does not: `send_event`
|
||||||
/// sending, timing out, and splitting results. `noun` ("note", "metadata")
|
/// waits for each relay to become writable itself, and the per-send timeout
|
||||||
/// only shapes the user-facing timeout wording.
|
/// below already bounds the whole publish. Waiting for *all* relays first
|
||||||
|
/// would burn the full timeout whenever a single relay is unreachable.
|
||||||
|
/// `noun` ("note", "metadata") only shapes user-facing timeout wording.
|
||||||
pub(crate) async fn send_to_all_relays(
|
pub(crate) async fn send_to_all_relays(
|
||||||
client: &Client,
|
client: &Client,
|
||||||
relay_urls: Vec<String>,
|
relay_urls: Vec<String>,
|
||||||
|
|
@ -210,7 +206,6 @@ pub(crate) async fn send_to_all_relays(
|
||||||
noun: &str,
|
noun: &str,
|
||||||
) -> (Vec<String>, Vec<RelayFailure>) {
|
) -> (Vec<String>, Vec<RelayFailure>) {
|
||||||
let noun = noun.to_string();
|
let noun = noun.to_string();
|
||||||
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
|
||||||
|
|
|
||||||
|
|
@ -69,6 +69,30 @@ pub fn enabled_urls(settings: &Settings) -> Vec<String> {
|
||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Build a client pool over `relay_urls`, adding each relay and optionally
|
||||||
|
/// waiting (up to `wait`) for connections before returning.
|
||||||
|
///
|
||||||
|
/// Fails on the first relay that cannot be added. Callers that should tolerate
|
||||||
|
/// an unreachable relay instead of aborting add their relays themselves.
|
||||||
|
pub(crate) async fn open_pool(
|
||||||
|
keys: Keys,
|
||||||
|
relay_urls: &[String],
|
||||||
|
wait: Option<Duration>,
|
||||||
|
) -> Result<Client, AppError> {
|
||||||
|
let client = Client::new(keys);
|
||||||
|
for url in relay_urls {
|
||||||
|
client
|
||||||
|
.add_relay(url.as_str())
|
||||||
|
.await
|
||||||
|
.map_err(|e| AppError::network(format!("Could not add relay {url}: {e}")))?;
|
||||||
|
}
|
||||||
|
client.connect().await;
|
||||||
|
if let Some(timeout) = wait {
|
||||||
|
client.wait_for_connection(timeout).await;
|
||||||
|
}
|
||||||
|
Ok(client)
|
||||||
|
}
|
||||||
|
|
||||||
/// Result of testing a relay connection.
|
/// Result of testing a relay connection.
|
||||||
#[derive(Debug, Clone, Serialize)]
|
#[derive(Debug, Clone, Serialize)]
|
||||||
pub struct RelayTestResult {
|
pub struct RelayTestResult {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue