diff --git a/CHANGELOG.md b/CHANGELOG.md index cc58b63..1d430a7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- 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 6cd1754..92a20f1 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` | @@ -974,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 @@ -984,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 @@ -993,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 d09fe22..c506a5f 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. @@ -4278,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 @@ -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/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) 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 ///