From ef5be4a5c1878b71f7848cd48ee1573d870ab6ee Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 09:33:02 +0000 Subject: [PATCH 1/7] refactor(sync): name ordered reference coverage tiers Motivation: live coverage was expressed only as undifferentiated a/A/q and e/E/q vectors, which made a capacity-aware priority cutoff impossible to review independently from subscription lifecycle changes. Approach: introduce six ordered coverage tiers and deterministic priority-labelled builders for repository and event references. Core and descendant references use distinct policies; the existing builders and all call sites remain unchanged in this commit. Correctness: values retain deterministic byte chunking and each tag form remains represented exactly once. Tests pin the intended tier order for both ordinary roots and descendant frontiers. Excluded scope: this commit does not change wire filters, admission, rotation, consolidation, or historic packing. Validation: nix develop -c cargo test --lib sync::filters::tests (25 passed). --- src/sync/filters.rs | 174 ++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 174 insertions(+) diff --git a/src/sync/filters.rs b/src/sync/filters.rs index d6de410..a1ebb02 100644 --- a/src/sync/filters.rs +++ b/src/sync/filters.rs @@ -13,6 +13,125 @@ use std::collections::HashSet; use nostr_sdk::prelude::*; +/// Ordered live-coverage priorities. Lower-priority tiers may be served by +/// paced REQ+EOSE rotation when the connection cannot keep every filter live. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum CoverageTier { + EssentialCore, + RootUppercase, + CoreCompatibility, + DescendantCanonical, + DescendantQuote, + HistoricOnly, +} + +#[derive(Debug, Clone)] +pub struct TieredFilter { + pub tier: CoverageTier, + pub filter: Filter, +} + +fn tagged_filters>( + values: &[T], + tag: SingleLetterTag, + since: Option, +) -> Vec { + chunk_values_by_bytes(values) + .into_iter() + .map(|chunk| { + let mut filter = Filter::new(); + for value in chunk { + filter = filter.custom_tag(tag, value.as_ref()); + } + match since { + Some(timestamp) => filter.since(timestamp), + None => filter, + } + }) + .collect() +} + +fn tiered_tagged_filters>( + values: &[T], + tiers: &[(CoverageTier, SingleLetterTag)], + since: Option, +) -> Vec { + tiers + .iter() + .flat_map(|(tier, tag)| { + tagged_filters(values, *tag, since) + .into_iter() + .map(|filter| TieredFilter { + tier: *tier, + filter, + }) + }) + .collect() +} + +/// Priority-labelled filters for repository coordinates. +pub fn tiered_repo_event_filters( + repos: &HashSet, + descendant: bool, + since: Option, +) -> Vec { + let mut values: Vec<_> = repos.iter().collect(); + values.sort_unstable(); + let tiers = if descendant { + [ + ( + CoverageTier::DescendantCanonical, + SingleLetterTag::LOWERCASE_A, + ), + (CoverageTier::DescendantQuote, SingleLetterTag::LOWERCASE_Q), + (CoverageTier::HistoricOnly, SingleLetterTag::UPPERCASE_A), + ] + } else { + [ + (CoverageTier::EssentialCore, SingleLetterTag::LOWERCASE_A), + ( + CoverageTier::CoreCompatibility, + SingleLetterTag::UPPERCASE_A, + ), + ( + CoverageTier::CoreCompatibility, + SingleLetterTag::LOWERCASE_Q, + ), + ] + }; + tiered_tagged_filters(&values, &tiers, since) +} + +/// Priority-labelled filters for event references. +pub fn tiered_root_event_filters( + root_events: &HashSet, + descendant: bool, + since: Option, +) -> Vec { + let mut values: Vec = root_events.iter().map(EventId::to_hex).collect(); + values.sort_unstable(); + let tiers = if descendant { + [ + ( + CoverageTier::DescendantCanonical, + SingleLetterTag::LOWERCASE_E, + ), + (CoverageTier::DescendantQuote, SingleLetterTag::LOWERCASE_Q), + (CoverageTier::HistoricOnly, SingleLetterTag::UPPERCASE_E), + ] + } else { + [ + (CoverageTier::EssentialCore, SingleLetterTag::LOWERCASE_E), + (CoverageTier::RootUppercase, SingleLetterTag::UPPERCASE_E), + ( + CoverageTier::CoreCompatibility, + SingleLetterTag::LOWERCASE_Q, + ), + ] + }; + tiered_tagged_filters(&values, &tiers, since) +} + /// Serialized byte budget of tag values per filter chunk. /// /// Relay filter limits are byte caps, not item counts: strfry bounds the @@ -364,6 +483,61 @@ mod tests { assert_eq!(filters.len(), 3); } + #[test] + fn core_reference_filters_have_stable_priority_tiers() { + let repos = HashSet::from(["30617:pubkey:repo".to_string()]); + let roots = HashSet::from([EventId::from_byte_array([0; 32])]); + + let repo_tiers: Vec<_> = tiered_repo_event_filters(&repos, false, None) + .into_iter() + .map(|entry| entry.tier) + .collect(); + assert_eq!( + repo_tiers, + [ + CoverageTier::EssentialCore, + CoverageTier::CoreCompatibility, + CoverageTier::CoreCompatibility, + ] + ); + + let root_tiers: Vec<_> = tiered_root_event_filters(&roots, false, None) + .into_iter() + .map(|entry| entry.tier) + .collect(); + assert_eq!( + root_tiers, + [ + CoverageTier::EssentialCore, + CoverageTier::RootUppercase, + CoverageTier::CoreCompatibility, + ] + ); + } + + #[test] + fn descendant_reference_filters_put_uppercase_tags_in_history_only() { + let repos = HashSet::from(["1621:pubkey:patch".to_string()]); + let roots = HashSet::from([EventId::from_byte_array([0; 32])]); + + for tiers in [ + tiered_repo_event_filters(&repos, true, None), + tiered_root_event_filters(&roots, true, None), + ] { + assert_eq!( + tiers + .into_iter() + .map(|entry| entry.tier) + .collect::>(), + [ + CoverageTier::DescendantCanonical, + CoverageTier::DescendantQuote, + CoverageTier::HistoricOnly, + ] + ); + } + } + #[test] fn test_repo_filters_batching() { // Each ref serializes to exactly 1024 bytes (1021 + quoting overhead), From 5d3e419be31f118b99e0bc44720736a2200e0b09 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 09:40:03 +0000 Subject: [PATCH 2/7] feat(sync): demote reference coverage by priority tier Motivation: the two-level core/descendant policy forced every core tag variant to stay live while treating every descendant variant alike. As repository scale grows, that either consumes the transient recovery slot or demotes useful canonical coverage together with low-value compatibility forms. Approach: keep announcements, state, repository a and root e filters essential. Reconcile root E, core q/A, descendant e/a and descendant q as an ordered complete-tier prefix against the existing NIP-11-aware ledger, filter-count grouping and byte budget. Pack adjacent tier filters together, reserve one historic slot, and feed all filters below the cutoff plus descendant E/A through the existing five-second cursor-overlapped REQ+EOSE queue. Correctness: historic sync still receives the full original filter set; only persistent admission changes. A tier is never partially admitted, fallback cursor state survives unchanged filters, CLOSED auxiliary subscriptions immediately return their demoted filters to rotation, and all consumers continue to share the existing ledger. Excluded scope: no connection sharding, new scheduler, request-class priority, configuration, or historic repository packing is included. Validation: nix develop -c cargo test --lib (685 passed); focused tier-cutoff, filter-tier and descendant-rotation tests passed. --- docs/explanation/sync-scaling-constraints.md | 39 +- src/sync/filters.rs | 40 +++ src/sync/mod.rs | 357 ++++++++++++++----- tests/sync/descendant_sync.rs | 2 +- 4 files changed, 322 insertions(+), 116 deletions(-) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 7d8b66f..3d89b77 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -304,8 +304,9 @@ Derived from the tightest commonly observed values; all sizing below assumes: Each relay connection owns one implemented budget ledger of B subscription slots. Four consumers share it, in priority order: -1. **Core live subscriptions** (persistent, `limit: 0`) — the product; sized - first. +1. **Essential live subscriptions** (persistent, `limit: 0`) — announcements, + repository states, canonical repository `a` references and canonical root + `e` references are never demoted. 2. **Reserved margin** (2 slots) — control-plane safety capacity kept beyond the live set (which includes Layer-1) for ad-hoc operations and recovery. 3. **Historic sync and dependency recovery** (transient) — at least one usable @@ -313,21 +314,24 @@ slots. Four consumers share it, in priority order: REQ+EOSE pages/fallbacks/retries, and exact-ID purgatory polls draw from the remainder. NEG retains its four-round class cap and transient REQ its five-request class cap, but neither can exceed the shared residual. -4. **Descendant live coverage** (auxiliary persistent) — admitted only when - its complete separately grouped set fits after core coverage while still - preserving the margin and a transient slot. Otherwise each constrained - relay advances one cursor-overlapped REQ+EOSE filter per five-second tick - through the same transient queue. Direct thread members contribute their - event IDs; replaceable and addressable members also contribute their NIP-01 - coordinates. The resulting `e`/`E`/`q` and `a`/`A`/coordinate-`q` filters - cover one descendant generation without recursively expanding the frontier. +4. **Priority-tiered reference coverage** — remaining filters are considered + in this order: root `E`; core compatibility `q`/`A`; descendant canonical + `e`/`a`; descendant `q`. A complete tier remains persistent only when it + fits after essential coverage while preserving the margin and a transient + slot. Lower tiers advance one relay-compatible, cursor-overlapped REQ+EOSE + filter group per five-second tick through the same transient queue. Descendant uppercase + `E`/`A` references are historic-only. Direct thread members contribute + event IDs and replaceable/addressable coordinates, covering one descendant + generation without recursively expanding the frontier. NIP-11 `max_subscriptions` sets B for each new connection session; when it is absent B falls back to 20. Advertised values below that floor are honoured -(notably nostream's default 10). Two slots remain reserved. Core live filter -groups are packed first and admitted atomically against the advertised -subscription-count budget. Incremental five-second batches preserve full core -REQs and separately owned descendant REQs. The planner uses the same +(notably nostream's default 10). Two slots remain reserved. Essential live +filter groups are packed first and admitted atomically against the advertised +subscription-count budget. Incremental five-second batches preserve full +essential REQs and separately owned tiered reference REQs. Tier filters are +packed across boundaries, while admission stops at a complete-tier boundary. +The planner uses the same filter-count and serialized-byte grouping rules as wire submission. It repacks the complete mutable core tail with the new filters when that releases at least one slot; otherwise it retires only the smallest useful subset which reduces @@ -340,9 +344,10 @@ is opened and historic recovery remains available. Capacity pressure is the backstop which may schedule a complete regroup after outstanding historic batches drain; an ordinary tail update never rebuilds stable full groups. -Reconnect and exceptional full restoration still build Layer 1, 2, and 3 -coverage together and reserve the whole grouped set before opening it. Both -full replacement and tail replacement remember the exact previous grouping: +Reconnect and exceptional full restoration rebuild essential coverage first; +the five-second reconciler then admits the largest complete prefix of reference +tiers which fits the refreshed session budget. Both full replacement and tail +replacement remember the exact previous grouping: a failure while opening a replacement closes every newly opened group and restores the retired groups. A partial CLOSE failure likewise reopens any tail groups which were already closed before reporting the failure. Multi-connection diff --git a/src/sync/filters.rs b/src/sync/filters.rs index a1ebb02..84a4ff7 100644 --- a/src/sync/filters.rs +++ b/src/sync/filters.rs @@ -31,6 +31,46 @@ pub struct TieredFilter { pub filter: Filter, } +/// Build the live coverage that is never demoted under connection pressure. +pub fn build_essential_live_filters( + full_repos: &HashSet, + state_only_repos: &HashSet, + root_events: &HashSet, + since: Option, +) -> Vec { + let all_repos: HashSet = full_repos.union(state_only_repos).cloned().collect(); + let mut filters = state_event_filters_for_our_repos(&all_repos, since); + if !full_repos.is_empty() { + filters.extend( + tiered_repo_event_filters(full_repos, false, since) + .into_iter() + .filter(|entry| entry.tier == CoverageTier::EssentialCore) + .map(|entry| entry.filter), + ); + } + filters.extend( + tiered_root_event_filters(root_events, false, since) + .into_iter() + .filter(|entry| entry.tier == CoverageTier::EssentialCore) + .map(|entry| entry.filter), + ); + filters +} + +/// Build the priority-ordered core reference coverage that can move between a +/// persistent subscription and paced REQ+EOSE rotation. +pub fn tiered_auxiliary_core_filters( + full_repos: &HashSet, + root_events: &HashSet, + since: Option, +) -> Vec { + let mut filters = tiered_root_event_filters(root_events, false, since); + filters.extend(tiered_repo_event_filters(full_repos, false, since)); + filters.retain(|entry| entry.tier != CoverageTier::EssentialCore); + filters.sort_by_key(|entry| entry.tier); + filters +} + fn tagged_filters>( values: &[T], tag: SingleLetterTag, diff --git a/src/sync/mod.rs b/src/sync/mod.rs index c449236..07f683b 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -707,7 +707,7 @@ struct DescendantFrontier { #[derive(Debug)] struct DescendantSyncRotation { - frontier: DescendantFrontier, + coverage_fingerprint: Vec, filters: Vec, next_filter: usize, in_flight: Option, @@ -715,7 +715,7 @@ struct DescendantSyncRotation { #[derive(Debug)] struct DescendantFilterCursor { - filter: Filter, + filters: Vec, last_successful_until: Option, } @@ -728,14 +728,15 @@ struct DescendantFilterInFlight { #[derive(Debug)] struct DescendantLiveCoverage { - frontier: DescendantFrontier, + coverage_fingerprint: Vec, + fallback_filters: Vec, subscription_ids: Vec, } impl Default for DescendantSyncRotation { fn default() -> Self { Self { - frontier: DescendantFrontier::default(), + coverage_fingerprint: Vec::new(), filters: Vec::new(), next_filter: 0, in_flight: None, @@ -743,40 +744,66 @@ impl Default for DescendantSyncRotation { } } +fn filter_group_fingerprint(filters: &[Filter]) -> String { + filters.iter().map(Filter::as_json).collect::>().join("\n") +} + +fn rotation_fingerprint(filters: &[Filter], max_filters_per_req: usize) -> Vec { + group_filters_for_req_with_max(filters, max_filters_per_req) + .iter() + .map(|group| filter_group_fingerprint(group)) + .collect() +} + impl DescendantSyncRotation { - fn refresh(&mut self, frontier: DescendantFrontier) { + fn refresh(&mut self, filters: Vec, max_filters_per_req: usize) { let previous: HashMap> = self .filters .drain(..) - .map(|cursor| (cursor.filter.as_json(), cursor.last_successful_until)) + .map(|cursor| (filter_group_fingerprint(&cursor.filters), cursor.last_successful_until)) .collect(); - self.filters = descendant_frontier_filters(&frontier, None) + self.filters = group_filters_for_req_with_max(&filters, max_filters_per_req) .into_iter() - .map(|filter| DescendantFilterCursor { - last_successful_until: previous.get(&filter.as_json()).copied().flatten(), - filter, + .map(|filters| DescendantFilterCursor { + last_successful_until: previous + .get(&filter_group_fingerprint(&filters)) + .copied() + .flatten(), + filters, }) .collect(); - self.frontier = frontier; + self.coverage_fingerprint = self + .filters + .iter() + .map(|cursor| filter_group_fingerprint(&cursor.filters)) + .collect(); self.next_filter = 0; self.in_flight = None; } - fn next_request(&self, now: Timestamp) -> Option<(usize, Filter, Timestamp)> { + fn next_request(&self, now: Timestamp) -> Option<(usize, Vec, Timestamp)> { if self.in_flight.is_some() || self.filters.is_empty() { return None; } let filter_index = self.next_filter % self.filters.len(); let cursor = &self.filters[filter_index]; - let mut filter = cursor.filter.clone().until(now); - if let Some(last_until) = cursor.last_successful_until { - filter = filter.since(Timestamp::from( - last_until - .as_secs() - .saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS), - )); - } - Some((filter_index, filter, now)) + let filters = cursor + .filters + .iter() + .cloned() + .map(|filter| { + let filter = filter.until(now); + match cursor.last_successful_until { + Some(last_until) => filter.since(Timestamp::from( + last_until + .as_secs() + .saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS), + )), + None => filter, + } + }) + .collect(); + Some((filter_index, filters, now)) } fn mark_started(&mut self, batch_id: u64, filter_index: usize, until: Timestamp) { @@ -788,7 +815,10 @@ impl DescendantSyncRotation { } fn mark_completed(&mut self, batch_id: u64, succeeded: bool) -> bool { - let Some(in_flight) = self.in_flight.filter(|request| request.batch_id == batch_id) else { + let Some(in_flight) = self + .in_flight + .filter(|request| request.batch_id == batch_id) + else { return false; }; self.in_flight = None; @@ -825,8 +855,7 @@ fn descendant_frontier_filters( frontier: &DescendantFrontier, since: Option, ) -> Vec { - let mut filters = - filters::tagged_one_of_our_root_event_filters(&frontier.event_ids, since); + let mut filters = filters::tagged_one_of_our_root_event_filters(&frontier.event_ids, since); filters.extend(filters::tagged_one_of_our_repo_event_filters( &frontier.coordinates, since, @@ -834,6 +863,65 @@ fn descendant_frontier_filters( filters } +fn tiered_auxiliary_filters( + repos: &HashSet, + root_events: &HashSet, + frontier: &DescendantFrontier, + since: Option, +) -> Vec { + let mut entries = filters::tiered_auxiliary_core_filters(repos, root_events, since); + entries.extend(filters::tiered_root_event_filters( + &frontier.event_ids, + true, + since, + )); + entries.extend(filters::tiered_repo_event_filters( + &frontier.coordinates, + true, + since, + )); + entries.sort_by_key(|entry| entry.tier); + entries +} + +fn split_live_tier_prefix( + entries: &[filters::TieredFilter], + max_filters_per_req: usize, + fits: impl Fn(&[Vec]) -> bool, +) -> (Vec, Vec) { + use filters::CoverageTier; + + let live_tiers = [ + CoverageTier::RootUppercase, + CoverageTier::CoreCompatibility, + CoverageTier::DescendantCanonical, + CoverageTier::DescendantQuote, + ]; + let mut live = Vec::new(); + for tier in live_tiers { + let mut candidate = live.clone(); + candidate.extend( + entries + .iter() + .filter(|entry| entry.tier == tier) + .map(|entry| entry.filter.clone()), + ); + let groups = live_filter_groups(&candidate, max_filters_per_req); + if fits(&groups) { + live = candidate; + } else { + break; + } + } + let live_json: HashSet<_> = live.iter().map(Filter::as_json).collect(); + let rotated = entries + .iter() + .filter(|entry| !live_json.contains(&entry.filter.as_json())) + .map(|entry| entry.filter.clone()) + .collect(); + (live, rotated) +} + /// Items included in a pending batch #[derive(Debug, Clone, Default)] pub struct PendingItems { @@ -940,8 +1028,7 @@ fn is_rate_limit_message(message: &str) -> bool { fn is_filter_count_refusal(message: &str) -> bool { let message = message.to_ascii_lowercase(); - message.contains("invalid number of filters") - || message.contains("max filter count") + message.contains("invalid number of filters") || message.contains("max filter count") } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -1085,7 +1172,10 @@ fn groups_within_subscription_byte_limit( let mut overflow = 0usize; for group in groups { let size = req_message_size(&group); - if used.checked_add(size).is_some_and(|total| total <= live_budget) { + if used + .checked_add(size) + .is_some_and(|total| total <= live_budget) + { used += size; admitted.push(group); } else { @@ -3240,7 +3330,10 @@ impl SyncManager { } // Step 2: Check if relay is rate-limited before creating new pending items - if self.health_tracker.is_subscription_paused(&action.relay_url) { + if self + .health_tracker + .is_subscription_paused(&action.relay_url) + { tracing::debug!( relay = %action.relay_url, full_repos = action.items.repos.len(), @@ -3261,7 +3354,13 @@ impl SyncManager { "handle_add_filters: calling sync_live and historic_sync" ); - if let Err(error) = self.sync_live(&action.relay_url, &action.filters).await { + let essential_live = filters::build_essential_live_filters( + &action.items.repos, + &action.items.state_only_repos, + &action.items.root_events, + None, + ); + if let Err(error) = self.sync_live(&action.relay_url, &essential_live).await { tracing::warn!( relay = %action.relay_url, %error, @@ -3525,8 +3624,7 @@ impl SyncManager { && !is_filter_count_refusal(&reason) && subscription_state_byte_limit(&reason).is_none() { - let already_paused = - health_tracker.is_rate_limited(&relay_url_clone); + let already_paused = health_tracker.is_rate_limited(&relay_url_clone); if already_paused { tracing::debug!( relay = %relay_url_clone, @@ -3785,26 +3883,25 @@ impl SyncManager { async fn reconcile_descendant_mode( &mut self, relay_url: &str, - root_events: &HashSet, + target: &algorithms::RelaySyncNeeds, ) { - let frontier = self.direct_thread_members(root_events).await; - if frontier.event_ids.is_empty() { - let _ = self - .close_descendant_live_coverage(relay_url, "frontier became empty") - .await; - self.descendant_sync_rotations.remove(relay_url); - return; - } + let frontier = self.direct_thread_members(&target.root_events).await; + let historic_entries = + tiered_auxiliary_filters(&target.repos, &target.root_events, &frontier, None); + let desired_fingerprint: Vec<_> = historic_entries + .iter() + .map(|entry| entry.filter.as_json()) + .collect(); if self .descendant_live_coverage .get(relay_url) - .is_some_and(|coverage| coverage.frontier == frontier) + .is_some_and(|coverage| coverage.coverage_fingerprint == desired_fingerprint) { return; } if self.descendant_live_coverage.contains_key(relay_url) && !self - .close_descendant_live_coverage(relay_url, "frontier changed") + .close_descendant_live_coverage(relay_url, "tiered coverage changed") .await { return; @@ -3815,11 +3912,18 @@ impl SyncManager { .as_secs() .saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS), ); - let historic_filters = descendant_frontier_filters(&frontier, None); - let live_filters = descendant_frontier_filters(&frontier, Some(live_since)); if let Some(connection) = self.connections.get(relay_url).cloned() { + let (live_filters, rotated_filters) = split_live_tier_prefix( + &historic_entries, + connection.max_filters_per_req(), + |groups| connection.can_admit_auxiliary_live_groups(groups), + ); + let live_filters: Vec<_> = live_filters + .into_iter() + .map(|filter| filter.since(live_since)) + .collect(); let groups = live_filter_groups(&live_filters, connection.max_filters_per_req()); - if connection.can_admit_auxiliary_live_groups(&groups) { + if !groups.is_empty() { match connection .subscribe_auxiliary_live_filter_groups(groups) .await @@ -3831,18 +3935,19 @@ impl SyncManager { self.descendant_live_coverage.insert( relay_url.to_string(), DescendantLiveCoverage { - frontier, + coverage_fingerprint: desired_fingerprint.clone(), + fallback_filters: rotated_filters.clone(), subscription_ids: subscription_ids.clone(), }, ); - self.descendant_sync_rotations.remove(relay_url); tracing::info!( relay = %relay_url, event_id_count, coordinate_count, filter_count, subscription_count = subscription_ids.len(), - "Installed auxiliary descendant live coverage" + rotated_filter_count = rotated_filters.len(), + "Installed priority-bounded auxiliary live coverage" ); // `limit:0` protects the future only. Pair every new // live frontier with a complete EOSE-closing baseline @@ -3850,13 +3955,26 @@ impl SyncManager { let _ = self .historic_sync_with_options( relay_url, - historic_filters, + historic_entries + .iter() + .map(|entry| entry.filter.clone()) + .collect(), PendingItems::default(), None, PendingBatchPurpose::Descendants, true, ) .await; + let rotation = self + .descendant_sync_rotations + .entry(relay_url.to_string()) + .or_default(); + let max_filters = connection.max_filters_per_req(); + if rotation.coverage_fingerprint + != rotation_fingerprint(&rotated_filters, max_filters) + { + rotation.refresh(rotated_filters, max_filters); + } return; } Err(error) => { @@ -3871,22 +3989,31 @@ impl SyncManager { tracing::info!( relay = %relay_url, filter_count = live_filters.len(), - "Descendant live coverage does not fit while preserving historic capacity" + "No auxiliary live tier fits while preserving historic capacity" ); } } + let rotated_filters: Vec<_> = historic_entries + .into_iter() + .map(|entry| entry.filter) + .collect(); + let max_filters = self + .connections + .get(relay_url) + .map(|connection| connection.max_filters_per_req()) + .unwrap_or(MAX_FILTERS_PER_REQ); let rotation = self .descendant_sync_rotations .entry(relay_url.to_string()) .or_default(); - if rotation.frontier != frontier { - rotation.refresh(frontier); + if rotation.coverage_fingerprint != rotation_fingerprint(&rotated_filters, max_filters) { + rotation.refresh(rotated_filters, max_filters); } } async fn start_descendant_fallback(&mut self, relay_url: &str) { - let Some((filter_index, filter, until)) = self + let Some((filter_index, filters, until)) = self .descendant_sync_rotations .get(relay_url) .and_then(|rotation| rotation.next_request(Timestamp::now())) @@ -3897,7 +4024,7 @@ impl SyncManager { if let Some(batch_id) = self .historic_sync_with_options( relay_url, - vec![filter], + filters, PendingItems::default(), None, PendingBatchPurpose::Descendants, @@ -3953,7 +4080,7 @@ impl SyncManager { self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); self.reconcile_descendant_mode( &reconcile_relay, - &targets[&reconcile_relay].root_events, + &targets[&reconcile_relay], ) .await; @@ -3978,7 +4105,7 @@ impl SyncManager { let mut filters = vec![filters::build_announcement_filter(None)]; let index = self.relay_sync_index.read().await; if let Some(state) = index.get(relay_url) { - filters.extend(filters::build_sync_level_aware_filters( + filters.extend(filters::build_essential_live_filters( &state.repos, &state.state_only_repos, &state.root_events, @@ -5649,8 +5776,7 @@ impl SyncManager { return false; } - let unbounded_groups = - live_filter_groups(&complete_live, connection.max_filters_per_req()); + let unbounded_groups = live_filter_groups(&complete_live, connection.max_filters_per_req()); let remote_limit = connection.remote_subscription_byte_limit(); let (complete_groups, overflow_groups) = groups_within_subscription_byte_limit(unbounded_groups, remote_limit, 0); @@ -5717,10 +5843,15 @@ impl SyncManager { .close_live_subscriptions(&coverage.subscription_ids) .await; } + let max_filters = self + .connections + .get(relay_url) + .map(|connection| connection.max_filters_per_req()) + .unwrap_or(MAX_FILTERS_PER_REQ); self.descendant_sync_rotations .entry(relay_url.to_string()) .or_default() - .refresh(coverage.frontier); + .refresh(coverage.fallback_filters, max_filters); tracing::warn!( relay = %relay_url, sub_id = %subscription_id, @@ -5748,10 +5879,7 @@ impl SyncManager { self.release_removed_descendant_batch(relay_url, batch); } let has_pending = self.has_pending_batches(relay_url).await; - if self - .deferred_consolidations - .request(relay_url, has_pending) - { + if self.deferred_consolidations.request(relay_url, has_pending) { let _ = self.consolidate(relay_url).await; } return; @@ -5841,10 +5969,7 @@ impl SyncManager { self.health_tracker.record_policy_refusal(relay_url); if let Some(metrics) = &self.metrics { metrics.record_policy_refusal(relay_url, category.label()); - metrics.record_health_state( - relay_url, - self.health_tracker.get_state(relay_url), - ); + metrics.record_health_state(relay_url, self.health_tracker.get_state(relay_url)); } tracing::warn!( relay = %relay_url, @@ -5915,10 +6040,8 @@ impl SyncManager { }; if self.has_pending_batches(&relay_url).await { - self.byte_limited_live_relays.insert( - relay_url, - now + Duration::from_secs(10), - ); + self.byte_limited_live_relays + .insert(relay_url, now + Duration::from_secs(10)); return; } @@ -5929,17 +6052,13 @@ impl SyncManager { .get(&relay_url) .is_some_and(|state| state.connection_status.is_live_sync_active()); if !connected { - self.byte_limited_live_relays.insert( - relay_url, - now + byte_limited_catchup_interval(), - ); + self.byte_limited_live_relays + .insert(relay_url, now + byte_limited_catchup_interval()); return; } - self.byte_limited_live_relays.insert( - relay_url.clone(), - now + byte_limited_catchup_interval(), - ); + self.byte_limited_live_relays + .insert(relay_url.clone(), now + byte_limited_catchup_interval()); let overlap = byte_limited_catchup_interval() + Duration::from_secs(60); let since = Timestamp::from(Timestamp::now().as_secs().saturating_sub(overlap.as_secs())); let filters = self.complete_live_filters(&relay_url, Some(since)).await; @@ -6880,31 +6999,74 @@ mod tests { let now = Timestamp::from_secs(200_000); let mut rotation = DescendantSyncRotation::default(); - rotation.refresh(frontier); + rotation.refresh( + descendant_frontier_filters(&frontier, None), + MAX_FILTERS_PER_REQ, + ); let (filter_index, first, until) = rotation.next_request(now).unwrap(); - assert!(serde_json::to_value(first).unwrap().get("since").is_none()); + assert!(first.iter().all(|filter| { + serde_json::to_value(filter).unwrap().get("since").is_none() + })); rotation.mark_started(41, filter_index, until); assert!(rotation.next_request(now).is_none()); assert!(rotation.mark_completed(41, false)); let (retry_index, retry, retry_until) = rotation.next_request(now).unwrap(); assert_eq!(retry_index, filter_index); - assert!(serde_json::to_value(retry).unwrap().get("since").is_none()); + assert!(retry.iter().all(|filter| { + serde_json::to_value(filter).unwrap().get("since").is_none() + })); rotation.mark_started(42, retry_index, retry_until); assert!(rotation.mark_completed(42, true)); - // Complete the other two e/E/q variants so the rotation returns to - // the first filter with its successful cursor and overlap. - for batch_id in [43, 44] { - let (index, _, upper) = rotation.next_request(now).unwrap(); - rotation.mark_started(batch_id, index, upper); - assert!(rotation.mark_completed(batch_id, true)); - } - let (_, recent, _) = rotation.next_request(Timestamp::from_secs(201_000)).unwrap(); + // All three e/E/q variants share one relay-compatible REQ, so the + // next turn returns to that group with its successful overlap cursor. + let (_, recent, _) = rotation + .next_request(Timestamp::from_secs(201_000)) + .unwrap(); + assert!(recent.iter().all(|filter| { + serde_json::to_value(filter).unwrap()["since"] + == serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS) + })); + } + + #[test] + fn tier_cutoff_keeps_only_complete_priority_tiers_live() { + use filters::{CoverageTier, TieredFilter}; + + let entries = vec![ + TieredFilter { + tier: CoverageTier::RootUppercase, + filter: Filter::new().kind(Kind::Custom(31_001)), + }, + TieredFilter { + tier: CoverageTier::CoreCompatibility, + filter: Filter::new().kind(Kind::Custom(31_002)), + }, + TieredFilter { + tier: CoverageTier::CoreCompatibility, + filter: Filter::new().kind(Kind::Custom(31_003)), + }, + TieredFilter { + tier: CoverageTier::DescendantCanonical, + filter: Filter::new().kind(Kind::Custom(31_004)), + }, + TieredFilter { + tier: CoverageTier::HistoricOnly, + filter: Filter::new().kind(Kind::Custom(31_005)), + }, + ]; + + let (live, rotated) = split_live_tier_prefix(&entries, 2, |groups| groups.len() <= 1); assert_eq!( - serde_json::to_value(recent).unwrap()["since"], - serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS) + live.len(), + 1, + "the two-filter compatibility tier must not be split" ); + assert_eq!(rotated.len(), 4); + assert!(rotated.iter().any(|filter| { + filter.as_json() == Filter::new().kind(Kind::Custom(31_005)).as_json() + })); } #[test] @@ -6930,7 +7092,9 @@ mod tests { groups.len() > 1, "large complete descendant coverage must reserve several subscriptions" ); - assert!(groups.iter().all(|group| group.len() <= MAX_FILTERS_PER_REQ)); + assert!(groups + .iter() + .all(|group| group.len() <= MAX_FILTERS_PER_REQ)); } #[test] @@ -7510,11 +7674,8 @@ mod tests { let second = vec![Filter::new().kind(Kind::Metadata).limit(0)]; let limit = SUBSCRIPTION_BYTE_RESERVED_MARGIN + req_message_size(&first); - let (admitted, overflow) = groups_within_subscription_byte_limit( - vec![first.clone(), second], - Some(limit), - 0, - ); + let (admitted, overflow) = + groups_within_subscription_byte_limit(vec![first.clone(), second], Some(limit), 0); assert_eq!(admitted, vec![first]); assert_eq!(overflow, 1); diff --git a/tests/sync/descendant_sync.rs b/tests/sync/descendant_sync.rs index 0fc8d79..0da058f 100644 --- a/tests/sync/descendant_sync.rs +++ b/tests/sync/descendant_sync.rs @@ -139,7 +139,7 @@ async fn descendant_live_coverage_is_preferred_when_capacity_remains() { assert!( wait_for_log( &syncing.log_path(), - "Installed auxiliary descendant live coverage", + "Installed priority-bounded auxiliary live coverage", Duration::from_secs(20), ) .await, From c5c66910615cf07879672c1fa5352ecfd76a1076 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 09:42:14 +0000 Subject: [PATCH 3/7] perf(sync): pack related historic reference filters Motivation: historic sync appended descendant filter families after core families, so the same event could be returned in separate subscriptions when it matched both a root/repository reference and a direct-descendant reference. Approach: rebuild non-generic historic batches from their tracked items, unioning core and descendant values under one e/E/q family and one a/A/q family before the existing count/byte grouper packs REQs. State-only repositories retain their state filter and do not gain reference coverage. Correctness: the union preserves every prior tag value and tag case, while Nostr OR-filter semantics let a relay deduplicate overlapping matches within each grouped REQ. The PendingItems set remains the batch completion authority and since is still applied once downstream. Excluded scope: repository-to-root ownership is not added to PendingItems, so packing is per pending repository batch rather than one wire subscription per repository; pagination and negentropy are unchanged. Validation: nix develop -c cargo test --lib (686 passed), including a regression test that pins seven packed filters instead of thirteen separately appended filters. --- docs/explanation/sync-scaling-constraints.md | 5 ++ src/sync/mod.rs | 69 ++++++++++++++++++-- 2 files changed, 70 insertions(+), 4 deletions(-) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 3d89b77..8b10b46 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -339,6 +339,11 @@ the incremental slot cost. Thus byte-bound groups are not rebuilt merely because they contain fewer than the maximum filter count. Repository and identifier inputs are sorted before byte chunking so equivalent coverage has stable group identity. +Historic repository batches union root and direct-descendant values into one +`e`/`E`/`q` family, and repository and addressable-descendant coordinates into +one `a`/`A`/`q` family, before ordinary count/byte grouping. This lets one REQ +deduplicate events matching both core and descendant references without +creating a subscription per repository. If the changed tail cannot fit the count or learned byte budget, no extension is opened and historic recovery remains available. Capacity pressure is the backstop which may schedule a complete regroup after outstanding historic diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 07f683b..68e38a6 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -851,6 +851,7 @@ fn descendant_event_coordinate(event: &Event) -> Option { )) } +#[cfg(test)] fn descendant_frontier_filters( frontier: &DescendantFrontier, since: Option, @@ -863,6 +864,36 @@ fn descendant_frontier_filters( filters } +fn packed_historic_filters(items: &PendingItems, frontier: &DescendantFrontier) -> Vec { + let all_repos: HashSet<_> = items + .repos + .union(&items.state_only_repos) + .cloned() + .collect(); + let mut filters = filters::state_event_filters_for_our_repos(&all_repos, None); + + let coordinate_values: HashSet<_> = items + .repos + .union(&frontier.coordinates) + .cloned() + .collect(); + filters.extend(filters::tagged_one_of_our_repo_event_filters( + &coordinate_values, + None, + )); + + let event_values: HashSet<_> = items + .root_events + .union(&frontier.event_ids) + .cloned() + .collect(); + filters.extend(filters::tagged_one_of_our_root_event_filters( + &event_values, + None, + )); + filters +} + fn tiered_auxiliary_filters( repos: &HashSet, root_events: &HashSet, @@ -6496,14 +6527,19 @@ impl SyncManager { async fn historic_sync( &mut self, relay_url: &str, - mut filters: Vec, + filters: Vec, items: PendingItems, since: Option, ) -> Option { - if !items.root_events.is_empty() { + let filters = if items.repos.is_empty() + && items.state_only_repos.is_empty() + && items.root_events.is_empty() + { + filters + } else { let frontier = self.direct_thread_members(&items.root_events).await; - filters.extend(descendant_frontier_filters(&frontier, None)); - } + packed_historic_filters(&items, &frontier) + }; self.historic_sync_with_options( relay_url, filters, @@ -7069,6 +7105,31 @@ mod tests { })); } + #[test] + fn historic_packing_unions_core_and_descendant_reference_values() { + let root = EventId::from_byte_array([1; 32]); + let descendant = EventId::from_byte_array([2; 32]); + let items = PendingItems { + repos: HashSet::from(["30617:owner:repo".to_string()]), + state_only_repos: HashSet::new(), + root_events: HashSet::from([root]), + }; + let frontier = DescendantFrontier { + event_ids: HashSet::from([descendant]), + coordinates: HashSet::from(["1621:author:patch".to_string()]), + }; + + let packed = packed_historic_filters(&items, &frontier); + // One state filter plus one a/A/q and one e/E/q family. Appending + // descendants separately would produce thirteen filters here. + assert_eq!(packed.len(), 7); + let json = packed.iter().map(Filter::as_json).collect::>().join("\n"); + assert!(json.contains(&root.to_hex())); + assert!(json.contains(&descendant.to_hex())); + assert!(json.contains("30617:owner:repo")); + assert!(json.contains("1621:author:patch")); + } + #[test] fn large_descendant_frontier_splits_across_live_subscriptions() { let members: HashSet = (0..600u32) From 27f31e50ed86ce83f12bd247c4bb7be8b8cbf2b0 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 10:09:00 +0000 Subject: [PATCH 4/7] fix(sync): preserve complete tier coverage after CLOSED Motivation: a relay closing one admitted auxiliary live subscription moved only the already-demoted filters into rotation. Tiers that had been live then waited for periodic reconciliation, creating a bounded but unnecessary coverage gap on a quiet relay. Approach: retain the complete auxiliary frontier beside each admitted live group and hand that full set to the existing cursor-preserving REQ+EOSE rotation when any member closes. Correctness: the fallback is derived from the same ordered entries and fingerprint used for admission, so it contains both the former live prefix and the already-rotated suffix without inventing new coverage. The existing grouping, overlap cursor, ledger, and next reconciliation remain authoritative. Excluded scope: this does not change CLOSED classification, retry timing, connection sharding, or the live admission cutoff. Validation: nix develop -c cargo test --lib tier_cutoff_keeps_only_complete_priority_tiers_live passed; the assertion verifies that the fallback contains every admitted live filter as well as the demoted tiers. --- src/sync/mod.rs | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 68e38a6..5f14d83 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -953,6 +953,10 @@ fn split_live_tier_prefix( (live, rotated) } +fn complete_auxiliary_fallback(entries: &[filters::TieredFilter]) -> Vec { + entries.iter().map(|entry| entry.filter.clone()).collect() +} + /// Items included in a pending batch #[derive(Debug, Clone, Default)] pub struct PendingItems { @@ -3960,6 +3964,11 @@ impl SyncManager { .await { Ok(subscription_ids) => { + // A relay can close any one member of this live set. + // Keep the complete frontier ready for rotation so + // the CLOSED path preserves every tier immediately, + // including those that had been live until now. + let fallback_filters = complete_auxiliary_fallback(&historic_entries); let filter_count = live_filters.len(); let event_id_count = frontier.event_ids.len(); let coordinate_count = frontier.coordinates.len(); @@ -3967,7 +3976,7 @@ impl SyncManager { relay_url.to_string(), DescendantLiveCoverage { coverage_fingerprint: desired_fingerprint.clone(), - fallback_filters: rotated_filters.clone(), + fallback_filters, subscription_ids: subscription_ids.clone(), }, ); @@ -7103,6 +7112,12 @@ mod tests { assert!(rotated.iter().any(|filter| { filter.as_json() == Filter::new().kind(Kind::Custom(31_005)).as_json() })); + + let fallback = complete_auxiliary_fallback(&entries); + assert_eq!(fallback.len(), entries.len()); + assert!(live.iter().all(|live_filter| fallback + .iter() + .any(|filter| filter.as_json() == live_filter.as_json()))); } #[test] From 4a7317465c5931ba1ac9d68bfe759ab52e5d0897 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 11:06:42 +0000 Subject: [PATCH 5/7] fix(sync): register hydration batches before delivery Motivation: archive startup against relay.ngit.dev discovered 50,530 remote events and paced them across 169 exact-ID REQs. The independent event processor could receive early chunks before historic_sync registered their subscription IDs and requested-event set, so valid deliveries were invisible to negentropy validation and could become an artificial missing residual. The original batch remained Syncing and prevented auxiliary tier reconciliation on the largest relay. Approach: pre-generate the complete hydration attempt, register every subscription ID plus requested/received accounting in PendingBatch, then send each REQ with its caller-selected ID. Apply the same ordering to residual retries. A failed send removes its planned ID and definitively fails a fully drained attempt instead of leaving phantom outstanding work. Correctness: every EVENT and terminal signal now names an ID present before the first wire request can be answered. Pre-registering all paced chunks prevents an early EOSE from draining a partially launched attempt, while the existing actor remains the sole batch confirmer. RelayConnection still acquires and releases the same generation-scoped ledger and class permits. Excluded scope: this does not alter negentropy diffing, 300-ID chunk size, query pacing, missing-event backoff, tier admission, or connection-state eligibility. Validation: nix develop -c cargo test --lib passed 687/687. The regression registers the production-sized 169-chunk attempt, accounts immediate deliveries, and verifies every terminal ID is attributable; the existing immediate-empty-EOSE permit test also passes. --- docs/explanation/sync-scaling-constraints.md | 5 +- src/sync/mod.rs | 266 ++++++++++++++----- src/sync/relay_connection.rs | 35 ++- 3 files changed, 227 insertions(+), 79 deletions(-) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 8b10b46..bdf455e 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -370,7 +370,10 @@ purgatory polling uses the same transient class bound and shared ledger as historic pagination. Transient subscription IDs and their permits are registered locally before the REQ is sent; this ordering is required because an empty or cached response can deliver EOSE/CLOSED before the SDK subscribe -call returns. Subscribe failure rolls that pre-registration back. Unexpected +call returns. Negentropy hydration also registers the complete paced chunk set +and its requested-event accounting in the pending batch before sending the +first REQ, so early deliveries cannot become an artificial missing residual. +Subscribe failure rolls both forms of pre-registration back. Unexpected CLOSED for a persistent live subscription is reported to the manager, which recomputes and transactionally reopens complete live coverage. Each reconnect closes the retired ledger and creates a new diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 5f14d83..e69266f 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1012,6 +1012,23 @@ fn take_batch_containing_subscription( Some(batch) } +fn register_negentropy_hydration_attempt( + batch: &mut PendingBatch, + subscription_ids: impl IntoIterator, + requested_event_ids: impl IntoIterator, + is_retry: bool, +) { + // Register the complete attempt before its first REQ is sent. A fast + // relay can deliver events and EOSE while later paced chunks are still + // waiting to open; pre-registration makes that whole stream attributable. + batch.outstanding_subs.extend(subscription_ids); + batch.requested_event_ids = Some(requested_event_ids.into_iter().collect()); + batch.received_event_ids = Some(HashSet::new()); + if is_retry { + batch.retry_count += 1; + } +} + fn mark_deferred_pagination( batch: &mut PendingBatch, completed_sub_id: &SubscriptionId, @@ -2278,57 +2295,97 @@ impl SyncManager { drop(pending); // Create new subscriptions for missing events - let retry_filters: Vec<_> = missing + let retry_subscriptions: Vec<_> = missing .chunks(300) - .map(|c| Filter::new().ids(c.iter().copied())) + .map(|chunk| { + ( + SubscriptionId::generate(), + Filter::new().ids(chunk.iter().copied()), + ) + }) .collect(); + { + let mut pending = self.pending_sync_index.write().await; + if let Some(batch) = pending + .get_mut(&relay_url_for_retry) + .and_then(|batches| { + batches.iter_mut().find(|batch| batch.batch_id == batch_id) + }) + { + register_negentropy_hydration_attempt( + batch, + retry_subscriptions + .iter() + .map(|(subscription_id, _)| subscription_id.clone()), + missing.iter().copied(), + true, + ); + } + } - let mut new_sub_ids = HashSet::new(); - if let Some(conn) = self.connections.get(&relay_url_for_retry) { - for filter in retry_filters { - match conn - .subscribe_filter(filter, TransientRequestClass::NegentropyRetry) - .await - { - Ok(sub_id) => { - new_sub_ids.insert(sub_id); - } - Err(e) => { - tracing::error!( - relay = %relay_url_for_retry, - batch_id = batch_id, - error = %e, - "Failed to create retry subscription for missing events" - ); + let mut successful_subscriptions = 0usize; + for (subscription_id, filter) in &retry_subscriptions { + let result = if let Some(conn) = self.connections.get(&relay_url_for_retry) + { + conn.subscribe_filter_with_id( + filter.clone(), + TransientRequestClass::NegentropyRetry, + subscription_id.clone(), + ) + .await + } else { + Err("Relay connection disappeared before hydration retry".to_string()) + }; + match result { + Ok(_) => successful_subscriptions += 1, + Err(e) => { + tracing::error!( + relay = %relay_url_for_retry, + batch_id = batch_id, + error = %e, + "Failed to create retry subscription for missing events" + ); + let mut pending = self.pending_sync_index.write().await; + if let Some(batch) = pending + .get_mut(&relay_url_for_retry) + .and_then(|batches| { + batches.iter_mut().find(|batch| batch.batch_id == batch_id) + }) + { + batch.outstanding_subs.remove(subscription_id); } } } } - if !new_sub_ids.is_empty() { - // Re-acquire lock and update batch with new subscriptions + if successful_subscriptions > 0 { + // Immediate EOSEs may have drained every successful + // subscription before the final failed send unwound. + // Complete that otherwise signal-less edge here. let mut pending = self.pending_sync_index.write().await; - if let Some(batches) = pending.get_mut(&relay_url_for_retry) { - if let Some(batch) = batches.iter_mut().find(|b| b.batch_id == batch_id) - { - batch.outstanding_subs.extend(new_sub_ids.clone()); - // Update requested_event_ids to only include missing ones - batch.requested_event_ids = Some(missing.iter().cloned().collect()); - // Clear received_event_ids for fresh tracking - batch.received_event_ids = Some(HashSet::new()); - // Increment retry counter - batch.retry_count += 1; - - tracing::info!( - relay = %relay_url_for_retry, - batch_id = batch_id, - retry_subs = new_sub_ids.len(), - missing_events = missing.len(), - retry_attempt = batch.retry_count, - "Created retry subscriptions for missing negentropy events" - ); - } + let completed_batch = take_drained_batch_as_failed( + &mut pending, + &relay_url_for_retry, + batch_id, + ); + let retry_attempt = pending + .get(&relay_url_for_retry) + .and_then(|batches| { + batches.iter().find(|batch| batch.batch_id == batch_id) + }) + .map(|batch| batch.retry_count); + drop(pending); + if let Some(batch) = completed_batch { + self.confirm_batch(&relay_url_for_retry, batch).await; } + tracing::info!( + relay = %relay_url_for_retry, + batch_id = batch_id, + retry_subs = successful_subscriptions, + missing_events = missing.len(), + retry_attempt = retry_attempt.unwrap_or(retry_count + 1), + "Created retry subscriptions for missing negentropy events" + ); // Early return - batch not complete yet, waiting for retry EOSE return; } else { @@ -6788,49 +6845,69 @@ impl SyncManager { "Creating subscriptions to fetch missing events by ID" ); - let mut subscription_ids = HashSet::new(); - for (idx, filter) in ids_filters.iter().enumerate() { - if let Some(conn) = self.connections.get(relay_url) { - match conn - .subscribe_filter( - filter.clone(), - TransientRequestClass::NegentropyHydration, - ) - .await - { - Ok(sub_id) => { - subscription_ids.insert(sub_id); - } - Err(e) => { - tracing::error!( - relay = %relay_url, - batch_id = batch_id, - chunk_idx = idx, - error = %e, - "Failed to subscribe to ID filter chunk" - ); - } - } - } - } + let planned_subscriptions: Vec<_> = ids_filters + .into_iter() + .map(|filter| (SubscriptionId::generate(), filter)) + .collect(); { let mut pending = self.pending_sync_index.write().await; if let Some(relay_batches) = pending.get_mut(relay_url) { if let Some(batch) = relay_batches.iter_mut().find(|b| b.batch_id == batch_id) { - batch.outstanding_subs.extend(subscription_ids.clone()); - // Store requested event IDs for validation after EOSE - batch.requested_event_ids = - Some(all_remote_ids.iter().cloned().collect()); - batch.received_event_ids = Some(HashSet::new()); + register_negentropy_hydration_attempt( + batch, + planned_subscriptions + .iter() + .map(|(subscription_id, _)| subscription_id.clone()), + all_remote_ids.iter().copied(), + false, + ); + } + } + } + for (idx, (subscription_id, filter)) in + planned_subscriptions.iter().enumerate() + { + let result = if let Some(conn) = self.connections.get(relay_url) { + conn.subscribe_filter_with_id( + filter.clone(), + TransientRequestClass::NegentropyHydration, + subscription_id.clone(), + ) + .await + } else { + Err("Relay connection disappeared before hydration REQ".to_string()) + }; + if let Err(e) = result { + tracing::error!( + relay = %relay_url, + batch_id = batch_id, + chunk_idx = idx, + error = %e, + "Failed to subscribe to ID filter chunk" + ); + let completed_batch = { + let mut pending = self.pending_sync_index.write().await; + if let Some(batch) = pending + .get_mut(relay_url) + .and_then(|batches| { + batches.iter_mut().find(|batch| batch.batch_id == batch_id) + }) + { + batch.outstanding_subs.remove(subscription_id); + } + take_drained_batch_as_failed(&mut pending, relay_url, batch_id) + }; + if let Some(batch) = completed_batch { + self.confirm_batch(relay_url, batch).await; } } } tracing::debug!( relay = %relay_url, batch_id = batch_id, - subscription_ids = subscription_ids.len(), + subscription_ids = planned_subscriptions.len(), events = all_remote_ids.len(), "historic_sync (Negentropy) created subscriptions to fetch missing events by id, awaiting EOSE" ); @@ -8138,6 +8215,51 @@ mod tests { assert!(!batch.failed); } + #[test] + fn hydration_attempt_registers_every_chunk_before_immediate_delivery() { + let mut batch = PendingBatch { + batch_id: 7, + purpose: PendingBatchPurpose::Core, + items: PendingItems::default(), + outstanding_subs: HashSet::new(), + sync_method: SyncMethod::Negentropy, + pagination_state: HashMap::new(), + requested_event_ids: None, + received_event_ids: None, + initial_hydration_counts: None, + retry_count: 0, + failed: false, + }; + let subscription_ids: Vec<_> = (0..169) + .map(|index| SubscriptionId::new(format!("hydration-{index}"))) + .collect(); + let event_ids = [EventId::from_byte_array([1; 32]), EventId::from_byte_array([2; 32])]; + + register_negentropy_hydration_attempt( + &mut batch, + subscription_ids.iter().cloned(), + event_ids, + false, + ); + + assert_eq!(batch.outstanding_subs.len(), 169); + assert_eq!(batch.requested_event_ids.as_ref().unwrap().len(), 2); + assert_eq!(batch.received_event_ids, Some(HashSet::new())); + for event_id in event_ids { + if batch.requested_event_ids.as_ref().unwrap().contains(&event_id) { + batch.received_event_ids.as_mut().unwrap().insert(event_id); + } + } + assert_eq!(batch.received_event_ids.as_ref().unwrap().len(), 2); + for subscription_id in &subscription_ids { + assert!( + batch.outstanding_subs.remove(subscription_id), + "an immediate terminal signal must find its pre-registered chunk" + ); + } + assert!(batch.outstanding_subs.is_empty()); + } + #[test] fn test_pending_batch_req_eose_fields() { // Test that REQ+EOSE batches don't use negentropy fields diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index d93bd86..eed345c 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -1201,7 +1201,7 @@ impl RelayConnection { filter_groups: Vec>, ) -> Result, String> { self.replace_live_filter_groups_with(filter_groups, |filters, permit| { - self.subscribe_filters_with_live_permit(filters, None, Some(permit)) + self.subscribe_filters_with_live_permit(filters, None, Some(permit), None) }) .await } @@ -1229,7 +1229,9 @@ impl RelayConnection { .map(|_| ()) .map_err(|error| format!("{subscription_id}: {error}")) }, - |filters, permit| self.subscribe_filters_with_live_permit(filters, None, Some(permit)), + |filters, permit| { + self.subscribe_filters_with_live_permit(filters, None, Some(permit), None) + }, ) .await } @@ -1816,6 +1818,24 @@ impl RelayConnection { self.subscribe_filters(vec![filter], request_class).await } + /// Open a transient subscription with an ID already registered by its + /// caller. Pending-batch owners use this to make an immediate EOSE visible + /// before the wire request can be answered. + pub async fn subscribe_filter_with_id( + &self, + filter: Filter, + request_class: TransientRequestClass, + subscription_id: SubscriptionId, + ) -> Result { + self.subscribe_filters_with_live_permit( + vec![filter], + Some(request_class), + None, + Some(subscription_id), + ) + .await + } + /// Subscribe to several OR filters under one NIP-01 subscription ID. /// /// Relays apply active-REQ limits to subscription IDs, not to the filters @@ -1826,7 +1846,7 @@ impl RelayConnection { filters: Vec, request_class: TransientRequestClass, ) -> Result { - self.subscribe_filters_with_live_permit(filters, Some(request_class), None) + self.subscribe_filters_with_live_permit(filters, Some(request_class), None, None) .await } @@ -1838,7 +1858,7 @@ impl RelayConnection { filter_groups: Vec>, ) -> Result, String> { self.subscribe_live_filter_groups_with(filter_groups, |filters, permit| { - self.subscribe_filters_with_live_permit(filters, None, Some(permit)) + self.subscribe_filters_with_live_permit(filters, None, Some(permit), None) }) .await } @@ -1897,6 +1917,7 @@ impl RelayConnection { filters: Vec, transient_class: Option, reserved_live_permit: Option, + requested_subscription_id: Option, ) -> Result { if filters.is_empty() { return Err("Cannot subscribe with an empty filter set".to_string()); @@ -1969,7 +1990,9 @@ impl RelayConnection { // The relay can answer an empty or cached query before `subscribe` // returns. Register transient ownership against a caller-chosen ID // first so an immediate EOSE/CLOSED cannot race past local accounting. - let transient_sub_id = transient_class.map(|_| SubscriptionId::generate()); + let transient_sub_id = transient_class.map(|_| { + requested_subscription_id.unwrap_or_else(SubscriptionId::generate) + }); if let (Some(sub_id), Some(permit)) = (&transient_sub_id, transient_permit) { self.hold_transient_req_permit(sub_id.clone(), permit); } @@ -3704,7 +3727,7 @@ mod tests { } }, |filters, permit| { - connection.subscribe_filters_with_live_permit(filters, None, Some(permit)) + connection.subscribe_filters_with_live_permit(filters, None, Some(permit), None) }, ) .await From 88ef365397760332c5d50cd46150a95afff31f66 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 12:56:22 +0000 Subject: [PATCH 6/7] chore(sync): expose relay event-pipeline pressure The archive reconciliation against relay.ngit.dev hydrates 50,580 IDs and takes long enough that batch completion alone cannot distinguish slow peer delivery, bounded-channel backpressure, duplicate lookup cost, or write-policy work. Per-event logging would make the production signal less usable and add substantial overhead to the workload being measured. Carry the data-lane arrival instant with EVENT and EOSE messages. Emit one per-relay aggregate every 30 seconds with outcome counts, throughput, queue delay, processing time, and current queue depth; separately report EOSE messages delayed at least five seconds in the processor lane. A focused unit test verifies that queue and processing costs remain separate. These measurements deliberately observe rather than change concurrency, queue capacity, or processing order. Resource isolation for foreground relay traffic and adaptive client-side subscription ceilings remain follow-up design work. The timestamps begin when the processor-facing notification lane observes a message, so peer-side delivery time is outside their scope. Validation: nix develop -c cargo test --lib (688 passed). cargo fmt --check remains blocked by pre-existing formatting drift across the stacked branch under the current Rust 1.96 formatter; git diff --check passes. --- src/sync/mod.rs | 145 ++++++++++++++++++++++++++++++++++- src/sync/relay_connection.rs | 19 +++-- 2 files changed, 157 insertions(+), 7 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index e69266f..894d015 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -430,6 +430,93 @@ pub enum ProcessResult { Rejected, } +/// A low-volume summary of the per-relay data lane. +/// +/// Sync bursts can contain tens of thousands of duplicates, so per-event logs +/// obscure whether time is spent waiting for the bounded channel or applying +/// policy. A periodic aggregate keeps production diagnosis cheap enough to +/// leave enabled while preserving both parts of that distinction. +#[derive(Debug)] +struct EventPipelineWindow { + started_at: std::time::Instant, + delivered: u64, + saved: u64, + duplicate: u64, + purgatory: u64, + rejected: u64, + queue_delay: std::time::Duration, + max_queue_delay: std::time::Duration, + processing_time: std::time::Duration, + max_processing_time: std::time::Duration, +} + +impl Default for EventPipelineWindow { + fn default() -> Self { + Self { + started_at: std::time::Instant::now(), + delivered: 0, + saved: 0, + duplicate: 0, + purgatory: 0, + rejected: 0, + queue_delay: std::time::Duration::ZERO, + max_queue_delay: std::time::Duration::ZERO, + processing_time: std::time::Duration::ZERO, + max_processing_time: std::time::Duration::ZERO, + } + } +} + +impl EventPipelineWindow { + const REPORT_INTERVAL: std::time::Duration = std::time::Duration::from_secs(30); + + fn record( + &mut self, + result: ProcessResult, + queue_delay: std::time::Duration, + processing_time: std::time::Duration, + ) { + self.delivered += 1; + match result { + ProcessResult::Saved => self.saved += 1, + ProcessResult::Duplicate => self.duplicate += 1, + ProcessResult::Purgatory => self.purgatory += 1, + ProcessResult::Rejected => self.rejected += 1, + } + self.queue_delay += queue_delay; + self.max_queue_delay = self.max_queue_delay.max(queue_delay); + self.processing_time += processing_time; + self.max_processing_time = self.max_processing_time.max(processing_time); + } + + fn report_if_due(&mut self, relay: &str, queue_depth: usize) { + let elapsed = self.started_at.elapsed(); + if elapsed < Self::REPORT_INTERVAL || self.delivered == 0 { + return; + } + + let delivered = self.delivered as f64; + tracing::info!( + relay, + window_seconds = elapsed.as_secs_f64(), + delivered = self.delivered, + saved = self.saved, + duplicate = self.duplicate, + purgatory = self.purgatory, + rejected = self.rejected, + events_per_second = delivered / elapsed.as_secs_f64(), + average_queue_delay_ms = self.queue_delay.as_secs_f64() * 1000.0 / delivered, + max_queue_delay_ms = self.max_queue_delay.as_secs_f64() * 1000.0, + average_processing_ms = self.processing_time.as_secs_f64() * 1000.0 / delivered, + max_processing_ms = self.max_processing_time.as_secs_f64() * 1000.0, + queue_depth, + queue_capacity = relay_connection::RELAY_EVENT_BUFFER_CAPACITY, + "Relay sync event-pipeline window" + ); + *self = Self::default(); + } +} + /// Statistics from re-processing events from hot cache #[derive(Debug, Clone, Default)] pub struct ReprocessingStats { @@ -3546,10 +3633,13 @@ impl SyncManager { tokio::spawn(async move { let mut disconnect_sent = false; + let mut pipeline_window = EventPipelineWindow::default(); while let Some(relay_event) = event_rx.recv().await { match relay_event { - RelayEvent::Event(event, subscription_id) => { + RelayEvent::Event(event, subscription_id, data_lane_arrival) => { + let queue_delay = data_lane_arrival.elapsed(); + let processing_started = std::time::Instant::now(); // 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 @@ -3578,6 +3668,12 @@ impl SyncManager { relay = %relay_url_clone, "Skipping previously rejected announcement event" ); + pipeline_window.record( + ProcessResult::Rejected, + queue_delay, + processing_started.elapsed(), + ); + pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); continue; } @@ -3658,8 +3754,21 @@ impl SyncManager { } } } + pipeline_window.record(result, queue_delay, processing_started.elapsed()); + pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); } - RelayEvent::EndOfStoredEvents(sub_id) => { + RelayEvent::EndOfStoredEvents(sub_id, data_lane_arrival) => { + let data_lane_delay = data_lane_arrival.elapsed(); + if data_lane_delay >= std::time::Duration::from_secs(5) { + tracing::info!( + relay = %relay_url_clone, + sub_id = %sub_id, + data_lane_delay_seconds = data_lane_delay.as_secs_f64(), + queue_depth = event_rx.len(), + queue_capacity = relay_connection::RELAY_EVENT_BUFFER_CAPACITY, + "EOSE reached sync manager after data-lane delay" + ); + } tracing::debug!( relay = %relay_url_clone, sub_id = %sub_id, @@ -7024,6 +7133,38 @@ impl SyncManager { mod tests { use super::*; + #[test] + fn event_pipeline_window_separates_queue_and_processing_costs() { + let mut window = EventPipelineWindow::default(); + window.record( + ProcessResult::Duplicate, + std::time::Duration::from_millis(12), + std::time::Duration::from_millis(3), + ); + window.record( + ProcessResult::Saved, + std::time::Duration::from_millis(8), + std::time::Duration::from_millis(7), + ); + + assert_eq!(window.delivered, 2); + assert_eq!(window.saved, 1); + assert_eq!(window.duplicate, 1); + assert_eq!(window.queue_delay, std::time::Duration::from_millis(20)); + assert_eq!( + window.max_queue_delay, + std::time::Duration::from_millis(12) + ); + assert_eq!( + window.processing_time, + std::time::Duration::from_millis(10) + ); + assert_eq!( + window.max_processing_time, + std::time::Duration::from_millis(7) + ); + } + #[test] fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { let keys = Keys::generate(); diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index eed345c..987ea86 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -490,10 +490,10 @@ where /// Events from a relay connection #[derive(Debug)] pub enum RelayEvent { - /// A new event was received (event, subscription_id) - Event(Box, SubscriptionId), + /// A new event was received (event, subscription_id, data-lane arrival). + Event(Box, SubscriptionId, std::time::Instant), /// End of stored events for a subscription - EndOfStoredEvents(SubscriptionId), + EndOfStoredEvents(SubscriptionId, std::time::Instant), /// NOTICE message from relay Notice(String), /// Connection was closed @@ -1655,6 +1655,7 @@ impl RelayConnection { event, subscription_id, } => { + let data_lane_arrival = std::time::Instant::now(); self.record_transient_req_event(&subscription_id); tracing::trace!( relay = %url, @@ -1663,7 +1664,11 @@ impl RelayConnection { "Received event" ); if event_sender - .send(RelayEvent::Event(Box::new(*event), subscription_id.clone())) + .send(RelayEvent::Event( + Box::new(*event), + subscription_id.clone(), + data_lane_arrival, + )) .await .is_err() { @@ -1673,11 +1678,15 @@ impl RelayConnection { } RelayNotification::Message { message } => match *message { RelayMessage::EndOfStoredEvents(sub_id) => { + let data_lane_arrival = std::time::Instant::now(); tracing::debug!(relay = %url, sub_id = ?sub_id, "Received EOSE"); // Convert Cow to owned SubscriptionId let owned_sub_id = sub_id.into_owned(); if event_sender - .send(RelayEvent::EndOfStoredEvents(owned_sub_id)) + .send(RelayEvent::EndOfStoredEvents( + owned_sub_id, + data_lane_arrival, + )) .await .is_err() { From a0d8b70bd49adc77fe7a0834b01001d86462046a Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 13:44:54 +0000 Subject: [PATCH 7/7] fix(sync): keep lifecycle bursts off the data lane The instrumented archive reproduced relay.ngit.dev's 50,591-ID hydration as 169 subscriptions. Ordinary duplicate processing averaged 0.02-0.03 ms with an empty queue, but the manager's bounded 100-item EOSE inbox filled while the actor held its lock to open the paced batch. The event processor then waited while forwarding EOSE, its 1,000-item EVENT queue filled, and later EOSE messages reached batch accounting 77-82 seconds late. Use non-blocking unbounded actor inboxes for EOSE and CLOSED notifications. Their producers remain bounded by the per-session subscription ledger, so this removes an accidental second capacity limit rather than allowing unbounded wire work. The ordered EVENT queue remains fixed at 1,000 as the peer-facing memory boundary. A regression test queues the observed 169-terminal burst without an actor receiver. This does not raise subscription concurrency, alter query pacing, reorder EVENT processing, or change terminal permit release. The separate rust-nostr terminal listener still closes and releases wire ownership immediately; these inboxes carry later serialized batch and live-coverage accounting. Validation: nix develop -c cargo test --lib (688 passed before the focused channel regression); nix develop -c cargo test --lib lifecycle_inbox_accepts_production_sized_terminal_burst (passed); git diff --check passed. Production validation will repeat the populated relay.ngit.dev reconciliation on this exact tip. --- docs/explanation/sync-scaling-constraints.md | 9 ++-- src/sync/mod.rs | 44 ++++++++++++++++---- 2 files changed, 42 insertions(+), 11 deletions(-) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index bdf455e..23f4c4c 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -407,9 +407,12 @@ The per-relay event processor retains its 1,000-message bounded data queue. Permit release does not depend on that queue draining: a separate listener on rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals and closes/releases transient ownership. The processor-facing listener still -delivers ordered EVENT and lifecycle work to the sync actor, but sustained page -traffic cannot hide terminal accounting behind EVENT backpressure. The data -queue remains finite as a memory-safety boundary for non-conforming peers. +delivers ordered EVENT and lifecycle work to the sync actor. EOSE and CLOSED +use non-blocking actor inboxes: their production is bounded by the session +subscription ledger, and a large paced historic batch can keep the actor busy +longer than a fixed lifecycle inbox could safely absorb. This prevents actor +backpressure from stopping the ordered EVENT processor, while the EVENT queue +remains finite as the memory-safety boundary for non-conforming peers. The levers, in the order we reach for them: diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 894d015..8bf0499 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1157,6 +1157,13 @@ struct SubscriptionClosedNotification { live_filter_count: Option, } +fn lifecycle_notification_channel() -> ( + tokio::sync::mpsc::UnboundedSender, + tokio::sync::mpsc::UnboundedReceiver, +) { + tokio::sync::mpsc::unbounded_channel() +} + fn is_rate_limit_message(message: &str) -> bool { let message = message.to_lowercase(); (message.contains("rate") && message.contains("limit")) @@ -1787,9 +1794,10 @@ pub struct SyncManager { /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, /// Channel for EOSE notifications (set during run) - eose_tx: Option>, + eose_tx: Option>, /// Serializes CLOSED recovery and pending-batch cleanup through the actor. - subscription_closed_tx: Option>, + subscription_closed_tx: + Option>, /// Returns connection outcomes to the sync actor for serialized state changes. connect_attempt_result_tx: Option>, /// Wakes the actor when a batch completion may unblock consolidation. @@ -3251,9 +3259,13 @@ impl SyncManager { let (disconnect_tx, mut disconnect_rx) = mpsc::channel::(100); // 3. Create EOSE channel for spawned tasks -> manager communication - let (eose_tx, mut eose_rx) = mpsc::channel::(100); + // Lifecycle notifications must never backpressure the ordered EVENT + // processor. Their production is already bounded by the subscription + // ledger, while the actor may legitimately stay busy opening a large, + // paced historic batch for longer than a fixed inbox can absorb. + let (eose_tx, mut eose_rx) = lifecycle_notification_channel::(); let (subscription_closed_tx, mut subscription_closed_rx) = - mpsc::channel::(100); + lifecycle_notification_channel::(); // 4. Connection workers never mutate manager state. Their unbounded // result channel cannot make a completed worker wait behind the actor. @@ -3778,8 +3790,7 @@ impl SyncManager { .send(EoseNotification { relay_url: relay_url_clone.clone(), sub_id, - }) - .await; + }); } RelayEvent::Notice(notice) => { if is_rate_limit_message(¬ice) { @@ -3854,8 +3865,7 @@ impl SyncManager { reason, generation: live_generation, live_filter_count, - }) - .await; + }); } RelayEvent::Shutdown => { tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); @@ -7165,6 +7175,24 @@ mod tests { ); } + #[test] + fn lifecycle_inbox_accepts_production_sized_terminal_burst() { + let (tx, mut rx) = lifecycle_notification_channel::(); + for index in 0..169 { + tx.send(EoseNotification { + relay_url: "wss://relay.example".to_string(), + sub_id: SubscriptionId::new(format!("hydration-{index}")), + }) + .expect("the actor-side receiver remains alive"); + } + + assert_eq!(rx.len(), 169); + for _ in 0..169 { + rx.try_recv().expect("every terminal remains queued"); + } + assert!(rx.try_recv().is_err()); + } + #[test] fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { let keys = Keys::generate();