From bf6be68d79ae117c3e1e555796c8ceecc0058ded Mon Sep 17 00:00:00 2001 From: Michael Neale Date: Tue, 2 Jun 2026 15:05:25 +1000 Subject: [PATCH] =?UTF-8?q?fix(serverless):=20persistent=20connection=20po?= =?UTF-8?q?ol=20=E2=80=94=20stop=20rate-limit=20storm?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root cause of 'everything is rate-limited' / 'noting too much': the serverless transport opened a FRESH WebSocket per query and per publish (connect→send→close). get_channels fires ~10 queries, create channel 2 publishes, add-agent several more — each a new connection. Public relays (damus, nos.lol) aggressively rate-limit connection storms + event bursts, so every action failed. Fix (validated against damus/nostr-tools idiomatic patterns + NIP-01/65): - ws_pool.rs: one persistent WebSocket per relay, reused for all queries/publishes; background reader dispatches by sub/event id; transparent reconnect; NIP-42 AUTH answered on the same socket. - Don't block on slow relays: publish returns on FIRST acceptance (FuturesUnordered); connect has a 4s timeout; failed relays enter a 30s cooldown so a dead relay isn't re-dialed every query. - Multi-relay publish (success if any accepts) confirmed idiomatic — not a mistake; it's the resilience model and hides per-relay rate caps. - Pool cleared on workspace switch. Verified live: 5-relay burst (15 queries + 6 publishes) that previously hung 3+min now completes in 17s with zero rate-limits; create→join roundtrip passes on 2 relays. 426 unit tests pass, clippy clean. New live test: serverless_burst_no_rate_limit. --- desktop/src-tauri/src/app_state.rs | 5 + .../src-tauri/src/commands/channels_tests.rs | 65 +++ desktop/src-tauri/src/commands/messages.rs | 2 +- desktop/src-tauri/src/commands/workspace.rs | 4 + desktop/src-tauri/src/lib.rs | 1 + desktop/src-tauri/src/relay.rs | 2 +- desktop/src-tauri/src/ws_pool.rs | 406 ++++++++++++++++++ desktop/src-tauri/src/ws_relay.rs | 285 +----------- 8 files changed, 505 insertions(+), 265 deletions(-) create mode 100644 desktop/src-tauri/src/ws_pool.rs diff --git a/desktop/src-tauri/src/app_state.rs b/desktop/src-tauri/src/app_state.rs index b2a5cd000..ac4a4313c 100644 --- a/desktop/src-tauri/src/app_state.rs +++ b/desktop/src-tauri/src/app_state.rs @@ -22,6 +22,10 @@ pub struct AppState { /// writes are performed over the plain WebSocket (REQ/EVENT) instead. /// Set by `apply_workspace`. See docs/SPROUT_LITE_MODE.md. pub serverless: std::sync::atomic::AtomicBool, + /// Persistent WebSocket pool for serverless mode — one long-lived + /// connection per relay, reused for all queries/publishes (avoids the + /// connect-per-op storm that public relays rate-limit). + pub relay_pool: std::sync::Arc, pub managed_agents_store_lock: Mutex<()>, pub channel_templates_store_lock: Mutex<()>, pub managed_agent_processes: Mutex>, @@ -75,6 +79,7 @@ pub fn build_app_state() -> AppState { .unwrap_or_else(|_| reqwest::Client::new()), relay_url_override: Mutex::new(None), serverless: std::sync::atomic::AtomicBool::new(false), + relay_pool: std::sync::Arc::new(crate::ws_pool::RelayPool::new()), managed_agents_store_lock: Mutex::new(()), channel_templates_store_lock: Mutex::new(()), managed_agent_processes: Mutex::new(HashMap::new()), diff --git a/desktop/src-tauri/src/commands/channels_tests.rs b/desktop/src-tauri/src/commands/channels_tests.rs index 7d5b737f0..00bea5e83 100644 --- a/desktop/src-tauri/src/commands/channels_tests.rs +++ b/desktop/src-tauri/src/commands/channels_tests.rs @@ -345,3 +345,68 @@ async fn serverless_multi_relay_fanout() { eprintln!("✅ multi-relay fanout OK: published to 2, read+deduped from 2"); } + +// ── Connection-reuse burst test (the rate-limit fix) ───────────────────────── +// +// Reproduces the real-app failure: a single AppState (one pooled connection) +// firing many ops in quick succession — like get_channels (~10 queries) + +// create channel (2 publishes) + several messages. Before the connection pool, +// each op opened a fresh WebSocket and public relays rate-limited the storm +// ("you are noting too much"). With the pool, all ops reuse ONE socket per +// relay, so the burst goes through. +// +// Run with: +// cargo test --manifest-path desktop/src-tauri/Cargo.toml \ +// serverless_burst_no_rate_limit -- --ignored --nocapture +#[tokio::test] +#[ignore = "network: hits a live public relay (RELAY_URL, default damus)"] +async fn serverless_burst_no_rate_limit() { + use crate::app_state::build_app_state; + use crate::relay::{query_relay, submit_event}; + use std::sync::atomic::Ordering; + + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + + let keys = nostr::Keys::generate(); + let pk = keys.public_key().to_hex(); + let relay = std::env::var("RELAY_URL").unwrap_or_else(|_| "wss://relay.damus.io".to_string()); + + let state = build_app_state(); + *state.keys.lock().unwrap() = keys.clone(); + *state.relay_url_override.lock().unwrap() = Some(relay.clone()); + state.serverless.store(true, Ordering::Relaxed); + + // Burst of ~15 sequential queries (mimics get_channels firing many REQs). + for i in 0..15 { + query_relay( + &state, + &[serde_json::json!({"kinds":[39000],"#d":[format!("burst-{i}")],"limit":1})], + ) + .await + .unwrap_or_else(|e| panic!("query {i} failed (rate-limited?): {e}")); + } + eprintln!("✅ 15 sequential queries reused one pooled connection"); + + // Burst of publishes (mimics create channel + messages). All must be + // accepted — a rate-limit rejection here is the bug we fixed. + for i in 0..6 { + let channel = uuid::Uuid::new_v4(); + let builder = crate::events::build_message( + channel, + &format!("burst message {i} from {}", &pk[..8]), + None, + &[], + &[], + ) + .expect("build message"); + let resp = submit_event(builder, &state) + .await + .unwrap_or_else(|e| panic!("publish {i} failed (rate-limited?): {e}")); + assert!( + !resp.message.contains("rate-limit"), + "publish {i} was rate-limited — connection pool not reusing socket: {}", + resp.message + ); + } + eprintln!("✅ 6 sequential publishes reused one pooled connection — no rate limit"); +} diff --git a/desktop/src-tauri/src/commands/messages.rs b/desktop/src-tauri/src/commands/messages.rs index 745d3d6cc..ee05a752b 100644 --- a/desktop/src-tauri/src/commands/messages.rs +++ b/desktop/src-tauri/src/commands/messages.rs @@ -117,7 +117,7 @@ async fn send_encrypted_message( let mut published = 0; let mut last_err = None; for wrap in &wraps { - match crate::ws_relay::publish_signed_event_ws(wrap, &keys, &relay_urls).await { + match crate::ws_relay::publish_signed_event_ws(state, wrap, &keys, &relay_urls).await { Ok(()) => published += 1, Err(e) => last_err = Some(e), } diff --git a/desktop/src-tauri/src/commands/workspace.rs b/desktop/src-tauri/src/commands/workspace.rs index facb0c688..712e38a10 100644 --- a/desktop/src-tauri/src/commands/workspace.rs +++ b/desktop/src-tauri/src/commands/workspace.rs @@ -54,6 +54,10 @@ pub fn apply_workspace( std::sync::atomic::Ordering::Relaxed, ); + // Drop any pooled relay connections from the previous workspace so we don't + // reuse a socket authed to a different relay/identity. + state.relay_pool.clear(); + if let Some(keys) = parsed_keys { let mut keys_guard = state.keys.lock().map_err(|e| e.to_string())?; *keys_guard = keys; diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index c6b6e17c7..f97732830 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -12,6 +12,7 @@ mod prevent_sleep; mod relay; mod templates; mod util; +mod ws_pool; mod ws_relay; use app_state::{build_app_state, resolve_persisted_identity, AppState}; diff --git a/desktop/src-tauri/src/relay.rs b/desktop/src-tauri/src/relay.rs index 370dde221..b7116f074 100644 --- a/desktop/src-tauri/src/relay.rs +++ b/desktop/src-tauri/src/relay.rs @@ -284,7 +284,7 @@ pub async fn sync_managed_agent_profile( .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .collect(); - return crate::ws_relay::publish_signed_event_ws(&event, agent_keys, &relay_urls) + return crate::ws_relay::publish_signed_event_ws(state, &event, agent_keys, &relay_urls) .await .map_err(|e| { format!("Created the agent, but could not sync its profile metadata: {e}") diff --git a/desktop/src-tauri/src/ws_pool.rs b/desktop/src-tauri/src/ws_pool.rs new file mode 100644 index 000000000..85261ef2e --- /dev/null +++ b/desktop/src-tauri/src/ws_pool.rs @@ -0,0 +1,406 @@ +//! Persistent WebSocket connection pool for serverless mode. +//! +//! **Why this exists:** the original serverless transport opened a *fresh* +//! WebSocket per query and per publish (connect → send → close). A single +//! `get_channels` fires ~10 queries, creating a channel fires 2 publishes, +//! adding an agent fires several more — each a brand-new connection. Public +//! relays (damus, nos.lol, …) aggressively rate-limit connection storms and +//! event bursts ("rate-limited: you are noting too much"), so *every* action +//! started failing. +//! +//! This pool keeps **one long-lived connection per relay URL** and multiplexes +//! all queries and publishes over it. A background reader task per connection +//! dispatches incoming messages to in-flight requests by subscription id +//! (queries) or event id (OK acks). Connections are re-established on demand +//! after a drop. NIP-42 AUTH challenges are answered automatically. + +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; + +use futures_util::stream::{SplitSink, SplitStream}; +use futures_util::{SinkExt, StreamExt}; +use nostr::Keys; +use tokio::net::TcpStream; +use tokio::sync::{mpsc, oneshot, Mutex}; +use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, WebSocketStream}; + +use crate::relay::SubmitEventResponse; + +type WsStream = WebSocketStream>; +type WsWrite = SplitSink; +type WsRead = SplitStream; + +/// Max time to wait for a relay's EOSE before returning whatever arrived. +/// Reachable public relays respond in well under this; the cap stops a slow +/// relay from stalling the merge. +const QUERY_TIMEOUT: Duration = Duration::from_secs(6); +const PUBLISH_TIMEOUT: Duration = Duration::from_secs(8); +/// Max time to establish a WebSocket before treating the relay as unreachable. +const CONNECT_TIMEOUT: Duration = Duration::from_secs(4); +/// How long to skip a relay after a connection failure (so we don't re-dial a +/// dead relay on every single query in a burst). +const FAILED_COOLDOWN: Duration = Duration::from_secs(30); + +/// A pending query: collects events for one subscription until EOSE. +struct PendingQuery { + events_tx: mpsc::UnboundedSender, + done_tx: Option>, +} + +/// A pending publish: waits for the OK ack for one event id. +struct PendingPublish { + event_id: String, + ok_tx: oneshot::Sender, +} + +/// Shared dispatch state for one connection's reader task. +#[derive(Default)] +struct Dispatch { + queries: HashMap, + publishes: HashMap, +} + +/// A single persistent relay connection. +struct Conn { + write: Mutex, + dispatch: Arc>, + keys: Keys, + relay_url: String, + /// Set false when the reader task exits (connection dropped); the pool then + /// re-establishes on next use. + alive: Arc, +} + +impl Conn { + async fn send(&self, text: String) -> Result<(), String> { + let mut w = self.write.lock().await; + w.send(Message::Text(text.into())) + .await + .map_err(|e| format!("relay send failed: {e}")) + } +} + +/// One persistent connection per relay URL. +#[derive(Default)] +pub struct RelayPool { + /// Map of live connections. A std Mutex (not tokio) so `clear()` can run + /// from the synchronous `apply_workspace` command. We never hold this lock + /// across an `.await`; a `connect_lock` serializes dials instead. + conns: std::sync::Mutex>>, + /// Serializes connection establishment so concurrent first-use of the same + /// relay doesn't open multiple sockets (thundering herd). + connect_lock: Mutex<()>, + /// Relays that recently failed to connect, with the time they may be + /// retried. Lets a burst of queries skip a dead relay instead of each + /// paying the full connect timeout. + failed: std::sync::Mutex>, +} + +impl RelayPool { + pub fn new() -> Self { + Self::default() + } + + fn live(&self, relay_url: &str) -> Option> { + let conns = self.conns.lock().unwrap(); + conns.get(relay_url).and_then(|c| { + if c.alive.load(std::sync::atomic::Ordering::Relaxed) { + Some(c.clone()) + } else { + None + } + }) + } + + /// Whether `relay_url` is in the failure cooldown window. + fn in_cooldown(&self, relay_url: &str) -> bool { + let mut failed = self.failed.lock().unwrap(); + match failed.get(relay_url) { + Some(until) if std::time::Instant::now() < *until => true, + Some(_) => { + failed.remove(relay_url); + false + } + None => false, + } + } + + /// Get (or establish) the persistent connection for `relay_url`. + async fn get(&self, relay_url: &str, keys: &Keys) -> Result, String> { + if let Some(c) = self.live(relay_url) { + return Ok(c); + } + // Skip relays that recently failed to connect (avoids paying the connect + // timeout for every query in a burst to a dead relay). + if self.in_cooldown(relay_url) { + return Err(format!("relay {relay_url} in failure cooldown")); + } + // Serialize dials to the same relay (the std lock is never held across + // the await below). + let _dial = self.connect_lock.lock().await; + if let Some(c) = self.live(relay_url) { + return Ok(c); + } + match connect(relay_url, keys).await { + Ok(conn) => { + self.failed.lock().unwrap().remove(relay_url); + self.conns + .lock() + .unwrap() + .insert(relay_url.to_string(), conn.clone()); + Ok(conn) + } + Err(e) => { + self.failed.lock().unwrap().insert( + relay_url.to_string(), + std::time::Instant::now() + FAILED_COOLDOWN, + ); + Err(e) + } + } + } + + /// Run a REQ over the pooled connection, collecting events until EOSE (or + /// timeout). Reuses the persistent socket — no per-call connect. + pub async fn query( + &self, + relay_url: &str, + keys: &Keys, + filters: &[serde_json::Value], + ) -> Result, String> { + let conn = self.get(relay_url, keys).await?; + let sub_id = format!("q-{}", uuid::Uuid::new_v4()); + let (events_tx, mut events_rx) = mpsc::unbounded_channel(); + let (done_tx, done_rx) = oneshot::channel(); + + { + let mut d = conn.dispatch.lock().await; + d.queries.insert( + sub_id.clone(), + PendingQuery { + events_tx, + done_tx: Some(done_tx), + }, + ); + } + + let mut req = vec![ + serde_json::Value::String("REQ".into()), + serde_json::Value::String(sub_id.clone()), + ]; + req.extend(filters.iter().cloned()); + let req_json = serde_json::Value::Array(req).to_string(); + if let Err(e) = conn.send(req_json).await { + conn.dispatch.lock().await.queries.remove(&sub_id); + return Err(e); + } + + // Collect events until EOSE (done) or timeout. + let mut events = Vec::new(); + let _ = tokio::time::timeout(QUERY_TIMEOUT, async { + // Drain events while waiting for the done signal. + tokio::pin!(done_rx); + loop { + tokio::select! { + Some(ev) = events_rx.recv() => events.push(ev), + _ = &mut done_rx => { + // Drain any straggler events already queued. + while let Ok(ev) = events_rx.try_recv() { + events.push(ev); + } + break; + } + } + } + }) + .await; + + // Best-effort CLOSE + cleanup. + let _ = conn + .send(serde_json::json!(["CLOSE", sub_id]).to_string()) + .await; + conn.dispatch.lock().await.queries.remove(&sub_id); + Ok(events) + } + + /// Publish a signed event over the pooled connection and await OK. + pub async fn publish( + &self, + relay_url: &str, + keys: &Keys, + event: &nostr::Event, + ) -> Result { + let conn = self.get(relay_url, keys).await?; + let event_id = event.id.to_hex(); + let (ok_tx, ok_rx) = oneshot::channel(); + + { + let mut d = conn.dispatch.lock().await; + d.publishes.insert( + event_id.clone(), + PendingPublish { + event_id: event_id.clone(), + ok_tx, + }, + ); + } + + let event_json = serde_json::json!(["EVENT", event]).to_string(); + if let Err(e) = conn.send(event_json).await { + conn.dispatch.lock().await.publishes.remove(&event_id); + return Err(e); + } + + match tokio::time::timeout(PUBLISH_TIMEOUT, ok_rx).await { + Ok(Ok(resp)) => Ok(resp), + // Channel dropped (reader died) — clean up and report. + Ok(Err(_)) => { + conn.dispatch.lock().await.publishes.remove(&event_id); + Err("relay connection lost during publish".to_string()) + } + // Many relays are slow/silent on OK — treat timeout as best-effort + // success so writes don't spuriously fail in the UI. + Err(_) => { + conn.dispatch.lock().await.publishes.remove(&event_id); + Ok(SubmitEventResponse { + event_id, + accepted: true, + message: "published (no OK received before timeout)".to_string(), + }) + } + } + } + + /// Drop all connections (e.g. on workspace switch). Synchronous so it can be + /// called from `apply_workspace`. Dropping a `Conn` drops its write half and + /// the reader task's stream, closing the socket. + pub fn clear(&self) { + self.conns.lock().unwrap().clear(); + self.failed.lock().unwrap().clear(); + } +} + +/// Establish a connection and spawn its background reader task. +async fn connect(relay_url: &str, keys: &Keys) -> Result, String> { + let (ws, _) = tokio::time::timeout(CONNECT_TIMEOUT, connect_async(relay_url)) + .await + .map_err(|_| format!("relay {relay_url} connection timed out"))? + .map_err(|e| format!("relay connection failed: {e}"))?; + let (write, read) = ws.split(); + let dispatch = Arc::new(Mutex::new(Dispatch::default())); + let alive = Arc::new(std::sync::atomic::AtomicBool::new(true)); + + let conn = Arc::new(Conn { + write: Mutex::new(write), + dispatch: dispatch.clone(), + keys: keys.clone(), + relay_url: relay_url.to_string(), + alive: alive.clone(), + }); + + // Background reader: dispatch incoming messages to pending requests. + let reader_conn = conn.clone(); + tokio::spawn(async move { + run_reader(read, reader_conn, dispatch, alive).await; + }); + + Ok(conn) +} + +/// Reader loop: routes EVENT/EOSE/OK/AUTH/CLOSED to the right pending request. +async fn run_reader( + mut read: WsRead, + conn: Arc, + dispatch: Arc>, + alive: Arc, +) { + while let Some(msg) = read.next().await { + let text = match msg { + Ok(Message::Text(t)) => t, + Ok(Message::Close(_)) | Err(_) => break, + Ok(_) => continue, + }; + let Ok(arr) = serde_json::from_str::(&text) else { + continue; + }; + let Some(arr) = arr.as_array() else { continue }; + let Some(tag) = arr.first().and_then(|v| v.as_str()) else { + continue; + }; + + match tag { + "EVENT" => { + // ["EVENT", , ] + let sub = arr.get(1).and_then(|v| v.as_str()).unwrap_or(""); + if let Some(ev) = arr + .get(2) + .and_then(|v| serde_json::from_value::(v.clone()).ok()) + { + let d = dispatch.lock().await; + if let Some(q) = d.queries.get(sub) { + let _ = q.events_tx.send(ev); + } + } + } + "EOSE" => { + let sub = arr.get(1).and_then(|v| v.as_str()).unwrap_or(""); + let mut d = dispatch.lock().await; + if let Some(q) = d.queries.get_mut(sub) { + if let Some(done) = q.done_tx.take() { + let _ = done.send(()); + } + } + } + "CLOSED" => { + let sub = arr.get(1).and_then(|v| v.as_str()).unwrap_or(""); + let mut d = dispatch.lock().await; + if let Some(q) = d.queries.get_mut(sub) { + if let Some(done) = q.done_tx.take() { + let _ = done.send(()); + } + } + } + "OK" => { + // ["OK", , , ] + let id = arr.get(1).and_then(|v| v.as_str()).unwrap_or(""); + let accepted = arr.get(2).and_then(|v| v.as_bool()).unwrap_or(false); + let message = arr + .get(3) + .and_then(|v| v.as_str()) + .unwrap_or("") + .to_string(); + let mut d = dispatch.lock().await; + if let Some(p) = d.publishes.remove(id) { + let _ = p.ok_tx.send(SubmitEventResponse { + event_id: p.event_id, + accepted, + message, + }); + } + } + "AUTH" => { + // NIP-42 challenge — sign and answer over the same socket. + if let Some(challenge) = arr.get(1).and_then(|v| v.as_str()) { + if let Ok(auth_json) = + crate::ws_relay::build_auth_message(&conn.keys, &conn.relay_url, challenge) + { + let _ = conn.send(auth_json).await; + } + } + } + _ => {} + } + } + + // Connection closed: mark dead so the pool reconnects, and fail any + // in-flight requests so callers don't hang. + alive.store(false, std::sync::atomic::Ordering::Relaxed); + let mut d = dispatch.lock().await; + for (_, q) in d.queries.drain() { + if let Some(done) = q.done_tx { + let _ = done.send(()); + } + } + d.publishes.clear(); // oneshot senders dropped → callers get RecvError +} diff --git a/desktop/src-tauri/src/ws_relay.rs b/desktop/src-tauri/src/ws_relay.rs index 39a9a958a..6ec4e72e4 100644 --- a/desktop/src-tauri/src/ws_relay.rs +++ b/desktop/src-tauri/src/ws_relay.rs @@ -14,18 +14,12 @@ //! (channels, DMs, agents) is transport-agnostic: it calls `query_relay` / //! `submit_event`, which dispatch here when `state.is_serverless()`. -use std::time::Duration; - -use futures_util::{SinkExt, StreamExt}; +use futures_util::StreamExt; use nostr::EventBuilder; -use tokio_tungstenite::{connect_async, tungstenite::Message}; use crate::app_state::AppState; use crate::relay::SubmitEventResponse; -const QUERY_TIMEOUT: Duration = Duration::from_secs(10); -const PUBLISH_TIMEOUT: Duration = Duration::from_secs(10); - /// Execute one or more filters as a single `REQ` and collect matching events /// until the relay sends `EOSE`. Mirrors `relay::query_relay` but over a plain /// WebSocket against a generic relay. @@ -71,101 +65,7 @@ async fn query_relay_ws_one( let guard = state.keys.lock().map_err(|e| e.to_string())?; guard.clone() }; - - let (ws, _) = connect_async(relay_url) - .await - .map_err(|e| format!("relay connection failed: {e}"))?; - let (mut write, mut read) = ws.split(); - - // Build ["REQ", , , , ...] - let sub_id = format!("q-{}", uuid::Uuid::new_v4()); - let mut req = vec![ - serde_json::Value::String("REQ".into()), - serde_json::Value::String(sub_id.clone()), - ]; - req.extend(filters.iter().cloned()); - let req_json = serde_json::Value::Array(req).to_string(); - - write - .send(Message::Text(req_json.into())) - .await - .map_err(|e| format!("failed to send REQ: {e}"))?; - - let mut events: Vec = Vec::new(); - - let collect = tokio::time::timeout(QUERY_TIMEOUT, async { - loop { - let msg = match read.next().await { - Some(Ok(m)) => m, - Some(Err(e)) => return Err(format!("WS read error: {e}")), - None => return Err("relay closed during query".to_string()), - }; - let Message::Text(text) = msg else { continue }; - let Ok(arr) = serde_json::from_str::(&text) else { - continue; - }; - let Some(arr) = arr.as_array() else { continue }; - let Some(tag) = arr.first().and_then(|v| v.as_str()) else { - continue; - }; - - // For EVENT/EOSE/CLOSED, arr[1] is the subscription id. - let sub_matches = arr.get(1).and_then(|v| v.as_str()) == Some(sub_id.as_str()); - - match tag { - "EVENT" if sub_matches => { - // ["EVENT", , ] - if let Some(ev) = arr - .get(2) - .and_then(|v| serde_json::from_value::(v.clone()).ok()) - { - events.push(ev); - } - } - "EOSE" if sub_matches => return Ok(()), - "CLOSED" if sub_matches => { - let reason = arr - .get(2) - .and_then(|v| v.as_str()) - .unwrap_or("subscription closed by relay"); - return Err(format!("relay closed subscription: {reason}")); - } - "AUTH" => { - // Relay wants NIP-42 auth. Sign and send, then keep reading. - if let Some(challenge) = arr.get(1).and_then(|v| v.as_str()) { - if let Ok(auth_json) = build_auth_message(&keys, relay_url, challenge) { - let _ = write.send(Message::Text(auth_json.into())).await; - } - } - } - _ => {} - } - } - }) - .await; - - // Best-effort CLOSE so we don't leave a dangling sub. - let close_json = serde_json::json!(["CLOSE", sub_id]).to_string(); - let _ = write.send(Message::Text(close_json.into())).await; - let _ = write.close().await; - - match collect { - Ok(Ok(())) => Ok(events), - Ok(Err(e)) => { - // A relay that CLOSED for auth reasons but we still got some events: - // return what we have rather than failing hard. - if !events.is_empty() { - Ok(events) - } else { - Err(e) - } - } - Err(_) => { - // Timed out waiting for EOSE — return whatever arrived. Many public - // relays are slow to EOSE; partial results are better than nothing. - Ok(events) - } - } + state.relay_pool.query(relay_url, &keys, filters).await } /// Publish a signed event over a plain WebSocket and wait for the relay's @@ -184,23 +84,23 @@ pub async fn submit_event_ws( .sign_with_keys(&keys) .map_err(|e| format!("failed to sign event: {e}"))?; - // Publish to all relays, succeeding if any accepts. Public relays often - // rate-limit bursts ("noting too much"); since the same event going to N - // relays is idempotent (dedup by id), a transient rate-limit on one relay - // is fine as long as another accepts. If ALL relays rate-limit, retry once - // after a short backoff — adding an agent fires several writes in quick - // succession, which is the usual trigger. + // Publish to all relays concurrently and return as soon as ONE accepts — + // don't block on slow/unreachable relays (a stalled public relay must not + // freeze the UI). The same event to N relays is idempotent (dedup by id), + // so the first acceptance is authoritative. Public relays rate-limit bursts + // per connection; with multiple relays another usually accepts. If all + // relays rate-limit, retry once after a short backoff. for attempt in 0..2 { - let futures = relay_urls + let mut futures: futures_util::stream::FuturesUnordered<_> = relay_urls .iter() - .map(|url| submit_event_ws_one(&event, &keys, url)); - let results = futures_util::future::join_all(futures).await; + .map(|url| submit_event_ws_one(state, &event, &keys, url)) + .collect(); let mut last_err = None; - let mut all_rate_limited = !results.is_empty(); - for r in results { + let mut all_rate_limited = true; + while let Some(r) = futures.next().await { match r { - Ok(resp) if resp.accepted => return Ok(resp), + Ok(resp) if resp.accepted => return Ok(resp), // first acceptance wins Ok(resp) => { if !resp.message.contains("rate-limit") { all_rate_limited = false; @@ -214,9 +114,7 @@ pub async fn submit_event_ws( } } - // Only the rate-limit case is worth retrying; other rejections won't - // change on retry. - if attempt == 0 && all_rate_limited { + if attempt == 0 && all_rate_limited && !relay_urls.is_empty() { tokio::time::sleep(std::time::Duration::from_millis(1200)).await; continue; } @@ -226,89 +124,12 @@ pub async fn submit_event_ws( } async fn submit_event_ws_one( + state: &AppState, event: &nostr::Event, keys: &nostr::Keys, relay_url: &str, ) -> Result { - let event_id = event.id.to_hex(); - let event_json = serde_json::json!(["EVENT", event]).to_string(); - - let (ws, _) = connect_async(relay_url) - .await - .map_err(|e| format!("relay connection failed: {e}"))?; - let (mut write, mut read) = ws.split(); - - write - .send(Message::Text(event_json.clone().into())) - .await - .map_err(|e| format!("failed to send EVENT: {e}"))?; - - let result = tokio::time::timeout(PUBLISH_TIMEOUT, async { - loop { - let msg = match read.next().await { - Some(Ok(m)) => m, - Some(Err(e)) => return Err(format!("WS read error: {e}")), - None => return Err("relay closed during publish".to_string()), - }; - let Message::Text(text) = msg else { continue }; - let Ok(arr) = serde_json::from_str::(&text) else { - continue; - }; - let Some(arr) = arr.as_array() else { continue }; - let Some(tag) = arr.first().and_then(|v| v.as_str()) else { - continue; - }; - - match tag { - // ["OK", , , ] - "OK" if arr.get(1).and_then(|v| v.as_str()) == Some(event_id.as_str()) => { - let accepted = arr.get(2).and_then(|v| v.as_bool()).unwrap_or(false); - let message = arr - .get(3) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - return Ok(SubmitEventResponse { - event_id: event_id.clone(), - accepted, - message, - }); - } - "AUTH" => { - if let Some(challenge) = arr.get(1).and_then(|v| v.as_str()) { - if let Ok(auth_json) = build_auth_message(keys, relay_url, challenge) { - let _ = write.send(Message::Text(auth_json.into())).await; - // Re-send the event after authenticating. - let _ = write.send(Message::Text(event_json.clone().into())).await; - } - } - } - _ => {} - } - } - }) - .await; - - let _ = write.close().await; - - match result { - Ok(Ok(resp)) => { - if !resp.accepted { - return Err(format!("relay rejected event: {}", resp.message)); - } - Ok(resp) - } - Ok(Err(e)) => Err(e), - Err(_) => { - // Many relays accept silently or are slow to OK. Treat a timeout as - // best-effort success so writes don't spuriously fail in the UI. - Ok(SubmitEventResponse { - event_id, - accepted: true, - message: "published (no OK received before timeout)".to_string(), - }) - } - } + state.relay_pool.publish(relay_url, keys, event).await } /// Publish an already-signed event over a plain WebSocket and wait for `OK`. @@ -318,85 +139,23 @@ async fn submit_event_ws_one( /// agent-profile sync, where the event is signed by the agent's keys rather /// than the user's identity key. pub async fn publish_signed_event_ws( + state: &AppState, event: &nostr::Event, keys: &nostr::Keys, relay_urls: &[String], ) -> Result<(), String> { - let futures = relay_urls - .iter() - .map(|url| publish_signed_event_ws_one(event, keys, url)); - let results = futures_util::future::join_all(futures).await; let mut last_err = None; - for r in results { - match r { - Ok(()) => return Ok(()), + for url in relay_urls { + match state.relay_pool.publish(url, keys, event).await { + Ok(_) => return Ok(()), Err(e) => last_err = Some(e), } } Err(last_err.unwrap_or_else(|| "all relays failed".to_string())) } -async fn publish_signed_event_ws_one( - event: &nostr::Event, - keys: &nostr::Keys, - relay_url: &str, -) -> Result<(), String> { - let event_id = event.id.to_hex(); - let event_json = serde_json::json!(["EVENT", event]).to_string(); - - let (ws, _) = connect_async(relay_url) - .await - .map_err(|e| format!("relay connection failed: {e}"))?; - let (mut write, mut read) = ws.split(); - - write - .send(Message::Text(event_json.clone().into())) - .await - .map_err(|e| format!("failed to send EVENT: {e}"))?; - - let result = tokio::time::timeout(PUBLISH_TIMEOUT, async { - loop { - let msg = match read.next().await { - Some(Ok(m)) => m, - Some(Err(e)) => return Err(format!("WS read error: {e}")), - None => return Err("relay closed during publish".to_string()), - }; - let Message::Text(text) = msg else { continue }; - let Ok(arr) = serde_json::from_str::(&text) else { - continue; - }; - let Some(arr) = arr.as_array() else { continue }; - let Some(tag) = arr.first().and_then(|v| v.as_str()) else { - continue; - }; - match tag { - "OK" if arr.get(1).and_then(|v| v.as_str()) == Some(event_id.as_str()) => { - return Ok(()); - } - "AUTH" => { - if let Some(challenge) = arr.get(1).and_then(|v| v.as_str()) { - if let Ok(auth_json) = build_auth_message(keys, relay_url, challenge) { - let _ = write.send(Message::Text(auth_json.into())).await; - let _ = write.send(Message::Text(event_json.clone().into())).await; - } - } - } - _ => {} - } - } - }) - .await; - - let _ = write.close().await; - - match result { - Ok(Ok(())) | Err(_) => Ok(()), // best-effort: tolerate slow/silent relays - Ok(Err(e)) => Err(e), - } -} - /// Build a NIP-42 `["AUTH", ]` message string. -fn build_auth_message( +pub(crate) fn build_auth_message( keys: &nostr::Keys, relay_url: &str, challenge: &str,