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:
DanConwayDev
2026-08-08 11:08:53 +00:00
parent c5e52c5a1e
commit 35894d7ee0
4 changed files with 342 additions and 13 deletions
+6
View File
@@ -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
+19 -1
View File
@@ -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
View File
@@ -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";
+94 -1
View File
@@ -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());