diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index e77fb45..3a39e26 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -198,11 +198,11 @@ Consequences: - Correct against any relay whose effective per-filter cap is ≥ 200. The smallest audited defaults are 250 (haven/Badger) and 300 (rnostr); the - gap below 250 is deliberate margin, because only events processed as - Saved or Duplicate count toward the threshold — events routed to - purgatory or rejected consume the relay's allowance without being - counted, so a threshold equal to a relay's cap would mistake a full page - for an exhausted filter. + current gap below 250 is no longer required for event-accounting safety: + every raw delivery matching a tracked filter now counts before + deduplication or write-policy processing. Purgatory-routed, rejected, and + repeated events therefore consume both the relay's allowance and our page + count, and the `until` cursor is derived from that same raw stream. - A relay capping a filter below 200, or enforcing an aggregate per-REQ cap, silently truncates history. No audited implementation has an aggregate cap. **Known live exception (2026-08-06): Ditto Relay applies @@ -219,16 +219,13 @@ Consequences: raised from its original ultra-conservative 75 once the audit established the real floor; filters with 75–199 results no longer pay the extra page. -- Planned design (accepted 2026-08-06, not implemented): keep omitting +- Implementation in progress (design accepted 2026-08-06): keep omitting `limit` — an explicit limit would cap the relays that serve unbounded - pages — fix the counting, and adapt the threshold per relay: - 1. **Count raw delivered events.** Today only events processed as - Saved or Duplicate count toward the threshold and the `until` - cursor, so purgatory-routed and rejected events consume relay - allowance invisibly; this is the sole reason thresholds need - margin. Counting every delivered event that matches the filter - (and cursoring on them) makes a truncated page count exactly the - relay's page size. + pages — count raw deliveries, and adapt the threshold per relay: + 1. **Count raw delivered events (implemented).** Every delivered event + that matches a tracked filter is counted before deduplication and write + policy, and the cursor uses the same stream. Purgatory-routed, rejected, + and repeated events can no longer consume relay allowance invisibly. 2. **Adaptive per-relay threshold:** `estimated_cap = max(largest observed page, advertised default_limit if present)`; diff --git a/src/sync/mod.rs b/src/sync/mod.rs index cf3b8f7..e6d197b 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -597,21 +597,17 @@ const MAX_CONCURRENT_CONNECT_ATTEMPTS: usize = 8; /// always applied per filter, never in aggregate across a REQ, with the /// smallest finite default caps at 250 (haven/Badger) and 300 (rnostr). /// -/// The threshold must sit below the smallest cap we may meet, with margin: -/// only events processed as Saved or Duplicate are counted here, so events -/// routed to purgatory or rejected consume the relay's allowance without -/// being counted, and a threshold equal to a relay's cap would mistake a -/// full page for an exhausted filter. 200 keeps a 50-event margin under the -/// tightest audited default while sparing filters with fewer than 200 -/// results the redundant final page. A relay capped below this threshold -/// silently truncates history. Known live exception (2026-08-06): Ditto -/// Relay applies a 100-event default to filters that omit `limit` — which -/// ours do — while advertising only its larger explicit-request cap in -/// NIP-11. The accepted mitigation (not yet implemented) is to count raw -/// delivered events instead of only Saved/Duplicate ones and adapt the -/// threshold per relay from observed page sizes and advertised -/// `default_limit`, with a floor of 90. See "Per-query result limits and -/// the pagination model" in docs/explanation/sync-scaling-constraints.md. +/// Every raw delivery matching a tracked filter counts, before write-policy +/// processing. Purgatory-routed, rejected, and repeated events therefore +/// consume both the relay's allowance and our page count, and the cursor is +/// derived from that same raw stream. 200 remains a temporary static floor: +/// a relay capped below it silently truncates history. Known live exception +/// (2026-08-06): Ditto Relay applies a 100-event default to filters that omit +/// `limit` — which ours do — while advertising only its larger +/// explicit-request cap in NIP-11. The remaining accepted mitigation is to +/// adapt the threshold per relay from observed raw page sizes and advertised +/// `default_limit`, with a floor of 90. See "Per-query result limits and the +/// pagination model" in docs/explanation/sync-scaling-constraints.md. const PAGINATION_THRESHOLD: usize = 200; /// Conservative number of OR filters carried by one NIP-01 REQ. @@ -2798,6 +2794,23 @@ impl SyncManager { while let Some(relay_event) = event_rx.recv().await { match relay_event { RelayEvent::Event(event, subscription_id) => { + // Count raw deliveries before deduplication or write policy. Relays spend + // their result allowance on every matching delivery, including events we + // route to purgatory, reject, or have already stored; pagination must use + // that same stream for both its page count and `until` cursor. + { + let mut pending = pending_sync_index.write().await; + if let Some(batches) = pending.get_mut(&relay_url_clone) { + for batch in batches.iter_mut() { + if let Some(state) = + batch.pagination_state.get_mut(&subscription_id) + { + state.record_event(&event); + } + } + } + } + // Skip events we've already rejected (announcements only) if (event.kind == Kind::GitRepoAnnouncement || event.kind == Kind::RepoState) @@ -2869,19 +2882,13 @@ impl SyncManager { } } - // Track pagination state for this subscription (REQ+EOSE) - // and received event IDs for negentropy batches + // Track received event IDs for negentropy batches. Unlike REQ+EOSE + // pagination above, negentropy completion is concerned with events that + // were actually saved or already present locally. if result == ProcessResult::Saved || result == ProcessResult::Duplicate { let mut pending = pending_sync_index.write().await; if let Some(batches) = pending.get_mut(&relay_url_clone) { for batch in batches.iter_mut() { - // Track pagination state (REQ+EOSE path) - if let Some(state) = - batch.pagination_state.get_mut(&subscription_id) - { - state.record_event(&event); - } - // Track received event IDs (negentropy path) // Only track if this batch has requested_event_ids set // and the subscription is one we're waiting on @@ -5832,6 +5839,86 @@ mod tests { ); } + #[test] + fn raw_pagination_counts_purgatory_and_rejected_deliveries() { + let keys = Keys::generate(); + let filter = Filter::new().kinds([Kind::GitRepoAnnouncement, Kind::RepoState]); + let mut pagination = PaginationState::new(vec![filter]); + let deliveries = [ + ( + EventBuilder::new(Kind::GitRepoAnnouncement, "purgatory") + .custom_created_at(Timestamp::from_secs(20)) + .finalize(&keys) + .expect("build purgatory-routed event"), + ProcessResult::Purgatory, + ), + ( + EventBuilder::new(Kind::RepoState, "rejected") + .custom_created_at(Timestamp::from_secs(10)) + .finalize(&keys) + .expect("build rejected event"), + ProcessResult::Rejected, + ), + ]; + + for (event, policy_result) in deliveries { + // The production handler records here, before it knows this result. + pagination.record_event(&event); + assert!(matches!( + policy_result, + ProcessResult::Purgatory | ProcessResult::Rejected + )); + } + + assert_eq!(pagination.filters[0].event_count, 2); + assert_eq!( + pagination.filters[0].min_created_at, + Some(Timestamp::from_secs(10)) + ); + } + + #[test] + fn raw_pagination_counts_repeat_deliveries_at_the_boundary() { + let keys = Keys::generate(); + let filter = Filter::new().kind(Kind::TextNote); + let mut pagination = PaginationState::new(vec![filter]); + let event = EventBuilder::new(Kind::TextNote, "same relay delivery") + .custom_created_at(Timestamp::from_secs(42)) + .finalize(&keys) + .expect("build repeated event"); + + for _ in 0..PAGINATION_THRESHOLD { + // Repeat deliveries each consume a result slot even though the second and later + // process as Duplicate after the raw-delivery accounting point. + pagination.record_event(&event); + } + + assert_eq!(pagination.filters[0].event_count, PAGINATION_THRESHOLD); + assert_eq!(pagination.next_page_filters().len(), 1); + } + + #[test] + fn raw_pagination_counts_an_event_for_each_overlapping_filter() { + let keys = Keys::generate(); + let event = EventBuilder::new(Kind::TextNote, "overlap") + .custom_created_at(Timestamp::from_secs(42)) + .finalize(&keys) + .expect("build overlapping event"); + let mut pagination = PaginationState::new(vec![ + Filter::new().kind(Kind::TextNote), + Filter::new().author(keys.public_key()), + ]); + + pagination.record_event(&event); + + assert_eq!(pagination.filters[0].event_count, 1); + assert_eq!(pagination.filters[1].event_count, 1); + assert_eq!( + pagination.filters[0].min_created_at, + pagination.filters[1].min_created_at + ); + } + #[test] fn deferred_consolidation_runs_only_after_final_batch_completion() { let relay_url = "wss://relay.example";