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
|
### 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
|
- 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
|
production 34-filter repository-sync live set reached roughly 1.2 MiB, so
|
||||||
rust-nostr's newly introduced default repeatedly closed part of persistent
|
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
|
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
|
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
|
(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
|
the existing subscriptions are consolidated into the byte- and filter-count
|
||||||
bounded REQ groups first. If the consolidated live set still cannot fit,
|
bounded REQ groups first. If the consolidated live set still cannot fit,
|
||||||
partial coverage is not opened, historic work is deferred, and a warning
|
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
|
generation; queued or late borrowers therefore fail before sending on the new
|
||||||
SDK session and cannot inflate or bypass its capacity.
|
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.
|
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
|
Permit release does not depend on that queue draining: a separate listener on
|
||||||
rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals
|
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(
|
fn select_purgatory_dependency_events(
|
||||||
mut events: Vec<Event>,
|
mut events: Vec<Event>,
|
||||||
attempts: &mut HashMap<EventId, Instant>,
|
attempts: &mut HashMap<EventId, Instant>,
|
||||||
@@ -744,6 +752,14 @@ fn is_rate_limit_message(message: &str) -> bool {
|
|||||||
|| message.contains("throttl")
|
|| 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)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
struct ConnectAttemptToken(u64);
|
struct ConnectAttemptToken(u64);
|
||||||
|
|
||||||
@@ -816,6 +832,40 @@ const MAX_FILTERS_PER_REQ: usize = 10;
|
|||||||
/// REQ when chunks are full. See
|
/// REQ when chunks are full. See
|
||||||
/// docs/explanation/sync-scaling-constraints.md.
|
/// docs/explanation/sync-scaling-constraints.md.
|
||||||
const REQ_MESSAGE_BYTE_BUDGET: usize = 96 * 1024;
|
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.
|
/// Pack filters into REQ-sized groups.
|
||||||
///
|
///
|
||||||
@@ -1165,7 +1215,12 @@ async fn run_health_and_metrics_checker(
|
|||||||
// 3. Check for rate limit recovery
|
// 3. Check for rate limit recovery
|
||||||
manager.check_rate_limit_recovery().await;
|
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() {
|
if let Some(naughty_list) = manager.health_tracker.naughty_list() {
|
||||||
let recovered = naughty_list.expire_old_entries();
|
let recovered = naughty_list.expire_old_entries();
|
||||||
for url in recovered {
|
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 {
|
if let Some(ref metrics) = manager.metrics {
|
||||||
// Get all tracked relay URLs
|
// Get all tracked relay URLs
|
||||||
let relay_urls: Vec<String> = {
|
let relay_urls: Vec<String> = {
|
||||||
@@ -1267,6 +1322,9 @@ pub struct SyncManager {
|
|||||||
connect_attempt_semaphore: Arc<Semaphore>,
|
connect_attempt_semaphore: Arc<Semaphore>,
|
||||||
/// Relays whose subscription consolidation waits for in-flight batches to drain.
|
/// Relays whose subscription consolidation waits for in-flight batches to drain.
|
||||||
deferred_consolidations: DeferredConsolidations,
|
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)
|
/// Channel for disconnect notifications (set during run)
|
||||||
disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>,
|
disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>,
|
||||||
/// Channel for EOSE notifications (set during run)
|
/// Channel for EOSE notifications (set during run)
|
||||||
@@ -1366,6 +1424,7 @@ impl SyncManager {
|
|||||||
in_flight_connect_attempts: HashMap::new(),
|
in_flight_connect_attempts: HashMap::new(),
|
||||||
connect_attempt_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECT_ATTEMPTS)),
|
connect_attempt_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECT_ATTEMPTS)),
|
||||||
deferred_consolidations: DeferredConsolidations::default(),
|
deferred_consolidations: DeferredConsolidations::default(),
|
||||||
|
byte_limited_live_relays: HashMap::new(),
|
||||||
disconnect_tx: None,
|
disconnect_tx: None,
|
||||||
eose_tx: None,
|
eose_tx: None,
|
||||||
subscription_closed_tx: None,
|
subscription_closed_tx: None,
|
||||||
@@ -2958,12 +3017,15 @@ impl SyncManager {
|
|||||||
"handle_add_filters: calling sync_live and historic_sync"
|
"handle_add_filters: calling sync_live and historic_sync"
|
||||||
);
|
);
|
||||||
|
|
||||||
if self
|
if let Err(error) = self
|
||||||
.sync_live(&action.relay_url, &action.filters)
|
.sync_live(&action.relay_url, &action.filters)
|
||||||
.await
|
.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)
|
self.historic_sync(&action.relay_url, action.filters, action.items, None)
|
||||||
.await;
|
.await;
|
||||||
@@ -3201,7 +3263,9 @@ impl SyncManager {
|
|||||||
reason = %reason,
|
reason = %reason,
|
||||||
"Relay closed a subscription (not a connection close)"
|
"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 =
|
let already_paused =
|
||||||
health_tracker.is_rate_limited(&relay_url_clone);
|
health_tracker.is_rate_limited(&relay_url_clone);
|
||||||
if already_paused {
|
if already_paused {
|
||||||
@@ -3398,6 +3462,18 @@ impl SyncManager {
|
|||||||
filters
|
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
|
/// Quick reconnect - for disconnections < 15 minutes
|
||||||
///
|
///
|
||||||
/// Re-establishes subscriptions after a brief disconnection by:
|
/// 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
|
// Replace L1+L2+L3 as one reserved transaction. The connection-level
|
||||||
// opener rolls every successful group back if a later group fails.
|
// opener rolls every successful group back if a later group fails.
|
||||||
let connection = match self.connections.get(relay_url) {
|
let connection = match self.connections.get(relay_url) {
|
||||||
Some(conn) => conn,
|
Some(conn) => conn.clone(),
|
||||||
None => {
|
None => {
|
||||||
tracing::debug!(
|
tracing::debug!(
|
||||||
relay = %relay_url,
|
relay = %relay_url,
|
||||||
@@ -5100,7 +5176,10 @@ impl SyncManager {
|
|||||||
return false;
|
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
|
if connection
|
||||||
.replace_live_filter_groups(complete_groups)
|
.replace_live_filter_groups(complete_groups)
|
||||||
.await
|
.await
|
||||||
@@ -5112,11 +5191,23 @@ impl SyncManager {
|
|||||||
);
|
);
|
||||||
return false;
|
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!(
|
tracing::info!(
|
||||||
relay = %relay_url,
|
relay = %relay_url,
|
||||||
since = %since,
|
since = %since,
|
||||||
|
overflow_groups,
|
||||||
"Consolidation complete - filter count reset"
|
"Consolidation complete - filter count reset"
|
||||||
);
|
);
|
||||||
true
|
true
|
||||||
@@ -5129,6 +5220,27 @@ impl SyncManager {
|
|||||||
reason: &str,
|
reason: &str,
|
||||||
live_generation: Option<u64>,
|
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) {
|
if is_rate_limit_message(reason) {
|
||||||
let removed_batch = {
|
let removed_batch = {
|
||||||
let mut pending = self.pending_sync_index.write().await;
|
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
|
/// Check for relays that should be disconnected
|
||||||
///
|
///
|
||||||
/// This method is called periodically by run_disconnect_checker.
|
/// This method is called periodically by run_disconnect_checker.
|
||||||
@@ -5450,7 +5611,7 @@ impl SyncManager {
|
|||||||
/// # Returns
|
/// # Returns
|
||||||
/// Vec of subscription IDs for the live subscriptions, or empty if connection not found
|
/// Vec of subscription IDs for the live subscriptions, or empty if connection not found
|
||||||
async fn sync_live(
|
async fn sync_live(
|
||||||
&self,
|
&mut self,
|
||||||
relay_url: &str,
|
relay_url: &str,
|
||||||
filters: &[Filter],
|
filters: &[Filter],
|
||||||
) -> Result<Vec<SubscriptionId>, String> {
|
) -> Result<Vec<SubscriptionId>, String> {
|
||||||
@@ -5459,7 +5620,7 @@ impl SyncManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let connection = match self.connections.get(relay_url) {
|
let connection = match self.connections.get(relay_url) {
|
||||||
Some(conn) => conn,
|
Some(conn) => conn.clone(),
|
||||||
None => {
|
None => {
|
||||||
tracing::debug!(relay = %relay_url, "No connection found for live sync");
|
tracing::debug!(relay = %relay_url, "No connection found for live sync");
|
||||||
return Err(format!("No connection found for live sync on {relay_url}"));
|
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 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
|
connection
|
||||||
.subscribe_live_filter_groups(filter_groups)
|
.subscribe_live_filter_groups(filter_groups)
|
||||||
.await
|
.await
|
||||||
@@ -6495,6 +6677,36 @@ mod tests {
|
|||||||
assert!(!is_rate_limit_message("blocked: unsupported filter"));
|
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]
|
#[test]
|
||||||
fn rate_limited_closed_removes_only_its_pending_batch_for_retry() {
|
fn rate_limited_closed_removes_only_its_pending_batch_for_retry() {
|
||||||
let relay_url = "wss://limited.example";
|
let relay_url = "wss://limited.example";
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ use std::time::Duration;
|
|||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
|
|
||||||
use super::health::RATE_LIMIT_COOLDOWN_SECS;
|
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::nostr::SharedDatabase;
|
||||||
use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource};
|
use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSource};
|
||||||
|
|
||||||
@@ -270,6 +270,7 @@ type TransientReqPermitMap = std::sync::Arc<
|
|||||||
struct HeldTransientPermits {
|
struct HeldTransientPermits {
|
||||||
_class_cap: tokio::sync::OwnedSemaphorePermit,
|
_class_cap: tokio::sync::OwnedSemaphorePermit,
|
||||||
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
|
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
|
||||||
|
_byte_limit_gate: Option<tokio::sync::OwnedSemaphorePermit>,
|
||||||
generation: u64,
|
generation: u64,
|
||||||
request_class: TransientRequestClass,
|
request_class: TransientRequestClass,
|
||||||
opened_at: std::time::Instant,
|
opened_at: std::time::Instant,
|
||||||
@@ -478,6 +479,13 @@ pub struct RelayConnection {
|
|||||||
transient_req_permits_held: TransientReqPermitMap,
|
transient_req_permits_held: TransientReqPermitMap,
|
||||||
/// Ledger slots held for persistent subscriptions until CLOSE/teardown.
|
/// Ledger slots held for persistent subscriptions until CLOSE/teardown.
|
||||||
live_req_permits_held: LiveReqPermitMap,
|
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.
|
/// 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.
|
/// 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://)
|
/// Normalize a relay URL to include a scheme (wss:// or ws://)
|
||||||
///
|
///
|
||||||
/// If the URL already has a scheme, it's returned as-is.
|
/// 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(
|
live_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new(
|
||||||
std::collections::HashMap::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()),
|
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
|
||||||
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::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(
|
live_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new(
|
||||||
std::collections::HashMap::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()),
|
query_start_pacer: std::sync::Arc::new(QueryStartPacer::default()),
|
||||||
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
|
background_query_pacer: std::sync::Arc::new(BackgroundQueryPacer::default()),
|
||||||
}
|
}
|
||||||
@@ -918,11 +951,23 @@ impl RelayConnection {
|
|||||||
&self,
|
&self,
|
||||||
request_class: TransientRequestClass,
|
request_class: TransientRequestClass,
|
||||||
) -> Result<HeldTransientPermits, String> {
|
) -> 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 {
|
loop {
|
||||||
let ledger_slot = self.acquire_subscription_slots(1).await?;
|
let ledger_slot = self.acquire_subscription_slots(1).await?;
|
||||||
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
|
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
|
||||||
return Ok(HeldTransientPermits {
|
return Ok(HeldTransientPermits {
|
||||||
_class_cap: class_cap,
|
_class_cap: class_cap,
|
||||||
|
_byte_limit_gate: byte_limit_gate,
|
||||||
generation: ledger_slot.generation,
|
generation: ledger_slot.generation,
|
||||||
_ledger_slot: ledger_slot.permit,
|
_ledger_slot: ledger_slot.permit,
|
||||||
request_class,
|
request_class,
|
||||||
@@ -942,6 +987,7 @@ impl RelayConnection {
|
|||||||
if let Ok(ledger_slot) = self.try_acquire_subscription_slot() {
|
if let Ok(ledger_slot) = self.try_acquire_subscription_slot() {
|
||||||
return Ok(HeldTransientPermits {
|
return Ok(HeldTransientPermits {
|
||||||
_class_cap: class_cap,
|
_class_cap: class_cap,
|
||||||
|
_byte_limit_gate: byte_limit_gate,
|
||||||
generation: ledger_slot.generation,
|
generation: ledger_slot.generation,
|
||||||
_ledger_slot: ledger_slot.permit,
|
_ledger_slot: ledger_slot.permit,
|
||||||
request_class,
|
request_class,
|
||||||
@@ -988,6 +1034,19 @@ impl RelayConnection {
|
|||||||
.collect()
|
.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) {
|
async fn unsubscribe_live(&self) {
|
||||||
let ids: Vec<_> = self
|
let ids: Vec<_> = self
|
||||||
.live_req_permits_held
|
.live_req_permits_held
|
||||||
@@ -1253,6 +1312,13 @@ impl RelayConnection {
|
|||||||
if is_query_rate_limit_message(&msg) {
|
if is_query_rate_limit_message(&msg) {
|
||||||
self.record_query_rate_limit();
|
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) {
|
if is_rate_limit_message(&msg) {
|
||||||
// The sync actor emits one canonical signal and owns
|
// The sync actor emits one canonical signal and owns
|
||||||
// cooldown deduplication for a rate-limit episode.
|
// cooldown deduplication for a rate-limit episode.
|
||||||
@@ -2956,6 +3022,7 @@ mod tests {
|
|||||||
opened_at: std::time::Instant::now(),
|
opened_at: std::time::Instant::now(),
|
||||||
last_event_at: None,
|
last_event_at: None,
|
||||||
delivered_events: 0,
|
delivered_events: 0,
|
||||||
|
_byte_limit_gate: None,
|
||||||
_class_cap: connection
|
_class_cap: connection
|
||||||
.transient_req_permits
|
.transient_req_permits
|
||||||
.clone()
|
.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]
|
#[tokio::test]
|
||||||
async fn session_reset_cannot_be_inflated_by_an_old_borrower() {
|
async fn session_reset_cannot_be_inflated_by_an_old_borrower() {
|
||||||
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
|
let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());
|
||||||
|
|||||||
Reference in New Issue
Block a user