fix(sync): pace background query starts proactively

The subscription ledger bounds simultaneous work but allows completed historic operations to turn over as fast as relays answer. Waiting for a rate-limit refusal before pacing makes every startup impose an avoidable burst even though historic completeness is not latency-sensitive.

Add one per-connection background gate that spaces historic pages, pagination, hydration and retry pages, exact-ID dependency fetches, and NIP-77 round starts by one second from session startup. Persistent live subscriptions bypass this gate and retain priority. The existing reactive gate still covers every application-visible start after a refusal.

The gate runs before class and ledger permit acquisition, so paced work cannot pin scarce local capacity while sleeping. Session reset clears its timestamp. Correctness assumes delayed historic and dependency work remains re-derivable through the existing pending/retry machinery.

SDK-managed NIP-77 NEG-MSG continuations remain outside application control; reactive refusal classification and session REQ fallback are deliberately retained. Configurability and different priority levels within background work are excluded.

Validation: nix develop -c cargo test --lib (658 passed), including paused-time pacing boundaries and the real query-limited LocalRelay scenario; rustfmt --edition 2021 --check src/sync/relay_connection.rs; git diff --check.
This commit is contained in:
DanConwayDev
2026-08-08 09:02:25 +00:00
parent 21f3a46ca4
commit 55205b1f6b
5 changed files with 117 additions and 24 deletions
+7 -3
View File
@@ -14,14 +14,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
a finite DoS backstop while avoiding an unusually restrictive default that a finite DoS backstop while avoiding an unusually restrictive default that
also charges every SDK-managed NIP-77 `NEG-MSG` continuation. Re-evaluate also charges every SDK-managed NIP-77 `NEG-MSG` continuation. Re-evaluate
the 10× allowance after upstream revises that NIP-77 accounting. the 10× allowance after upstream revises that NIP-77 accounting.
- Proactively space non-urgent historic sync, pagination, hydration, retry,
and dependency query starts at one per second per relay connection. Live
subscriptions bypass the background gate, while the existing adaptive
refusal handling remains the backstop for tighter or internally generated
traffic.
- Pace query starts after a relay reports its per-connection query-rate budget - Pace query starts after a relay reports its per-connection query-rate budget
is exhausted, preventing fixed cooldown recovery from replaying the same is exhausted, preventing fixed cooldown recovery from replaying the same
fast historic-sync burst indefinitely. If an SDK-managed NIP-77 exchange fast historic-sync burst indefinitely. If an SDK-managed NIP-77 exchange
exhausts that budget, use paced REQs for the rest of the connection session exhausts that budget, use paced REQs for the rest of the connection session
because the application cannot pace individual `NEG-MSG` frames. This is a because the application cannot pace individual `NEG-MSG` frames. This is a
reactive compatibility backstop for an unadvertised limit; proactive pacing reactive compatibility backstop for an unadvertised limit, layered over the
of lower-priority historic work remains separate so live coverage is not proactive background pacing while leaving live coverage prioritised.
delayed.
## [2.1.0] - 2026-08-07 ## [2.1.0] - 2026-08-07
+9 -7
View File
@@ -1062,13 +1062,15 @@ Probing -> previous state: Recovery REQs succeed
Probing -> RateLimited: Recovery REQ is rate limited again Probing -> RateLimited: Recovery REQ is rate limited again
``` ```
The subscription ledger limits simultaneous work, while some relays also The subscription ledger limits simultaneous work, while a proactive
limit completed query operations per minute. A `too many queries` refusal per-connection gate spaces non-urgent historic, dependency, pagination,
therefore makes already-queued starts unwind during the 65-second cooldown so hydration, retry, and NIP-77 round starts at one per second. Persistent live
their work can be re-derived in priority order, then activates per-connection subscriptions bypass that background gate. Some relays also limit completed
query-start pacing: 600 ms between starts initially, doubling for a distinct query operations per minute. A `too many queries` refusal therefore makes
later episode up to 10 seconds. The pacing gate is shared by live and transient already-queued starts unwind during the 65-second cooldown so their work can be
REQs, exact-ID fetches and NIP-77 round starts, and resets with the relay re-derived in priority order, then activates a shared reactive gate for live
and background operations: 600 ms between starts initially, doubling for a
distinct later episode up to 10 seconds. Both gates reset with the relay
connection session. Any query-rate refusal disables NIP-77 for the rest of the connection session. Any query-rate refusal disables NIP-77 for the rest of the
session and falls back to paced REQs, because rust-nostr owns the internal session and falls back to paced REQs, because rust-nostr owns the internal
`NEG-MSG` frames and the application cannot guarantee their pacing. `NEG-MSG` frames and the application cannot guarantee their pacing.
+3 -1
View File
@@ -150,7 +150,9 @@ After a `too many queries` response, `Activated adaptive query-start pacing`
reports the learned per-connection interval. It begins at 600 ms and doubles reports the learned per-connection interval. It begins at 600 ms and doubles
only when a distinct later episode proves that pace too fast. Queued starts only when a distinct later episode proves that pace too fast. Queued starts
unwind during the existing 65-second cooldown before paced recovery; unwind during the existing 65-second cooldown before paced recovery;
reconnecting starts a fresh session without pacing. reconnecting clears that reactive lesson. Independent proactive pacing still
spaces background historic and dependency starts at one per second; persistent
live subscriptions bypass the background gate.
`Falling back to paced REQs for the query-limited connection session` means `Falling back to paced REQs for the query-limited connection session` means
NIP-77 remains skipped until reconnect because SDK-managed `NEG-MSG` traffic NIP-77 remains skipped until reconnect because SDK-managed `NEG-MSG` traffic
cannot be passed individually through the learned gate. cannot be passed individually through the learned gate.
+21 -13
View File
@@ -203,13 +203,13 @@ whether it applies independently to each filter or to the merged REQ, which is
why the source audit above remains necessary. Relay-specific extensions can why the source audit above remains necessary. Relay-specific extensions can
add fields, but clients cannot assume common names or semantics. add fields, but clients cannot assume common names or semantics.
The implemented query-start pacer is therefore a reactive compatibility Two gates cover the distinct concerns. A proactive background gate spaces
backstop: it remains inactive until an explicit `too many queries` response, historic, dependency, pagination, hydration, retry, and NIP-77 round starts at
then spaces all application-visible starts and selects REQ fallback because one per second from session startup; persistent live subscriptions bypass it.
the application cannot pace individual `NEG-MSG` frames. Proactively being A reactive compatibility gate remains inactive until an explicit
gentle with non-urgent historic work is desirable but separate. It requires `too many queries` response, then spaces every application-visible start and
request-class priority—live coverage first, dependency recovery next, bulk selects REQ fallback because the application cannot pace individual
history last—rather than enabling this shared gate from connection startup. `NEG-MSG` frames.
The client encodes this model in per-connection `RelayPaginationSession` The client encodes this model in per-connection `RelayPaginationSession`
state (`src/sync/mod.rs`). After EOSE it learns the largest raw page seen state (`src/sync/mod.rs`). After EOSE it learns the largest raw page seen
@@ -400,15 +400,23 @@ So concurrency is not a free scaling axis; it is the residual of the ledger:
- Rounds queue behind a per-connection semaphore; each completion releases - Rounds queue behind a per-connection semaphore; each completion releases
the next. No timed batches or sleeps — throughput degrades smoothly instead the next. No timed batches or sleeps — throughput degrades smoothly instead
of bursting into rejections. of bursting into rejections.
- The ledger bounds simultaneous resource use, not query starts over time. If - The ledger bounds simultaneous resource use, not query starts over time.
a relay returns `rate-limited: too many queries`, the current connection Non-urgent historic REQs, pagination and hydration/retry pages, exact-ID
dependency fetches, and NIP-77 round starts proactively share a one-second
start interval from the beginning of each connection session. Persistent
live subscriptions bypass that background gate and retain priority. This
bounds application-visible starts, but cannot pace SDK-managed NIP-77
`NEG-MSG` continuations. If a relay returns
`rate-limited: too many queries`, the current connection
first rejects locally queued starts for the existing 65-second cooldown so first rejects locally queued starts for the existing 65-second cooldown so
their incomplete work can be re-derived in priority order, then learns a their incomplete work can be re-derived in priority order, then learns a
shared query-start interval: 600 ms initially (100 starts/minute, below our shared query-start interval: 600 ms initially (100 starts/minute, below the
own 120/minute serving limit), doubling on a later rate-limit episode up to rust-nostr 120/minute default), doubling on a later rate-limit episode up to
10 seconds. Live, transient REQ, exact-ID fetch and NIP-77 round starts all 10 seconds. Live, transient REQ, exact-ID fetch and NIP-77 round starts all
pass through that pacer. Relays that never report a query pass through that reactive pacer in addition to background work retaining
rate limit remain unpaced, and reconnecting resets the per-session lesson. its proactive spacing. Relays that never report a query rate limit still
receive proactively paced background work, while live starts remain
immediate; reconnecting resets both per-session gates.
After any query-budget refusal, NIP-77 is skipped for the rest of the session: After any query-budget refusal, NIP-77 is skipped for the rest of the session:
rust-nostr owns the internal `NEG-MSG` exchange, so the application cannot rust-nostr owns the internal `NEG-MSG` exchange, so the application cannot
guarantee that each charged frame passes through its pacing gate. Historic guarantee that each charged frame passes through its pacing gate. Historic
+77
View File
@@ -58,6 +58,15 @@ const QUERY_PACING_MAX_INTERVAL: Duration = Duration::from_secs(10);
/// refusal after a full rate-limit window means the learned pace was too fast. /// refusal after a full rate-limit window means the learned pace was too fast.
const QUERY_RATE_LIMIT_EPISODE_GAP: Duration = Duration::from_secs(60); const QUERY_RATE_LIMIT_EPISODE_GAP: Duration = Duration::from_secs(60);
/// Minimum spacing between non-urgent historic and dependency query starts.
///
/// Historic completeness matters, but does not require burst throughput. One
/// background start per second is deliberately gentle to public relays and
/// leaves query capacity for live coverage. This only controls
/// application-visible starts; rust-nostr owns NIP-77 `NEG-MSG` continuations,
/// so the reactive refusal handling remains necessary.
const BACKGROUND_QUERY_INTERVAL: Duration = Duration::from_secs(1);
fn effective_subscription_budget(advertised: Option<usize>) -> usize { fn effective_subscription_budget(advertised: Option<usize>) -> usize {
advertised.unwrap_or(FALLBACK_SUBSCRIPTION_BUDGET) advertised.unwrap_or(FALLBACK_SUBSCRIPTION_BUDGET)
} }
@@ -87,6 +96,39 @@ struct QueryStartPacer {
state: std::sync::Mutex<QueryPacingState>, state: std::sync::Mutex<QueryPacingState>,
} }
/// Serialises background query starts from the beginning of each session.
/// Live subscriptions bypass this gate and therefore retain priority.
#[derive(Debug, Default)]
struct BackgroundQueryPacer {
gate: tokio::sync::Mutex<()>,
last_start: std::sync::Mutex<Option<tokio::time::Instant>>,
}
impl BackgroundQueryPacer {
async fn wait_for_start(&self) {
let _gate = self.gate.lock().await;
let deadline = self
.last_start
.lock()
.expect("background query pacing state poisoned")
.map(|last| last + BACKGROUND_QUERY_INTERVAL);
if let Some(deadline) = deadline {
tokio::time::sleep_until(deadline).await;
}
*self
.last_start
.lock()
.expect("background query pacing state poisoned") = Some(tokio::time::Instant::now());
}
fn reset(&self) {
*self
.last_start
.lock()
.expect("background query pacing state poisoned") = None;
}
}
impl QueryStartPacer { impl QueryStartPacer {
fn record_rate_limit(&self) -> (Duration, bool) { fn record_rate_limit(&self) -> (Duration, bool) {
let now = tokio::time::Instant::now(); let now = tokio::time::Instant::now();
@@ -438,6 +480,8 @@ pub struct RelayConnection {
live_req_permits_held: LiveReqPermitMap, live_req_permits_held: LiveReqPermitMap,
/// Learned per-session spacing for query starts after a query-rate refusal. /// Learned per-session spacing for query starts after a query-rate refusal.
query_start_pacer: std::sync::Arc<QueryStartPacer>, query_start_pacer: std::sync::Arc<QueryStartPacer>,
/// Proactive spacing for non-urgent historic and dependency query starts.
background_query_pacer: std::sync::Arc<BackgroundQueryPacer>,
} }
impl RelayConnection { impl RelayConnection {
@@ -454,6 +498,10 @@ impl RelayConnection {
}) })
} }
async fn await_background_query_start(&self) {
self.background_query_pacer.wait_for_start().await;
}
fn record_query_rate_limit(&self) { fn record_query_rate_limit(&self) {
let (interval, new_episode) = self.query_start_pacer.record_rate_limit(); let (interval, new_episode) = self.query_start_pacer.record_rate_limit();
if new_episode { if new_episode {
@@ -548,6 +596,7 @@ impl RelayConnection {
std::collections::HashMap::new(), std::collections::HashMap::new(),
)), )),
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()), query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
} }
} }
@@ -607,6 +656,7 @@ impl RelayConnection {
std::collections::HashMap::new(), std::collections::HashMap::new(),
)), )),
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()), query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
} }
} }
@@ -785,6 +835,7 @@ impl RelayConnection {
pub fn reset_subscription_budget(&self, advertised: Option<usize>) { pub fn reset_subscription_budget(&self, advertised: Option<usize>) {
self.clear_subscription_permits(); self.clear_subscription_permits();
self.query_start_pacer.reset(); self.query_start_pacer.reset();
self.background_query_pacer.reset();
self.nip77_query_rate_limited self.nip77_query_rate_limited
.store(false, std::sync::atomic::Ordering::Relaxed); .store(false, std::sync::atomic::Ordering::Relaxed);
let budget = effective_subscription_budget(advertised); let budget = effective_subscription_budget(advertised);
@@ -1385,6 +1436,9 @@ impl RelayConnection {
// Acquire pacing before subscription/class permits: a learned remote // Acquire pacing before subscription/class permits: a learned remote
// query-rate window must not turn local capacity into queued sleepers. // query-rate window must not turn local capacity into queued sleepers.
self.await_query_start().await?; self.await_query_start().await?;
if transient_class.is_some() {
self.await_background_query_start().await;
}
// Transient (auto-close) subscriptions share a bounded number of // Transient (auto-close) subscriptions share a bounded number of
// per-connection slots so historic bursts queue instead of // per-connection slots so historic bursts queue instead of
@@ -1519,6 +1573,7 @@ impl RelayConnection {
// nostr-sdk. It must share the same permit bound as historic pages, // nostr-sdk. It must share the same permit bound as historic pages,
// fallbacks and retries rather than escaping the connection budget. // fallbacks and retries rather than escaping the connection budget.
self.await_query_start().await?; self.await_query_start().await?;
self.await_background_query_start().await;
if self.historic_capacity_consumed_by_live() { if self.historic_capacity_consumed_by_live() {
tracing::warn!( tracing::warn!(
relay = %self.url, relay = %self.url,
@@ -1928,6 +1983,7 @@ impl RelayConnection {
// shared with live subscriptions (see MAX_CONCURRENT_NEG_DIFFS). // shared with live subscriptions (see MAX_CONCURRENT_NEG_DIFFS).
// The permit is held for the whole round, including the timeout. // The permit is held for the whole round, including the timeout.
self.await_query_start().await?; self.await_query_start().await?;
self.await_background_query_start().await;
if self.historic_capacity_consumed_by_live() { if self.historic_capacity_consumed_by_live() {
tracing::warn!( tracing::warn!(
relay = %self.url, relay = %self.url,
@@ -3203,6 +3259,27 @@ mod tests {
.expect("paced query start should be admitted"); .expect("paced query start should be admitted");
} }
#[tokio::test(start_paused = true)]
async fn background_query_pacer_spaces_starts_from_session_beginning() {
let pacer = std::sync::Arc::new(BackgroundQueryPacer::default());
pacer.wait_for_start().await;
let waiting = {
let pacer = std::sync::Arc::clone(&pacer);
tokio::spawn(async move { pacer.wait_for_start().await })
};
tokio::task::yield_now().await;
assert!(!waiting.is_finished());
tokio::time::advance(BACKGROUND_QUERY_INTERVAL - Duration::from_millis(1)).await;
tokio::task::yield_now().await;
assert!(!waiting.is_finished());
tokio::time::advance(Duration::from_millis(1)).await;
waiting.await.expect("background query should be admitted");
pacer.reset();
pacer.wait_for_start().await;
}
#[tokio::test(start_paused = true)] #[tokio::test(start_paused = true)]
async fn query_pacer_deduplicates_bursts_and_slows_new_episodes() { async fn query_pacer_deduplicates_bursts_and_slows_new_episodes() {
let pacer = QueryStartPacer::default(); let pacer = QueryStartPacer::default();