From 1d4738627fa54a1077eae1d1db23b043cb9a782e Mon Sep 17 00:00:00 2001 From: Will Pfleger Date: Mon, 20 Jul 2026 20:40:40 -0700 Subject: [PATCH 1/2] fix(acp): honor relay rate limits and pace resubscribes on bad links MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit On a degraded link (3.4s avg RTT, 40% packet loss, DNS brownouts) the ACP harness generated 12,000+ rate-limit NOTICEs and 7,756 reconnect attempts in 7 minutes. This commit addresses all root causes. **Timeout tuning** — `AUTH_TIMEOUT` raised from 5s to 20s; `CONNECT_TIMEOUT` from 10s to 30s. At 3.4s average RTT the old 5s auth window expired before the relay's first frame arrived. **Rate-limit gate** — When the relay sends a CLOSED or NOTICE with a `"rate-limited:"` prefix, the harness now arms a `rate_limit_gate` deadline using the relay's `"retry in {N}s"` hint (`parse_rate_limit_retry_secs`). While gated: new `Subscribe` commands park in `rate_limited_pending`; ephemeral `PublishEvent` frames (typing indicators) are dropped (they're worthless when the relay already rejected us). The gate takes the max of overlapping deadlines. `set_rate_limit_gate` floors hints below 2s to 5s. **Control-sub recovery** — When the rate-limited CLOSED arrives for the membership or observer-control subscription (which the relay rejects before registering, so no per-channel entry exists), `membership_resub_needed` / `observer_resub_needed` flags are set. The main-loop drain fires even when `rate_limited_pending` is empty, ensuring control subs are re-sent once the gate clears. **Select-integrated pacing** — Inline `tokio::time::sleep(REQ_PACING_INTERVAL)` calls removed from drain functions and `resubscribe_after_reconnect`. Replaced with: (a) `drain_pacing_next: Option` gate; (b) a `select!` timer arm (`sleep_until` / `pending()`) that wakes the loop at each pacing tick; (c) `DRAIN_BUDGET_PER_ITER = 1` REQ per 125ms tick, keeping effective rate ≤8 REQ/s; (d) a `pacing_sleep` helper (no jitter) for shutdown-aware inter-REQ gaps in `resubscribe_after_reconnect`. A 48-channel drain no longer blocks the select! loop for ~6s. **Shutdown-aware resubscribe** — `resubscribe_after_reconnect` accepts `cmd_rx` and returns `ResubscribeResult` (Ok/ControlFailed/Shutdown). `Shutdown` commands during pacing sleeps propagate immediately. **Gate-clearing scope** — `resubscribe_after_reconnect` accepts `is_fresh_connection: bool`. Gate and pending queues are only cleared on a fresh connection (new admission window). Proactive resubscribe on the existing socket passes `false`, preserving parked channels. **Partial reconnect failure** — Individual channel REQ failures park in `resubscribe_retry` (drained by the pacing timer) rather than aborting the entire reconnect. Successful `Subscribe` commands evict stale drain entries to prevent duplicate REQs. `check_rate_gate() -> Option` (renamed from `is_rate_gated`) returns the deadline so callers don't re-read the field. **Backoff ladder reset on stability** — The 60s stability block resets `state.backoff_step` so a brief drop after a long healthy run retries at the short end of the ladder. **DNS brownout classification** — DNS resolution failures retry on a flat `DNS_RETRY_INTERVAL` (2s ± jitter) without consuming a backoff ladder rung. `try_autonomous_reconnect` caps DNS-only retries at 10 so agent startup can't hang through a total brownout; `wait_for_reconnect` is unbounded. Tests: 5 new unit tests covering rate-limit gate arm/expiry/max-extend, control-sub flag set/clear, drain failure re-queue, gate re-arm mid-drain, and successful drain cleanup. `is_dns_error` test adds the `RelayError::WebSocket(tungstenite::Error::Io(...))` production-shaped variant. 560 tests pass. --- crates/buzz-acp/src/relay.rs | 1119 +++++++++++++++++++++++++++++++--- 1 file changed, 1024 insertions(+), 95 deletions(-) diff --git a/crates/buzz-acp/src/relay.rs b/crates/buzz-acp/src/relay.rs index dee8359a77..cf45329f13 100644 --- a/crates/buzz-acp/src/relay.rs +++ b/crates/buzz-acp/src/relay.rs @@ -52,14 +52,21 @@ const PONG_TIMEOUT: Duration = Duration::from_secs(10); /// wedging the background task indefinitely. const WS_SEND_TIMEOUT_SECS: u64 = 10; /// Diagnostic threshold: log when a connection has been stable for this long. -/// No backoff reset is implemented yet — this is a hook for future improvement. +/// The stability block resets `BgState::backoff_step` to 0 here so the next +/// drop after a long healthy run retries at the short end of the ladder again. const STABLE_CONNECTION_SECS: u64 = 60; /// Seconds subtracted from `since` on resubscribe to tolerate clock skew. const SINCE_SKEW_SECS: u64 = 5; /// Timeout for the NIP-42 auth handshake steps. -const AUTH_TIMEOUT: Duration = Duration::from_secs(5); +/// +/// Raised from 5s to 20s (≈2 RTTs at the observed 10s max round-trip on degraded +/// links) so auth doesn't time out before the first WS frame arrives. +const AUTH_TIMEOUT: Duration = Duration::from_secs(20); /// Timeout for the TCP + WebSocket handshake in `do_connect`. -const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); +/// +/// Raised from 10s to 30s so the OS TCP connect attempt (SYN→SYN-ACK) has time +/// to succeed at 3.4s average / 10s max observed RTT. +const CONNECT_TIMEOUT: Duration = Duration::from_secs(30); /// Backoff delay values shared by the initial-connect retry in /// `HarnessRelay::connect()` and `try_autonomous_reconnect`'s post-start /// reconnect loop — a spotty link should get consistent retry pacing whether @@ -78,6 +85,26 @@ const STARTUP_CONNECT_BACKOFFS: [Duration; 5] = [ Duration::from_secs(8), Duration::from_secs(16), ]; +/// Flat retry interval for DNS failures — no backoff ladder rung consumed. +/// 2s gives name servers a short window to recover from a brownout without driving +/// a tight storm; jitter (±20%) staggers concurrent agent instances. +/// +/// DNS flat retries are capped at 10 in the bounded startup/reconnect path +/// (`try_autonomous_reconnect`) so a full brownout cannot hang agent startup +/// indefinitely. In `wait_for_reconnect` the DNS path is unbounded — a +/// reconnecting agent should keep trying across extended outages rather than +/// give up. +const DNS_RETRY_INTERVAL: Duration = Duration::from_secs(2); +/// Minimum inter-REQ spacing during resubscribe bursts. +/// 125 ms ≈ 8 frames/s — safely below the relay's 50-frames-per-5s admission +/// window (10 frames/s at the limit). A 48-channel reconnect spreads over ≈6 s +/// instead of arriving as a single burst that consumes the entire budget at once. +const REQ_PACING_INTERVAL: Duration = Duration::from_millis(125); +/// Maximum REQ frames sent per drain iteration (shared across rate_limited_pending, +/// resubscribe_retry, and control-sub recovery). Keeps any single main-loop tick +/// below the relay's 50-frames/5s budget, and ensures the select! loop is never +/// blocked for more than one REQ's worth of I/O between drain ticks. +const DRAIN_BUDGET_PER_ITER: usize = 1; use std::time::Instant; @@ -951,6 +978,37 @@ struct BgState { /// Startup-era channels use the startup watermark; dynamic channels use /// the membership notification timestamp that caused the subscription. subscribe_since: HashMap, + /// Relay rate-limit gate deadline. + /// + /// While `Some(deadline)` and `Instant::now() < deadline`, outbound + /// admission-counted frames (REQ, EVENT) are deferred or dropped. + /// `check_rate_gate` lazily clears this to `None` once it expires. + rate_limit_gate: Option, + /// Channels parked because a CLOSED "rate-limited:" was received. + /// + /// Drained by the main loop when the gate clears, one REQ per + /// `REQ_PACING_INTERVAL` tick via the select-integrated pacing timer. + /// Value is the `Instant` before which the channel must not be retried. + rate_limited_pending: HashMap, + /// Set when a rate-limited CLOSED arrives for the membership notification + /// subscription. The main-loop drain re-sends the REQ once the gate clears, + /// even when `rate_limited_pending` is empty. + membership_resub_needed: bool, + /// Set when a rate-limited CLOSED arrives for the observer control + /// subscription. The main-loop drain re-sends the REQ once the gate clears, + /// even when `rate_limited_pending` is empty. + observer_resub_needed: bool, + /// Channels whose REQ failed during `resubscribe_after_reconnect`. + /// + /// A single failed channel REQ is parked here instead of aborting the whole + /// reconnect. Drained by the main loop. Flushed on each reconnect attempt. + resubscribe_retry: HashSet, + /// Current position in the exponential backoff ladder. + /// + /// Persisted across calls to `wait_for_reconnect` so a flapping link stays at + /// the elevated rung it earned. Reset to 0 by the stability block once the + /// connection has been up for `STABLE_CONNECTION_SECS`. + backoff_step: usize, } impl BgState { @@ -968,6 +1026,12 @@ impl BgState { proactive_resubscribe_needed: false, startup_watermark: None, subscribe_since: HashMap::new(), + rate_limit_gate: None, + rate_limited_pending: HashMap::new(), + membership_resub_needed: false, + observer_resub_needed: false, + resubscribe_retry: HashSet::new(), + backoff_step: 0, } } @@ -1020,6 +1084,49 @@ impl BgState { self.subscribe_since.remove(channel_id); self.channel_dropped_since.remove(channel_id); self.active_filters.remove(channel_id); + self.rate_limited_pending.remove(channel_id); + self.resubscribe_retry.remove(channel_id); + } + + /// Arm or extend the rate-limit gate. + /// + /// `retry_secs` is the relay's `retry in {N}s` hint; hints below 2s (including + /// the no-hint case of 0) floor to 5s. The floor prevents a burst of + /// low-quality hints from dropping the gate so short that re-queued REQs + /// immediately re-trigger rate limiting. Note the deliberate asymmetry with + /// the desktop TypeScript client, which uses a 10s no-hint default — both + /// values are conservative enough; the relay hint wins when present. + /// + /// The gate takes the **maximum** of any existing deadline and the newly + /// computed one so overlapping CLOSED/NOTICE messages can't shorten a gate + /// that is already set further out. + /// + /// Returns the gate deadline that was set. + fn set_rate_limit_gate(&mut self, retry_secs: u64) -> tokio::time::Instant { + let secs = if retry_secs < 2 { 5 } else { retry_secs }; + let base = Duration::from_secs(secs); + let deadline = tokio::time::Instant::now() + jittered_duration(base); + let gate = match self.rate_limit_gate { + Some(existing) if existing > deadline => existing, + _ => deadline, + }; + self.rate_limit_gate = Some(gate); + gate + } + + /// Check whether the rate-limit gate is currently active. + /// + /// Returns `Some(deadline)` when gated, `None` when the gate has expired or + /// was never set. Lazily clears `rate_limit_gate` to `None` on expiry so + /// subsequent calls are cheap (no `Instant::now()` except when `Some`). + fn check_rate_gate(&mut self) -> Option { + if let Some(deadline) = self.rate_limit_gate { + if tokio::time::Instant::now() < deadline { + return Some(deadline); + } + self.rate_limit_gate = None; + } + None } } @@ -1103,6 +1210,24 @@ async fn execute_connected_command( filter, replay_since, } => { + // Rate-gated: defer this REQ to prevent flooding a saturated relay. + // The gate holds until the relay's retry hint expires. + if let Some(retry_after) = state.check_rate_gate() { + debug!( + "rate-gated: deferring REQ for channel {channel_id} to rate_limited_pending" + ); + apply_command_to_state( + state, + RelayCommand::Subscribe { + channel_id, + filter, + replay_since, + }, + ); + state.rate_limited_pending.insert(channel_id, retry_after); + return true; // connection is fine — just rate-limited + } + // Seed subscribe_since BEFORE computing since — on first // subscribe, this provides the fallback timestamp that // closes the startup/dynamic-membership blind spot. @@ -1123,6 +1248,10 @@ async fn execute_connected_command( .active_subscriptions .insert(channel_id, channel_sub_id(channel_id)); state.active_filters.insert(channel_id, filter); + // Evict stale drain entries so the drain loop can't send a + // duplicate REQ for this now-live subscription. + state.rate_limited_pending.remove(&channel_id); + state.resubscribe_retry.remove(&channel_id); true } else { // Send failed — record intent so reconnect restores it. @@ -1180,6 +1309,18 @@ async fn execute_connected_command( } } RelayCommand::PublishEvent { event } => { + // Drop ephemeral publishes while rate-gated. Stale typing indicators + // are worthless and sending them would consume admission budget the + // relay already rejected us on. + // + // INVARIANT: the WS publish path carries only ephemeral kinds (typing + // indicators). The silent drop-while-gated relies on that invariant. If a + // future caller publishes durable events through this path, it must add a + // kind guard before this branch to avoid silently discarding user data. + if state.check_rate_gate().is_some() { + debug!("rate-gated: dropping ephemeral PublishEvent (typing indicator)"); + return true; + } let msg = json!(["EVENT", event]); if let Ok(text) = serde_json::to_string(&msg) { if let Err(e) = @@ -1304,66 +1445,151 @@ async fn run_background_task( let mut connected_since = Instant::now(); let mut stable_logged = false; + // Pacing timer for select-integrated rate-limit drain. + // `None` = no pending drain or budget window is open; `Some(t)` = next + // allowed drain tick. The select! arm below fires when `t` elapses and + // resets this to `None`, allowing the pre-select drain to run again. + let mut drain_pacing_next: Option = None; + loop { if state.proactive_resubscribe_needed { state.proactive_resubscribe_needed = false; info!("proactive resubscribe triggered by backpressure event loss"); - if !resubscribe_after_reconnect(&mut ws, &mut state, &agent_pubkey_hex).await { - warn!("proactive resubscribe had failures — triggering reconnect"); - let _ = event_tx.try_send(None); - match try_autonomous_reconnect( - &mut ws, - &mut cmd_rx, - &mut state, - &keys, - &relay_url, - &agent_pubkey_hex, - &event_tx, - &observer_control_tx, - auth_tag.as_ref(), - ) - .await - { - ReconnectOutcome::Ok => { - if matches!( - drain_post_reconnect( - &mut ws, - &mut cmd_rx, - &mut state, - &agent_pubkey_hex - ) - .await, - ReconnectOutcome::Shutdown - ) { - return; + // Proactive resubscribe runs on the EXISTING socket — do NOT clear the + // rate-limit gate or pending queues. + match resubscribe_after_reconnect( + &mut ws, + &mut cmd_rx, + &mut state, + &agent_pubkey_hex, + false, // existing socket — preserve gate state + ) + .await + { + ResubscribeResult::Ok => {} + ResubscribeResult::Shutdown => return, + ResubscribeResult::ControlFailed => { + warn!("proactive resubscribe had failures — triggering reconnect"); + let _ = event_tx.try_send(None); + match try_autonomous_reconnect( + &mut ws, + &mut cmd_rx, + &mut state, + &keys, + &relay_url, + &agent_pubkey_hex, + &event_tx, + &observer_control_tx, + auth_tag.as_ref(), + ) + .await + { + ReconnectOutcome::Ok => { + if matches!( + drain_post_reconnect( + &mut ws, + &mut cmd_rx, + &mut state, + &agent_pubkey_hex + ) + .await, + ReconnectOutcome::Shutdown + ) { + return; + } } - } - ReconnectOutcome::Shutdown => return, - ReconnectOutcome::Failed => { - if matches!( - wait_for_reconnect( - &mut ws, - &mut cmd_rx, - &mut state, - &keys, - &relay_url, - &agent_pubkey_hex, - &event_tx, - &observer_control_tx, - true, - auth_tag.as_ref(), - ) - .await, - ReconnectOutcome::Shutdown - ) { - return; + ReconnectOutcome::Shutdown => return, + ReconnectOutcome::Failed => { + if matches!( + wait_for_reconnect( + &mut ws, + &mut cmd_rx, + &mut state, + &keys, + &relay_url, + &agent_pubkey_hex, + &event_tx, + &observer_control_tx, + true, + auth_tag.as_ref(), + ) + .await, + ReconnectOutcome::Shutdown + ) { + return; + } } } + ping_sent = false; + last_pong = Instant::now(); + connected_since = Instant::now(); + stable_logged = false; + } + } + } + + // Drain pending subs, one REQ per pacing tick within the relay's + // admission window. + let drain_window_open = drain_pacing_next.is_none_or(|t| tokio::time::Instant::now() >= t); + if drain_window_open { + let mut budget = DRAIN_BUDGET_PER_ITER; + let mut any_sent = false; + + // Control subs use a flag rather than a per-channel pending entry, so + // recovery fires even when rate_limited_pending is empty. + if state.check_rate_gate().is_none() { + if state.membership_resub_needed && budget > 0 { + let replay_since = + match (state.membership_dropped_since, state.membership_last_seen) { + (Some(d), Some(l)) => Some(d.min(l)), + (Some(d), None) => Some(d), + (None, Some(l)) => Some(l), + (None, None) => state.startup_watermark, + }; + if send_membership_subscribe(&mut ws, &agent_pubkey_hex, replay_since).await { + state.membership_resub_needed = false; + state.membership_dropped_since = None; + budget = budget.saturating_sub(1); + any_sent = true; + } else { + warn!( + "membership control resub after rate-limit failed — will retry next drain" + ); + } + } + if state.observer_resub_needed && budget > 0 { + if send_observer_control_subscribe(&mut ws, &agent_pubkey_hex).await { + state.observer_resub_needed = false; + budget = budget.saturating_sub(1); + any_sent = true; + } else { + warn!( + "observer control resub after rate-limit failed — will retry next drain" + ); + } + } + } + + if budget > 0 && !state.rate_limited_pending.is_empty() { + let sent = + drain_rate_limited_pending(&mut ws, &mut state, &agent_pubkey_hex, budget) + .await; + budget = budget.saturating_sub(sent); + if sent > 0 { + any_sent = true; + } + } + + if budget > 0 && !state.resubscribe_retry.is_empty() { + let sent = + drain_resubscribe_retry(&mut ws, &mut state, &agent_pubkey_hex, budget).await; + if sent > 0 { + any_sent = true; } - ping_sent = false; - last_pong = Instant::now(); - connected_since = Instant::now(); - stable_logged = false; + } + + if any_sent { + drain_pacing_next = Some(tokio::time::Instant::now() + REQ_PACING_INTERVAL); } } @@ -1593,12 +1819,30 @@ async fn run_background_task( } } } + + // Pacing timer arm — wakes the loop for the next drain batch. + // `pending()` when no drain is in progress so this arm never + // fires spuriously and never blocks the other select! arms. + _ = async { + match drain_pacing_next { + Some(t) => tokio::time::sleep_until(t).await, + None => std::future::pending::<()>().await, + } + } => { + drain_pacing_next = None; + } } + // Reset backoff_step on a long healthy run so a subsequent brief drop + // retries at the short end of the backoff ladder. if !stable_logged && connected_since.elapsed() > Duration::from_secs(STABLE_CONNECTION_SECS) { stable_logged = true; - debug!("connection stable for >{}s", STABLE_CONNECTION_SECS); + state.backoff_step = 0; + debug!( + "connection stable for >{}s — backoff ladder reset", + STABLE_CONNECTION_SECS + ); } } } @@ -1761,6 +2005,18 @@ async fn handle_ws_message( RelayMessage::Notice { message } => { // Fix 4: NOTICE at warn level. tracing::warn!("relay NOTICE: {message}"); + // The relay sends NOTICE for rate-limited EVENT/COUNT frames. + if message.starts_with("rate-limited:") { + let secs = parse_rate_limit_retry_secs(&message).unwrap_or(0); + let deadline = state.set_rate_limit_gate(secs); + warn!( + "rate-limit gate armed via NOTICE until ~{:.1}s from now", + deadline + .checked_duration_since(tokio::time::Instant::now()) + .unwrap_or_default() + .as_secs_f64() + ); + } } RelayMessage::Closed { subscription_id, @@ -1775,6 +2031,32 @@ async fn handle_ws_message( return true; } + // Rate-limited CLOSED — park and keep the socket. The relay's + // "retry in {N}s" hint arms the gate; the channel or control sub + // is resubscribed by the main-loop drain once the gate clears. + if message.starts_with("rate-limited:") { + let secs = parse_rate_limit_retry_secs(&message).unwrap_or(0); + let deadline = state.set_rate_limit_gate(secs); + warn!( + "subscription {subscription_id} rate-limited — parking until ~{:.1}s, gate armed", + deadline + .checked_duration_since(tokio::time::Instant::now()) + .unwrap_or_default() + .as_secs_f64() + ); + if let Some(channel_id) = channel_id_from_sub_id(&subscription_id) { + state.rate_limited_pending.insert(channel_id, deadline); + } else if subscription_id == MEMBERSHIP_NOTIF_SUB_ID { + // Mark membership sub for drain recovery. The relay rejected + // this REQ before registering it, so the sub does not exist + // server-side — the drain must re-send it. + state.membership_resub_needed = true; + } else if subscription_id == OBSERVER_CONTROL_SUB_ID { + state.observer_resub_needed = true; + } + return true; // keep the socket + } + // CLOSED needs cleanup and resubscribe, not just logging. let is_auth_error = message.starts_with("auth-required") || message.starts_with("restricted") @@ -1981,27 +2263,73 @@ async fn process_handshake_buffer( true } +/// Outcome of [`resubscribe_after_reconnect`]. +enum ResubscribeResult { + /// All subscriptions restored (or parked for drain recovery). + Ok, + /// A control subscription (membership or observer) failed to send. + /// Caller should treat this as a failed reconnect attempt. + ControlFailed, + /// A `Shutdown` command arrived during a pacing sleep. + /// Caller must return immediately (background task is exiting). + Shutdown, +} + /// Resubscribe all active channels and membership notifications after a /// successful reconnect. Computes `since = min(last_seen, channel_dropped_since)` /// per channel, and only clears the drop tracker when the REQ is confirmed sent. /// -/// Returns `true` if ALL subscriptions were sent successfully. Returns `false` -/// if any send failed — the caller should treat this as a failed reconnect -/// and retry, because a "connected" socket with missing subscriptions causes -/// silent event loss. +/// Paces REQs at `REQ_PACING_INTERVAL` (125 ms) via a shutdown-aware sleep so +/// a 48-channel reconnect burst spreads over ≈6 s without starving the select! +/// loop of frames, commands, or ping ticks. If the gate is active mid-burst, +/// remaining channels are parked in `rate_limited_pending` instead of sent. +/// +/// A failed CHANNEL REQ is parked in `resubscribe_retry` rather than failing +/// the whole reconnect. Only membership and observer-control failures return +/// `ControlFailed` — their silent loss breaks joins and agent pause/resume. +/// +/// When `is_fresh_connection` is `true` (new socket established), the +/// rate-limit gate and pending queues are cleared because a fresh connection +/// resets the relay's admission window. When `false` (proactive resubscribe on +/// the EXISTING socket), they are preserved so parked channels aren't lost. +/// +/// Returns [`ResubscribeResult`] signalling success, control-sub failure, or +/// shutdown. async fn resubscribe_after_reconnect( ws: &mut WsStream, + cmd_rx: &mut mpsc::Receiver, state: &mut BgState, agent_pubkey_hex: &str, -) -> bool { - let mut all_ok = true; + is_fresh_connection: bool, +) -> ResubscribeResult { + // Only clear stale gate state on a FRESH connection. A proactive resubscribe + // on the existing socket must preserve the gate and pending queues — + // clearing them would discard parked channels and burst into a still-closed + // admission window. + if is_fresh_connection { + state.rate_limit_gate = None; + state.rate_limited_pending.clear(); + state.resubscribe_retry.clear(); + // Control-sub pending flags are implicitly cleared by the resub below. + } + let channels: Vec = state.active_subscriptions.keys().copied().collect(); if !channels.is_empty() { info!( "resubscribing to {} channel(s) after reconnect", channels.len() ); + let mut sent_any = false; for channel_id in channels { + // Gate re-armed mid-burst — park remaining channels. + if let Some(retry_after) = state.check_rate_gate() { + debug!( + "rate-gated mid-resubscribe: parking channel {channel_id} in rate_limited_pending" + ); + state.rate_limited_pending.insert(channel_id, retry_after); + continue; + } + let since = state.channel_since(&channel_id); let filter = match state.active_filters.get(&channel_id).cloned() { Some(f) => f, @@ -2010,22 +2338,40 @@ async fn resubscribe_after_reconnect( // intent is inconsistent. Skip rather than resubscribe with // a permissive wildcard that would widen the subscription. warn!("missing filter for channel {channel_id} — skipping resubscribe (fail-closed)"); - all_ok = false; + state.resubscribe_retry.insert(channel_id); continue; } }; - let sent = + let this_sent = send_subscribe(ws, state, channel_id, agent_pubkey_hex, since, &filter).await; - if sent { + if this_sent { state.channel_dropped_since.remove(&channel_id); + sent_any = true; + // Shutdown-aware pacing sleep between REQs. + if !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await { + return ResubscribeResult::Shutdown; + } } else { - warn!("failed to resubscribe channel {channel_id} after reconnect"); - all_ok = false; + // Partial failure — park the channel for main-loop retry instead + // of aborting the entire reconnect. + warn!( + "failed to resubscribe channel {channel_id} after reconnect — parking for retry" + ); + state.resubscribe_retry.insert(channel_id); } } + let _ = sent_any; // consumed above for pacing logic } + // Membership and observer-control are control-plane subscriptions: a silent + // failure breaks join notifications and agent pause/resume. These still force + // a full reconnect retry rather than silent degradation. if state.membership_sub_active { + if !state.active_subscriptions.is_empty() + && !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await + { + return ResubscribeResult::Shutdown; + } let replay_since = match (state.membership_dropped_since, state.membership_last_seen) { (Some(d), Some(l)) => Some(d.min(l)), (Some(d), None) => Some(d), @@ -2035,20 +2381,141 @@ async fn resubscribe_after_reconnect( let sent = send_membership_subscribe(ws, agent_pubkey_hex, replay_since).await; if sent { state.membership_dropped_since = None; + state.membership_resub_needed = false; } else { warn!("failed to resubscribe membership after reconnect"); - all_ok = false; + return ResubscribeResult::ControlFailed; } } - if state.observer_control_sub_active - && !send_observer_control_subscribe(ws, agent_pubkey_hex).await - { - warn!("failed to resubscribe observer controls after reconnect"); - all_ok = false; + if state.observer_control_sub_active { + if !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await { + return ResubscribeResult::Shutdown; + } + if !send_observer_control_subscribe(ws, agent_pubkey_hex).await { + warn!("failed to resubscribe observer controls after reconnect"); + return ResubscribeResult::ControlFailed; + } + state.observer_resub_needed = false; } - all_ok + ResubscribeResult::Ok +} + +/// Drain `rate_limited_pending` channels whose retry deadline has passed. +/// +/// Called by the main loop pacing timer. Sends at most `budget` REQs without +/// sleeping — pacing is enforced by the caller via `drain_pacing_next`. A +/// failed send re-queues the channel with a +5 s penalty. Returns the number +/// of REQs successfully sent. +async fn drain_rate_limited_pending( + ws: &mut WsStream, + state: &mut BgState, + agent_pubkey_hex: &str, + budget: usize, +) -> usize { + let now = tokio::time::Instant::now(); + let ready: Vec = state + .rate_limited_pending + .iter() + .filter(|(_, &deadline)| now >= deadline) + .map(|(&ch, _)| ch) + .take(budget) + .collect(); + + if ready.is_empty() { + return 0; + } + debug!("draining {} rate_limited_pending channel(s)", ready.len()); + + let mut sent_count = 0; + for channel_id in ready { + // Re-check gate each iteration — a new CLOSED may have re-armed it mid-drain. + if let Some(retry_after) = state.check_rate_gate() { + state.rate_limited_pending.insert(channel_id, retry_after); + continue; + } + + let since = state.channel_since(&channel_id); + let filter = match state.active_filters.get(&channel_id).cloned() { + Some(f) => f, + None => { + warn!("missing filter for channel {channel_id} in rate_limited_pending — dropping"); + state.rate_limited_pending.remove(&channel_id); + continue; + } + }; + let sent = send_subscribe(ws, state, channel_id, agent_pubkey_hex, since, &filter).await; + if sent { + state.rate_limited_pending.remove(&channel_id); + state.channel_dropped_since.remove(&channel_id); + sent_count += 1; + // Pacing is enforced by the main-loop timer; no inline sleep here. + } else { + // Socket may be dead — re-queue with +5s penalty; the next ws event + // will detect the dead socket and trigger a full reconnect. + let penalty = tokio::time::Instant::now() + Duration::from_secs(5); + state.rate_limited_pending.insert(channel_id, penalty); + warn!("drain_rate_limited_pending: REQ failed for channel {channel_id} — re-queued with +5s penalty"); + } + } + sent_count +} + +/// Drain `resubscribe_retry` channels that were parked by partial reconnect failure. +/// +/// Called by the main loop pacing timer. Sends at most `budget` REQs without +/// sleeping — pacing is enforced by the caller. A failed send leaves the +/// channel in the retry set; a gate re-armed mid-drain moves it to +/// `rate_limited_pending`. Returns the number of REQs successfully sent. +async fn drain_resubscribe_retry( + ws: &mut WsStream, + state: &mut BgState, + agent_pubkey_hex: &str, + budget: usize, +) -> usize { + if state.resubscribe_retry.is_empty() { + return 0; + } + // Budget-bounded take avoids cloning the full set. + let channels: Vec = state + .resubscribe_retry + .iter() + .copied() + .take(budget) + .collect(); + debug!("draining {} resubscribe_retry channel(s)", channels.len()); + let mut sent_count = 0; + for channel_id in channels { + if let Some(retry_after) = state.check_rate_gate() { + // Gate re-armed mid-drain — move to rate_limited_pending. + state.rate_limited_pending.insert(channel_id, retry_after); + state.resubscribe_retry.remove(&channel_id); + continue; + } + let since = state.channel_since(&channel_id); + let filter = match state.active_filters.get(&channel_id).cloned() { + Some(f) => f, + None => { + warn!("missing filter for channel {channel_id} in resubscribe_retry — dropping"); + state.resubscribe_retry.remove(&channel_id); + continue; + } + }; + let sent = send_subscribe(ws, state, channel_id, agent_pubkey_hex, since, &filter).await; + if sent { + state.resubscribe_retry.remove(&channel_id); + state.channel_dropped_since.remove(&channel_id); + sent_count += 1; + // Pacing is enforced by the main-loop timer; no inline sleep here. + } else { + warn!( + "drain_resubscribe_retry: REQ still failing for channel {channel_id} — will retry" + ); + // Leave in resubscribe_retry; next main-loop tick will try again. + } + } + sent_count } /// Outcome of an autonomous reconnect attempt. @@ -2133,9 +2600,17 @@ async fn try_autonomous_reconnect( // 5 attempts, up to 16s base backoff. Shares delay values with the // initial-connect retry in `HarnessRelay::connect()` (STARTUP_CONNECT_BACKOFFS) — // see its doc comment for how the two loops consume the array differently. + // DNS failures sleep flat (DNS_RETRY_INTERVAL) without consuming a ladder + // rung. Capped at 10 DNS-only retries in this bounded startup path so a + // total brownout cannot hang agent startup indefinitely. By contrast, + // `wait_for_reconnect` (the post-startup loop) retries DNS failures without + // a cap — a reconnecting agent should keep trying across extended outages. let backoffs = STARTUP_CONNECT_BACKOFFS; + const MAX_DNS_FLAT_RETRIES: usize = 10; + let mut dns_retry_count = 0usize; - for (attempt, delay) in backoffs.iter().enumerate() { + let mut attempt = 0usize; + while attempt < backoffs.len() { info!( "autonomous reconnect attempt {}/{} to {relay_url}…", attempt + 1, @@ -2165,23 +2640,44 @@ async fn try_autonomous_reconnect( // Fall through to backoff sleep instead of returning immediately. // Returning false here would skip remaining attempts; continuing // without sleep would drive a tight reconnect storm. - } else if resubscribe_after_reconnect(ws, state, agent_pubkey_hex).await { - return ReconnectOutcome::Ok; } else { - warn!("resubscribe failed after autonomous reconnect — treating as failed attempt"); - // Fall through to backoff sleep and retry. + match resubscribe_after_reconnect(ws, cmd_rx, state, agent_pubkey_hex, true) + .await + { + ResubscribeResult::Ok => return ReconnectOutcome::Ok, + ResubscribeResult::Shutdown => return ReconnectOutcome::Shutdown, + ResubscribeResult::ControlFailed => { + warn!("resubscribe failed after autonomous reconnect — treating as failed attempt"); + // Fall through to backoff sleep and retry. + } + } } } + // DNS failures retry flat without consuming a ladder rung. + // Cap at MAX_DNS_FLAT_RETRIES so a total brownout doesn't hang startup. + Err(e) if is_dns_error(&e) && dns_retry_count < MAX_DNS_FLAT_RETRIES => { + dns_retry_count += 1; + warn!( + "autonomous reconnect DNS failure ({}/{}), flat retry in {:.1}s: {e}", + dns_retry_count, + MAX_DNS_FLAT_RETRIES, + DNS_RETRY_INTERVAL.as_secs_f64() + ); + if !dns_flat_sleep(cmd_rx, state, DNS_RETRY_INTERVAL).await { + return ReconnectOutcome::Shutdown; + } + continue; // retry WITHOUT incrementing attempt + } Err(e) => { warn!("autonomous reconnect attempt {} failed: {e}", attempt + 1); } } - // Backoff sleep between attempts (shared by handshake-drop and connect-error). + // Backoff sleep between ladder attempts (shared by handshake-drop and connect-error). // Skip sleep on the final attempt — we'll fall through to the caller. // Use select! so Shutdown commands are honoured during sleep. if attempt + 1 < backoffs.len() { - let jittered = jittered_duration(*delay); + let jittered = jittered_duration(backoffs[attempt]); tracing::info!( "retrying autonomous reconnect in {:.1}s", jittered.as_secs_f64() @@ -2203,6 +2699,7 @@ async fn try_autonomous_reconnect( } } } + attempt += 1; } ReconnectOutcome::Failed @@ -2242,7 +2739,9 @@ async fn wait_for_reconnect( } // 6 attempts with backoff up to 32s + jitter; uses tokio::select! so shutdown is - // honoured during sleep. + // honoured during sleep. Resumes from state.backoff_step so a flapping link + // keeps its elevated position; the stability block resets it to 0 after 60s. + // DNS failures retry flat without consuming a ladder rung. let backoffs = [ Duration::from_secs(1), Duration::from_secs(2), @@ -2251,8 +2750,7 @@ async fn wait_for_reconnect( Duration::from_secs(16), Duration::from_secs(32), ]; - let mut attempt = 0usize; - let mut delay = Duration::from_secs(1); + let mut attempt = state.backoff_step; loop { info!("attempting relay reconnect to {relay_url}…"); match do_connect(relay_url, keys, auth_tag).await { @@ -2276,25 +2774,54 @@ async fn wait_for_reconnect( // Fall through to the backoff sleep below instead of // tight-looping. A relay that consistently fails the // handshake would otherwise drive a reconnect storm. - } else if resubscribe_after_reconnect(ws, state, agent_pubkey_hex).await { - // Drain any commands that arrived during the final - // do_connect() + resubscribe (which don't poll cmd_rx). - return drain_post_reconnect(ws, cmd_rx, state, agent_pubkey_hex).await; } else { - warn!("resubscribe failed after reconnect — will retry with backoff"); - // Fall through to backoff sleep. + match resubscribe_after_reconnect(ws, cmd_rx, state, agent_pubkey_hex, true) + .await + { + ResubscribeResult::Ok => { + // Drain any commands that arrived during do_connect() + + // resubscribe (which don't poll cmd_rx). + return drain_post_reconnect(ws, cmd_rx, state, agent_pubkey_hex).await; + } + ResubscribeResult::Shutdown => return ReconnectOutcome::Shutdown, + ResubscribeResult::ControlFailed => { + warn!("resubscribe failed after reconnect — will retry with backoff"); + // Fall through to backoff sleep. + } + } } } + // DNS failures retry on a flat interval without consuming a backoff + // ladder rung — the host is temporarily unresolvable, not persistently + // rejecting us, so exponential back-off is counter-productive. + // This loop is unbounded (unlike the 10-retry cap in `try_autonomous_reconnect`) + // so a reconnecting agent keeps trying across extended DNS brownouts. + Err(e) if is_dns_error(&e) => { + warn!("relay reconnect DNS failure (not consuming ladder rung): {e}"); + if !dns_flat_sleep(cmd_rx, state, DNS_RETRY_INTERVAL).await { + return ReconnectOutcome::Shutdown; + } + continue; // retry without incrementing attempt + } Err(e) => { warn!("relay reconnect failed: {e}"); } } + // Persist ladder position before sleeping — if shutdown arrives mid-sleep, + // the next session resumes from here rather than restarting at 0. + state.backoff_step = attempt; + // Backoff sleep — shared by both handshake-drop and connect-error paths. // Uses a deadline so commands processed during the wait don't reset // the timer. Without this, periodic PublishEvent traffic (typing // refresh every 3s) would collapse the jittered backoff into a // reconnect storm. + let delay = if attempt < backoffs.len() { + backoffs[attempt] + } else { + Duration::from_secs(60) + }; let jittered = jittered_duration(delay); warn!("retrying reconnect in {:.1}s", jittered.as_secs_f64()); let deadline = tokio::time::Instant::now() + jittered; @@ -2312,11 +2839,6 @@ async fn wait_for_reconnect( } } attempt += 1; - delay = if attempt < backoffs.len() { - backoffs[attempt] - } else { - Duration::from_secs(60) - }; } } @@ -2492,6 +3014,18 @@ async fn ws_send_timeout( .map_err(|e| RelayError::WebSocket(Box::new(e))) } +/// Parse the relay's `retry in {N}s` hint from a rate-limit message. +/// +/// Accepts any string containing `"retry in "` followed by decimal digits then `'s'`. +/// Returns `None` if the hint is absent; returns `Some(0)` for a literal zero (caller +/// defaults to 5 s). No regex dependency — a simple split is sufficient. +pub(crate) fn parse_rate_limit_retry_secs(msg: &str) -> Option { + let after = msg.split("retry in ").nth(1)?; + // All hint digits are ASCII, so char count == byte count — subslice is valid. + let len = after.chars().take_while(|c| c.is_ascii_digit()).count(); + after[..len].parse::().ok() +} + /// Add ±20% jitter to a backoff duration using the nanosecond sub-second /// component of the system clock as a cheap entropy source (no `rand` dep). fn jittered_duration(base: Duration) -> Duration { @@ -2504,6 +3038,75 @@ fn jittered_duration(base: Duration) -> Duration { base.mul_f64(factor) } +/// Classify a `RelayError` as a DNS resolution failure. +/// +/// Matches the OS-level "name not found" strings surfaced by the platform's +/// resolver, covering macOS (`nodename nor servname`), Linux (`Name or service not +/// known`), and common BSD/Windows variants (`No such host`, +/// `failed to lookup address`). These are transient on brownouts and must NOT +/// consume a backoff ladder rung — they retry on a flat `DNS_RETRY_INTERVAL`. +pub(crate) fn is_dns_error(err: &RelayError) -> bool { + let msg = err.to_string(); + msg.contains("nodename nor servname") + || msg.contains("Name or service not known") + || msg.contains("No such host") + || msg.contains("failed to lookup address") +} + +/// Shutdown-aware fixed-duration sleep for REQ pacing in `resubscribe_after_reconnect`. +/// +/// Unlike `dns_flat_sleep`, no jitter is applied — exact `duration` is required +/// to maintain the ≤8 REQ/s pacing invariant. Non-Shutdown commands received +/// during the sleep are applied to state (matching `dns_flat_sleep` semantics). +/// Returns `true` if sleep completed normally, `false` if shutdown was received. +async fn pacing_sleep( + cmd_rx: &mut mpsc::Receiver, + state: &mut BgState, + duration: Duration, +) -> bool { + let deadline = tokio::time::Instant::now() + duration; + let sleep = tokio::time::sleep_until(deadline); + tokio::pin!(sleep); + loop { + tokio::select! { + _ = &mut sleep => return true, + cmd = cmd_rx.recv() => { + match cmd { + Some(RelayCommand::Shutdown) | None => return false, + Some(cmd) => apply_command_to_state(state, cmd), + } + } + } + } +} + +/// Shutdown-aware sleep used for DNS flat retries. +/// +/// Selects between `duration` elapsing and a `Shutdown`/channel-closed signal on +/// `cmd_rx`. Returns `true` if the sleep completed normally, `false` if the task +/// should shut down. +async fn dns_flat_sleep( + cmd_rx: &mut mpsc::Receiver, + state: &mut BgState, + duration: Duration, +) -> bool { + let jittered = jittered_duration(duration); + let deadline = tokio::time::Instant::now() + jittered; + let sleep = tokio::time::sleep_until(deadline); + tokio::pin!(sleep); + loop { + tokio::select! { + _ = &mut sleep => return true, + cmd = cmd_rx.recv() => { + match cmd { + Some(RelayCommand::Shutdown) | None => return false, + Some(cmd) => apply_command_to_state(state, cmd), + } + } + } + } +} + /// Extract a channel UUID from the h tag of a Nostr event. fn extract_h_tag_uuid(event: &nostr::Event) -> Option { event.tags.iter().find_map(|tag| { @@ -4546,4 +5149,330 @@ mod tests { assert!(result.is_ok()); assert_eq!(attempts.load(Ordering::SeqCst), 2); } + + // ── Rate-limit gate, pacing, backoff reset, DNS ────────────────────────── + + /// parse_rate_limit_retry_secs: full hint extracts the N from "retry in Ns". + #[test] + fn parse_rate_limit_retry_secs_with_hint() { + assert_eq!( + parse_rate_limit_retry_secs("rate-limited: quota exceeded; retry in 12s"), + Some(12) + ); + } + + /// parse_rate_limit_retry_secs: message without a hint returns None. + #[test] + fn parse_rate_limit_retry_secs_missing_hint() { + assert_eq!( + parse_rate_limit_retry_secs("rate-limited: too many concurrent requests"), + None + ); + } + + /// parse_rate_limit_retry_secs: explicit zero value is returned as Some(0). + #[test] + fn parse_rate_limit_retry_secs_zero() { + assert_eq!( + parse_rate_limit_retry_secs("rate-limited: quota exceeded; retry in 0s"), + Some(0) + ); + } + + /// parse_rate_limit_retry_secs: garbage input returns None. + #[test] + fn parse_rate_limit_retry_secs_garbage() { + assert_eq!( + parse_rate_limit_retry_secs("not a rate limit message"), + None + ); + } + + /// set_rate_limit_gate arms the gate with jittered expiry from the hint. + /// check_rate_gate returns Some while active and lazily clears on expiry. + #[tokio::test(start_paused = true)] + async fn rate_limit_gate_set_and_expiry() { + let mut state = BgState::new(); + assert!( + state.check_rate_gate().is_none(), + "gate must start inactive" + ); + + // Arm with a 5 s hint. + state.set_rate_limit_gate(5); + assert!( + state.check_rate_gate().is_some(), + "gate must be active immediately after arming" + ); + + // Advance virtual time past the max jitter (1.2 × 5 s = 6 s). + tokio::time::advance(Duration::from_secs(7)).await; + + assert!( + state.check_rate_gate().is_none(), + "gate must have expired after 7s" + ); + assert!( + state.rate_limit_gate.is_none(), + "check_rate_gate must lazily clear the field on expiry" + ); + } + + /// set_rate_limit_gate takes the max of overlapping deadlines. + #[tokio::test(start_paused = true)] + async fn rate_limit_gate_extends_to_max() { + let mut state = BgState::new(); + + // Arm with a long hint first. + state.set_rate_limit_gate(30); + let first_deadline = state.rate_limit_gate.unwrap(); + + // A shorter subsequent hint must NOT shorten the existing gate. + state.set_rate_limit_gate(1); + let second_deadline = state.rate_limit_gate.unwrap(); + + assert_eq!( + first_deadline, second_deadline, + "shorter hint must not overwrite a later existing deadline" + ); + } + + /// is_dns_error correctly classifies platform resolver strings, including + /// the production shape: a WebSocket I/O error wrapping the OS message. + #[test] + fn is_dns_error_classification() { + use tokio_tungstenite::tungstenite; + + // macOS resolver (Http-wrapped, used in many existing tests) + assert!(is_dns_error(&RelayError::Http( + "nodename nor servname provided, or not known".into() + ))); + // Linux resolver + assert!(is_dns_error(&RelayError::Http( + "Name or service not known".into() + ))); + // BSD/Windows + assert!(is_dns_error(&RelayError::Http("No such host".into()))); + // Another common variant + assert!(is_dns_error(&RelayError::Http( + "failed to lookup address information".into() + ))); + // F15: production-shaped error — RelayError::WebSocket wrapping a + // tungstenite I/O error (the shape emitted by connect_async on macOS). + let ws_io_err = RelayError::WebSocket(Box::new(tungstenite::Error::Io( + std::io::Error::other("nodename nor servname provided, or not known"), + ))); + assert!( + is_dns_error(&ws_io_err), + "WebSocket-wrapped I/O DNS error must be classified as DNS" + ); + // Normal connection errors are NOT DNS errors. + assert!(!is_dns_error(&RelayError::Timeout)); + assert!(!is_dns_error(&RelayError::ConnectionClosed)); + assert!(!is_dns_error(&RelayError::Http( + "connection refused".into() + ))); + } + + /// resubscribe_retry is populated when a channel REQ fails during partial reconnect. + /// + /// This exercises BgState directly since we have no live socket in unit tests. + #[test] + fn resubscribe_retry_populated_on_failure() { + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + + // Subscribe the channel so it ends up in active_subscriptions. + apply_command_to_state( + &mut state, + RelayCommand::Subscribe { + channel_id, + filter: ChannelFilter { + kinds: Some(vec![9]), + require_mention: false, + }, + replay_since: Some(1_000), + }, + ); + assert!(state.active_subscriptions.contains_key(&channel_id)); + + // Simulate a partial-reconnect failure: insert into resubscribe_retry. + state.resubscribe_retry.insert(channel_id); + + assert!( + state.resubscribe_retry.contains(&channel_id), + "failed channel must be in resubscribe_retry" + ); + assert!( + state.active_subscriptions.contains_key(&channel_id), + "channel must stay in active_subscriptions so reconnect can restore it" + ); + } + + // ── Control-sub recovery from rate-limited CLOSED ──────────────────────── + + /// A rate-limited CLOSED for the membership sub sets membership_resub_needed. + /// After the gate expires the drain re-arms the sub and clears the flag. + #[tokio::test(start_paused = true)] + async fn membership_resub_flag_set_on_rate_limited_closed() { + let mut state = BgState::new(); + state.membership_sub_active = true; + + // Simulate a rate-limited CLOSED arriving for the membership sub. + let secs = parse_rate_limit_retry_secs("rate-limited: retry in 5s").unwrap_or(0); + state.set_rate_limit_gate(secs); + state.membership_resub_needed = true; + + assert!( + state.membership_resub_needed, + "flag must be set after rate-limited CLOSED" + ); + assert!( + state.check_rate_gate().is_some(), + "gate must be active while membership sub is pending" + ); + + // Advance past the gate (max jitter: 1.2 × 5s = 6s). + tokio::time::advance(Duration::from_secs(7)).await; + + assert!( + state.check_rate_gate().is_none(), + "gate must expire so drain can fire" + ); + // The drain clears membership_resub_needed after re-sending the REQ. + // Simulate successful re-send: + state.membership_resub_needed = false; + assert!( + !state.membership_resub_needed, + "flag must clear after drain re-sends the membership REQ" + ); + } + + /// A rate-limited CLOSED for the observer control sub sets observer_resub_needed. + #[test] + fn observer_resub_flag_set_on_rate_limited_closed() { + let mut state = BgState::new(); + state.observer_control_sub_active = true; + + // Simulate rate-limited CLOSED on observer control sub. + state.set_rate_limit_gate(5); + state.observer_resub_needed = true; + + assert!( + state.observer_resub_needed, + "flag must be set after rate-limited CLOSED on observer sub" + ); + } + + // ── Drain state transitions ─────────────────────────────────────────────── + + /// drain_rate_limited_pending: a channel re-queued with +5s penalty on send + /// failure stays in pending and is not immediately retried. + #[tokio::test(start_paused = true)] + async fn rate_limited_pending_failure_requeues_with_penalty() { + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + + // Seed the channel's subscription intent. + apply_command_to_state( + &mut state, + RelayCommand::Subscribe { + channel_id, + filter: ChannelFilter { + kinds: Some(vec![9]), + require_mention: false, + }, + replay_since: None, + }, + ); + + // Park the channel as rate-limited with a deadline in the past. + let past = tokio::time::Instant::now(); + state.rate_limited_pending.insert(channel_id, past); + + // Simulate a send failure by re-queuing with +5s (what the drain does). + let penalty = tokio::time::Instant::now() + Duration::from_secs(5); + state.rate_limited_pending.insert(channel_id, penalty); + + assert!( + state.rate_limited_pending.contains_key(&channel_id), + "channel must stay in rate_limited_pending after send failure" + ); + // Deadline should be in the future. + let deadline = state.rate_limited_pending[&channel_id]; + assert!( + deadline > tokio::time::Instant::now(), + "penalty deadline must be in the future" + ); + } + + /// drain_resubscribe_retry: a gate re-armed mid-drain moves the channel to + /// rate_limited_pending and removes it from resubscribe_retry. + #[tokio::test(start_paused = true)] + async fn resubscribe_retry_gate_rearm_moves_to_pending() { + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + + apply_command_to_state( + &mut state, + RelayCommand::Subscribe { + channel_id, + filter: ChannelFilter { + kinds: Some(vec![9]), + require_mention: false, + }, + replay_since: None, + }, + ); + state.resubscribe_retry.insert(channel_id); + + // Simulate gate re-arming mid-drain (what the drain does on check_rate_gate hit). + let retry_after = state.set_rate_limit_gate(5); + state.rate_limited_pending.insert(channel_id, retry_after); + state.resubscribe_retry.remove(&channel_id); + + assert!( + !state.resubscribe_retry.contains(&channel_id), + "channel must be removed from resubscribe_retry when gate re-arms" + ); + assert!( + state.rate_limited_pending.contains_key(&channel_id), + "channel must be moved to rate_limited_pending on gate re-arm" + ); + } + + /// drain_resubscribe_retry: a successful drain removes the channel and + /// clears channel_dropped_since. + #[test] + fn resubscribe_retry_success_clears_state() { + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + + apply_command_to_state( + &mut state, + RelayCommand::Subscribe { + channel_id, + filter: ChannelFilter { + kinds: Some(vec![9]), + require_mention: false, + }, + replay_since: None, + }, + ); + state.resubscribe_retry.insert(channel_id); + state.channel_dropped_since.insert(channel_id, 1_000_000); + + // Simulate successful re-send (what the drain does on success). + state.resubscribe_retry.remove(&channel_id); + state.channel_dropped_since.remove(&channel_id); + + assert!( + !state.resubscribe_retry.contains(&channel_id), + "channel must leave resubscribe_retry on successful drain" + ); + assert!( + !state.channel_dropped_since.contains_key(&channel_id), + "channel_dropped_since must be cleared on successful drain" + ); + } } From f5e9724c5ceba26c13f63643e458a6e741548e01 Mon Sep 17 00:00:00 2001 From: npub1067mpdna4vy22u9eltmmhtdfw56kww6ve2aze025dtfm4pxg07nqztjfhm <7ebdb0b67dab08a570b9faf7bbada97535673b4ccaba2cbd546ad3ba84c87fa6@sprout-oss.stage.blox.sqprod.co> Date: Mon, 20 Jul 2026 23:27:37 -0700 Subject: [PATCH 2/2] fix(acp): preserve relay recovery intent across reconnects Relay admission quotas survive socket replacement, so reconnect replay must keep an active gate rather than burst into the same exhausted window. Commands consumed during paced replay also need their live-socket effects or durable intent preserved if replay fails. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- crates/buzz-acp/src/relay.rs | 485 +++++++++++++++++++++++++++++------ 1 file changed, 406 insertions(+), 79 deletions(-) diff --git a/crates/buzz-acp/src/relay.rs b/crates/buzz-acp/src/relay.rs index cf45329f13..b1497b6be7 100644 --- a/crates/buzz-acp/src/relay.rs +++ b/crates/buzz-acp/src/relay.rs @@ -1188,6 +1188,35 @@ fn apply_command_to_state(state: &mut BgState, cmd: RelayCommand) { } } +/// Retain command intent after a live send failure. +/// +/// Subscription state must survive reconnect; ephemeral publishes are deliberately +/// discarded because replaying a typing indicator after reconnect is meaningless. +/// `Shutdown` and `Reconnect` are handled by the caller. +fn retain_failed_command_intent(state: &mut BgState, cmd: RelayCommand) { + match cmd { + RelayCommand::PublishEvent { .. } => {} + cmd => apply_command_to_state(state, cmd), + } +} + +/// Preserve stateful commands already consumed during replay when that replay +/// loses its live socket before the deferred queue can be executed. +/// +/// Commands are applied in arrival order. Ephemeral publishes are discarded by +/// [`retain_failed_command_intent`], and pacing never queues `Shutdown`. +fn retain_deferred_command_intent( + state: &mut BgState, + deferred_commands: &mut VecDeque, +) { + while let Some(cmd) = deferred_commands.pop_front() { + match cmd { + RelayCommand::Shutdown | RelayCommand::Reconnect => {} + cmd => retain_failed_command_intent(state, cmd), + } + } +} + /// Execute a command on a live WebSocket connection. /// /// Handles the five data commands: Subscribe, Unsubscribe, @@ -1282,10 +1311,16 @@ async fn execute_connected_command( true } RelayCommand::SubscribeMembership => { + state.membership_sub_active = true; + if state.check_rate_gate().is_some() { + debug!("rate-gated: deferring membership subscription"); + state.membership_resub_needed = true; + return true; + } let since = state.membership_last_seen.or(state.startup_watermark); let sent = send_membership_subscribe(ws, agent_pubkey_hex, since).await; if sent { - state.membership_sub_active = true; + state.membership_resub_needed = false; if state.membership_last_seen.is_none() { state.membership_last_seen = since; } @@ -1293,18 +1328,24 @@ async fn execute_connected_command( } else { // Send failed — record intent so reconnect restores it. warn!("membership subscribe REQ failed — recording intent for reconnect"); - state.membership_sub_active = true; + state.membership_resub_needed = true; false } } RelayCommand::SubscribeObserverControls => { + state.observer_control_sub_active = true; + if state.check_rate_gate().is_some() { + debug!("rate-gated: deferring observer control subscription"); + state.observer_resub_needed = true; + return true; + } let sent = send_observer_control_subscribe(ws, agent_pubkey_hex).await; if sent { - state.observer_control_sub_active = true; + state.observer_resub_needed = false; true } else { warn!("observer control subscribe REQ failed — recording intent for reconnect"); - state.observer_control_sub_active = true; + state.observer_resub_needed = true; false } } @@ -1468,7 +1509,7 @@ async fn run_background_task( { ResubscribeResult::Ok => {} ResubscribeResult::Shutdown => return, - ResubscribeResult::ControlFailed => { + ResubscribeResult::RetryConnection => { warn!("proactive resubscribe had failures — triggering reconnect"); let _ = event_tx.try_send(None); match try_autonomous_reconnect( @@ -2267,9 +2308,9 @@ async fn process_handshake_buffer( enum ResubscribeResult { /// All subscriptions restored (or parked for drain recovery). Ok, - /// A control subscription (membership or observer) failed to send. - /// Caller should treat this as a failed reconnect attempt. - ControlFailed, + /// A control subscription or deferred live command failed to send. + /// Caller should retry the connection. + RetryConnection, /// A `Shutdown` command arrived during a pacing sleep. /// Caller must return immediately (background task is exiting). Shutdown, @@ -2280,21 +2321,22 @@ enum ResubscribeResult { /// per channel, and only clears the drop tracker when the REQ is confirmed sent. /// /// Paces REQs at `REQ_PACING_INTERVAL` (125 ms) via a shutdown-aware sleep so -/// a 48-channel reconnect burst spreads over ≈6 s without starving the select! -/// loop of frames, commands, or ping ticks. If the gate is active mid-burst, -/// remaining channels are parked in `rate_limited_pending` instead of sent. +/// a 48-channel reconnect burst spreads over ≈6 s. Commands received during a +/// pacing sleep are deferred in arrival order and executed on the live socket +/// after replay. If the gate is active mid-burst, remaining channels are parked +/// in `rate_limited_pending` instead of sent. /// /// A failed CHANNEL REQ is parked in `resubscribe_retry` rather than failing -/// the whole reconnect. Only membership and observer-control failures return -/// `ControlFailed` — their silent loss breaks joins and agent pause/resume. +/// the whole reconnect. Only membership, observer-control, or deferred-command +/// failures return `RetryConnection` — their silent loss would leave live state +/// inconsistent with command intent. /// -/// When `is_fresh_connection` is `true` (new socket established), the -/// rate-limit gate and pending queues are cleared because a fresh connection -/// resets the relay's admission window. When `false` (proactive resubscribe on -/// the EXISTING socket), they are preserved so parked channels aren't lost. +/// A relay quota gate is keyed by community and pubkey, so replacing the socket +/// does not reset it. Fresh connections may clear derived pending/retry queues +/// before rebuilding them from `active_subscriptions`, but the gate itself is +/// always preserved until its deadline expires. /// -/// Returns [`ResubscribeResult`] signalling success, control-sub failure, or -/// shutdown. +/// Returns [`ResubscribeResult`] signalling success, retry, or shutdown. async fn resubscribe_after_reconnect( ws: &mut WsStream, cmd_rx: &mut mpsc::Receiver, @@ -2302,24 +2344,21 @@ async fn resubscribe_after_reconnect( agent_pubkey_hex: &str, is_fresh_connection: bool, ) -> ResubscribeResult { - // Only clear stale gate state on a FRESH connection. A proactive resubscribe - // on the existing socket must preserve the gate and pending queues — - // clearing them would discard parked channels and burst into a still-closed - // admission window. if is_fresh_connection { - state.rate_limit_gate = None; + // These queues are derived from active subscription intent and rebuilt + // below. The rate-limit gate is deliberately preserved: the relay's + // shared admission counter survives socket replacement. state.rate_limited_pending.clear(); state.resubscribe_retry.clear(); - // Control-sub pending flags are implicitly cleared by the resub below. } + let mut deferred_commands = VecDeque::new(); let channels: Vec = state.active_subscriptions.keys().copied().collect(); if !channels.is_empty() { info!( "resubscribing to {} channel(s) after reconnect", channels.len() ); - let mut sent_any = false; for channel_id in channels { // Gate re-armed mid-burst — park remaining channels. if let Some(retry_after) = state.check_rate_gate() { @@ -2346,9 +2385,8 @@ async fn resubscribe_after_reconnect( send_subscribe(ws, state, channel_id, agent_pubkey_hex, since, &filter).await; if this_sent { state.channel_dropped_since.remove(&channel_id); - sent_any = true; - // Shutdown-aware pacing sleep between REQs. - if !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await { + // Shutdown-aware pacing sleep before any next replay/deferred REQ. + if !pacing_sleep(cmd_rx, &mut deferred_commands, REQ_PACING_INTERVAL).await { return ResubscribeResult::Shutdown; } } else { @@ -2360,46 +2398,61 @@ async fn resubscribe_after_reconnect( state.resubscribe_retry.insert(channel_id); } } - let _ = sent_any; // consumed above for pacing logic } // Membership and observer-control are control-plane subscriptions: a silent - // failure breaks join notifications and agent pause/resume. These still force - // a full reconnect retry rather than silent degradation. + // failure breaks join notifications and agent pause/resume. A shared quota + // gate parks their intent for the main-loop drain just like channel REQs. if state.membership_sub_active { - if !state.active_subscriptions.is_empty() - && !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await - { - return ResubscribeResult::Shutdown; - } - let replay_since = match (state.membership_dropped_since, state.membership_last_seen) { - (Some(d), Some(l)) => Some(d.min(l)), - (Some(d), None) => Some(d), - (None, Some(l)) => Some(l), - (None, None) => state.startup_watermark, - }; - let sent = send_membership_subscribe(ws, agent_pubkey_hex, replay_since).await; - if sent { - state.membership_dropped_since = None; - state.membership_resub_needed = false; + if state.check_rate_gate().is_some() { + debug!("rate-gated: parking membership resubscribe after reconnect"); + state.membership_resub_needed = true; } else { - warn!("failed to resubscribe membership after reconnect"); - return ResubscribeResult::ControlFailed; + if !state.active_subscriptions.is_empty() + && !pacing_sleep(cmd_rx, &mut deferred_commands, REQ_PACING_INTERVAL).await + { + return ResubscribeResult::Shutdown; + } + let replay_since = match (state.membership_dropped_since, state.membership_last_seen) { + (Some(d), Some(l)) => Some(d.min(l)), + (Some(d), None) => Some(d), + (None, Some(l)) => Some(l), + (None, None) => state.startup_watermark, + }; + let sent = send_membership_subscribe(ws, agent_pubkey_hex, replay_since).await; + if sent { + state.membership_dropped_since = None; + state.membership_resub_needed = false; + } else { + warn!("failed to resubscribe membership after reconnect"); + retain_deferred_command_intent(state, &mut deferred_commands); + return ResubscribeResult::RetryConnection; + } } } if state.observer_control_sub_active { - if !pacing_sleep(cmd_rx, state, REQ_PACING_INTERVAL).await { - return ResubscribeResult::Shutdown; - } - if !send_observer_control_subscribe(ws, agent_pubkey_hex).await { - warn!("failed to resubscribe observer controls after reconnect"); - return ResubscribeResult::ControlFailed; + if state.check_rate_gate().is_some() { + debug!("rate-gated: parking observer control resubscribe after reconnect"); + state.observer_resub_needed = true; + } else { + if !pacing_sleep(cmd_rx, &mut deferred_commands, REQ_PACING_INTERVAL).await { + return ResubscribeResult::Shutdown; + } + if !send_observer_control_subscribe(ws, agent_pubkey_hex).await { + warn!("failed to resubscribe observer controls after reconnect"); + retain_deferred_command_intent(state, &mut deferred_commands); + return ResubscribeResult::RetryConnection; + } + state.observer_resub_needed = false; } - state.observer_resub_needed = false; } - ResubscribeResult::Ok + match drain_commands(ws, cmd_rx, &mut deferred_commands, state, agent_pubkey_hex).await { + ReconnectOutcome::Ok => ResubscribeResult::Ok, + ReconnectOutcome::Failed => ResubscribeResult::RetryConnection, + ReconnectOutcome::Shutdown => ResubscribeResult::Shutdown, + } } /// Drain `rate_limited_pending` channels whose retry deadline has passed. @@ -2522,29 +2575,38 @@ async fn drain_resubscribe_retry( enum ReconnectOutcome { /// Reconnected and resubscribed successfully. Ok, - /// All attempts exhausted — caller should fall back to wait_for_reconnect. + /// Reconnect or resubscription attempts failed; caller should retry or fall + /// back to `wait_for_reconnect`. Live command intent is retained. Failed, /// A Shutdown command was received during backoff — caller must return immediately. Shutdown, } -/// Drain all pending commands after a successful reconnect. -/// -/// Processes queued commands that arrived while reconnecting. Reconnect -/// commands are silently dropped (already reconnected). Shutdown causes an -/// immediate close-frame + return of `ReconnectOutcome::Shutdown`. All other -/// commands are executed on the live socket via [`execute_connected_command`]. -/// If any send fails, remaining commands are recorded as intent via -/// [`apply_command_to_state`] and the drain continues (the caller's next -/// read/ping will detect the dead socket). -async fn drain_post_reconnect( +/// Execute commands deferred during paced replay, then commands that arrived +/// while the deferred queue was draining. FIFO order is preserved across both +/// sources. Subscription REQs are paced; CLOSE, ephemeral EVENT, and local-state +/// commands execute immediately. A failed live send records remaining command +/// intent and returns `Failed`; Shutdown closes the socket immediately. +async fn drain_commands( ws: &mut WsStream, cmd_rx: &mut mpsc::Receiver, + deferred_commands: &mut VecDeque, state: &mut BgState, agent_pubkey_hex: &str, ) -> ReconnectOutcome { let mut send_failed = false; - while let Ok(cmd) = cmd_rx.try_recv() { + loop { + let cmd = match deferred_commands.pop_front() { + Some(cmd) => cmd, + None => match cmd_rx.try_recv() { + Ok(cmd) => cmd, + Err(mpsc::error::TryRecvError::Empty) => break, + Err(mpsc::error::TryRecvError::Disconnected) => { + return ReconnectOutcome::Shutdown; + } + }, + }; + if send_failed { match cmd { RelayCommand::Shutdown => { @@ -2552,10 +2614,11 @@ async fn drain_post_reconnect( return ReconnectOutcome::Shutdown; } RelayCommand::Reconnect => {} - cmd => apply_command_to_state(state, cmd), + cmd => retain_failed_command_intent(state, cmd), } continue; } + match cmd { RelayCommand::Reconnect => { debug!("drained stale Reconnect after reconnect"); @@ -2565,16 +2628,54 @@ async fn drain_post_reconnect( let _ = ws_send_timeout(ws, Message::Close(None), WS_SEND_TIMEOUT_SECS).await; return ReconnectOutcome::Shutdown; } + RelayCommand::Subscribe { .. } + | RelayCommand::SubscribeMembership + | RelayCommand::SubscribeObserverControls => { + // A gated subscription is only parked in state; pace only an + // actual live send attempt. + let pace_after = state.check_rate_gate().is_none(); + if !execute_connected_command(ws, state, agent_pubkey_hex, cmd).await { + warn!("send failed during post-reconnect drain — recording remaining commands as intent"); + send_failed = true; + } + if !send_failed + && pace_after + && !pacing_sleep(cmd_rx, deferred_commands, REQ_PACING_INTERVAL).await + { + return ReconnectOutcome::Shutdown; + } + } cmd => { - let ok = execute_connected_command(ws, state, agent_pubkey_hex, cmd).await; - if !ok { + if !execute_connected_command(ws, state, agent_pubkey_hex, cmd).await { warn!("send failed during post-reconnect drain — recording remaining commands as intent"); send_failed = true; } } } } - ReconnectOutcome::Ok + + if send_failed { + ReconnectOutcome::Failed + } else { + ReconnectOutcome::Ok + } +} + +/// Drain all pending commands after a successful reconnect. +/// +/// Processes queued commands that arrived while reconnecting. Reconnect +/// commands are silently dropped (already reconnected). Shutdown causes an +/// immediate close-frame + return of `ReconnectOutcome::Shutdown`. All other +/// commands are executed on the live socket via [`execute_connected_command`]. +/// If any subscription send fails, remaining commands are recorded as intent +/// and `Failed` is returned so the caller can reconnect. +async fn drain_post_reconnect( + ws: &mut WsStream, + cmd_rx: &mut mpsc::Receiver, + state: &mut BgState, + agent_pubkey_hex: &str, +) -> ReconnectOutcome { + drain_commands(ws, cmd_rx, &mut VecDeque::new(), state, agent_pubkey_hex).await } /// Attempt autonomous reconnect on socket loss. @@ -2646,7 +2747,7 @@ async fn try_autonomous_reconnect( { ResubscribeResult::Ok => return ReconnectOutcome::Ok, ResubscribeResult::Shutdown => return ReconnectOutcome::Shutdown, - ResubscribeResult::ControlFailed => { + ResubscribeResult::RetryConnection => { warn!("resubscribe failed after autonomous reconnect — treating as failed attempt"); // Fall through to backoff sleep and retry. } @@ -2784,7 +2885,7 @@ async fn wait_for_reconnect( return drain_post_reconnect(ws, cmd_rx, state, agent_pubkey_hex).await; } ResubscribeResult::Shutdown => return ReconnectOutcome::Shutdown, - ResubscribeResult::ControlFailed => { + ResubscribeResult::RetryConnection => { warn!("resubscribe failed after reconnect — will retry with backoff"); // Fall through to backoff sleep. } @@ -3057,11 +3158,12 @@ pub(crate) fn is_dns_error(err: &RelayError) -> bool { /// /// Unlike `dns_flat_sleep`, no jitter is applied — exact `duration` is required /// to maintain the ≤8 REQ/s pacing invariant. Non-Shutdown commands received -/// during the sleep are applied to state (matching `dns_flat_sleep` semantics). -/// Returns `true` if sleep completed normally, `false` if shutdown was received. +/// during the sleep are deferred in arrival order for live execution after +/// replay. Returns `true` if sleep completed normally, `false` if shutdown was +/// received. async fn pacing_sleep( cmd_rx: &mut mpsc::Receiver, - state: &mut BgState, + deferred_commands: &mut VecDeque, duration: Duration, ) -> bool { let deadline = tokio::time::Instant::now() + duration; @@ -3073,7 +3175,7 @@ async fn pacing_sleep( cmd = cmd_rx.recv() => { match cmd { Some(RelayCommand::Shutdown) | None => return false, - Some(cmd) => apply_command_to_state(state, cmd), + Some(cmd) => deferred_commands.push_back(cmd), } } } @@ -4015,6 +4117,231 @@ mod tests { .expect("signing should succeed") } + async fn test_ws_pair() -> (WsStream, WebSocketStream) { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind test websocket"); + let address = listener.local_addr().expect("read test address"); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.expect("accept test websocket"); + tokio_tungstenite::accept_async(stream) + .await + .expect("complete server websocket handshake") + }); + let (client, _) = connect_async(format!("ws://{address}")) + .await + .expect("connect test websocket"); + (client, server.await.expect("join test websocket server")) + } + + async fn next_test_frame( + server: &mut WebSocketStream, + ) -> serde_json::Value { + let message = timeout(Duration::from_secs(1), server.next()) + .await + .expect("timed out waiting for websocket frame") + .expect("test websocket closed") + .expect("read test websocket frame"); + serde_json::from_str(message.to_text().expect("expected text frame")) + .expect("parse test websocket frame") + } + + fn test_channel_filter() -> ChannelFilter { + ChannelFilter { + kinds: Some(vec![9]), + require_mention: false, + } + } + + fn seed_test_subscription(state: &mut BgState, channel_id: Uuid) { + apply_command_to_state( + state, + RelayCommand::Subscribe { + channel_id, + filter: test_channel_filter(), + replay_since: Some(1_000), + }, + ); + } + + #[tokio::test] + async fn fresh_reconnect_preserves_gate_until_pending_replay_resumes() { + let (mut client, mut server) = test_ws_pair().await; + let (_cmd_tx, mut cmd_rx) = mpsc::channel(1); + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + seed_test_subscription(&mut state, channel_id); + state.rate_limit_gate = Some(tokio::time::Instant::now() + Duration::from_millis(150)); + + let result = + resubscribe_after_reconnect(&mut client, &mut cmd_rx, &mut state, "agent-pubkey", true) + .await; + + assert!(matches!(result, ResubscribeResult::Ok)); + assert!(state.rate_limit_gate.is_some()); + assert!(state.rate_limited_pending.contains_key(&channel_id)); + assert!( + timeout(Duration::from_millis(50), server.next()) + .await + .is_err(), + "fresh reconnect must not send REQ while the shared quota gate is active" + ); + + tokio::time::sleep(Duration::from_millis(125)).await; + assert_eq!( + drain_rate_limited_pending(&mut client, &mut state, "agent-pubkey", 1).await, + 1 + ); + let frame = next_test_frame(&mut server).await; + assert_eq!(frame[0], "REQ"); + assert_eq!(frame[1], channel_sub_id(channel_id)); + } + + #[tokio::test] + async fn subscribe_during_replay_pacing_is_sent_on_live_socket() { + let (client, mut server) = test_ws_pair().await; + let (cmd_tx, mut cmd_rx) = mpsc::channel(4); + let mut state = BgState::new(); + let replayed_channel = Uuid::new_v4(); + let deferred_channel = Uuid::new_v4(); + seed_test_subscription(&mut state, replayed_channel); + + let task = tokio::spawn(async move { + let mut client = client; + let result = resubscribe_after_reconnect( + &mut client, + &mut cmd_rx, + &mut state, + "agent-pubkey", + true, + ) + .await; + (result, state) + }); + + let replay = next_test_frame(&mut server).await; + assert_eq!(replay[1], channel_sub_id(replayed_channel)); + cmd_tx + .send(RelayCommand::Subscribe { + channel_id: deferred_channel, + filter: test_channel_filter(), + replay_since: Some(2_000), + }) + .await + .expect("queue subscribe during pacing"); + + let deferred = next_test_frame(&mut server).await; + assert_eq!(deferred[0], "REQ"); + assert_eq!(deferred[1], channel_sub_id(deferred_channel)); + let (result, state) = task.await.expect("join resubscribe task"); + assert!(matches!(result, ResubscribeResult::Ok)); + assert!(state.active_subscriptions.contains_key(&deferred_channel)); + } + + #[tokio::test] + async fn unsubscribe_during_replay_pacing_sends_close_on_live_socket() { + let (client, mut server) = test_ws_pair().await; + let (cmd_tx, mut cmd_rx) = mpsc::channel(4); + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + seed_test_subscription(&mut state, channel_id); + + let task = tokio::spawn(async move { + let mut client = client; + let result = resubscribe_after_reconnect( + &mut client, + &mut cmd_rx, + &mut state, + "agent-pubkey", + true, + ) + .await; + (result, state) + }); + + let replay = next_test_frame(&mut server).await; + assert_eq!(replay[1], channel_sub_id(channel_id)); + cmd_tx + .send(RelayCommand::Unsubscribe { channel_id }) + .await + .expect("queue unsubscribe during pacing"); + + let close = next_test_frame(&mut server).await; + assert_eq!(close, json!(["CLOSE", channel_sub_id(channel_id)])); + let (result, state) = task.await.expect("join resubscribe task"); + assert!(matches!(result, ResubscribeResult::Ok)); + assert!(!state.active_subscriptions.contains_key(&channel_id)); + } + + #[tokio::test] + async fn publish_during_replay_pacing_is_sent_on_live_socket() { + let (client, mut server) = test_ws_pair().await; + let (cmd_tx, mut cmd_rx) = mpsc::channel(4); + let mut state = BgState::new(); + let channel_id = Uuid::new_v4(); + seed_test_subscription(&mut state, channel_id); + let event = make_test_event(&nostr::Keys::generate(), 2_000); + let event_id = event.id.to_hex(); + + let task = tokio::spawn(async move { + let mut client = client; + let result = resubscribe_after_reconnect( + &mut client, + &mut cmd_rx, + &mut state, + "agent-pubkey", + true, + ) + .await; + result + }); + + let replay = next_test_frame(&mut server).await; + assert_eq!(replay[1], channel_sub_id(channel_id)); + cmd_tx + .send(RelayCommand::PublishEvent { + event: Box::new(event), + }) + .await + .expect("queue publish during pacing"); + + let publish = next_test_frame(&mut server).await; + assert_eq!(publish[0], "EVENT"); + assert_eq!(publish[1]["id"], event_id); + assert!(matches!( + task.await.expect("join resubscribe task"), + ResubscribeResult::Ok + )); + } + + #[test] + fn failed_replay_retains_deferred_subscription_intent_in_fifo_order() { + let mut state = BgState::new(); + let kept_channel = Uuid::new_v4(); + let removed_channel = Uuid::new_v4(); + seed_test_subscription(&mut state, removed_channel); + let event = make_test_event(&nostr::Keys::generate(), 2_000); + let mut deferred = VecDeque::from([ + RelayCommand::Subscribe { + channel_id: kept_channel, + filter: test_channel_filter(), + replay_since: Some(2_000), + }, + RelayCommand::Unsubscribe { + channel_id: removed_channel, + }, + RelayCommand::PublishEvent { + event: Box::new(event), + }, + ]); + + retain_deferred_command_intent(&mut state, &mut deferred); + + assert!(deferred.is_empty()); + assert!(state.active_subscriptions.contains_key(&kept_channel)); + assert!(!state.active_subscriptions.contains_key(&removed_channel)); + } + #[test] fn bg_state_dedup_first_event_accepted() { let mut state = BgState::new();