From e713750aa77ebfef6cafd304474517732d8366ab Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 6 Aug 2026 11:10:59 +0000 Subject: [PATCH] fix(sync): adapt pagination to relay page size Static threshold 200 silently truncated historic filters on Ditto-shaped relays whose omitted-limit pages contain 100 events. NIP-11 max_limit cannot correct that because it describes explicit limits, while GRASP deliberately omits limit to preserve unbounded relay responses. Learn the largest raw page per relay connection session and combine it with NIP-11 default_limit, then paginate at max(90, floor(90% of the estimate)). Refetch NIP-11 on each successful connection, reset learning on disconnect, and ignore max_limit. Treat default_limit only as a hint: one suspiciously short page receives an inclusive-cursor verification request; any unseen event discards the hint session-wide. Preserve verification state across rate-limit deferral and across every filter in grouped REQs. Correctness assumes the audited relays apply result caps per filter rather than across the merged REQ, and that a cap of at least 90 covers the observed interoperability floor. Relays capped below 90 remain a documented residual risk. Explicit limits, configurable thresholds, aggregate-cap support, and the pre-existing transient-REQ concurrency defect are deliberately excluded. Validation: nix develop -c cargo test --lib (623 passed); nix develop -c cargo test --test sync adaptive_pagination -- --nocapture (3 passed: Ditto-shaped, honest hint, lying-high hint); nix develop -c cargo check --workspace --all-targets passed. The required standalone req-concurrency scenario failed twice with its documented proxy-rejection failure. Under this diff the same too many concurrent REQs signature repeated during phase-one pagination, a noisier form which is reported rather than masked. --- Cargo.toml | 2 +- docs/explanation/sync-scaling-constraints.md | 54 ++- src/sync/mod.rs | 329 ++++++++++++++++--- src/sync/relay_connection.rs | 54 +++ tests/common/mock_relay.rs | 94 +++++- tests/sync.rs | 1 + tests/sync/adaptive_pagination.rs | 166 ++++++++++ 7 files changed, 608 insertions(+), 92 deletions(-) create mode 100644 tests/sync/adaptive_pagination.rs diff --git a/Cargo.toml b/Cargo.toml index 48eed7f..575e676 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -69,6 +69,7 @@ async-trait = "0.1" # Temporary directories (used for GRASP-06 empty-repo synthesis) tempfile = "3" +reqwest = { version = "0.13", default-features = false, features = ["native-tls"] } # Git (for future use) # git-http-backend = "0.3" @@ -78,7 +79,6 @@ tempfile = "3" grasp-audit = { path = "grasp-audit", version = "0.2.0" } tempfile = "3" tokio = { version = "1.35", features = ["full", "test-util"] } -reqwest = { version = "0.13", default-features = false, features = ["native-tls"] } tokio-tungstenite = "0.28.0" [lib] diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 3a39e26..bb7d0a5 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -190,43 +190,37 @@ whether it applies independently to each filter or to the merged REQ, which is why the source audit above remains necessary. Relay-specific extensions can add fields, but clients cannot assume common names or semantics. -The client encodes this model in `PAGINATION_THRESHOLD` (200, -`src/sync/mod.rs`): after EOSE, a filter that delivered ≥ 200 counted -events is treated as possibly-truncated and fetched again with `until` set -to its oldest seen `created_at`; below 200 it is treated as exhausted. +The client encodes this model in per-connection `RelayPaginationSession` +state (`src/sync/mod.rs`). After EOSE it learns the largest raw page seen +from that relay and computes `max(90, floor(0.9 × estimated_cap))`, where +`estimated_cap` also includes an advertised NIP-11 `default_limit` while +that hint remains trusted. A filter meeting the adaptive threshold is +fetched again with `until` set to its oldest raw `created_at`. Consequences: -- Correct against any relay whose effective per-filter cap is ≥ 200. The - smallest audited defaults are 250 (haven/Badger) and 300 (rnostr); the - 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 - a 100-event default to filters that omit `limit`** — which ours - currently do — while accepting explicit limits up to its advertised - `max_limit` (1000). Against a Ditto relay, filters with more than 100 - results are silently truncated until historic filters carry an explicit - `limit`. Note that NIP-11 `max_limit` cannot reveal this: it advertises - the largest *accepted* request, not the default applied when `limit` is - omitted. Beyond Ditto, the remaining risk is a deliberately - restrictive, non-default strfry or rnostr configuration. -- The price of the floor is one redundant page for any filter whose result - count lands between 200 and the relay's actual cap. The threshold was - 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. -- Implementation in progress (design accepted 2026-08-06): keep omitting +- Every raw delivery matching a tracked filter 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. +- Ditto's 100-event omitted-limit default is now above the adaptive floor + and is learned from its first page even though it advertises only + `max_limit: 1000`. `max_limit` never raises the threshold because it + describes explicit limits, not the omitted-limit filters sent here. +- A relay capping a filter below 90 can still silently truncate history. + No such deployment was found in the audit. No audited implementation + enforces an aggregate cap across filters in one REQ. +- Larger learned pages raise the threshold and avoid redundant requests. + The 0.9 slack can still produce one final verification-shaped page when + a result count falls near the learned cap; this is the deliberate cost + of tolerating relay-side page shrinkage. +- Implemented design (accepted 2026-08-06): keep omitting `limit` — an explicit limit would cap the relays that serve unbounded 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:** + 2. **Adaptive per-relay threshold (implemented):** `estimated_cap = max(largest observed page, advertised default_limit if present)`; `threshold = max(90, floor(0.9 × estimated_cap))`. Observed pages @@ -237,7 +231,7 @@ Consequences: Learned state is per connection session and NIP-11 is refetched on reconnect, so an operator lowering their cap cannot strand a stale threshold. - 3. **NIP-11 fields:** `default_limit` ("maximum returned events if + 3. **NIP-11 fields (implemented):** `default_limit` ("maximum returned events if you send a filter without a limit") is the standard field for exactly this and is used as a hint when advertised — though rarely: neither nos.lol nor relay.ditto.pub advertises it (checked diff --git a/src/sync/mod.rs b/src/sync/mod.rs index e6d197b..9efc6ee 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -402,6 +402,12 @@ pub struct FilterPaginationState { pub min_created_at: Option, /// Original filter to reconstruct for next page pub original_filter: Filter, + /// IDs delivered on this page, retained only long enough to verify a NIP-11 hint. + page_event_ids: HashSet, + /// IDs from the page which caused a one-page NIP-11 hint verification. + verification_baseline: Option>, + /// Whether the verification page delivered an ID absent from its triggering page. + verification_productive: bool, } /// Pagination state for every OR filter carried by one subscription. @@ -419,6 +425,9 @@ impl PaginationState { event_count: 0, min_created_at: None, original_filter, + page_event_ids: HashSet::new(), + verification_baseline: None, + verification_productive: false, }) .collect(), } @@ -431,6 +440,14 @@ impl PaginationState { .match_event(event, MatchEventOptions::new()) { state.event_count += 1; + state.page_event_ids.insert(event.id); + if state + .verification_baseline + .as_ref() + .is_some_and(|baseline| !baseline.contains(&event.id)) + { + state.verification_productive = true; + } match state.min_created_at { None => state.min_created_at = Some(event.created_at), Some(min) if event.created_at < min => { @@ -442,21 +459,141 @@ impl PaginationState { } } - fn next_page_filters(self) -> Vec { + fn 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())) - }) - }) + .iter() + .map(|state| state.original_filter.clone()) .collect() } + + fn next_page(mut self, session: &mut RelayPaginationSession) -> Option { + let mut next = Vec::new(); + for mut state in self.filters.drain(..) { + session.observe_page(state.event_count); + let completed_verification = state.verification_baseline.is_some(); + + let continue_page = if completed_verification { + session.complete_hint_verification(state.verification_productive); + state.verification_productive && state.event_count >= session.pagination_threshold() + } else if state.event_count >= session.pagination_threshold() { + true + } else if state.event_count >= PAGINATION_THRESHOLD_FLOOR + && (session.begin_hint_verification() || session.hint_verification_in_progress()) + { + state.verification_baseline = Some(std::mem::take(&mut state.page_event_ids)); + true + } else { + false + }; + if completed_verification { + state.verification_baseline = None; + } + + let Some(min_created_at) = continue_page.then_some(state.min_created_at).flatten() + else { + continue; + }; + state.original_filter = state + .original_filter + .until(Timestamp::from(min_created_at.as_secs())); + state.event_count = 0; + state.min_created_at = None; + state.page_event_ids.clear(); + state.verification_productive = false; + next.push(state); + } + (!next.is_empty()).then_some(Self { filters: next }) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum PaginationHint { + Absent, + Unverified(usize), + Verifying(usize), + Verified(usize), + Discarded, +} + +#[derive(Debug, Clone)] +struct RelayPaginationSession { + largest_raw_page: usize, + hint: PaginationHint, +} + +impl Default for RelayPaginationSession { + fn default() -> Self { + Self::new(None) + } +} + +impl RelayPaginationSession { + fn new(advertised_default_limit: Option) -> Self { + Self { + largest_raw_page: 0, + hint: advertised_default_limit + .map(PaginationHint::Unverified) + .unwrap_or(PaginationHint::Absent), + } + } + + fn observe_page(&mut self, raw_count: usize) { + self.largest_raw_page = self.largest_raw_page.max(raw_count); + } + + fn active_hint(&self) -> Option { + match self.hint { + PaginationHint::Unverified(limit) | PaginationHint::Verified(limit) => Some(limit), + PaginationHint::Absent | PaginationHint::Verifying(_) | PaginationHint::Discarded => { + None + } + } + } + + fn estimated_cap(&self) -> usize { + self.largest_raw_page.max(self.active_hint().unwrap_or(0)) + } + + fn pagination_threshold(&self) -> usize { + // Ten percent slack absorbs relay-side shrinkage such as expired-event filtering. + // The floor remains below Ditto's observed 100-event omitted-limit page, the smallest + // live default found in the relay audit. These are interoperability constants, not + // operator policy, so deliberately do not enlarge the four-source config surface. + PAGINATION_THRESHOLD_FLOOR.max( + self.estimated_cap() + .saturating_mul(PAGINATION_THRESHOLD_PERCENT) + / 100, + ) + } + + fn begin_hint_verification(&mut self) -> bool { + match self.hint { + PaginationHint::Unverified(limit) => { + self.hint = PaginationHint::Verifying(limit); + true + } + _ => false, + } + } + + fn hint_verification_in_progress(&self) -> bool { + matches!(self.hint, PaginationHint::Verifying(_)) + } + + fn complete_hint_verification(&mut self, productive: bool) { + match (self.hint, productive) { + (PaginationHint::Verifying(_) | PaginationHint::Verified(_), true) => { + // Several filters can share the grouped verification page. Any one of them + // finding an unseen event disproves the relay-wide hint, even if an earlier + // exhausted filter provisionally marked it verified. + self.hint = PaginationHint::Discarded; + } + (PaginationHint::Verifying(limit), false) => { + self.hint = PaginationHint::Verified(limit); + } + _ => {} + } + } } /// A batch of items pending confirmation @@ -562,7 +699,9 @@ struct ConnectAttemptToken(u64); #[derive(Debug)] enum ConnectAttemptOutcome { - Connected, + Connected { + advertised_default_limit: Option, + }, Failed(String), } @@ -584,31 +723,27 @@ const CONSOLIDATION_THRESHOLD: usize = 70; /// exhaust network resources while keeping the sync actor responsive. const MAX_CONCURRENT_CONNECT_ATTEMPTS: usize = 8; -/// Per-filter threshold for historic REQ+EOSE pagination. -/// -/// After EOSE, a filter that delivered at least this many counted events is -/// treated as possibly truncated by the relay's per-filter result cap and is -/// fetched again with `until` set to its oldest seen `created_at`; below the -/// threshold it is treated as exhausted. +/// Adaptive per-relay threshold for historic REQ+EOSE pagination. /// /// NIP-01 guarantees none of this; the model is empirical. A source audit of /// nine relay implementations (2026-08-06, recorded with citations in /// docs/explanation/sync-scaling-constraints.md) found result limits are -/// always applied per filter, never in aggregate across a REQ, with the -/// smallest finite default caps at 250 (haven/Badger) and 300 (rnostr). +/// always applied per filter, never in aggregate across a REQ. Ditto is the +/// smallest live omitted-limit default found, at 100 events. /// /// 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; +/// derived from that same raw stream. Each connection session learns its +/// largest raw page and combines it with NIP-11 `default_limit` when present: +/// `threshold = max(90, floor(0.9 * max(observed, default_limit)))`. A hint is +/// verified once before it may stop pagination; a productive verification +/// page discards it for that session. `max_limit` is intentionally ignored +/// because these filters omit `limit`. State and NIP-11 data reset on every +/// reconnect. See "Per-query result limits and the pagination model" in +/// docs/explanation/sync-scaling-constraints.md. +const PAGINATION_THRESHOLD_FLOOR: usize = 90; +const PAGINATION_THRESHOLD_PERCENT: usize = 90; /// Conservative number of OR filters carried by one NIP-01 REQ. /// @@ -1035,6 +1170,8 @@ pub struct SyncManager { rejected_events_index: Arc, /// Active relay connections - keyed by relay URL connections: HashMap, + /// Adaptive pagination learning for each relay's current connection session. + pagination_sessions: HashMap, /// Event-directed relay targets rejected by the outbound target policy. /// /// Rejected URLs stay in `repo_sync_index` (they come from stored events), @@ -1147,6 +1284,7 @@ impl SyncManager { pending_sync_index: Arc::new(RwLock::new(HashMap::new())), rejected_events_index, connections: HashMap::new(), + pagination_sessions: HashMap::new(), rejected_relay_targets: HashSet::new(), dependency_refetch_attempts: Arc::new(std::sync::Mutex::new(HashMap::new())), missing_event_recovery: Arc::new(std::sync::Mutex::new( @@ -1284,10 +1422,16 @@ impl SyncManager { "EOSE processed for subscription" ); - // Check for pagination: if this subscription hit the threshold, fetch next page + // Check for pagination using this relay connection session's observed page sizes and + // verified NIP-11 default-limit hint. if let Some(pagination_state) = batch.pagination_state.remove(&sub_id) { - let next_filters = pagination_state.next_page_filters(); - if !next_filters.is_empty() { + let next_page = pagination_state.next_page( + self.pagination_sessions + .entry(relay_url.to_string()) + .or_default(), + ); + if let Some(next_page) = next_page { + let next_filters = next_page.filters(); let relay_url_for_pagination = relay_url.to_string(); let batch_id = batch.batch_id; tracing::info!( @@ -1321,7 +1465,7 @@ impl SyncManager { relay_url_for_pagination, batch_id, deferred_sub_id, - next_filters, + next_page, ); return; } @@ -1339,10 +1483,7 @@ impl SyncManager { { batch.outstanding_subs.insert(new_sub_id.clone()); next_page_started = true; - batch.pagination_state.insert( - new_sub_id.clone(), - PaginationState::new(next_filters), - ); + batch.pagination_state.insert(new_sub_id.clone(), next_page); tracing::info!( relay = %relay_url_for_pagination, new_sub_id = %new_sub_id, @@ -1763,9 +1904,10 @@ impl SyncManager { relay_url: String, batch_id: u64, deferred_sub_id: SubscriptionId, - next_filters: Vec, + next_page: PaginationState, ) { tokio::spawn(async move { + let next_filters = next_page.filters(); tracing::info!( relay = %relay_url, batch_id, @@ -1805,9 +1947,7 @@ impl SyncManager { 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::new(next_filters)); + batch.pagination_state.insert(new_sub_id.clone(), next_page); tracing::info!( relay = %relay_url, new_sub_id = %new_sub_id, @@ -3406,7 +3546,9 @@ impl SyncManager { }; let outcome = tokio::select! { result = connection.connect(timeout) => match result { - Ok(()) => ConnectAttemptOutcome::Connected, + Ok(()) => ConnectAttemptOutcome::Connected { + advertised_default_limit: connection.fetch_default_limit().await, + }, Err(error) => ConnectAttemptOutcome::Failed(error), }, _ = shutdown_rx.recv() => { @@ -3457,7 +3599,16 @@ impl SyncManager { } match result.outcome { - ConnectAttemptOutcome::Connected => { + ConnectAttemptOutcome::Connected { + advertised_default_limit, + } => { + // A session begins only after a successful WebSocket connection. Replacing this + // entry on every connection result resets learned caps and re-applies the NIP-11 + // hint fetched for that exact session. + self.pagination_sessions.insert( + result.relay_url.clone(), + RelayPaginationSession::new(advertised_default_limit), + ); self.health_tracker.record_success(&result.relay_url); if let Some(ref metrics) = self.metrics { metrics.record_connection_attempt(&result.relay_url, true); @@ -4080,6 +4231,9 @@ impl SyncManager { /// - Unexpected disconnects: Updates state to Disconnected, keeps RelayConnection for reconnect /// - Intentional disconnects: Completes cleanup of Disconnecting relays (removes from indices) async fn handle_disconnect(&mut self, relay_url: &str) { + // Learned page sizes and NIP-11 hints belong to the ended WebSocket session. + self.pagination_sessions.remove(relay_url); + // Check if this was an intentional disconnect (Disconnecting status) let was_intentional = { let index = self.relay_sync_index.read().await; @@ -5805,7 +5959,7 @@ mod tests { let mut pagination = PaginationState::new(vec![metadata_filter.clone(), note_filter.clone()]); - for created_at in 1..=PAGINATION_THRESHOLD { + for created_at in 1..=100 { let event = EventBuilder::new(Kind::Metadata, created_at.to_string()) .custom_created_at(Timestamp::from_secs(created_at as u64)) .finalize(&keys) @@ -5818,7 +5972,10 @@ mod tests { .expect("build text note"); pagination.record_event(¬e); - let next_filters = pagination.next_page_filters(); + let next_filters = pagination + .next_page(&mut RelayPaginationSession::default()) + .expect("full observed page should paginate") + .filters(); assert_eq!(next_filters.len(), 1); assert_eq!(next_filters[0].until, Some(Timestamp::from_secs(1))); assert!( @@ -5887,14 +6044,19 @@ mod tests { .finalize(&keys) .expect("build repeated event"); - for _ in 0..PAGINATION_THRESHOLD { + for _ in 0..PAGINATION_THRESHOLD_FLOOR { // 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); + assert_eq!( + pagination.filters[0].event_count, + PAGINATION_THRESHOLD_FLOOR + ); + assert!(pagination + .next_page(&mut RelayPaginationSession::default()) + .is_some()); } #[test] @@ -5919,6 +6081,75 @@ mod tests { ); } + #[test] + fn adaptive_threshold_selects_from_observation_and_hint() { + let mut observed_only = RelayPaginationSession::new(None); + observed_only.observe_page(100); + assert_eq!(observed_only.pagination_threshold(), 90); + + let hint_only = RelayPaginationSession::new(Some(500)); + assert_eq!(hint_only.pagination_threshold(), 450); + + let mut hint_vs_observed = RelayPaginationSession::new(Some(500)); + hint_vs_observed.observe_page(600); + assert_eq!(hint_vs_observed.pagination_threshold(), 540); + + let mut sub_floor = RelayPaginationSession::new(Some(50)); + sub_floor.observe_page(80); + assert_eq!(sub_floor.pagination_threshold(), 90); + } + + #[test] + fn productive_verification_page_discards_the_hint() { + let keys = Keys::generate(); + let mut first_page = PaginationState::new(vec![Filter::new().kind(Kind::TextNote)]); + for created_at in 100..200 { + let event = EventBuilder::new(Kind::TextNote, created_at.to_string()) + .custom_created_at(Timestamp::from_secs(created_at)) + .finalize(&keys) + .expect("build first-page event"); + first_page.record_event(&event); + } + let mut session = RelayPaginationSession::new(Some(1000)); + let mut verification = first_page + .next_page(&mut session) + .expect("a suspiciously short hinted page needs verification"); + assert_eq!(session.hint, PaginationHint::Verifying(1000)); + + let unseen = EventBuilder::new(Kind::TextNote, "older unseen event") + .custom_created_at(Timestamp::from_secs(99)) + .finalize(&keys) + .expect("build productive verification event"); + for _ in 0..100 { + verification.record_event(&unseen); + } + + assert!(verification.next_page(&mut session).is_some()); + assert_eq!(session.hint, PaginationHint::Discarded); + assert_eq!(session.pagination_threshold(), 90); + } + + #[test] + fn empty_verification_page_confirms_the_hint_without_another_page() { + let keys = Keys::generate(); + let mut first_page = PaginationState::new(vec![Filter::new().kind(Kind::TextNote)]); + for created_at in 100..200 { + let event = EventBuilder::new(Kind::TextNote, created_at.to_string()) + .custom_created_at(Timestamp::from_secs(created_at)) + .finalize(&keys) + .expect("build first-page event"); + first_page.record_event(&event); + } + let mut session = RelayPaginationSession::new(Some(1000)); + let verification = first_page + .next_page(&mut session) + .expect("a suspiciously short hinted page needs verification"); + + assert!(verification.next_page(&mut session).is_none()); + assert_eq!(session.hint, PaginationHint::Verified(1000)); + assert_eq!(session.pagination_threshold(), 900); + } + #[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 00ded76..5a1a76e 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -30,6 +30,9 @@ use crate::outbound::{OutboundTargetKind, OutboundTargetPolicy, RelayTargetSourc /// exchange alive indefinitely by continuing to send reconciliation messages. const NEGENTROPY_DIFF_TIMEOUT: Duration = Duration::from_secs(15); +/// NIP-11 is advisory and must not hold up a connected relay indefinitely. +const NIP11_FETCH_TIMEOUT: Duration = Duration::from_secs(3); + /// Cooldown schedule for transient negentropy failures. /// /// Indexed by the number of consecutive failed attempts (capped at the last @@ -411,6 +414,57 @@ impl RelayConnection { Ok(()) } + /// Fetch the omitted-limit page-size hint advertised by this connection's relay. + /// + /// This is deliberately fetched for every WebSocket session rather than cached on the + /// long-lived `RelayConnection`: operators can change their cap between reconnects. Only + /// NIP-11 `default_limit` is returned. `max_limit` describes explicit `limit` values and is + /// not applicable to the historic filters we intentionally send without one. + pub async fn fetch_default_limit(&self) -> Option { + let mut document_url = reqwest::Url::parse(&self.url).ok()?; + let http_scheme = match document_url.scheme() { + "ws" => "http", + "wss" => "https", + _ => return None, + }; + document_url.set_scheme(http_scheme).ok()?; + + let fetch = async { + let response = reqwest::Client::new() + .get(document_url) + .header(reqwest::header::ACCEPT, "application/nostr+json") + .send() + .await?; + if !response.status().is_success() { + tracing::debug!( + relay = %self.url, + status = %response.status(), + "NIP-11 fetch returned a non-success status" + ); + return Ok(None); + } + response.text().await.map(Some) + }; + let body = match tokio::time::timeout(NIP11_FETCH_TIMEOUT, fetch).await { + Ok(Ok(Some(body))) => body, + Ok(Ok(None)) => return None, + Ok(Err(error)) => { + tracing::debug!(relay = %self.url, error = %error, "NIP-11 fetch failed"); + return None; + } + Err(_) => { + tracing::debug!(relay = %self.url, "NIP-11 fetch timed out"); + return None; + } + }; + let document = nostr::nips::nip11::RelayInformationDocument::from_json(body).ok()?; + document + .limitation + .and_then(|limitation| limitation.default_limit) + .and_then(|limit| usize::try_from(limit).ok()) + .filter(|limit| *limit > 0) + } + /// Run the event loop, sending events through the provided channel /// /// This method blocks and processes notifications from the relay using diff --git a/tests/common/mock_relay.rs b/tests/common/mock_relay.rs index 7e3caff..5f63e49 100644 --- a/tests/common/mock_relay.rs +++ b/tests/common/mock_relay.rs @@ -38,7 +38,7 @@ use std::sync::Arc; use http_body_util::Full; use hyper::body::Bytes; -use hyper::header::{CONNECTION, SEC_WEBSOCKET_ACCEPT, SEC_WEBSOCKET_KEY, UPGRADE}; +use hyper::header::{ACCEPT, CONNECTION, SEC_WEBSOCKET_ACCEPT, SEC_WEBSOCKET_KEY, UPGRADE}; use hyper::server::conn::http1; use hyper::service::service_fn; use hyper::{Request, Response, StatusCode}; @@ -66,25 +66,55 @@ pub struct MockRelay { relay: LocalRelay, } +#[derive(Clone, Copy)] +struct PaginationConfig { + page_size: usize, + advertised_default_limit: Option, + advertised_max_limit: Option, +} + impl MockRelay { /// Start a mock relay on a random free port. /// /// 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 + Self::start_with_rate_limit(RateLimit::default(), None).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() - }) + Self::start_with_rate_limit( + RateLimit { + max_reqs, + ..RateLimit::default() + }, + None, + ) .await } - async fn start_with_rate_limit(rate_limit: RateLimit) -> Self { + /// Start a relay whose omitted-limit pages and NIP-11 hints can be varied independently. + pub async fn start_with_pagination( + page_size: usize, + advertised_default_limit: Option, + advertised_max_limit: Option, + ) -> Self { + Self::start_with_rate_limit( + RateLimit::default(), + Some(PaginationConfig { + page_size, + advertised_default_limit, + advertised_max_limit, + }), + ) + .await + } + + async fn start_with_rate_limit( + rate_limit: RateLimit, + pagination: Option, + ) -> 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"); @@ -100,7 +130,7 @@ impl MockRelay { let listener = TcpListener::from_std(std_listener).expect("Failed to convert to tokio listener"); - Self::start_with_listener(listener, port, rate_limit).await + Self::start_with_listener(listener, port, rate_limit, pagination).await } /// Start a mock relay on a specific port. @@ -109,13 +139,25 @@ impl MockRelay { let listener = TcpListener::bind(addr) .await .expect("Failed to bind to address"); - Self::start_with_listener(listener, port, RateLimit::default()).await + Self::start_with_listener(listener, port, RateLimit::default(), None).await } /// Internal method to start the relay with an existing listener. - async fn start_with_listener(listener: TcpListener, port: u16, rate_limit: RateLimit) -> Self { + async fn start_with_listener( + listener: TcpListener, + port: u16, + rate_limit: RateLimit, + pagination: Option, + ) -> Self { // Create a simple relay with no write policy (accepts all events) - let relay = LocalRelayBuilder::default().rate_limit(rate_limit).build(); + let mut builder = LocalRelayBuilder::default().rate_limit(rate_limit); + if let Some(config) = pagination { + builder = builder + .default_filter_limit(config.page_size) + .max_filter_limit(config.page_size.max(1000)) + .max_query_results(config.page_size.max(1000)); + } + let relay = builder.build(); // Create shutdown channel let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); @@ -130,12 +172,15 @@ impl MockRelay { match accept_result { Ok((stream, remote_addr)) => { let relay = server_relay.clone(); + let pagination = pagination; let io = TokioIo::new(stream); tokio::spawn(async move { let service = service_fn(move |req| { let relay = relay.clone(); - async move { handle_request(req, relay, remote_addr).await } + async move { + handle_request(req, relay, remote_addr, pagination).await + } }); if let Err(e) = http1::Builder::new() @@ -215,6 +260,7 @@ async fn handle_request( req: Request, relay: LocalRelay, addr: SocketAddr, + pagination: Option, ) -> Result>, hyper::Error> { // Check for WebSocket upgrade request let is_websocket = req @@ -257,6 +303,30 @@ async fn handle_request( } } + if req + .headers() + .get(ACCEPT) + .and_then(|value| value.to_str().ok()) + .is_some_and(|value| value.contains("application/nostr+json")) + { + let limitation = pagination.map(|config| { + serde_json::json!({ + "default_limit": config.advertised_default_limit, + "max_limit": config.advertised_max_limit, + }) + }); + let document = serde_json::json!({ + "name": "pagination test relay", + "supported_nips": [1, 11], + "limitation": limitation, + }); + return Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", "application/nostr+json") + .body(Full::new(Bytes::from(document.to_string()))) + .unwrap()); + } + // Non-WebSocket request - return simple response Ok(Response::builder() .status(StatusCode::OK) diff --git a/tests/sync.rs b/tests/sync.rs index 6ddabf8..6e3b20b 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -31,6 +31,7 @@ mod common; // Include sync test submodules (located in tests/sync/) mod sync { + pub mod adaptive_pagination; pub mod catchup; pub mod discovery; pub mod historic_recovery; diff --git a/tests/sync/adaptive_pagination.rs b/tests/sync/adaptive_pagination.rs new file mode 100644 index 0000000..b614e00 --- /dev/null +++ b/tests/sync/adaptive_pagination.rs @@ -0,0 +1,166 @@ +//! Adaptive historic-pagination integration scenarios. +//! +//! These drive the real REQ+EOSE sync path against a configurable rust-nostr +//! LocalRelay. Its omitted-limit page size models Ditto independently from +//! the NIP-11 values it advertises, which covers missing, honest, and wrong-high +//! `default_limit` documents without adding an explicit `limit` to our filters. + +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::purgatory_helpers::{ + create_state_event, create_test_repo_with_commit, push_to_relay, CommitVariant, +}; +use crate::common::sync_helpers::{repo_coord, wait_for_event_on_relay, TestClient}; +use crate::common::{port, MockRelay, TestRelay}; + +const SEED_BATCH: usize = 50; + +async fn seed_issues( + source: &MockRelay, + repo_keys: &Keys, + coordinate: &str, + count: usize, +) -> Vec { + let base_created_at = Timestamp::now().as_secs() - count as u64 - 10; + let mut issues = Vec::with_capacity(count); + + for batch_start in (0..count).step_by(SEED_BATCH) { + let client = TestClient::new(source.url(), Keys::generate()) + .await + .expect("connect seeding client"); + for index in batch_start..(batch_start + SEED_BATCH).min(count) { + let issue = EventBuilder::new(Kind::GitIssue, format!("Historic issue {index}")) + .tags(vec![Tag::custom("a", vec![coordinate.to_string()])]) + .custom_created_at(Timestamp::from_secs(base_created_at + index as u64)) + .finalize(repo_keys) + .expect("build historic issue"); + client + .send_event(&issue) + .await + .expect("seed historic issue"); + issues.push(issue); + } + client.disconnect().await; + } + + issues +} + +async fn run_pagination_scenario( + page_size: usize, + advertised_default_limit: Option, + advertised_max_limit: Option, + event_count: usize, +) -> (TestRelay, MockRelay) { + let reservation = port::reserve_port(); + let syncing_domain = format!("127.0.0.1:{}", reservation.port()); + let source = + MockRelay::start_with_pagination(page_size, advertised_default_limit, advertised_max_limit) + .await; + let repo_keys = Keys::generate(); + let identifier = "adaptive-pagination"; + let coordinate = repo_coord(&repo_keys, identifier); + let issues = seed_issues(&source, &repo_keys, &coordinate, event_count).await; + let oldest = issues.first().expect("at least one event").id; + + let syncing = TestRelay::start_on_reservation_with_options( + reservation, + Some(source.url().to_string()), + true, + ) + .await; + + // Admit one real repository locally so its self-subscriber installs the + // Layer-2 historic filter that matches the already-seeded issues. + let git_temp_dir = tempfile::tempdir().expect("create pagination git repo"); + let commit_hash = create_test_repo_with_commit(git_temp_dir.path(), CommitVariant::StateTest) + .expect("create pagination git history"); + let npub = repo_keys.public_key().to_bech32().expect("npub"); + let clone_urls = vec![format!("http://{syncing_domain}/{npub}/{identifier}.git")]; + let relay_urls = vec![source.url().to_string(), syncing.url().to_string()]; + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "pagination repository") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom("clone", clone_urls.clone()), + Tag::custom("relays", relay_urls.clone()), + ]) + .finalize(&repo_keys) + .expect("build repository announcement"); + let state = create_state_event( + &repo_keys, + identifier, + &[("main", &commit_hash)], + &[], + &clone_urls.iter().map(String::as_str).collect::>(), + &relay_urls.iter().map(String::as_str).collect::>(), + ) + .expect("build repository state"); + let client = TestClient::new(syncing.url(), repo_keys.clone()) + .await + .expect("connect announcement client"); + client + .send_event(&announcement) + .await + .expect("submit repository announcement"); + client + .send_event(&state) + .await + .expect("submit repository state"); + client.disconnect().await; + push_to_relay(git_temp_dir.path(), &syncing.domain(), &npub, identifier) + .expect("push repository data to syncing relay"); + + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(oldest), + Duration::from_secs(30), + ) + .await, + "the oldest of {event_count} issues must survive pagination" + ); + + (syncing, source) +} + +#[tokio::test] +async fn ditto_shaped_omitted_limit_pages_reach_the_oldest_event() { + // Ditto advertises only the explicit-request maximum (1000), while an + // omitted-limit filter receives 100-event pages. + let (syncing, source) = run_pagination_scenario(100, None, Some(1000), 320).await; + syncing.stop().await; + source.stop().await; +} + +#[tokio::test] +async fn honest_default_limit_stops_after_its_verification_page() { + // The relay honestly advertises a 500-event default but has only 100 + // matches. The short first page gets exactly one verification request; + // its inclusive cursor repeats only the boundary event, then stops. + let (syncing, source) = run_pagination_scenario(500, Some(500), Some(1000), 100).await; + let log_path = format!( + "/tmp/relay-{}.log", + syncing.domain().split(':').next_back().unwrap() + ); + let log = std::fs::read_to_string(&log_path).expect("read syncing relay log"); + assert_eq!( + log.matches("Grouped subscription hit pagination threshold") + .count(), + 1, + "an honest hint should require only its single verification page" + ); + syncing.stop().await; + source.stop().await; +} + +#[tokio::test] +async fn wrong_high_default_limit_is_discarded_after_productive_verification() { + // The document claims 1000, but omitted-limit pages contain only 100. + // The first verification page is productive, so learned-only threshold + // selection must rescue the remaining history. + let (syncing, source) = run_pagination_scenario(100, Some(1000), Some(1000), 320).await; + syncing.stop().await; + source.stop().await; +}