Merge #ea2c33a2: Enforce relay cooldowns for queued sync requests

nostr:nevent1qqsw5tpn5g7nvs4jp8vz85p6dfxxlfndnrqzkqq4xsfxcgydlgvgdfgpz3mhxue69uhhyetvv9ujumn8d96zuer9wcf6qtsd

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

PR description:

A relay rate-limit NOTICE paused scheduler admission, but requests already waiting for pacing or subscription capacity could still be sent during the cooldown. A full event-processing queue could also delay handling the NOTICE.

This shares cooldown state between the scheduler and connection, handles rate-limit signals on the independent control stream, and checks admission again before handing ordinary REQs, direct fetches, and new NIP-77 rounds to the SDK. Refused work releases its permits for existing recovery paths. Repeated NOTICEs move to debug logging, leaving one warning per cooldown episode.

The wire regression queues requests before a real NOTICE, blocks the event-processing queue, verifies that refused requests never reach the relay or leak permits, and observes a successful request after cooldown is cleared. Validation: 401 sync-related library tests passed; the final focused regression, Clippy with warnings denied, formatting, and diff checks passed.

This does not cancel requests already handed to the SDK or SDK-owned protocol continuations. It does not establish why Primal’s concurrency limit was exceeded. No deployment has been performed.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-09-21 11:13:30 +00:00
4 changed files with 278 additions and 33 deletions
+7
View File
@@ -7,6 +7,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
### Fixed
- Fix queued sync requests bypassing relay rate-limit backoff and continuing
to send requests during the cooldown. Rate-limit signals now take effect even
when sync event processing is blocked, and repeated signals no longer flood
warning logs.
## [3.0.3] - 2026-09-18
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,
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
+11 -20
View File
@@ -4801,15 +4801,15 @@ impl SyncManager {
}
RelayEvent::Notice(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,
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());
+247 -9
View File
@@ -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<RelayHealthTracker>,
/// 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<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> {
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<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]
async fn query_rate_closed_paces_later_wire_requests() {
let relay = TestRelay::start(LocalRelayBuilder::default().queries_per_minute(1)).await;