fix(sync): enforce relay cooldown after queued request waits

A concurrency-limit NOTICE paused scheduler admission but requests already waiting for pacing or subscription permits could still reach the relay. Processing the NOTICE through the event queue could also delay that pause under backpressure.

Share relay health between the scheduler and connection. Record rate-limit signals on the independent control stream and recheck admission after waits before handing REQs, direct fetches, or new NIP-77 rounds to the SDK. Release refused reservations normally, retain existing incomplete-work recovery, and avoid counting queued NIP-77 refusals as failed rounds. Leave filter-count and byte-budget refusals on their specialized recovery paths.

Only the control stream records wire cooldown signals, avoiding delayed duplicate episodes from the processor queue. Keep one warning per episode and repeated NOTICE details at debug level. Already handed-off requests, SDK-owned continuations, and investigation of the original remote concurrency mismatch are outside this change.

Validation: 401 sync-related library tests passed; the final wire regression passed across background pacing, transient-class and shared-ledger waits, direct fetches, NIP-77, and live requests. It blocks the processor queue, verifies no refused request reaches the wire or leaks permits, and observes a successful post-cooldown request. cargo clippy --lib -- -D warnings, cargo fmt --all -- --check, and git diff --check passed. No deployment was performed.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-09-18 14:56:50 +00:00
parent 5ef7480a01
commit fa919afc23
4 changed files with 277 additions and 33 deletions
+6
View File
@@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased] ## [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 ## [3.0.3] - 2026-09-18
This patch improves repository-state correctness, startup recovery, and relay This patch improves repository-state correctness, startup recovery, and relay
+13 -4
View File
@@ -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, IDs as metric labels. Live subscriptions are ledgered first, so NEG,
transient REQ, transient REQ,
and purgatory exact-ID polling share only the remaining capacity. and purgatory exact-ID polling share only the remaining capacity.
- Permit acquisition checks relay health first: while a rate-limit or - The scheduler and outbound connection share the relay's rate-limit state.
transient-failure cooldown is active, queued rounds take the REQ+EOSE NOTICE and rate-limited CLOSED messages update that state on the independent
fallback path (which is itself budget-accounted) instead of firing into a control stream, even when event delivery to the processor is blocked. New
relay that just complained. 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 - The reactive machinery (escalating cooldown, rate-limit detection in both
NOTICE and subscription-specific CLOSED messages, and per-batch fallback) NOTICE and subscription-specific CLOSED messages, and per-batch fallback)
remains the backstop for relays whose limits are below our floors. A remains the backstop for relays whose limits are below our floors. A
+11 -20
View File
@@ -4801,15 +4801,15 @@ impl SyncManager {
} }
RelayEvent::Notice(notice) => { RelayEvent::Notice(notice) => {
if is_rate_limit_message(&notice) { if is_rate_limit_message(&notice) {
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, relay = %relay_url_clone,
notice = %notice, notice = %notice,
"Rate limiting NOTICE detected from relay" "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 // Update metrics with new health state
if let Some(ref metrics) = metrics_clone { if let Some(ref metrics) = metrics_clone {
let state = health_tracker.get_state(&relay_url_clone); let state = health_tracker.get_state(&relay_url_clone);
@@ -4843,21 +4843,11 @@ impl SyncManager {
&& !is_filter_count_refusal(&reason) && !is_filter_count_refusal(&reason)
&& subscription_state_byte_limit(&reason).is_none() && subscription_state_byte_limit(&reason).is_none()
{ {
let already_paused = health_tracker.is_rate_limited(&relay_url_clone); tracing::debug!(
if already_paused { relay = %relay_url_clone,
tracing::debug!( reason = %reason,
relay = %relay_url_clone, "Rate limiting CLOSED detected from relay"
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);
if let Some(ref metrics) = metrics_clone { if let Some(ref metrics) = metrics_clone {
metrics.record_health_state( metrics.record_health_state(
&relay_url_clone, &relay_url_clone,
@@ -5490,7 +5480,8 @@ impl SyncManager {
keys, keys,
source, source,
policy, policy,
); )
.with_health_tracker(Arc::clone(&self.health_tracker));
self.connections.insert(relay_url.clone(), connection); self.connections.insert(relay_url.clone(), connection);
if nip65_discovery_only { if nip65_discovery_only {
self.nip65_discovery_only_relays.insert(relay_url.clone()); self.nip65_discovery_only_relays.insert(relay_url.clone());
+247 -9
View File
@@ -21,7 +21,7 @@ use std::future::Future;
use std::time::Duration; use std::time::Duration;
use tokio::sync::mpsc; 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 super::{is_filter_count_refusal, is_rate_limit_message, subscription_state_byte_limit};
use crate::nostr::SharedDatabase; use crate::nostr::SharedDatabase;
use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource}; use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource};
@@ -586,6 +586,8 @@ pub struct RelayConnection {
source: RelayTargetSource, source: RelayTargetSource,
/// Outbound target policy applied before every event-directed dial /// Outbound target policy applied before every event-directed dial
policy: OutboundTargetPolicy, policy: OutboundTargetPolicy,
/// Shared with the scheduler so queued work observes the same cooldown.
health_tracker: std::sync::Arc<RelayHealthTracker>,
/// The underlying nostr-sdk client /// The underlying nostr-sdk client
client: Client, client: Client,
/// Whether a NIP-42 authenticator was attached to the client. Without /// Whether a NIP-42 authenticator was attached to the client. Without
@@ -638,7 +640,38 @@ pub struct RelayConnection {
} }
impl RelayConnection { impl RelayConnection {
pub(super) fn with_health_tracker(
mut self,
health_tracker: std::sync::Arc<RelayHealthTracker>,
) -> 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> { async fn await_query_start(&self) -> Result<(), String> {
self.ensure_query_admitted()?;
self.query_start_pacer self.query_start_pacer
.wait_for_start() .wait_for_start()
.await .await
@@ -738,6 +771,7 @@ impl RelayConnection {
source, source,
policy, policy,
client, client,
health_tracker: std::sync::Arc::new(RelayHealthTracker::with_defaults()),
has_authenticator, has_authenticator,
database: None, database: None,
nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new( nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
@@ -812,6 +846,7 @@ impl RelayConnection {
source, source,
policy, policy,
client, client,
health_tracker: std::sync::Arc::new(RelayHealthTracker::with_defaults()),
has_authenticator, has_authenticator,
database: Some(database), database: Some(database),
nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new( nip77_unsupported_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(
@@ -1696,8 +1731,8 @@ impl RelayConnection {
let mut notifications = relay.notifications(); let mut notifications = relay.notifications();
// EVENT forwarding can block on the bounded processor queue. Subscribe // EVENT forwarding can block on the bounded processor queue. Subscribe
// independently to terminal messages so EOSE/CLOSED can still return // independently to control messages so rate-limit cooldowns take
// transient ledger slots while the data lane drains. Relay // effect and EOSE/CLOSED return slots while the data lane drains. Relay
// notifications are broadcast, so this does not consume messages from // notifications are broadcast, so this does not consume messages from
// the processor-facing stream below. // the processor-facing stream below.
let terminal_connection = self.clone(); let terminal_connection = self.clone();
@@ -1708,6 +1743,9 @@ impl RelayConnection {
while let Some(notification) = terminals.next().await { while let Some(notification) = terminals.next().await {
match notification { match notification {
RelayNotification::Message { message } => match *message { RelayNotification::Message { message } => match *message {
RelayMessage::Notice(message) => {
terminal_connection.record_subscription_rate_limit(&message);
}
RelayMessage::EndOfStoredEvents(sub_id) => { RelayMessage::EndOfStoredEvents(sub_id) => {
terminal_connection terminal_connection
.close_and_release_transient_req_permit(&sub_id) .close_and_release_transient_req_permit(&sub_id)
@@ -1717,6 +1755,7 @@ impl RelayConnection {
subscription_id, subscription_id,
message, message,
} => { } => {
terminal_connection.record_subscription_rate_limit(&message);
let subscription_id = subscription_id.into_owned(); let subscription_id = subscription_id.into_owned();
// An auth-required CLOSED is only a retry signal // An auth-required CLOSED is only a retry signal
// when an authenticator can actually answer the // when an authenticator can actually answer the
@@ -1794,9 +1833,6 @@ impl RelayConnection {
} }
} }
RelayMessage::Notice(msg) => { RelayMessage::Notice(msg) => {
if is_query_rate_limit_message(&msg) {
self.record_query_rate_limit();
}
// Check if this is a negentropy-related notice // Check if this is a negentropy-related notice
let is_negentropy_notice = msg.contains("envelope") let is_negentropy_notice = msg.contains("envelope")
|| msg.contains("NEG-") || msg.contains("NEG-")
@@ -1847,9 +1883,6 @@ impl RelayConnection {
} else { } else {
self.release_live_req_permit(&subscription_id) 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) { if let Some(limit) = self.record_subscription_byte_limit(&msg) {
tracing::warn!( tracing::warn!(
relay = %url, relay = %url,
@@ -2113,6 +2146,9 @@ impl RelayConnection {
} }
// Transient permits are acquired through the same session helper; a // Transient permits are acquired through the same session helper; a
// reset closes queued acquisitions and the helper checks generation. // 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(); let retained_filters = filters.clone();
// The relay can answer an empty or cached query before `subscribe` // The relay can answer an empty or cached query before `subscribe`
// returns. Register transient ownership against a caller-chosen ID // returns. Register transient ownership against a caller-chosen ID
@@ -2239,6 +2275,7 @@ impl RelayConnection {
let subscription_id = SubscriptionId::generate(); let subscription_id = SubscriptionId::generate();
let mut notifications = relay.notifications(); let mut notifications = relay.notifications();
let fetch = async { let fetch = async {
self.ensure_query_admitted()?;
let mut stream = relay let mut stream = relay
.stream_events(filter) .stream_events(filter)
.with_id(subscription_id.clone()) .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 // Use dry_run to only identify differences without downloading events
self.ensure_current_session(&ledger_slot)?; self.ensure_current_session(&ledger_slot)?;
let sync_opts = SyncOptions::default().dry_run(); let sync_opts = SyncOptions::default().dry_run();
let client = self.client.clone(); let client = self.client.clone();
let sync_task = async move { let sync_task = async move {
self.ensure_query_admitted()?;
match client.sync(filter).opts(sync_opts).await { match client.sync(filter).opts(sync_opts).await {
Ok(output) => { Ok(output) => {
if !output.failed.is_empty() { if !output.failed.is_empty() {
@@ -3433,6 +3474,203 @@ mod tests {
relay.shutdown(); 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<dyn Future<Output = Result<(), String>> + 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] #[tokio::test]
async fn query_rate_closed_paces_later_wire_requests() { async fn query_rate_closed_paces_later_wire_requests() {
let relay = TestRelay::start(LocalRelayBuilder::default().queries_per_minute(1)).await; let relay = TestRelay::start(LocalRelayBuilder::default().queries_per_minute(1)).await;