diff --git a/CHANGELOG.md b/CHANGELOG.md index 48019ac..6a37b0a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 9f9da5f..b785d87 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -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 diff --git a/src/sync/mod.rs b/src/sync/mod.rs index b7146b2..75fef12 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -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, attempts: &mut HashMap, @@ -744,6 +752,14 @@ fn is_rate_limit_message(message: &str) -> bool { || message.contains("throttl") } +fn subscription_state_byte_limit(message: &str) -> Option { + 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>, + limit: Option, + already_used: usize, +) -> (Vec>, 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 = { @@ -1267,6 +1322,9 @@ pub struct SyncManager { connect_attempt_semaphore: Arc, /// 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, /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, /// 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, ) { + 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, 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"; diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 72f40ac..35380aa 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -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, 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, + /// 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, /// Learned per-session spacing for query starts after a query-rate refusal. query_start_pacer: std::sync::Arc, /// 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 { + 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 { + 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 { + 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());