diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 7d8b66f..23f4c4c 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 @@ -335,14 +339,20 @@ 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 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 @@ -360,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 @@ -394,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/filters.rs b/src/sync/filters.rs index d6de410..84a4ff7 100644 --- a/src/sync/filters.rs +++ b/src/sync/filters.rs @@ -13,6 +13,165 @@ 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, +} + +/// 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, + 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 +523,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), diff --git a/src/sync/mod.rs b/src/sync/mod.rs index c449236..8bf0499 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 { @@ -707,7 +794,7 @@ struct DescendantFrontier { #[derive(Debug)] struct DescendantSyncRotation { - frontier: DescendantFrontier, + coverage_fingerprint: Vec, filters: Vec, next_filter: usize, in_flight: Option, @@ -715,7 +802,7 @@ struct DescendantSyncRotation { #[derive(Debug)] struct DescendantFilterCursor { - filter: Filter, + filters: Vec, last_successful_until: Option, } @@ -728,14 +815,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 +831,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 +902,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; @@ -821,12 +938,12 @@ fn descendant_event_coordinate(event: &Event) -> Option { )) } +#[cfg(test)] 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 +951,99 @@ 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, + 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) +} + +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 { @@ -889,6 +1099,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, @@ -930,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")) @@ -940,8 +1174,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 +1318,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 { @@ -1558,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. @@ -2153,57 +2390,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 { @@ -2982,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. @@ -3240,7 +3521,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 +3545,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, @@ -3355,10 +3645,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 @@ -3387,6 +3680,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; } @@ -3467,8 +3766,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, @@ -3478,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) { @@ -3525,8 +3836,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, @@ -3555,8 +3865,7 @@ impl SyncManager { reason, generation: live_generation, live_filter_count, - }) - .await; + }); } RelayEvent::Shutdown => { tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); @@ -3785,26 +4094,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,34 +4123,47 @@ 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 { 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(); self.descendant_live_coverage.insert( relay_url.to_string(), DescendantLiveCoverage { - frontier, + coverage_fingerprint: desired_fingerprint.clone(), + fallback_filters, 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 +4171,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 +4205,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 +4240,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 +4296,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 +4321,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 +5992,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 +6059,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 +6095,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 +6185,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 +6256,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 +6268,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; @@ -6377,14 +6712,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, @@ -6624,49 +6964,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" ); @@ -6783,6 +7143,56 @@ 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 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(); @@ -6880,31 +7290,105 @@ 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() + })); + + 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] + 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] @@ -6930,7 +7414,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 +7996,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); @@ -7901,6 +8384,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..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 @@ -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 } @@ -1653,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, @@ -1661,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() { @@ -1671,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() { @@ -1816,6 +1827,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 +1855,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 +1867,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 +1926,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 +1999,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 +3736,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 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,