mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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:
+7
-3
@@ -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
|
||||||
|
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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();
|
||||||
|
|||||||
Reference in New Issue
Block a user