From 8fadb9beeacdf7cd5941f0183f0ff2d6f5c335e1 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 1 Aug 2026 13:43:05 +0100 Subject: [PATCH 1/3] refactor(sync): support grouped relay subscriptions Allow callers to send multiple NIP-01 filters under one subscription ID while retaining subscribe_filter as the single-filter compatibility wrapper. NIP-01 explicitly permits multiple filters in one REQ. Isolating that wire-level capability keeps the later proactive-sync batching change focused on policy rather than transport mechanics. Reject empty filter sets and surface immediate relay-pool send failures to callers. --- src/sync/relay_connection.rs | 39 ++++++++++++++++++++++++++++++------ 1 file changed, 33 insertions(+), 6 deletions(-) diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 09b6df8..ad0df4f 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -423,30 +423,57 @@ impl RelayConnection { filter: Filter, auto_close: bool, ) -> Result { - // DEBUG TRACING: Log the filter being subscribed to + self.subscribe_filters(vec![filter], auto_close).await + } + + /// Subscribe to several OR filters under one NIP-01 subscription ID. + /// + /// Relays apply active-REQ limits to subscription IDs, not to the filters + /// inside a REQ. Grouping related filters preserves NIP-01 semantics while + /// avoiding one persistent subscription per GRASP tag variant. + pub async fn subscribe_filters( + &self, + filters: Vec, + auto_close: bool, + ) -> Result { + if filters.is_empty() { + return Err("Cannot subscribe with an empty filter set".to_string()); + } + tracing::debug!( relay = %self.url, - filter = ?filter, + filter_count = filters.len(), + filters = ?filters, auto_close = auto_close, - "subscribe_filter called with filter" + "subscribe_filters called" ); let output = if auto_close { self.client - .subscribe(filter) + .subscribe(filters) .close_on( SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::ExitOnEOSE), ) .await } else { - self.client.subscribe(filter).await + self.client.subscribe(filters).await } .map_err(|e| format!("Failed to subscribe on {}: {}", self.url, e))?; + if !output.failed.is_empty() { + let failures = output + .failed + .values() + .cloned() + .collect::>() + .join("; "); + return Err(format!("Failed to subscribe on {}: {}", self.url, failures)); + } + tracing::debug!( relay = %self.url, subscription_id = %output.value, - "subscribe_filter succeeded" + "subscribe_filters succeeded" ); Ok(output.value) From 86b61aea4b8921ba5f9c9331016815cdc2ae057f Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 1 Aug 2026 13:44:18 +0100 Subject: [PATCH 2/3] fix(sync): batch proactive filters under relay REQ limits Public relays cap active subscription IDs. Opening separate live and historic REQs for every GRASP state, a, A, q, e, and E filter can reject the tail of the desired coverage. Carry up to ten compatible GRASP filters under each NIP-01 subscription ID. Preserve a separate event count and oldest timestamp for every filter so each filter advances only when its own result page may be full. This pagination model assumes relay result limits are independent per filter, the effective per-filter cap is at least 75, and no additional total-result cap can starve filters within a grouped REQ. Record those assumptions in the implementation and design documentation. Keep 300-ID negentropy fetch filters in separate REQs because their message-size budget is distinct from semantic GRASP filter batching. --- CHANGELOG.md | 3 + docs/explanation/grasp-02-proactive-sync.md | 70 +++- src/sync/mod.rs | 405 ++++++++++++-------- tests/common/mock_relay.rs | 26 +- tests/sync/live_sync.rs | 84 +++- 5 files changed, 398 insertions(+), 190 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cc58b63..cd846c1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Fixed proactive sync losing repository events when public relays cap active + subscriptions. Compatible GRASP filters now share bounded NIP-01 REQs while + retaining per-filter history pagination. - Removed superseded same-author repository states from purgatory after their replacement is promoted when the locally available Git data cannot reconstruct them. Reconstructable rollback states and other maintainers' diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 6cd1754..c369c84 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -24,6 +24,8 @@ Key Architectural Points: - **Clear separation** between Live sync (using `limit:0`) and Historic Sync (handled via negentropy falling back to REQ+EOSE with 'until' based pagination support) - **Discovery management**: The nature of discovery inherently leads to a drip feed of root_events (e.g., Repo Announcements, Issues, Patches and PRs) that require additional subscriptions. Without careful management this can lead to large numbers of subscriptions and potentially rate limiting. Mitigation strategies: - Self-subscriber waits for 5s to batch updates before creating new filters / subscriptions, allowing time for most events to be received from outstanding subscriptions from connected relays + - Up to ten compatible OR filters share each NIP-01 REQ, bounding + relay-visible active subscriptions without broadening any filter - PendingBatch tracks each new set of filters that may require pagination until they are complete - Websocket handshakes run in at most eight bounded workers outside the sync actor; only the actor applies their results, and subscriptions start only @@ -168,15 +170,18 @@ pub enum SyncMethod { /// Key: relay URL pub type PendingSyncIndex = Arc>>>; -/// Pagination state for a subscription in non-Negentropy historic sync +/// Pagination state for one filter inside a grouped subscription +#[derive(Debug, Clone)] +pub struct FilterPaginationState { + pub event_count: usize, + pub min_created_at: Option, + pub original_filter: Filter, +} + +/// Per-filter progress for every OR filter carried by one subscription #[derive(Debug, Clone)] pub struct PaginationState { - /// Number of events received for this subscription - pub event_count: usize, - /// Smallest created_at timestamp seen (for pagination with `until`) - pub min_created_at: Option, - /// Original filter to reconstruct for next page - pub original_filter: Filter, + pub filters: Vec, } pub struct PendingBatch { @@ -205,12 +210,18 @@ pub struct PendingItems { When a relay doesn't support NIP-77 Negentropy, historic sync falls back to traditional REQ+EOSE. To handle large result sets efficiently: -- **`PaginationState`** tracks per-subscription pagination progress +- **`PaginationState`** tracks pagination separately for each OR filter in a + grouped subscription - `event_count`: Number of events received so far - `min_created_at`: Smallest timestamp seen, used to set `until` for next page - `original_filter`: Base filter to reconstruct with updated `until` parameter -- **Automatic pagination**: When EOSE is received, if enough events were received to suggest more may exist, the system automatically issues a follow-up request with `until` set to `min_created_at` +- **Automatic pagination**: When EOSE is received, each filter that may have + more results is reconstructed with its own `until` timestamp; those next-page + filters remain grouped in one follow-up REQ - **Completion**: Pagination continues until an EOSE is received with fewer events than expected, indicating the end of results +- **Compatibility assumptions**: Relay result limits apply independently to + each filter, the effective per-filter limit is at least 75, and there is no + additional total-result cap across the grouped REQ --- @@ -707,8 +718,10 @@ fn compute_actions( ### Sync Primitives -- **`sync_live()`**: Creates subscriptions with `limit: 0` for ongoing event stream (not tracked in PendingSyncIndex) -- **`historic_sync()`**: Dispatches to negentropy or REQ+EOSE based on relay capability, creates PendingBatch, returns batch_id +- **`sync_live()`**: Groups compatible filters into bounded subscriptions with + `limit: 0` for the ongoing event stream (not tracked in PendingSyncIndex) +- **`historic_sync()`**: Dispatches to negentropy or grouped REQ+EOSE based on + relay capability, creates PendingBatch, and returns a batch ID ### Filter Processing @@ -843,8 +856,10 @@ When a relay doesn't support NIP-77 Negentropy, historic sync uses traditional R ### How Pagination Works -1. **Initial Request**: Send REQ with filters (may include `since` parameter) -2. **Track Events**: As events arrive, [`PaginationState`](src/sync/mod.rs:165) tracks: +1. **Initial Request**: Send bounded groups of OR filters in each REQ (filters + may include a `since` parameter) +2. **Track Events**: As events arrive, `PaginationState` tracks each matching + filter independently: - `event_count`: Number of events received - `min_created_at`: Smallest timestamp seen (oldest event) - `original_filter`: Base filter for reconstruction @@ -852,23 +867,40 @@ When a relay doesn't support NIP-77 Negentropy, historic sync uses traditional R 4. **Next Page**: If enough events were received (suggesting more exist): - Create new filter with `until: min_created_at` - Issue another REQ for events older than the oldest seen - - Reuse same subscription ID + - Group the next-page filters in a new subscription 5. **Completion**: Repeat until EOSE arrives with fewer events, indicating end of results +### Relay Compatibility Assumptions + +Per-filter completion is inferred from the number of returned events because +NIP-01 does not provide a pagination cursor or an explicit "filter exhausted" +signal. Grouped historic sync therefore assumes that a relay: + +- applies its result limit independently to every filter in the REQ; +- returns at least 75 events for a non-exhausted filter; and +- does not impose an additional total-result cap across the whole REQ that can + allow one filter to starve another. + +ngit-grasp's relay implementation has these semantics: it queries each filter +with its own limit before merging and deduplicating the results. Relays with a +smaller hidden per-filter cap or a shared total-result cap can cause historic +sync to conclude prematurely, so compatibility with those implementations is +not currently guaranteed. + ### Pagination State Lifecycle ```mermaid flowchart TB REQ[Send REQ with filters] --> TRACK[Initialize PaginationState] TRACK --> EVENT[Receive EVENT] - EVENT --> UPDATE[Update event_count and min_created_at] + EVENT --> UPDATE[Update each matching filter's count and oldest timestamp] UPDATE --> MORE{More events?} MORE --> |yes| EVENT MORE --> |no| EOSE[Receive EOSE] EOSE --> CHECK{event_count suggests more pages?} - CHECK --> |yes| NEXT[Create filter with until=min_created_at] - NEXT --> REQ2[Send next page REQ] - REQ2 --> RESET[Reset event_count, keep min_created_at] + CHECK --> |yes| NEXT[Create next filters with their own until timestamps] + NEXT --> REQ2[Send grouped next page REQ] + REQ2 --> RESET[Reset per-filter counters] RESET --> EVENT CHECK --> |no| DONE[Batch complete, confirm items] ``` @@ -880,7 +912,7 @@ flowchart TB | **Efficiency** | High (set reconciliation) | Lower (sequential pages) | | **Bandwidth** | Minimal (only missing items) | Higher (all matching events transferred) | | **Relay support** | Requires NIP-77 | Universal (standard Nostr) | -| **State tracking** | None needed | [`PaginationState`](src/sync/mod.rs:165) per subscription | +| **State tracking** | None needed | Per-filter state within each grouped subscription | | **Completion time** | Typically faster | Slower for large sets | | **Use cases** | Full sync, large event sets | Fallback, small gaps with `since` | diff --git a/src/sync/mod.rs b/src/sync/mod.rs index d09fe22..6c62883 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -371,15 +371,70 @@ pub struct ReprocessingStats { /// Pagination state for a subscription in non-Negentropy historic sync #[derive(Debug, Clone)] -pub struct PaginationState { - /// Number of events received for this subscription +pub struct FilterPaginationState { + /// Number of events received for this filter pub event_count: usize, - /// Smallest created_at timestamp seen (for pagination with `until`) + /// Smallest created_at timestamp seen for this filter pub min_created_at: Option, /// Original filter to reconstruct for next page pub original_filter: Filter, } +/// Pagination state for every OR filter carried by one subscription. +#[derive(Debug, Clone)] +pub struct PaginationState { + pub filters: Vec, +} + +impl PaginationState { + fn new(filters: Vec) -> Self { + Self { + filters: filters + .into_iter() + .map(|original_filter| FilterPaginationState { + event_count: 0, + min_created_at: None, + original_filter, + }) + .collect(), + } + } + + fn record_event(&mut self, event: &Event) { + for state in &mut self.filters { + if state + .original_filter + .match_event(event, MatchEventOptions::new()) + { + state.event_count += 1; + match state.min_created_at { + None => state.min_created_at = Some(event.created_at), + Some(min) if event.created_at < min => { + state.min_created_at = Some(event.created_at); + } + _ => {} + } + } + } + } + + fn next_page_filters(self) -> Vec { + self.filters + .into_iter() + .filter_map(|state| { + (state.event_count >= PAGINATION_THRESHOLD) + .then_some(state.min_created_at) + .flatten() + .map(|min_created_at| { + state + .original_filter + .until(Timestamp::from(min_created_at.as_secs())) + }) + }) + .collect() + } +} + /// A batch of items pending confirmation #[derive(Debug, Clone)] pub struct PendingBatch { @@ -503,10 +558,22 @@ const CONSOLIDATION_THRESHOLD: usize = 70; /// exhaust network resources while keeping the sync actor responsive. const MAX_CONCURRENT_CONNECT_ATTEMPTS: usize = 8; -/// Page size threshold for historic sync pagination (non-negentropy) -/// If a subscription receives >= 75 events, we fetch the next page +/// Per-filter threshold for historic REQ+EOSE pagination. +/// +/// Grouped pagination assumes that relays apply result limits independently to +/// each filter, return at least this many events for a non-exhausted filter, +/// and do not impose an additional total-result cap across the whole REQ. This +/// matches NIP-01's per-filter `limit` model and ngit-grasp's relay behavior. +/// Relays that violate these assumptions can make one filter appear exhausted +/// after another filter consumes the combined result allowance. const PAGINATION_THRESHOLD: usize = 75; +/// Conservative number of OR filters carried by one NIP-01 REQ. +/// +/// This keeps active subscription counts low without producing unusually +/// large REQ messages for relays that enforce their own per-REQ filter limits. +const MAX_FILTERS_PER_REQ: usize = 10; + fn reserve_connect_attempt( in_flight: &mut HashMap, next_token: &mut u64, @@ -560,6 +627,10 @@ fn should_consolidate(current_count: usize, new_count: usize, desired_baseline: > CONSOLIDATION_THRESHOLD } +fn grouped_subscription_count(filter_count: usize) -> usize { + filter_count.div_ceil(MAX_FILTERS_PER_REQ) +} + #[derive(Debug, Default)] struct DeferredConsolidations { relays: HashSet, @@ -1113,130 +1184,103 @@ impl SyncManager { // Check for pagination: if this subscription hit the threshold, fetch next page if let Some(pagination_state) = batch.pagination_state.remove(&sub_id) { - if pagination_state.event_count >= PAGINATION_THRESHOLD { - if let Some(min_created_at) = pagination_state.min_created_at { - tracing::info!( - relay = %relay_url, - sub_id = %sub_id, - batch_id = batch.batch_id, - event_count = pagination_state.event_count, - min_created_at = %min_created_at, - "Subscription hit pagination threshold, fetching next page" + let next_filters = pagination_state.next_page_filters(); + if !next_filters.is_empty() { + let relay_url_for_pagination = relay_url.to_string(); + let batch_id = batch.batch_id; + tracing::info!( + relay = %relay_url, + sub_id = %sub_id, + batch_id, + filter_count = next_filters.len(), + "Grouped subscription hit pagination threshold, fetching next page" + ); + + // A NOTICE can arrive immediately before this page's EOSE. + // Keep a sentinel in the batch and let a detached worker + // resume the exact grouped page after the cooldown. + if self.health_tracker.is_rate_limited(relay_url) { + let deferred_sub_id = mark_deferred_pagination(batch, &sub_id); + drop(pending); + + let Some(connection) = self.connections.get(&relay_url_for_pagination).cloned() + else { + tracing::error!( + relay = %relay_url_for_pagination, + batch_id, + "Cannot defer rate-limited pagination without a relay connection" + ); + return; + }; + Self::spawn_deferred_pagination( + connection, + self.health_tracker.clone(), + self.pending_sync_index.clone(), + relay_url_for_pagination, + batch_id, + deferred_sub_id, + next_filters, ); + return; + } - // Create next page filter: same as original but with .until(min_created_at) - // dont subtract 1 second to avoid duplicate events at the boundary - // as this would lead to missed events with the same created_at timestamp - let until_timestamp = Timestamp::from(min_created_at.as_secs()); - let mut next_filter = pagination_state.original_filter.clone(); - next_filter = next_filter.until(until_timestamp); + drop(pending); - // Store relay_url for spawning the subscription after releasing the lock - let relay_url_for_pagination = relay_url.to_string(); - let batch_id = batch.batch_id; - - // A NOTICE can arrive immediately before this page's EOSE. - // Never wait for the cooldown here: the caller owns the - // global SyncManager mutex, while the health checker that - // clears the cooldown needs that same mutex. Keep a - // sentinel in the batch and resume this exact page from a - // detached worker so generic history is not lost or - // restarted from page one. - if self.health_tracker.is_rate_limited(relay_url) { - let deferred_sub_id = mark_deferred_pagination(batch, &sub_id); - drop(pending); - - let Some(connection) = - self.connections.get(&relay_url_for_pagination).cloned() - else { + let mut next_page_started = false; + if let Some(conn) = self.connections.get(&relay_url_for_pagination) { + match conn.subscribe_filters(next_filters.clone(), true).await { + Ok(new_sub_id) => { + let mut pending = self.pending_sync_index.write().await; + if let Some(batches) = pending.get_mut(&relay_url_for_pagination) { + if let Some(batch) = + batches.iter_mut().find(|b| b.batch_id == batch_id) + { + batch.outstanding_subs.insert(new_sub_id.clone()); + next_page_started = true; + batch.pagination_state.insert( + new_sub_id.clone(), + PaginationState::new(next_filters), + ); + tracing::info!( + relay = %relay_url_for_pagination, + new_sub_id = %new_sub_id, + batch_id, + "Next grouped page subscription created" + ); + } + } + } + Err(error) => { tracing::error!( relay = %relay_url_for_pagination, batch_id, - "Cannot defer rate-limited pagination without a relay connection" + error = %error, + "Failed to create grouped pagination subscription" ); - return; - }; - Self::spawn_deferred_pagination( - connection, - self.health_tracker.clone(), - self.pending_sync_index.clone(), - relay_url_for_pagination, - batch_id, - deferred_sub_id, - next_filter, - until_timestamp, - ); - return; - } - - // Drop the lock before async operations - drop(pending); - - // Subscribe to next page and add to outstanding_subs - let mut next_page_started = false; - if let Some(conn) = self.connections.get(&relay_url_for_pagination) { - match conn.subscribe_filter(next_filter.clone(), true).await { - Ok(new_sub_id) => { - // Re-acquire lock to update the batch - let mut pending = self.pending_sync_index.write().await; - if let Some(batches) = pending.get_mut(&relay_url_for_pagination) { - if let Some(batch) = - batches.iter_mut().find(|b| b.batch_id == batch_id) - { - batch.outstanding_subs.insert(new_sub_id.clone()); - next_page_started = true; - // Initialize pagination state for new subscription - batch.pagination_state.insert( - new_sub_id.clone(), - PaginationState { - event_count: 0, - min_created_at: None, - original_filter: next_filter, - }, - ); - tracing::info!( - relay = %relay_url_for_pagination, - new_sub_id = %new_sub_id, - batch_id = batch_id, - until = %until_timestamp, - "Next page subscription created" - ); - } - } - } - Err(e) => { - tracing::error!( - relay = %relay_url_for_pagination, - batch_id = batch_id, - error = %e, - "Failed to create pagination subscription, continuing without next page" - ); - } } } - - if !next_page_started { - let completed_batch = { - let mut pending = self.pending_sync_index.write().await; - take_drained_batch_as_failed( - &mut pending, - &relay_url_for_pagination, - batch_id, - ) - }; - if let Some(batch) = completed_batch { - tracing::warn!( - relay = %relay_url_for_pagination, - batch_id, - "Pagination could not continue; completing drained batch as failed" - ); - self.confirm_batch(&relay_url_for_pagination, batch).await; - } - } - - // Early return since we've released and re-acquired locks - return; } + + if !next_page_started { + let completed_batch = { + let mut pending = self.pending_sync_index.write().await; + take_drained_batch_as_failed( + &mut pending, + &relay_url_for_pagination, + batch_id, + ) + }; + if let Some(batch) = completed_batch { + tracing::warn!( + relay = %relay_url_for_pagination, + batch_id, + "Pagination could not continue; completing drained batch as failed" + ); + self.confirm_batch(&relay_url_for_pagination, batch).await; + } + } + + return; } } @@ -1323,8 +1367,8 @@ impl SyncManager { let mut new_sub_ids = HashSet::new(); if let Some(conn) = self.connections.get(&relay_url_for_fallback) { - for filter in fallback_filters { - match conn.subscribe_filter(filter, true).await { + for filter_group in fallback_filters.chunks(MAX_FILTERS_PER_REQ) { + match conn.subscribe_filters(filter_group.to_vec(), true).await { Ok(sub_id) => { new_sub_ids.insert(sub_id); } @@ -1521,14 +1565,13 @@ impl SyncManager { relay_url: String, batch_id: u64, deferred_sub_id: SubscriptionId, - next_filter: Filter, - until_timestamp: Timestamp, + next_filters: Vec, ) { tokio::spawn(async move { tracing::info!( relay = %relay_url, batch_id, - until = %until_timestamp, + filter_count = next_filters.len(), "Rate limited during historic pagination; deferring the exact next page without blocking the sync actor" ); @@ -1557,23 +1600,20 @@ impl SyncManager { return; } - match connection.subscribe_filter(next_filter.clone(), true).await { + match connection + .subscribe_filters(next_filters.clone(), true) + .await + { Ok(new_sub_id) => { batch.outstanding_subs.remove(&deferred_sub_id); batch.outstanding_subs.insert(new_sub_id.clone()); - batch.pagination_state.insert( - new_sub_id.clone(), - PaginationState { - event_count: 0, - min_created_at: None, - original_filter: next_filter, - }, - ); + batch + .pagination_state + .insert(new_sub_id.clone(), PaginationState::new(next_filters)); tracing::info!( relay = %relay_url, new_sub_id = %new_sub_id, batch_id, - until = %until_timestamp, "Deferred pagination resumed after rate-limit cooldown" ); return; @@ -2125,8 +2165,11 @@ impl SyncManager { } // Step 3: Check if consolidation is needed BEFORE adding new filters - self.maybe_consolidate(&action.relay_url, action.filters.len()) - .await; + self.maybe_consolidate( + &action.relay_url, + grouped_subscription_count(action.filters.len()), + ) + .await; // Subscribe to each filter and collect subscription IDs tracing::info!( @@ -2293,15 +2336,7 @@ impl SyncManager { if let Some(state) = batch.pagination_state.get_mut(&subscription_id) { - state.event_count += 1; - // Track minimum created_at timestamp - match state.min_created_at { - None => state.min_created_at = Some(event.created_at), - Some(min) if event.created_at < min => { - state.min_created_at = Some(event.created_at); - } - _ => {} - } + state.record_event(&event); } // Track received event IDs (negentropy path) @@ -3989,7 +4024,7 @@ impl SyncManager { // Every connected relay carries one consolidated generic announcement // subscription in addition to its repository-specific desired filters. - 1 + desired_repo_filters + 1 + grouped_subscription_count(desired_repo_filters) } /// Check if incremental fragmentation exceeds the consolidation threshold. @@ -4372,15 +4407,17 @@ impl SyncManager { let mut sub_ids = Vec::new(); - for filter in filters.iter() { + for filter_group in filters.chunks(MAX_FILTERS_PER_REQ) { // Live subscriptions MUST use limit(0) to receive ONLY new events // This prevents fetching historic events that would be miscounted as "live" in metrics // The caller passes the same filters to both sync_live() and historic_sync() // Live subscriptions do NOT auto-close - we want them to stay open for new events - match connection - .subscribe_filter(filter.clone().limit(0), false) - .await - { + let grouped_filters = filter_group + .iter() + .cloned() + .map(|filter| filter.limit(0)) + .collect(); + match connection.subscribe_filters(grouped_filters, false).await { Ok(sub_id) => { sub_ids.push(sub_id); } @@ -4696,29 +4733,23 @@ impl SyncManager { let mut subscription_ids = HashSet::new(); let mut pagination_state = HashMap::new(); - // DEBUG TRACING: Log each filter in REQ+EOSE path - for (idx, filter) in filters_with_since.iter().enumerate() { + // Keep several OR filters under each relay-visible subscription. + for (idx, filter_group) in filters_with_since.chunks(MAX_FILTERS_PER_REQ).enumerate() { tracing::debug!( relay = %relay_url, batch_id = batch_id, - filter_idx = idx, - filter = ?filter, - "Subscribing to filter in REQ+EOSE path" + group_idx = idx, + filter_count = filter_group.len(), + filters = ?filter_group, + "Subscribing to grouped filters in REQ+EOSE path" ); if let Some(conn) = self.connections.get(relay_url) { - match conn.subscribe_filter(filter.clone(), true).await { + let grouped_filters = filter_group.to_vec(); + match conn.subscribe_filters(grouped_filters.clone(), true).await { Ok(sub_id) => { subscription_ids.insert(sub_id.clone()); - // Initialize pagination state for this subscription - pagination_state.insert( - sub_id, - PaginationState { - event_count: 0, - min_created_at: None, - original_filter: filter.clone(), - }, - ); + pagination_state.insert(sub_id, PaginationState::new(grouped_filters)); } Err(e) => { tracing::error!( @@ -5030,6 +5061,48 @@ mod tests { )); } + #[test] + fn grouped_pagination_advances_only_filters_that_fill_a_page() { + let keys = Keys::generate(); + let metadata_filter = Filter::new().kind(Kind::Metadata); + let note_filter = Filter::new().kind(Kind::TextNote); + let mut pagination = + PaginationState::new(vec![metadata_filter.clone(), note_filter.clone()]); + + for created_at in 1..=PAGINATION_THRESHOLD { + let event = EventBuilder::new(Kind::Metadata, created_at.to_string()) + .custom_created_at(Timestamp::from_secs(created_at as u64)) + .finalize(&keys) + .expect("build metadata event"); + pagination.record_event(&event); + } + let note = EventBuilder::text_note("one note") + .custom_created_at(Timestamp::from_secs(100)) + .finalize(&keys) + .expect("build text note"); + pagination.record_event(¬e); + + let next_filters = pagination.next_page_filters(); + assert_eq!(next_filters.len(), 1); + assert_eq!(next_filters[0].until, Some(Timestamp::from_secs(1))); + assert!( + next_filters[0] + .kinds + .as_ref() + .unwrap() + .contains(&Kind::Metadata), + "the full metadata filter should advance" + ); + assert!( + !next_filters[0] + .kinds + .as_ref() + .unwrap() + .contains(&Kind::TextNote), + "the partial text-note filter should not advance" + ); + } + #[test] fn deferred_consolidation_runs_only_after_final_batch_completion() { let relay_url = "wss://relay.example"; diff --git a/tests/common/mock_relay.rs b/tests/common/mock_relay.rs index 063287b..377784b 100644 --- a/tests/common/mock_relay.rs +++ b/tests/common/mock_relay.rs @@ -72,6 +72,19 @@ impl MockRelay { /// The relay accepts all events without validation and stores them /// in an in-memory database. pub async fn start() -> Self { + Self::start_with_rate_limit(RateLimit::default()).await + } + + /// Start a mock relay with a custom per-connection active REQ limit. + pub async fn start_with_max_reqs(max_reqs: usize) -> Self { + Self::start_with_rate_limit(RateLimit { + max_reqs, + ..RateLimit::default() + }) + .await + } + + async fn start_with_rate_limit(rate_limit: RateLimit) -> Self { // Create and bind listener (eliminates port race condition) let std_listener = std::net::TcpListener::bind("127.0.0.1:0").expect("Failed to bind to random port"); @@ -87,7 +100,7 @@ impl MockRelay { let listener = TcpListener::from_std(std_listener).expect("Failed to convert to tokio listener"); - Self::start_with_listener(listener, port).await + Self::start_with_listener(listener, port, rate_limit).await } /// Start a mock relay on a specific port. @@ -96,13 +109,13 @@ impl MockRelay { let listener = TcpListener::bind(addr) .await .expect("Failed to bind to address"); - Self::start_with_listener(listener, port).await + Self::start_with_listener(listener, port, RateLimit::default()).await } /// Internal method to start the relay with an existing listener. - async fn start_with_listener(listener: TcpListener, port: u16) -> Self { + async fn start_with_listener(listener: TcpListener, port: u16, rate_limit: RateLimit) -> Self { // Create a simple relay with no write policy (accepts all events) - let relay = LocalRelayBuilder::default().build(); + let relay = LocalRelayBuilder::default().rate_limit(rate_limit).build(); // Create shutdown channel let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); @@ -169,6 +182,11 @@ impl MockRelay { &self.url } + /// Get the relay domain as a host and port. + pub fn domain(&self) -> String { + format!("127.0.0.1:{}", self.port) + } + /// Stop the mock relay. pub async fn stop(mut self) { // Send shutdown signal diff --git a/tests/sync/live_sync.rs b/tests/sync/live_sync.rs index 5cfe0ad..738b6fe 100644 --- a/tests/sync/live_sync.rs +++ b/tests/sync/live_sync.rs @@ -22,7 +22,89 @@ use std::time::Duration; use nostr_sdk::prelude::*; -use crate::common::{sync_helpers::*, TestRelay}; +use crate::common::{sync_helpers::*, MockRelay, TestRelay}; + +/// A source relay's active-REQ cap must not silently remove one of the +/// repository filter variants. +/// +/// The generic announcement subscription occupies one active REQ. A single +/// repository then needs state, a, A, and q filters. Installing each filter as +/// a separate live subscription exceeds this source's four-REQ limit, leaving +/// q-tagged collaboration events permanently uncovered. +#[tokio::test] +async fn test_live_sync_batches_repo_filters_below_source_req_limit() { + let source = MockRelay::start_with_max_reqs(4).await; + let syncing = TestRelay::start_with_sync(None).await; + let keys = Keys::generate(); + let repo_id = "test-repo-bounded-reqs"; + let domains = [source.domain(), syncing.domain()]; + let domain_refs: Vec<&str> = domains.iter().map(String::as_str).collect(); + + let (announcement, _git_dir) = + setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await; + + let source_client = TestClient::new(source.url(), keys.clone()) + .await + .expect("connect to constrained source relay"); + source_client + .send_event(&announcement) + .await + .expect("publish announcement to constrained source relay"); + + wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(5)) + .await + .expect("syncing relay should connect to constrained source"); + + // Observe one q-tagged event completing the round trip before testing a + // second live event. This proves the relevant subscription is installed + // without relying on an arbitrary scheduling delay. + let readiness_issue = build_layer2_issue_with_q_tag( + &keys, + &repo_coord(&keys, repo_id), + "Subscription readiness probe", + ) + .expect("build q-tagged readiness issue"); + source_client + .send_event(&readiness_issue) + .await + .expect("publish q-tagged readiness issue"); + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(readiness_issue.id), + Duration::from_secs(5), + ) + .await, + "q-tagged readiness issue should sync before testing live delivery" + ); + + let issue = build_layer2_issue_with_q_tag( + &keys, + &repo_coord(&keys, repo_id), + "Issue behind the final repository filter", + ) + .expect("build q-tagged issue"); + source_client + .send_event(&issue) + .await + .expect("publish q-tagged issue"); + + let synced = wait_for_event_on_relay( + syncing.url(), + Filter::new().id(issue.id), + Duration::from_secs(5), + ) + .await; + + source_client.disconnect().await; + syncing.stop().await; + source.stop().await; + + assert!( + synced, + "q-tagged issue should sync even when the source permits only four active REQs" + ); +} /// Test 5: Live sync Layer 2 events /// From 8bce3970fcf0a6e6f757add2e342f723681c081b Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 1 Aug 2026 13:44:43 +0100 Subject: [PATCH 3/3] fix(sync): retry deliberately after relay rate limits Some relays repeat equivalent rate-limit notices throughout an active cooldown. Recomputing the deadline for every reminder lets periodic notices postpone recovery forever. Keep the first deadline while its cooldown is active. Once that deadline expires, treat a new rejection as a failed recovery probe and begin a fresh cooldown, whether the relay retained or extended its limit window. Do not clear an active cooldown merely because the WebSocket connected successfully: transport availability does not prove that the relay will accept a new REQ. Cover duplicate notices, post-deadline rejection, and connection-success behavior with unit tests. --- CHANGELOG.md | 3 + docs/explanation/grasp-02-proactive-sync.md | 11 ++- src/sync/health.rs | 94 +++++++++++++++++---- src/sync/mod.rs | 2 +- 4 files changed, 90 insertions(+), 20 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cd846c1..1d430a7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Fixed proactive sync losing repository events when public relays cap active subscriptions. Compatible GRASP filters now share bounded NIP-01 REQs while retaining per-filter history pagination. +- Fixed repeated relay rate-limit notices extending the cooldown indefinitely. + Notices during an active cooldown keep its original deadline, while a new + rejection after recovery begins a fresh cooldown. - Removed superseded same-author repository states from purgatory after their replacement is promoted when the locally available Git data cannot reconstruct them. Reconstructable rollback states and other maintainers' diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index c369c84..92a20f1 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -1006,7 +1006,9 @@ Degraded -> Dead: 24h+ of continuous failures Degraded -> Disconnected: Recovery (enters 5min stability period) Disconnected -> Healthy: Stable for 5 minutes after recovery Any -> RateLimited: NOTICE message from relay indicating rate limiting -RateLimited -> previous state: After 65-second cooldown expires +RateLimited -> Probing: After 65-second cooldown expires +Probing -> previous state: Recovery REQs succeed +Probing -> RateLimited: Recovery REQ is rate limited again ``` ### Backoff Configuration @@ -1016,7 +1018,9 @@ RateLimited -> previous state: After 65-second cooldown expires - **Default max**: 1 hour (configurable via `sync_max_backoff_secs`) - **Dead threshold**: 24 hours of continuous failures - **Dead retry interval**: Once per 24 hours -- **Rate limit cooldown**: Fixed 65 seconds (60s typical limit + 5s buffer) +- **Rate limit cooldown**: Fixed 65 seconds (60s typical limit + 5s buffer); + repeated notices during the same cooldown do not extend its deadline, while + a rejection after that deadline starts a new cooldown - **Stability period**: 5 minutes after recovery before marking as Healthy ### Special Behaviors @@ -1025,7 +1029,8 @@ RateLimited -> previous state: After 65-second cooldown expires - **Desired GRASP-02 sources**: Remain registered and retryable before their first successful historic batch; an initially empty or unavailable source cannot make a purgatory invitation permanently lose its sync path -- **Rate limiting**: Distinct from connection failures - triggered by relay NOTICE messages +- **Rate limiting**: Distinct from connection failures and therefore not cleared + by a successful WebSocket connection; it is triggered by relay NOTICE messages - **Connection timeout**: Set to `base_backoff_secs` to ensure retry timing works correctly - **Connection concurrency**: At most eight DNS/websocket attempts run at once; queued attempts do not start health backoff until a worker slot is available diff --git a/src/sync/health.rs b/src/sync/health.rs index 833918b..f82738b 100644 --- a/src/sync/health.rs +++ b/src/sync/health.rs @@ -105,7 +105,7 @@ impl RelayHealth { /// /// ## State Logic /// - /// 1. **RateLimited**: If rate_limited flag is set and cooldown hasn't expired + /// 1. **RateLimited**: If the rate-limit cooldown hasn't expired /// 2. **Dead**: 24+ hours of continuous failures /// 3. **Degraded**: Active connection failures OR in stability period after recovery /// 4. **Disconnected**: Not connected, but no recent failures or issues @@ -274,29 +274,41 @@ impl RelayHealthTracker { /// Record a successful connection to a relay /// - /// Clears failure counters and rate limiting. Sets connected = true. + /// Clears connection failure counters. Sets connected = true. + /// + /// A successful WebSocket connection does not prove that the relay will + /// accept a new REQ, so it deliberately leaves any active rate-limit + /// cooldown unchanged. pub fn record_success(&self, relay_url: &str) { let now = Instant::now(); let mut entry = self.health.entry(relay_url.to_string()).or_default(); let health = entry.value_mut(); let old_state = health.state(); + let active_rate_limit = health + .rate_limited + .then_some(health.next_retry_at) + .flatten() + .filter(|deadline| *deadline > now); - // Reset to healthy state + // Reset connection health. A live rate-limit cooldown is independent + // of whether the WebSocket handshake succeeded. health.connected = true; - health.rate_limited = false; + health.rate_limited = active_rate_limit.is_some(); health.consecutive_failures = 0; health.first_failure_time = None; health.last_failure_time = None; health.last_success_time = Some(now); health.last_attempt_time = Some(now); - health.next_retry_at = None; + health.next_retry_at = active_rate_limit; - if old_state != HealthState::Healthy { + let new_state = health.state(); + if old_state != new_state { tracing::info!( - "Relay {} recovered to healthy (was {:?})", + "Relay {} connection recovered ({:?} -> {:?})", relay_url, - old_state + old_state, + new_state ); } } @@ -382,6 +394,13 @@ impl RelayHealthTracker { let mut entry = self.health.entry(relay_url.to_string()).or_default(); let health = entry.value_mut(); + // A relay may repeat the same NOTICE throughout a cooldown. Ignore + // reminders for that episode, but treat a rejection after the deadline + // as a failed recovery probe and begin a new cooldown. + if health.rate_limited && health.next_retry_at.is_some_and(|deadline| now < deadline) { + return; + } + health.rate_limited = true; health.next_retry_at = Some(now + Duration::from_secs(RATE_LIMIT_COOLDOWN_SECS)); @@ -394,15 +413,14 @@ impl RelayHealthTracker { /// Clear rate limiting state for a specific relay /// - /// This only clears the rate_limited flag, without affecting connection status - /// or failure counters. Use this when rate limit cooldown has expired and we - /// want to allow new subscriptions. - /// - /// This is different from `record_success()` which resets all health state. + /// This clears the rate-limit episode without affecting connection status + /// or failure counters. Use this when the cooldown has expired and new + /// subscriptions may probe the relay again. pub fn clear_rate_limit(&self, relay_url: &str) { if let Some(mut entry) = self.health.get_mut(relay_url) { let health = entry.value_mut(); health.rate_limited = false; + health.next_retry_at = None; } } @@ -414,7 +432,7 @@ impl RelayHealthTracker { pub fn is_rate_limited(&self, relay_url: &str) -> bool { if let Some(entry) = self.health.get(relay_url) { let health = entry.value(); - health.rate_limited + health.is_rate_limited_now() } else { false } @@ -437,8 +455,8 @@ impl RelayHealthTracker { // Check if rate limited and cooldown has expired if health.rate_limited { - if let Some(next_retry) = health.next_retry_at { - if now > next_retry { + if let Some(deadline) = health.next_retry_at { + if now >= deadline { // Cooldown expired - clear rate limiting health.rate_limited = false; health.next_retry_at = None; @@ -749,4 +767,48 @@ mod tests { let health = tracker.get_health("wss://nonexistent.example.com"); assert!(health.is_none()); } + + #[test] + fn repeated_rate_limit_notice_does_not_extend_cooldown() { + let tracker = RelayHealthTracker::with_defaults(); + let relay = "wss://limited.example"; + + tracker.record_rate_limit(relay); + let first_deadline = tracker.get_health(relay).unwrap().next_retry_at; + tracker.record_rate_limit(relay); + + assert_eq!( + tracker.get_health(relay).unwrap().next_retry_at, + first_deadline + ); + } + + #[test] + fn rate_limit_notice_after_deadline_starts_new_cooldown() { + let tracker = RelayHealthTracker::with_defaults(); + let relay = "wss://limited.example"; + + tracker.record_rate_limit(relay); + tracker.health.get_mut(relay).unwrap().next_retry_at = + Some(Instant::now() - Duration::from_millis(1)); + + tracker.record_rate_limit(relay); + + assert!(tracker.get_health(relay).unwrap().is_rate_limited_now()); + } + + #[test] + fn connection_success_does_not_clear_rate_limit_cooldown() { + let tracker = RelayHealthTracker::with_defaults(); + let relay = "wss://limited.example"; + + tracker.record_rate_limit(relay); + let deadline = tracker.get_health(relay).unwrap().next_retry_at; + tracker.record_success(relay); + + let health = tracker.get_health(relay).unwrap(); + assert!(health.is_rate_limited_now()); + assert_eq!(health.next_retry_at, deadline); + assert!(health.connected); + } } diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 6c62883..c506a5f 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -4313,7 +4313,7 @@ impl SyncManager { /// Check for rate-limited relays that have exceeded cooldown /// - /// This method is called periodically by run_rate_limit_checker (every 1 second). + /// This method is called by the health and metrics checker every 2 seconds. /// For each relay in RateLimited state that has exceeded the 65-second cooldown: /// 1. Clears the rate limit state (sets to Healthy) /// 2. Recomputes required actions for that relay