diff --git a/CHANGELOG.md b/CHANGELOG.md index c16c06b..039886b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Enforce outbound relay rate-limit cooldowns after queued requests finish + waiting for pacing or subscription capacity. Process rate-limit signals even + when sync event delivery is blocked, and log one warning per cooldown episode. + ## [3.0.3] - 2026-09-18 This patch improves repository-state correctness, startup recovery, and relay diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 5bb77e0..5053491 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -503,10 +503,19 @@ So concurrency is not a free scaling axis; it is the residual of the ledger: IDs as metric labels. Live subscriptions are ledgered first, so NEG, transient REQ, and purgatory exact-ID polling share only the remaining capacity. -- Permit acquisition checks relay health first: while a rate-limit or - transient-failure cooldown is active, queued rounds take the REQ+EOSE - fallback path (which is itself budget-accounted) instead of firing into a - relay that just complained. +- The scheduler and outbound connection share the relay's rate-limit state. + NOTICE and rate-limited CLOSED messages update that state on the independent + control stream, even when event delivery to the processor is blocked. New + work checks it before pacing, then checks again after pacing and capacity + waits immediately before handing a REQ, direct fetch, or new NIP-77 round + to the SDK. Queued work rejected during the cooldown releases its permits + and returns an error for the existing incomplete-work recovery paths. + Filter-count and retained-subscription-byte refusals keep their specialized + recovery paths rather than activating this generic cooldown. Requests + already handed to the SDK and SDK-owned protocol continuations are not + cancelled by this check; existing live subscriptions remain open. + Negentropy's separate transient-failure cooldown still directs later rounds + to the budget-accounted REQ+EOSE fallback. - The reactive machinery (escalating cooldown, rate-limit detection in both NOTICE and subscription-specific CLOSED messages, and per-batch fallback) remains the backstop for relays whose limits are below our floors. A diff --git a/src/sync/mod.rs b/src/sync/mod.rs index fb4972f..85d2fe5 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -4801,15 +4801,15 @@ impl SyncManager { } RelayEvent::Notice(notice) => { if is_rate_limit_message(¬ice) { - tracing::warn!( + // The connection control lane records the cooldown; + // replaying it here could restart an expired episode + // after a delayed event queue drains. + tracing::debug!( relay = %relay_url_clone, notice = %notice, "Rate limiting NOTICE detected from relay" ); - // Mark relay as rate limited - health_tracker.record_rate_limit(&relay_url_clone); - // Update metrics with new health state if let Some(ref metrics) = metrics_clone { let state = health_tracker.get_state(&relay_url_clone); @@ -4843,21 +4843,11 @@ impl SyncManager { && !is_filter_count_refusal(&reason) && subscription_state_byte_limit(&reason).is_none() { - let already_paused = health_tracker.is_rate_limited(&relay_url_clone); - if already_paused { - tracing::debug!( - relay = %relay_url_clone, - reason = %reason, - "Repeated rate-limiting CLOSED during active cooldown" - ); - } else { - tracing::info!( - relay = %relay_url_clone, - reason = %reason, - "Rate limiting CLOSED detected from relay" - ); - } - health_tracker.record_rate_limit(&relay_url_clone); + tracing::debug!( + relay = %relay_url_clone, + reason = %reason, + "Rate limiting CLOSED detected from relay" + ); if let Some(ref metrics) = metrics_clone { metrics.record_health_state( &relay_url_clone, @@ -5490,7 +5480,8 @@ impl SyncManager { keys, source, policy, - ); + ) + .with_health_tracker(Arc::clone(&self.health_tracker)); self.connections.insert(relay_url.clone(), connection); if nip65_discovery_only { self.nip65_discovery_only_relays.insert(relay_url.clone()); diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 33022b9..8336239 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -21,7 +21,7 @@ use std::future::Future; use std::time::Duration; use tokio::sync::mpsc; -use super::health::RATE_LIMIT_COOLDOWN_SECS; +use super::health::{RelayHealthTracker, RATE_LIMIT_COOLDOWN_SECS}; use super::{is_filter_count_refusal, is_rate_limit_message, subscription_state_byte_limit}; use crate::nostr::SharedDatabase; use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource}; @@ -586,6 +586,8 @@ pub struct RelayConnection { source: RelayTargetSource, /// Outbound target policy applied before every event-directed dial policy: OutboundTargetPolicy, + /// Shared with the scheduler so queued work observes the same cooldown. + health_tracker: std::sync::Arc, /// The underlying nostr-sdk client client: Client, /// Whether a NIP-42 authenticator was attached to the client. Without @@ -638,7 +640,38 @@ pub struct RelayConnection { } impl RelayConnection { + pub(super) fn with_health_tracker( + mut self, + health_tracker: std::sync::Arc, + ) -> Self { + self.health_tracker = health_tracker; + self + } + + fn ensure_query_admitted(&self) -> Result<(), String> { + if self.health_tracker.is_rate_limited(&self.url) { + return Err(format!( + "Query start deferred for {}: rate-limit cooldown active", + self.url + )); + } + Ok(()) + } + + fn record_subscription_rate_limit(&self, message: &str) { + if is_rate_limit_message(message) + && !is_filter_count_refusal(message) + && subscription_state_byte_limit(message).is_none() + { + self.health_tracker.record_rate_limit(&self.url); + } + if is_query_rate_limit_message(message) { + self.record_query_rate_limit(); + } + } + async fn await_query_start(&self) -> Result<(), String> { + self.ensure_query_admitted()?; self.query_start_pacer .wait_for_start() .await @@ -738,6 +771,7 @@ impl RelayConnection { source, policy, client, + health_tracker: std::sync::Arc::new(RelayHealthTracker::with_defaults()), has_authenticator, database: None, nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new( @@ -812,6 +846,7 @@ impl RelayConnection { source, policy, client, + health_tracker: std::sync::Arc::new(RelayHealthTracker::with_defaults()), has_authenticator, database: Some(database), nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new( @@ -1696,8 +1731,8 @@ impl RelayConnection { let mut notifications = relay.notifications(); // EVENT forwarding can block on the bounded processor queue. Subscribe - // independently to terminal messages so EOSE/CLOSED can still return - // transient ledger slots while the data lane drains. Relay + // independently to control messages so rate-limit cooldowns take + // effect and EOSE/CLOSED return slots while the data lane drains. Relay // notifications are broadcast, so this does not consume messages from // the processor-facing stream below. let terminal_connection = self.clone(); @@ -1708,6 +1743,9 @@ impl RelayConnection { while let Some(notification) = terminals.next().await { match notification { RelayNotification::Message { message } => match *message { + RelayMessage::Notice(message) => { + terminal_connection.record_subscription_rate_limit(&message); + } RelayMessage::EndOfStoredEvents(sub_id) => { terminal_connection .close_and_release_transient_req_permit(&sub_id) @@ -1717,6 +1755,7 @@ impl RelayConnection { subscription_id, message, } => { + terminal_connection.record_subscription_rate_limit(&message); let subscription_id = subscription_id.into_owned(); // An auth-required CLOSED is only a retry signal // when an authenticator can actually answer the @@ -1794,9 +1833,6 @@ impl RelayConnection { } } RelayMessage::Notice(msg) => { - if is_query_rate_limit_message(&msg) { - self.record_query_rate_limit(); - } // Check if this is a negentropy-related notice let is_negentropy_notice = msg.contains("envelope") || msg.contains("NEG-") @@ -1847,9 +1883,6 @@ impl RelayConnection { } else { self.release_live_req_permit(&subscription_id) }; - if is_query_rate_limit_message(&msg) { - self.record_query_rate_limit(); - } if let Some(limit) = self.record_subscription_byte_limit(&msg) { tracing::warn!( relay = %url, @@ -2113,6 +2146,9 @@ impl RelayConnection { } // Transient permits are acquired through the same session helper; a // reset closes queued acquisitions and the helper checks generation. + // Recheck after every pacing/capacity wait, before registering ownership + // or handing a new request to the SDK. Rejected work drops its permits. + self.ensure_query_admitted()?; let retained_filters = filters.clone(); // The relay can answer an empty or cached query before `subscribe` // returns. Register transient ownership against a caller-chosen ID @@ -2239,6 +2275,7 @@ impl RelayConnection { let subscription_id = SubscriptionId::generate(); let mut notifications = relay.notifications(); let fetch = async { + self.ensure_query_admitted()?; let mut stream = relay .stream_events(filter) .with_id(subscription_id.clone()) @@ -2758,11 +2795,15 @@ impl RelayConnection { )); } + // A queued refusal is not a failed NIP-77 exchange. Return before + // entering the round's failure accounting. + self.ensure_query_admitted()?; // Use dry_run to only identify differences without downloading events self.ensure_current_session(&ledger_slot)?; let sync_opts = SyncOptions::default().dry_run(); let client = self.client.clone(); let sync_task = async move { + self.ensure_query_admitted()?; match client.sync(filter).opts(sync_opts).await { Ok(output) => { if !output.failed.is_empty() { @@ -3433,6 +3474,203 @@ mod tests { relay.shutdown(); } + #[tokio::test] + async fn notice_cooldown_rejects_queued_requests_before_sdk_send() { + use futures_util::SinkExt; + use tokio_tungstenite::tungstenite::Message; + + // Each request is polled into a known queue before a real wire NOTICE. + // The processor queue stays full: only the control lane can pause it. + for kind in [ + "historic", + "historic-capacity", + "fetch", + "negentropy", + "live", + ] { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let (notice_tx, notice_rx) = tokio::sync::oneshot::channel(); + let (request_tx, mut request_rx) = tokio::sync::mpsc::unbounded_channel(); + let server = tokio::spawn(async move { + let (socket, _) = listener.accept().await.unwrap(); + let mut ws = tokio_tungstenite::accept_async(socket).await.unwrap(); + notice_rx.await.unwrap(); + for message in ["queue blocker", "ERROR: too many concurrent REQs"] { + ws.send(Message::Text( + serde_json::json!(["NOTICE", message]).to_string().into(), + )) + .await + .unwrap(); + } + let mut requests = 0; + while let Some(frame) = ws.next().await { + let frame = frame.unwrap(); + if frame.is_close() { + break; + } + if !frame.is_text() { + continue; + } + let request: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + if request[0] == "REQ" || request[0] == "NEG-OPEN" { + requests += 1; + request_tx.send(()).unwrap(); + } + } + requests + }); + let health = std::sync::Arc::new(RelayHealthTracker::with_defaults()); + let connection = + permissive_connection(&url, Keys::generate()).with_health_tracker(health.clone()); + connection.connect(3).await.unwrap(); + let (event_tx, _event_rx) = tokio::sync::mpsc::channel(1); + event_tx + .send(RelayEvent::Notice("full queue".into())) + .await + .unwrap(); + let mut event_loop = Box::pin(connection.clone().run_event_loop(event_tx)); + assert!(futures_util::poll!(event_loop.as_mut()).is_pending()); + let event_loop = tokio::spawn(event_loop); + + let mut background_gate = Some(connection.background_query_pacer.gate.lock().await); + let query_gate = if kind == "live" { + Some(connection.query_start_pacer.gate.lock().await) + } else { + None + }; + let slots = connection + .subscription_usable_slots + .load(std::sync::atomic::Ordering::Relaxed); + let class_capacity = if kind == "historic-capacity" { + Some( + connection + .transient_req_permits + .acquire_many(MAX_CONCURRENT_TRANSIENT_REQS as u32) + .await + .unwrap(), + ) + } else { + None + }; + let capacity = if kind == "live" || kind == "historic-capacity" { + None + } else { + Some( + connection + .acquire_subscription_slots(slots as u32) + .await + .unwrap(), + ) + }; + if kind != "historic" { + drop(background_gate.take()); + } + let filter = Filter::new().kind(Kind::TextNote); + let mut request: std::pin::Pin< + Box> + Send + '_>, + > = match kind { + "historic" | "historic-capacity" => Box::pin(async { + connection + .subscribe_filter(filter, TransientRequestClass::HistoricPage) + .await + .map(|_| ()) + }), + "fetch" => Box::pin(async { + connection + .fetch_events(filter, Duration::from_secs(2)) + .await + .map(|_| ()) + }), + "negentropy" => { + Box::pin(async { connection.negentropy_sync_diff(filter).await.map(|_| ()) }) + } + "live" => Box::pin(async { + connection + .subscribe_live_filter_groups(vec![vec![filter]]) + .await + .map(|_| ()) + }), + _ => unreachable!(), + }; + assert!( + futures_util::poll!(request.as_mut()).is_pending(), + "{kind} must queue" + ); + notice_tx.send(()).unwrap(); + tokio::time::timeout(Duration::from_secs(3), async { + while !health.is_rate_limited(&connection.url) { + tokio::task::yield_now().await; + } + }) + .await + .expect("control lane must observe NOTICE behind blocked processor"); + assert!( + connection.query_start_pacer.interval().is_none(), + "concurrency refusal must not learn a query frequency" + ); + drop(capacity); + drop(class_capacity); + drop(query_gate); + drop(background_gate); + let result = tokio::time::timeout(Duration::from_secs(3), request).await; + let error = result + .expect("queued work must unwind") + .expect_err("cooldown must prevent send"); + assert!(error.contains("rate-limit cooldown"), "{kind}: {error}"); + assert_eq!( + connection + .nip77_transient_failures + .load(std::sync::atomic::Ordering::Relaxed), + 0, + "refused admission is not a failed NIP-77 exchange" + ); + assert!(connection + .transient_req_permits_held + .lock() + .unwrap() + .is_empty()); + assert_eq!( + connection.transient_req_permits.available_permits(), + MAX_CONCURRENT_TRANSIENT_REQS + ); + assert_eq!( + connection + .subscription_budget + .lock() + .unwrap() + .semaphore + .available_permits(), + slots + ); + + assert!( + request_rx.try_recv().is_err(), + "no request sent during cooldown" + ); + health.clear_rate_limit(&connection.url); + connection + .subscribe_live_filter_groups(vec![vec![Filter::new().kind(Kind::TextNote)]]) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(3), request_rx.recv()) + .await + .unwrap() + .unwrap(); + connection.disconnect().await; + let count = tokio::time::timeout(Duration::from_secs(3), server) + .await + .unwrap() + .unwrap(); + assert_eq!( + count, 1, + "only the post-cooldown request may reach the wire ({kind})" + ); + event_loop.abort(); + } + } + #[tokio::test] async fn query_rate_closed_paces_later_wire_requests() { let relay = TestRelay::start(LocalRelayBuilder::default().queries_per_minute(1)).await;