mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(serverless): persistent connection pool — stop rate-limit storm
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.
This commit is contained in:
@@ -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<crate::ws_pool::RelayPool>,
|
||||
pub managed_agents_store_lock: Mutex<()>,
|
||||
pub channel_templates_store_lock: Mutex<()>,
|
||||
pub managed_agent_processes: Mutex<HashMap<String, ManagedAgentProcess>>,
|
||||
@@ -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()),
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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};
|
||||
|
||||
@@ -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}")
|
||||
|
||||
@@ -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<MaybeTlsStream<TcpStream>>;
|
||||
type WsWrite = SplitSink<WsStream, Message>;
|
||||
type WsRead = SplitStream<WsStream>;
|
||||
|
||||
/// 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<nostr::Event>,
|
||||
done_tx: Option<oneshot::Sender<()>>,
|
||||
}
|
||||
|
||||
/// A pending publish: waits for the OK ack for one event id.
|
||||
struct PendingPublish {
|
||||
event_id: String,
|
||||
ok_tx: oneshot::Sender<SubmitEventResponse>,
|
||||
}
|
||||
|
||||
/// Shared dispatch state for one connection's reader task.
|
||||
#[derive(Default)]
|
||||
struct Dispatch {
|
||||
queries: HashMap<String, PendingQuery>,
|
||||
publishes: HashMap<String, PendingPublish>,
|
||||
}
|
||||
|
||||
/// A single persistent relay connection.
|
||||
struct Conn {
|
||||
write: Mutex<WsWrite>,
|
||||
dispatch: Arc<Mutex<Dispatch>>,
|
||||
keys: Keys,
|
||||
relay_url: String,
|
||||
/// Set false when the reader task exits (connection dropped); the pool then
|
||||
/// re-establishes on next use.
|
||||
alive: Arc<std::sync::atomic::AtomicBool>,
|
||||
}
|
||||
|
||||
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<HashMap<String, Arc<Conn>>>,
|
||||
/// 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<HashMap<String, std::time::Instant>>,
|
||||
}
|
||||
|
||||
impl RelayPool {
|
||||
pub fn new() -> Self {
|
||||
Self::default()
|
||||
}
|
||||
|
||||
fn live(&self, relay_url: &str) -> Option<Arc<Conn>> {
|
||||
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<Arc<Conn>, 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<Vec<nostr::Event>, 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<SubmitEventResponse, String> {
|
||||
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<Arc<Conn>, 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<Conn>,
|
||||
dispatch: Arc<Mutex<Dispatch>>,
|
||||
alive: Arc<std::sync::atomic::AtomicBool>,
|
||||
) {
|
||||
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::<serde_json::Value>(&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", <sub>, <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::<nostr::Event>(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", <event_id>, <accepted>, <message>]
|
||||
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
|
||||
}
|
||||
@@ -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", <sub>, <filter>, <filter>, ...]
|
||||
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<nostr::Event> = 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::<serde_json::Value>(&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", <sub>, <event>]
|
||||
if let Some(ev) = arr
|
||||
.get(2)
|
||||
.and_then(|v| serde_json::from_value::<nostr::Event>(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<SubmitEventResponse, 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::<serde_json::Value>(&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", <event_id>, <accepted: bool>, <message>]
|
||||
"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::<serde_json::Value>(&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", <event>]` message string.
|
||||
fn build_auth_message(
|
||||
pub(crate) fn build_auth_message(
|
||||
keys: &nostr::Keys,
|
||||
relay_url: &str,
|
||||
challenge: &str,
|
||||
|
||||
Reference in New Issue
Block a user