mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(sync): bound remote retained subscription state
Production sync to relay.ngit.dev disclosed a 1 MiB cumulative retained-REQ cap after the repository-scale live set crossed it. Treating that CLOSED response as a temporary rate-limit episode rebuilt the same impossible set after every cooldown; merely raising our own serving allowance does not protect outbound sync against third-party policy. Parse the disclosed byte limit and retain it across reconnects. Rebuild persistent groups within the learned cap while reserving one maximum-sized transient REQ, serialize transient requests on byte-limited connections, and run five-minute paced history with a one-minute overlap for filters outside persistent coverage. Historic work now continues when extending live coverage fails, so the capacity signal cannot stall the affected batch. Correctness assumes a relay that emits this rust-nostr CLOSED form enforces the numeric limit per connection and admits an individual request within the existing 96 KiB message budget. NIP-11 cannot advertise this cumulative limit, so the first refusal is unavoidable. Multi-connection sharding and a protocol capability extension are deliberately excluded; sharding would improve overflow latency, not eventual completeness. Validation: all 661 library tests passed, including cap parsing, live-budget arithmetic, and serialized transient occupancy; git diff --check passed; nix build .#ngit-grasp passed.
This commit is contained in:
@@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### Fixed
|
||||
|
||||
- Treat cumulative retained-subscription byte refusals as capacity signals,
|
||||
not temporary query-rate episodes. The sync client learns the disclosed cap,
|
||||
rebuilds persistent coverage within it while reserving one maximum transient
|
||||
REQ, serializes transient work against that reserve, and covers overflow with
|
||||
paced five-minute history plus overlap instead of retrying the same
|
||||
impossible live set.
|
||||
- Raise retained subscription state per connection from 1 MiB to 5 MiB. A
|
||||
production 34-filter repository-sync live set reached roughly 1.2 MiB, so
|
||||
rust-nostr's newly introduced default repeatedly closed part of persistent
|
||||
|
||||
@@ -305,7 +305,8 @@ consumers share it, in priority order:
|
||||
NIP-11 `max_subscriptions` sets B for each new connection session; when it is
|
||||
absent B falls back to 20. Advertised values below that floor are honoured
|
||||
(notably nostream's default 10). Two slots remain reserved. Live filter groups
|
||||
are packed first and admitted atomically: if the complete live set cannot fit,
|
||||
are packed first and admitted atomically against the advertised
|
||||
subscription-count budget: if the complete live set cannot fit,
|
||||
the existing subscriptions are consolidated into the byte- and filter-count
|
||||
bounded REQ groups first. If the consolidated live set still cannot fit,
|
||||
partial coverage is not opened, historic work is deferred, and a warning
|
||||
@@ -335,6 +336,23 @@ live coverage. Each reconnect closes the retired ledger and creates a new
|
||||
generation; queued or late borrowers therefore fail before sending on the new
|
||||
SDK session and cannot inflate or bypass its capacity.
|
||||
|
||||
Some relays additionally cap the cumulative serialized REQ state retained by
|
||||
one connection. NIP-11 has no field for this limit, so it cannot be negotiated
|
||||
before the first refusal. A CLOSED reason of the rust-nostr form `active
|
||||
subscriptions exceed max size N bytes` is treated as a durable capacity signal,
|
||||
not as a temporary query-rate episode. The connection remembers N across
|
||||
reconnects and rebuilds its persistent filter groups within that byte budget,
|
||||
reserving one maximum-sized transient REQ. Byte-limited sessions serialize
|
||||
transient REQs so actual relay occupancy cannot overdraw that reserve.
|
||||
|
||||
Persistent groups beyond the learned cap are not silently abandoned. One
|
||||
byte-limited relay is given a paced incremental historic catch-up every five
|
||||
minutes, with a one-minute overlap, through the same slot ledger and background
|
||||
query pacer as ordinary history. This preserves eventual completeness without
|
||||
recreating an impossible live set. The first capacity response remains
|
||||
unavoidable because the limit is not advertised; multi-connection sharding is
|
||||
still out of scope and would improve latency rather than correctness.
|
||||
|
||||
The per-relay event processor retains its 1,000-message bounded data queue.
|
||||
Permit release does not depend on that queue draining: a separate listener on
|
||||
rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals
|
||||
|
||||
+223
-11
@@ -88,6 +88,14 @@ fn dependency_relay_retention() -> Duration {
|
||||
}
|
||||
}
|
||||
|
||||
fn byte_limited_catchup_interval() -> Duration {
|
||||
if std::env::var("NGIT_TEST").as_deref() == Ok("1") {
|
||||
Duration::from_secs(2)
|
||||
} else {
|
||||
Duration::from_secs(5 * 60)
|
||||
}
|
||||
}
|
||||
|
||||
fn select_purgatory_dependency_events(
|
||||
mut events: Vec<Event>,
|
||||
attempts: &mut HashMap<EventId, Instant>,
|
||||
@@ -744,6 +752,14 @@ fn is_rate_limit_message(message: &str) -> bool {
|
||||
|| message.contains("throttl")
|
||||
}
|
||||
|
||||
fn subscription_state_byte_limit(message: &str) -> Option<usize> {
|
||||
let lower = message.to_ascii_lowercase();
|
||||
let marker = "active subscriptions exceed max size ";
|
||||
let tail = lower.split_once(marker)?.1;
|
||||
let digits = tail.split_whitespace().next()?;
|
||||
digits.parse().ok()
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
struct ConnectAttemptToken(u64);
|
||||
|
||||
@@ -816,6 +832,40 @@ const MAX_FILTERS_PER_REQ: usize = 10;
|
||||
/// REQ when chunks are full. See
|
||||
/// docs/explanation/sync-scaling-constraints.md.
|
||||
const REQ_MESSAGE_BYTE_BUDGET: usize = 96 * 1024;
|
||||
/// Leave room for one maximum-sized transient REQ when a relay discloses a
|
||||
/// cumulative retained-subscription byte cap. Byte-limited sessions serialize
|
||||
/// transient REQs through a matching connection gate.
|
||||
const SUBSCRIPTION_BYTE_RESERVED_MARGIN: usize = REQ_MESSAGE_BYTE_BUDGET + 256;
|
||||
|
||||
fn req_message_size(filters: &[Filter]) -> usize {
|
||||
ClientMessage::req(SubscriptionId::generate(), filters.to_vec())
|
||||
.as_json()
|
||||
.len()
|
||||
}
|
||||
|
||||
fn groups_within_subscription_byte_limit(
|
||||
groups: Vec<Vec<Filter>>,
|
||||
limit: Option<usize>,
|
||||
already_used: usize,
|
||||
) -> (Vec<Vec<Filter>>, usize) {
|
||||
let Some(limit) = limit else {
|
||||
return (groups, 0);
|
||||
};
|
||||
let live_budget = limit.saturating_sub(SUBSCRIPTION_BYTE_RESERVED_MARGIN);
|
||||
let mut admitted = Vec::new();
|
||||
let mut used = already_used;
|
||||
let mut overflow = 0usize;
|
||||
for group in groups {
|
||||
let size = req_message_size(&group);
|
||||
if used.checked_add(size).is_some_and(|total| total <= live_budget) {
|
||||
used += size;
|
||||
admitted.push(group);
|
||||
} else {
|
||||
overflow += 1;
|
||||
}
|
||||
}
|
||||
(admitted, overflow)
|
||||
}
|
||||
|
||||
/// Pack filters into REQ-sized groups.
|
||||
///
|
||||
@@ -1165,7 +1215,12 @@ async fn run_health_and_metrics_checker(
|
||||
// 3. Check for rate limit recovery
|
||||
manager.check_rate_limit_recovery().await;
|
||||
|
||||
// 4. Check for naughty list expiration
|
||||
// 4. Keep deliberately bounded live coverage complete through
|
||||
// paced incremental history rather than retrying an impossible
|
||||
// persistent set.
|
||||
manager.sync_due_byte_limited_relay().await;
|
||||
|
||||
// 5. Check for naughty list expiration
|
||||
if let Some(naughty_list) = manager.health_tracker.naughty_list() {
|
||||
let recovered = naughty_list.expire_old_entries();
|
||||
for url in recovered {
|
||||
@@ -1176,7 +1231,7 @@ async fn run_health_and_metrics_checker(
|
||||
}
|
||||
}
|
||||
|
||||
// 5. Update metrics with current health states and naughty list
|
||||
// 6. Update metrics with current health states and naughty list
|
||||
if let Some(ref metrics) = manager.metrics {
|
||||
// Get all tracked relay URLs
|
||||
let relay_urls: Vec<String> = {
|
||||
@@ -1267,6 +1322,9 @@ pub struct SyncManager {
|
||||
connect_attempt_semaphore: Arc<Semaphore>,
|
||||
/// Relays whose subscription consolidation waits for in-flight batches to drain.
|
||||
deferred_consolidations: DeferredConsolidations,
|
||||
/// Relays whose complete persistent filter set exceeds a learned remote
|
||||
/// byte cap, mapped to their next bounded catch-up deadline.
|
||||
byte_limited_live_relays: HashMap<String, Instant>,
|
||||
/// Channel for disconnect notifications (set during run)
|
||||
disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>,
|
||||
/// Channel for EOSE notifications (set during run)
|
||||
@@ -1366,6 +1424,7 @@ impl SyncManager {
|
||||
in_flight_connect_attempts: HashMap::new(),
|
||||
connect_attempt_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECT_ATTEMPTS)),
|
||||
deferred_consolidations: DeferredConsolidations::default(),
|
||||
byte_limited_live_relays: HashMap::new(),
|
||||
disconnect_tx: None,
|
||||
eose_tx: None,
|
||||
subscription_closed_tx: None,
|
||||
@@ -2958,12 +3017,15 @@ impl SyncManager {
|
||||
"handle_add_filters: calling sync_live and historic_sync"
|
||||
);
|
||||
|
||||
if self
|
||||
if let Err(error) = self
|
||||
.sync_live(&action.relay_url, &action.filters)
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
return;
|
||||
tracing::warn!(
|
||||
relay = %action.relay_url,
|
||||
%error,
|
||||
"Live coverage could not be extended; continuing bounded historic sync"
|
||||
);
|
||||
}
|
||||
self.historic_sync(&action.relay_url, action.filters, action.items, None)
|
||||
.await;
|
||||
@@ -3201,7 +3263,9 @@ impl SyncManager {
|
||||
reason = %reason,
|
||||
"Relay closed a subscription (not a connection close)"
|
||||
);
|
||||
if is_rate_limit_message(&reason) {
|
||||
if is_rate_limit_message(&reason)
|
||||
&& subscription_state_byte_limit(&reason).is_none()
|
||||
{
|
||||
let already_paused =
|
||||
health_tracker.is_rate_limited(&relay_url_clone);
|
||||
if already_paused {
|
||||
@@ -3398,6 +3462,18 @@ impl SyncManager {
|
||||
filters
|
||||
}
|
||||
|
||||
async fn desired_items_for_relay(&self, relay_url: &str) -> PendingItems {
|
||||
let index = self.repo_sync_index.read().await;
|
||||
let target = algorithms::derive_relay_targets(&index)
|
||||
.remove(relay_url)
|
||||
.unwrap_or_default();
|
||||
PendingItems {
|
||||
repos: target.repos,
|
||||
state_only_repos: target.state_only_repos,
|
||||
root_events: target.root_events,
|
||||
}
|
||||
}
|
||||
|
||||
/// Quick reconnect - for disconnections < 15 minutes
|
||||
///
|
||||
/// Re-establishes subscriptions after a brief disconnection by:
|
||||
@@ -5080,7 +5156,7 @@ impl SyncManager {
|
||||
// Replace L1+L2+L3 as one reserved transaction. The connection-level
|
||||
// opener rolls every successful group back if a later group fails.
|
||||
let connection = match self.connections.get(relay_url) {
|
||||
Some(conn) => conn,
|
||||
Some(conn) => conn.clone(),
|
||||
None => {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
@@ -5100,7 +5176,10 @@ impl SyncManager {
|
||||
return false;
|
||||
}
|
||||
|
||||
let complete_groups = live_filter_groups(&complete_live);
|
||||
let unbounded_groups = live_filter_groups(&complete_live);
|
||||
let remote_limit = connection.remote_subscription_byte_limit();
|
||||
let (complete_groups, overflow_groups) =
|
||||
groups_within_subscription_byte_limit(unbounded_groups, remote_limit, 0);
|
||||
if connection
|
||||
.replace_live_filter_groups(complete_groups)
|
||||
.await
|
||||
@@ -5112,11 +5191,23 @@ impl SyncManager {
|
||||
);
|
||||
return false;
|
||||
}
|
||||
self.sync_generic_history(relay_url, Some(since)).await;
|
||||
if overflow_groups > 0 {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url.to_string(),
|
||||
Instant::now() + byte_limited_catchup_interval(),
|
||||
);
|
||||
let items = self.desired_items_for_relay(relay_url).await;
|
||||
self.historic_sync(relay_url, complete_live, items, Some(since))
|
||||
.await;
|
||||
} else {
|
||||
self.byte_limited_live_relays.remove(relay_url);
|
||||
self.sync_generic_history(relay_url, Some(since)).await;
|
||||
}
|
||||
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
since = %since,
|
||||
overflow_groups,
|
||||
"Consolidation complete - filter count reset"
|
||||
);
|
||||
true
|
||||
@@ -5129,6 +5220,27 @@ impl SyncManager {
|
||||
reason: &str,
|
||||
live_generation: Option<u64>,
|
||||
) {
|
||||
if let Some(limit) = subscription_state_byte_limit(reason) {
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
limit,
|
||||
"Remote retained-subscription capacity exhausted; rebuilding bounded live coverage"
|
||||
);
|
||||
self.byte_limited_live_relays
|
||||
.insert(relay_url.to_string(), Instant::now());
|
||||
if live_generation.is_none() {
|
||||
let mut pending = self.pending_sync_index.write().await;
|
||||
take_batch_containing_subscription(&mut pending, relay_url, &subscription_id);
|
||||
}
|
||||
let has_pending = self.has_pending_batches(relay_url).await;
|
||||
if self
|
||||
.deferred_consolidations
|
||||
.request(relay_url, has_pending)
|
||||
{
|
||||
let _ = self.consolidate(relay_url).await;
|
||||
}
|
||||
return;
|
||||
}
|
||||
if is_rate_limit_message(reason) {
|
||||
let removed_batch = {
|
||||
let mut pending = self.pending_sync_index.write().await;
|
||||
@@ -5185,6 +5297,55 @@ impl SyncManager {
|
||||
}
|
||||
}
|
||||
|
||||
async fn sync_due_byte_limited_relay(&mut self) {
|
||||
let now = Instant::now();
|
||||
let due = self
|
||||
.byte_limited_live_relays
|
||||
.iter()
|
||||
.find_map(|(relay, deadline)| (*deadline <= now).then(|| relay.clone()));
|
||||
let Some(relay_url) = due else {
|
||||
return;
|
||||
};
|
||||
|
||||
if self.has_pending_batches(&relay_url).await {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url,
|
||||
now + Duration::from_secs(10),
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
let connected = self
|
||||
.relay_sync_index
|
||||
.read()
|
||||
.await
|
||||
.get(&relay_url)
|
||||
.is_some_and(|state| state.connection_status.is_live_sync_active());
|
||||
if !connected {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url,
|
||||
now + byte_limited_catchup_interval(),
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url.clone(),
|
||||
now + byte_limited_catchup_interval(),
|
||||
);
|
||||
let overlap = byte_limited_catchup_interval() + Duration::from_secs(60);
|
||||
let since = Timestamp::from(Timestamp::now().as_secs().saturating_sub(overlap.as_secs()));
|
||||
let filters = self.complete_live_filters(&relay_url, Some(since)).await;
|
||||
let items = self.desired_items_for_relay(&relay_url).await;
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
since = %since,
|
||||
"Starting paced catch-up for byte-limited persistent coverage"
|
||||
);
|
||||
self.historic_sync(&relay_url, filters, items, Some(since))
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Check for relays that should be disconnected
|
||||
///
|
||||
/// This method is called periodically by run_disconnect_checker.
|
||||
@@ -5450,7 +5611,7 @@ impl SyncManager {
|
||||
/// # Returns
|
||||
/// Vec of subscription IDs for the live subscriptions, or empty if connection not found
|
||||
async fn sync_live(
|
||||
&self,
|
||||
&mut self,
|
||||
relay_url: &str,
|
||||
filters: &[Filter],
|
||||
) -> Result<Vec<SubscriptionId>, String> {
|
||||
@@ -5459,7 +5620,7 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
let connection = match self.connections.get(relay_url) {
|
||||
Some(conn) => conn,
|
||||
Some(conn) => conn.clone(),
|
||||
None => {
|
||||
tracing::debug!(relay = %relay_url, "No connection found for live sync");
|
||||
return Err(format!("No connection found for live sync on {relay_url}"));
|
||||
@@ -5467,6 +5628,27 @@ impl SyncManager {
|
||||
};
|
||||
|
||||
let filter_groups = live_filter_groups(filters);
|
||||
let remote_limit = connection.remote_subscription_byte_limit();
|
||||
let (filter_groups, overflow_groups) = groups_within_subscription_byte_limit(
|
||||
filter_groups,
|
||||
remote_limit,
|
||||
connection.live_subscription_bytes(),
|
||||
);
|
||||
if overflow_groups > 0 {
|
||||
self.byte_limited_live_relays.insert(
|
||||
relay_url.to_string(),
|
||||
Instant::now() + byte_limited_catchup_interval(),
|
||||
);
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
remote_limit,
|
||||
overflow_groups,
|
||||
"Persistent coverage bounded by remote subscription-state limit; overflow will use paced catch-up"
|
||||
);
|
||||
}
|
||||
if filter_groups.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
connection
|
||||
.subscribe_live_filter_groups(filter_groups)
|
||||
.await
|
||||
@@ -6495,6 +6677,36 @@ mod tests {
|
||||
assert!(!is_rate_limit_message("blocked: unsupported filter"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn subscription_state_limit_parser_is_specific_and_extracts_bytes() {
|
||||
assert_eq!(
|
||||
subscription_state_byte_limit(
|
||||
"rate-limited: active subscriptions exceed max size 1048576 bytes"
|
||||
),
|
||||
Some(1_048_576)
|
||||
);
|
||||
assert_eq!(
|
||||
subscription_state_byte_limit("rate-limited: too many queries"),
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn learned_subscription_byte_limit_reserves_transient_capacity() {
|
||||
let first = vec![Filter::new().kind(Kind::TextNote).limit(0)];
|
||||
let second = vec![Filter::new().kind(Kind::Metadata).limit(0)];
|
||||
let limit = SUBSCRIPTION_BYTE_RESERVED_MARGIN + req_message_size(&first);
|
||||
|
||||
let (admitted, overflow) = groups_within_subscription_byte_limit(
|
||||
vec![first.clone(), second],
|
||||
Some(limit),
|
||||
0,
|
||||
);
|
||||
|
||||
assert_eq!(admitted, vec![first]);
|
||||
assert_eq!(overflow, 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rate_limited_closed_removes_only_its_pending_batch_for_retry() {
|
||||
let relay_url = "wss://limited.example";
|
||||
|
||||
@@ -22,7 +22,7 @@ use std::time::Duration;
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
use super::health::RATE_LIMIT_COOLDOWN_SECS;
|
||||
use super::is_rate_limit_message;
|
||||
use super::{is_rate_limit_message, subscription_state_byte_limit};
|
||||
use crate::nostr::SharedDatabase;
|
||||
use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource};
|
||||
|
||||
@@ -270,6 +270,7 @@ type TransientReqPermitMap = std::sync::Arc<
|
||||
struct HeldTransientPermits {
|
||||
_class_cap: tokio::sync::OwnedSemaphorePermit,
|
||||
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
|
||||
_byte_limit_gate: Option<tokio::sync::OwnedSemaphorePermit>,
|
||||
generation: u64,
|
||||
request_class: TransientRequestClass,
|
||||
opened_at: std::time::Instant,
|
||||
@@ -478,6 +479,13 @@ pub struct RelayConnection {
|
||||
transient_req_permits_held: TransientReqPermitMap,
|
||||
/// Ledger slots held for persistent subscriptions until CLOSE/teardown.
|
||||
live_req_permits_held: LiveReqPermitMap,
|
||||
/// Cumulative retained REQ bytes learned from a relay's CLOSED response.
|
||||
/// Zero means the relay has not exposed a limit. The value survives
|
||||
/// reconnects so a durable policy is not probed on every new session.
|
||||
remote_subscription_byte_limit: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
||||
/// Once a cumulative byte cap is learned, serialize transient REQs so the
|
||||
/// reserved byte margin and actual relay-side occupancy cannot diverge.
|
||||
byte_limited_transient_gate: std::sync::Arc<tokio::sync::Semaphore>,
|
||||
/// Learned per-session spacing for query starts after a query-rate refusal.
|
||||
query_start_pacer: std::sync::Arc<QueryStartPacer>,
|
||||
/// Proactive spacing for non-urgent historic and dependency query starts.
|
||||
@@ -519,6 +527,23 @@ impl RelayConnection {
|
||||
}
|
||||
}
|
||||
|
||||
fn record_subscription_byte_limit(&self, message: &str) -> Option<usize> {
|
||||
let limit = subscription_state_byte_limit(message)?;
|
||||
self.remote_subscription_byte_limit
|
||||
.store(limit, std::sync::atomic::Ordering::Relaxed);
|
||||
Some(limit)
|
||||
}
|
||||
|
||||
pub fn remote_subscription_byte_limit(&self) -> Option<usize> {
|
||||
match self
|
||||
.remote_subscription_byte_limit
|
||||
.load(std::sync::atomic::Ordering::Relaxed)
|
||||
{
|
||||
0 => None,
|
||||
limit => Some(limit),
|
||||
}
|
||||
}
|
||||
|
||||
/// Normalize a relay URL to include a scheme (wss:// or ws://)
|
||||
///
|
||||
/// If the URL already has a scheme, it's returned as-is.
|
||||
@@ -595,6 +620,10 @@ impl RelayConnection {
|
||||
live_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new(
|
||||
std::collections::HashMap::new(),
|
||||
)),
|
||||
remote_subscription_byte_limit: std::sync::Arc::new(
|
||||
std::sync::atomic::AtomicUsize::new(0),
|
||||
),
|
||||
byte_limited_transient_gate: std::sync::Arc::new(tokio::sync::Semaphore::new(1)),
|
||||
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
|
||||
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
|
||||
}
|
||||
@@ -655,6 +684,10 @@ impl RelayConnection {
|
||||
live_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new(
|
||||
std::collections::HashMap::new(),
|
||||
)),
|
||||
remote_subscription_byte_limit: std::sync::Arc::new(
|
||||
std::sync::atomic::AtomicUsize::new(0),
|
||||
),
|
||||
byte_limited_transient_gate: std::sync::Arc::new(tokio::sync::Semaphore::new(1)),
|
||||
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
|
||||
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
|
||||
}
|
||||
@@ -918,11 +951,23 @@ impl RelayConnection {
|
||||
&self,
|
||||
request_class: TransientRequestClass,
|
||||
) -> Result<HeldTransientPermits, String> {
|
||||
let byte_limit_gate = if self.remote_subscription_byte_limit().is_some() {
|
||||
Some(
|
||||
self.byte_limited_transient_gate
|
||||
.clone()
|
||||
.acquire_owned()
|
||||
.await
|
||||
.map_err(|_| format!("Transient byte-limit gate closed for {}", self.url))?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
loop {
|
||||
let ledger_slot = self.acquire_subscription_slots(1).await?;
|
||||
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
|
||||
return Ok(HeldTransientPermits {
|
||||
_class_cap: class_cap,
|
||||
_byte_limit_gate: byte_limit_gate,
|
||||
generation: ledger_slot.generation,
|
||||
_ledger_slot: ledger_slot.permit,
|
||||
request_class,
|
||||
@@ -942,6 +987,7 @@ impl RelayConnection {
|
||||
if let Ok(ledger_slot) = self.try_acquire_subscription_slot() {
|
||||
return Ok(HeldTransientPermits {
|
||||
_class_cap: class_cap,
|
||||
_byte_limit_gate: byte_limit_gate,
|
||||
generation: ledger_slot.generation,
|
||||
_ledger_slot: ledger_slot.permit,
|
||||
request_class,
|
||||
@@ -988,6 +1034,19 @@ impl RelayConnection {
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub fn live_subscription_bytes(&self) -> usize {
|
||||
self.live_req_permits_held
|
||||
.lock()
|
||||
.expect("live permit map poisoned")
|
||||
.values()
|
||||
.map(|held| {
|
||||
ClientMessage::req(SubscriptionId::generate(), held.filters.clone())
|
||||
.as_json()
|
||||
.len()
|
||||
})
|
||||
.sum()
|
||||
}
|
||||
|
||||
async fn unsubscribe_live(&self) {
|
||||
let ids: Vec<_> = self
|
||||
.live_req_permits_held
|
||||
@@ -1253,6 +1312,13 @@ impl RelayConnection {
|
||||
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,
|
||||
limit,
|
||||
"Learned remote cumulative subscription-state limit"
|
||||
);
|
||||
}
|
||||
if is_rate_limit_message(&msg) {
|
||||
// The sync actor emits one canonical signal and owns
|
||||
// cooldown deduplication for a rate-limit episode.
|
||||
@@ -2956,6 +3022,7 @@ mod tests {
|
||||
opened_at: std::time::Instant::now(),
|
||||
last_event_at: None,
|
||||
delivered_events: 0,
|
||||
_byte_limit_gate: None,
|
||||
_class_cap: connection
|
||||
.transient_req_permits
|
||||
.clone()
|
||||
@@ -3099,6 +3166,32 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn learned_byte_limit_serializes_transient_occupancy() {
|
||||
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
|
||||
connection
|
||||
.remote_subscription_byte_limit
|
||||
.store(1_048_576, std::sync::atomic::Ordering::Relaxed);
|
||||
let first = connection
|
||||
.acquire_transient_permits(TransientRequestClass::HistoricPage)
|
||||
.await
|
||||
.expect("first byte-limited transient");
|
||||
let waiter_connection = connection.clone();
|
||||
let mut waiter = tokio::spawn(async move {
|
||||
waiter_connection
|
||||
.acquire_transient_permits(TransientRequestClass::HistoricPage)
|
||||
.await
|
||||
});
|
||||
assert!(
|
||||
tokio::time::timeout(Duration::from_millis(50), &mut waiter)
|
||||
.await
|
||||
.is_err(),
|
||||
"a second transient must wait while the byte margin is occupied"
|
||||
);
|
||||
drop(first);
|
||||
drop(waiter.await.expect("transient waiter task").unwrap());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn session_reset_cannot_be_inflated_by_an_old_borrower() {
|
||||
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
|
||||
|
||||
Reference in New Issue
Block a user