From d203f9de5683b48ee8483681ae138956bc62a604 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 20:43:08 +0000 Subject: [PATCH 1/4] refactor(sync): share queued historic REQ submission Historic REQ+EOSE currently groups filters, waits for one shared-ledger transient permit at a time, and retains each permit until the relay terminates the subscription. Descendant fallback needs exactly that lifecycle; copying the loop would make capacity and pagination fixes diverge. Extract the grouped subscription loop behind one SyncManager helper. The caller still creates the batch at the same point, uses the same filter grouping and request class, and receives the same subscription and pagination maps. No scheduling, priority, grouping, or capacity policy changes. This deliberately does not introduce descendant work or a new queue. Correctness assumes the helper remains awaited serially just as the inlined loop was. Validated with cargo check --lib; workspace rustfmt was not applied because the installed formatter would rewrite unrelated baseline files. --- src/sync/mod.rs | 112 ++++++++++++++++++++++++++++-------------------- 1 file changed, 66 insertions(+), 46 deletions(-) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 6a9319b..ff79418 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -5856,6 +5856,69 @@ impl SyncManager { }) } + /// Submit grouped REQ+EOSE filters through the connection's transient + /// queue and return the subscriptions that actually started. + /// + /// Callers submit only the current group while the connection waits for a + /// shared-ledger permit, rather than pre-queuing an unbounded historic + /// batch. The connection owns each permit until EOSE/CLOSED (or watchdog + /// recovery). + async fn subscribe_historic_filter_groups( + &self, + relay_url: &str, + batch_id: u64, + filters: &[Filter], + ) -> ( + HashSet, + HashMap, + ) { + let mut subscription_ids = HashSet::new(); + let mut pagination_state = HashMap::new(); + let max_filters = self + .connections + .get(relay_url) + .map(RelayConnection::max_filters_per_req) + .unwrap_or(MAX_FILTERS_PER_REQ); + + for (group_idx, filter_group) in group_filters_for_req_with_max(filters, max_filters) + .into_iter() + .enumerate() + { + tracing::debug!( + relay = %relay_url, + batch_id, + group_idx, + filter_count = filter_group.len(), + filters = ?filter_group, + "Subscribing to grouped filters in REQ+EOSE path" + ); + + if let Some(connection) = self.connections.get(relay_url) { + match connection + .subscribe_filters(filter_group.clone(), TransientRequestClass::HistoricPage) + .await + { + Ok(subscription_id) => { + subscription_ids.insert(subscription_id.clone()); + pagination_state + .insert(subscription_id, PaginationState::new(filter_group)); + } + Err(error) => { + tracing::error!( + relay = %relay_url, + batch_id, + group_idx, + error = %error, + "Failed to subscribe to filter in historic_sync" + ); + } + } + } + } + + (subscription_ids, pagination_state) + } + /// Sync historical events and track in PendingSyncIndex /// /// This method handles historical synchronization for a set of filters, @@ -6162,53 +6225,10 @@ impl SyncManager { "Starting historic_sync with REQ+EOSE" ); - // Subscribe to each filter and collect subscription IDs - let mut subscription_ids = HashSet::new(); - let mut pagination_state = HashMap::new(); - // Keep several OR filters under each relay-visible subscription. - let max_filters = self - .connections - .get(relay_url) - .map(RelayConnection::max_filters_per_req) - .unwrap_or(MAX_FILTERS_PER_REQ); - for (idx, filter_group) in - group_filters_for_req_with_max(&filters_with_since, max_filters) - .into_iter() - .enumerate() - { - tracing::debug!( - relay = %relay_url, - batch_id = batch_id, - group_idx = idx, - filter_count = filter_group.len(), - filters = ?filter_group, - "Subscribing to grouped filters in REQ+EOSE path" - ); - - if let Some(conn) = self.connections.get(relay_url) { - let grouped_filters = filter_group; - match conn - .subscribe_filters( - grouped_filters.clone(), - TransientRequestClass::HistoricPage, - ) - .await - { - Ok(sub_id) => { - subscription_ids.insert(sub_id.clone()); - pagination_state.insert(sub_id, PaginationState::new(grouped_filters)); - } - Err(e) => { - tracing::error!( - relay = %relay_url, - error = %e, - "Failed to subscribe to filter in historic_sync" - ); - } - } - } - } + let (subscription_ids, pagination_state) = self + .subscribe_historic_filter_groups(relay_url, batch_id, &filters_with_since) + .await; if subscription_ids.is_empty() && !filters_with_since.is_empty() { tracing::warn!( From 5f84b8b55db6fd8de90f5bd94dcc52003e1d8d30 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 20:46:05 +0000 Subject: [PATCH 2/4] feat(sync): recover direct-member descendant history Repository collaboration events may reference only their immediate parent. Core Layer 3 filters follow repository root IDs, so a reply to an already discovered reply can remain absent forever when it omits repository and root tags. Derive a stable, non-recursive frontier from locally stored direct root-thread members and rotate EOSE-closing filters through the existing historic queue, pagination, shared ledger, and request pacing. Run a complete initial pass, then rolling 24-hour passes; reconnects, daily reconciliation, and failed auxiliary batches restart completeness. Batch purpose keeps auxiliary failures out of core health state. This deliberately provides the safe fallback baseline only: it does not retain descendant live subscriptions, recursively traverse threads, add configuration, persist cursors, or shard connections. Correctness assumes ordinary Layer 3 sync eventually stores direct root-thread members. Validated with cargo check --lib, the descendant rotation unit test, and historic_sync_recovers_one_generation_of_parent_only_descendants. The scenario proves parent-only recovery while asserting a further recursive grandchild remains out of scope. --- CHANGELOG.md | 10 + docs/explanation/grasp-02-proactive-sync.md | 23 ++ src/sync/algorithms.rs | 3 + src/sync/filters.rs | 3 +- src/sync/mod.rs | 277 +++++++++++++++++++- tests/sync.rs | 1 + tests/sync/descendant_sync.rs | 82 ++++++ 7 files changed, 390 insertions(+), 9 deletions(-) create mode 100644 tests/sync/descendant_sync.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index e349e73..c13e671 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- Recover repository-event descendants which reference a direct thread member + but omit the repository and root-event tags. One low-priority, EOSE-closing + historic filter starts globally every five seconds while rotating across + relays: a complete initial pass followed by rolling 24-hour passes. The + frontier is deliberately non-recursive and uses the existing subscription + ledger, pagination, and request pacing rather than adding retained + subscriptions or connection sharding. + ## [2.1.2] - 2026-08-08 ngit-grasp 2.1.2 is a patch release improving repository-event sync under diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index f9a894a..1a25699 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -11,6 +11,9 @@ Features: - Fetches all repository announcements from connected relays to discover new repos listing our service - Discovers and dynamically connects to new relays listed by repository announcements we have accepted (with optional bootstrap relay to get started) - Fetches events tagging repositories we are interested in, as well as events tagging Issues, Patches and PRs of these repositories +- Recovers one additional generation of events that tag those direct thread + members but omit repository and root-event tags, using scheduled history + queries rather than retained subscriptions - Supports live sync and historic sync (tries NIP-77 negentropy but falls back to REQ+EOSE with 'until' based pagination) - Plays nicely with other relays - connection backoff and rate-limiting detection with cooldown - Does a full reconciliation daily @@ -785,6 +788,26 @@ live in - **Function**: `build_root_event_tag_filters(root_events, since)` - **Only for `SyncLevel::Full` repos** — purgatory announcements (`StateOnly`) skip this layer +### Scheduled Direct-Member Descendants + +Some collaboration events reference only their immediate parent. Once the +ordinary Layer 3 filters have discovered an event which directly tags a +repository root, each source relay is also queried for events whose `e`, `E`, +or `q` tags reference that direct member. These filters are not retained live: + +- one filter starts globally on the existing five-second maintenance cadence; +- every request uses the existing historic REQ+EOSE path, including the + per-connection subscription ledger, proactive pacing, and pagination; +- the first stable frontier receives a complete historic pass, then subsequent + rotations use a rolling 24-hour `since` window; and +- recovered descendants never enter the frontier, so this is deliberately one + additional generation rather than recursive thread traversal. + +The rotation restarts with a complete pass after a reconnect or daily +reconciliation. This makes an interrupted filter recoverable without adding a +durable cursor, another scheduler, retained subscription pressure, or +multi-connection sharding. + ### Combined Layer 2+3 (SyncLevel-Aware) The `build_sync_level_aware_filters()` function combines both layers, partitioning repos by `SyncLevel`: diff --git a/src/sync/algorithms.rs b/src/sync/algorithms.rs index 7c7aa9e..103803b 100644 --- a/src/sync/algorithms.rs +++ b/src/sync/algorithms.rs @@ -485,6 +485,7 @@ mod tests { "wss://relay1.com".to_string(), vec![super::super::PendingBatch { batch_id: 1, + purpose: super::super::PendingBatchPurpose::Core, items: super::super::PendingItems { repos: vec!["repo1".to_string()].into_iter().collect(), state_only_repos: HashSet::new(), @@ -592,6 +593,7 @@ mod tests { relay_url.to_string(), vec![super::super::PendingBatch { batch_id: 1, + purpose: super::super::PendingBatchPurpose::Core, items: PendingItems { repos: HashSet::new(), state_only_repos: HashSet::from(["repo1".to_string()]), @@ -681,6 +683,7 @@ mod tests { "wss://relay1.com".to_string(), vec![super::super::PendingBatch { batch_id: 1, + purpose: super::super::PendingBatchPurpose::Core, items: super::super::PendingItems { repos: vec!["repo1".to_string()].into_iter().collect(), state_only_repos: HashSet::new(), diff --git a/src/sync/filters.rs b/src/sync/filters.rs index 8df8705..f24cacb 100644 --- a/src/sync/filters.rs +++ b/src/sync/filters.rs @@ -216,7 +216,8 @@ pub fn tagged_one_of_our_root_event_filters( ); let mut filters = Vec::new(); - let event_ids: Vec = root_events.iter().map(|id| id.to_hex()).collect(); + let mut event_ids: Vec = root_events.iter().map(|id| id.to_hex()).collect(); + event_ids.sort_unstable(); for (chunk_idx, chunk) in chunk_values_by_bytes(&event_ids).into_iter().enumerate() { // Lowercase 'e' tag - standard event reference diff --git a/src/sync/mod.rs b/src/sync/mod.rs index ff79418..0965cca 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -65,6 +65,7 @@ const MAX_PURGATORY_FILTER_ACTIONS_PER_TICK: usize = 1; const MAX_PURGATORY_DEPENDENCY_IDS_PER_QUERY: usize = 100; const SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS: usize = 20; const SEMANTIC_FALLBACK_MAX_DELIVERED_PERCENT: usize = 10; +const DESCENDANT_RECENT_WINDOW_SECS: u64 = 24 * 60 * 60; fn should_use_semantic_fallback(requested_count: usize, received_count: usize) -> bool { requested_count >= SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS @@ -662,6 +663,9 @@ impl RelayPaginationSession { pub struct PendingBatch { /// Unique ID for this batch - for debugging/logging pub batch_id: u64, + /// Why this batch exists. Auxiliary descendant discovery must not mutate + /// core historic-sync completion or health state. + pub purpose: PendingBatchPurpose, /// The items this batch is syncing pub items: PendingItems, /// Subscription IDs that must ALL receive EOSE before confirming (for ReqEose) @@ -688,6 +692,49 @@ pub struct PendingBatch { pub failed: bool, } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PendingBatchPurpose { + Core, + Announcements, + Descendants, +} + +#[derive(Debug)] +struct DescendantSyncRotation { + frontier: HashSet, + filters: Vec, + next_filter: usize, + historic: bool, +} + +impl Default for DescendantSyncRotation { + fn default() -> Self { + Self { + frontier: HashSet::new(), + filters: Vec::new(), + next_filter: 0, + historic: true, + } + } +} + +impl DescendantSyncRotation { + fn refresh(&mut self, members: HashSet, now: Timestamp) { + if self.historic + && !self.filters.is_empty() + && self.next_filter >= self.filters.len() + && self.frontier == members + { + self.historic = false; + } + self.frontier = members.clone(); + let since = (!self.historic) + .then(|| Timestamp::from(now.as_secs().saturating_sub(DESCENDANT_RECENT_WINDOW_SECS))); + self.filters = filters::tagged_one_of_our_root_event_filters(&members, since); + self.next_filter = 0; + } +} + /// Items included in a pending batch #[derive(Debug, Clone, Default)] pub struct PendingItems { @@ -1172,6 +1219,10 @@ async fn run_daily_timer( /// during negentropy reconciliation but failed to deliver on exact-ID fetches /// (see [`missing_events`]). Recovery attempts are backed off per relay, so /// the tick itself stays cheap when nothing is due. +/// +/// Finally, one relay at a time receives one low-priority repository-descendant +/// history query. Reusing this five-second maintenance cadence bounds the new +/// global query-start rate without another scheduler or configuration surface. async fn run_purgatory_announcement_sync( sync_manager: Arc>, mut shutdown_rx: broadcast::Receiver<()>, @@ -1187,6 +1238,7 @@ async fn run_purgatory_announcement_sync( let mut manager = sync_manager.lock().await; manager.sync_purgatory_announcements_to_index().await; manager.tick_missing_event_recovery().await; + manager.tick_descendant_sync().await; } _ = shutdown_rx.recv() => { tracing::debug!("Purgatory announcement sync timer received shutdown signal"); @@ -1415,6 +1467,11 @@ pub struct SyncManager { /// Relays whose complete persistent filter set exceeds a learned remote /// byte cap, mapped to their next bounded catch-up deadline. byte_limited_live_relays: HashMap, + /// Per-relay snapshots for low-priority, EOSE-closing descendant queries. + descendant_sync_rotations: HashMap, + /// Round-robin cursor so the maintenance loop starts at most one extra + /// query globally per tick. + descendant_relay_cursor: usize, /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, /// Channel for EOSE notifications (set during run) @@ -1515,6 +1572,8 @@ impl SyncManager { connect_attempt_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECT_ATTEMPTS)), deferred_consolidations: DeferredConsolidations::default(), byte_limited_live_relays: HashMap::new(), + descendant_sync_rotations: HashMap::new(), + descendant_relay_cursor: 0, disconnect_tx: None, eose_tx: None, subscription_closed_tx: None, @@ -2507,14 +2566,20 @@ impl SyncManager { /// # Arguments /// * `relay_url` - The relay URL the batch belongs to /// * `batch` - The completed batch to confirm - async fn confirm_batch(&self, relay_url: &str, batch: PendingBatch) { + async fn confirm_batch(&mut self, relay_url: &str, batch: PendingBatch) { let batch_id = batch.batch_id; let full_repos_count = batch.items.repos.len(); let state_only_repos_count = batch.items.state_only_repos.len(); let events_count = batch.items.root_events.len(); let sync_method = batch.sync_method; - let is_generic_filter = - full_repos_count == 0 && state_only_repos_count == 0 && events_count == 0; + let is_generic_filter = batch.purpose == PendingBatchPurpose::Announcements; + + if batch.failed && batch.purpose == PendingBatchPurpose::Descendants { + // Do not let an interrupted page turn the initial complete pass + // into a recent-only rotation. Starting this tiny auxiliary state + // afresh is simpler and safer than persisting per-filter cursors. + self.descendant_sync_rotations.remove(relay_url); + } let mut relay_index = self.relay_sync_index.write().await; @@ -2567,7 +2632,7 @@ impl SyncManager { } // Track if this batch failed (for ConnectedDegraded transition) - if batch.failed { + if batch.failed && batch.purpose != PendingBatchPurpose::Descendants { state.historic_sync_had_failures = true; // Failures unrelated to a relay's pending missing-event // recovery mean full recovery must not restore its health. @@ -2585,6 +2650,7 @@ impl SyncManager { tracing::info!( relay = %relay_url, batch_id = batch_id, + purpose = ?batch.purpose, sync_method = ?sync_method, full_repos_confirmed = full_repos_count, state_only_repos_confirmed = state_only_repos_count, @@ -2735,6 +2801,7 @@ impl SyncManager { async fn daily_sync(&mut self, relay_url: &str) { tracing::info!(relay = %relay_url, "Starting daily sync"); self.cancel_deferred_consolidation(relay_url, "daily sync reset"); + self.descendant_sync_rotations.remove(relay_url); // Get connection let connection = match self.connections.get(relay_url) { @@ -3528,16 +3595,151 @@ impl SyncManager { // Use historic_sync with empty PendingItems for generic filters // Generic filters (announcements) don't have associated repos or root_events let items = PendingItems::default(); - let _batch_id = self.historic_sync(relay_url, filters, items, since).await; + let _batch_id = self + .historic_sync_with_options( + relay_url, + filters, + items, + since, + PendingBatchPurpose::Announcements, + false, + ) + .await; } async fn sync_generic_history(&mut self, relay_url: &str, since: Option) { let filters = vec![filters::build_announcement_filter(None)]; let _batch_id = self - .historic_sync(relay_url, filters, PendingItems::default(), since) + .historic_sync_with_options( + relay_url, + filters, + PendingItems::default(), + since, + PendingBatchPurpose::Announcements, + false, + ) .await; } + /// Find the existing first-generation members of repository root threads. + /// + /// Only events that directly reference a root are admitted to this + /// frontier. Events fetched by the descendant rotation are deliberately + /// not fed back into it, keeping this feature non-recursive. + async fn direct_thread_members(&self, root_events: &HashSet) -> HashSet { + let mut members = HashSet::new(); + for filter in filters::tagged_one_of_our_root_event_filters(root_events, None) { + match self.database.query(filter).await { + Ok(events) => { + members.extend( + events + .iter() + .map(|event| event.id) + .filter(|event_id| !root_events.contains(event_id)), + ); + } + Err(error) => { + tracing::warn!( + error = %error, + root_event_count = root_events.len(), + "Failed to derive direct repository thread members" + ); + return HashSet::new(); + } + } + } + members + } + + /// Start at most one low-priority descendant query globally per tick. + /// + /// Each relay first receives a complete historic pass over a stable + /// snapshot of direct thread members. Later passes use a rolling 24-hour + /// window. Every request uses the ordinary historic REQ+EOSE pagination + /// path, so it draws from the existing ledger and closes at EOSE. + async fn tick_descendant_sync(&mut self) { + let targets = { + let index = self.repo_sync_index.read().await; + algorithms::derive_relay_targets(&index) + }; + let states = self.relay_sync_index.read().await; + let mut relay_urls: Vec = targets + .iter() + .filter(|(relay_url, needs)| { + !needs.root_events.is_empty() + && states.get(*relay_url).is_some_and(|state| { + matches!( + state.connection_status, + ConnectionStatus::Connected + | ConnectionStatus::ConnectedHistoricSyncFailures + ) + }) + && !self.health_tracker.is_subscription_paused(relay_url) + }) + .map(|(relay_url, _)| relay_url.clone()) + .collect(); + drop(states); + relay_urls.sort_unstable(); + if relay_urls.is_empty() { + return; + } + + let relay_url = relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); + self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); + let root_events = targets[&relay_url].root_events.clone(); + + let needs_snapshot = self + .descendant_sync_rotations + .get(&relay_url) + .is_none_or(|rotation| rotation.next_filter >= rotation.filters.len()); + if needs_snapshot { + let members = self.direct_thread_members(&root_events).await; + if members.is_empty() { + return; + } + + let rotation = self + .descendant_sync_rotations + .entry(relay_url.clone()) + .or_default(); + rotation.refresh(members, Timestamp::now()); + } + + let (filter, historic, filter_index, filter_count) = { + let rotation = self.descendant_sync_rotations.get(&relay_url).unwrap(); + ( + rotation.filters[rotation.next_filter].clone(), + rotation.historic, + rotation.next_filter, + rotation.filters.len(), + ) + }; + if self + .historic_sync_with_options( + &relay_url, + vec![filter], + PendingItems::default(), + None, + PendingBatchPurpose::Descendants, + true, + ) + .await + .is_some() + { + self.descendant_sync_rotations + .get_mut(&relay_url) + .unwrap() + .next_filter += 1; + tracing::debug!( + relay = %relay_url, + historic, + filter_index, + filter_count, + "Started scheduled repository descendant query" + ); + } + } + /// Build the complete persistent coverage for one connection. L1 must be /// admitted in the same transaction as rebuilt L2/L3 so a tight relay /// budget cannot leave a successful generic REQ hiding partial repo @@ -4538,6 +4740,10 @@ impl SyncManager { async fn handle_disconnect(&mut self, relay_url: &str) { // Learned page sizes and NIP-11 hints belong to the ended WebSocket session. self.pagination_sessions.remove(relay_url); + // A connection can end after a descendant REQ was sent but before its + // EOSE. Restart the bounded full rotation so that page is not treated + // as complete on the next session. + self.descendant_sync_rotations.remove(relay_url); // Check if this was an intentional disconnect (Disconnecting status) let was_intentional = { @@ -5945,6 +6151,26 @@ impl SyncManager { filters: Vec, items: PendingItems, since: Option, + ) -> Option { + self.historic_sync_with_options( + relay_url, + filters, + items, + since, + PendingBatchPurpose::Core, + false, + ) + .await + } + + async fn historic_sync_with_options( + &mut self, + relay_url: &str, + filters: Vec, + items: PendingItems, + since: Option, + purpose: PendingBatchPurpose, + force_req_eose: bool, ) -> Option { // DEBUG TRACING: Log all filters being passed to historic_sync tracing::debug!( @@ -5985,8 +6211,9 @@ impl SyncManager { }; // Check if we should use negentropy - let use_negentropy = - !self.config.sync_disable_negentropy && connection.supports_negentropy().await; + let use_negentropy = !force_req_eose + && !self.config.sync_disable_negentropy + && connection.supports_negentropy().await; // Generate batch ID let batch_id = self.next_batch_id(); @@ -6008,6 +6235,7 @@ impl SyncManager { // Create PendingBatch for negentropy (empty outstanding_subs and pagination_state) let batch = PendingBatch { batch_id, + purpose, items: items.clone(), outstanding_subs: HashSet::new(), sync_method: SyncMethod::Negentropy, @@ -6241,6 +6469,7 @@ impl SyncManager { // Create PendingBatch for REQ+EOSE let batch = PendingBatch { batch_id, + purpose, items, outstanding_subs: subscription_ids, sync_method: SyncMethod::ReqEose, @@ -6349,6 +6578,33 @@ mod tests { assert!(!should_use_semantic_fallback(19, 0)); } + #[test] + fn descendant_rotation_becomes_recent_only_after_a_stable_historic_pass() { + let member = EventId::from_byte_array([7; 32]); + let members = HashSet::from([member]); + let now = Timestamp::from_secs(200_000); + let mut rotation = DescendantSyncRotation::default(); + + rotation.refresh(members.clone(), now); + assert!(rotation.historic); + assert!(rotation + .filters + .iter() + .all(|filter| !serde_json::to_value(filter) + .unwrap() + .as_object() + .unwrap() + .contains_key("since"))); + + rotation.next_filter = rotation.filters.len(); + rotation.refresh(members, now); + assert!(!rotation.historic); + assert!(rotation.filters.iter().all(|filter| { + serde_json::to_value(filter).unwrap()["since"] + == serde_json::json!(now.as_secs() - DESCENDANT_RECENT_WINDOW_SECS) + })); + } + #[test] fn group_filters_for_req_respects_count_and_byte_budgets() { // Many small filters group by the count cap. @@ -6863,6 +7119,7 @@ mod tests { let still_pending = SubscriptionId::new("still-pending"); let make_batch = |batch_id, outstanding_subs| PendingBatch { batch_id, + purpose: PendingBatchPurpose::Core, items: PendingItems::default(), outstanding_subs, sync_method: SyncMethod::ReqEose, @@ -6965,6 +7222,7 @@ mod tests { let accepted = SubscriptionId::new("accepted"); let make_batch = |batch_id, subscription_id| PendingBatch { batch_id, + purpose: PendingBatchPurpose::Core, items: PendingItems::default(), outstanding_subs: HashSet::from([subscription_id]), sync_method: SyncMethod::ReqEose, @@ -7000,6 +7258,7 @@ mod tests { let completed_sub = SubscriptionId::new("completed-page"); let mut batch = PendingBatch { batch_id: 73, + purpose: PendingBatchPurpose::Core, items: PendingItems::default(), outstanding_subs: HashSet::new(), sync_method: SyncMethod::ReqEose, @@ -7318,6 +7577,7 @@ mod tests { // Test that PendingBatch properly tracks negentropy-specific fields let batch = PendingBatch { batch_id: 1, + purpose: PendingBatchPurpose::Core, items: PendingItems::default(), outstanding_subs: HashSet::new(), sync_method: SyncMethod::Negentropy, @@ -7341,6 +7601,7 @@ mod tests { // Test that REQ+EOSE batches don't use negentropy fields let batch = PendingBatch { batch_id: 1, + purpose: PendingBatchPurpose::Core, items: PendingItems::default(), outstanding_subs: HashSet::new(), sync_method: SyncMethod::ReqEose, diff --git a/tests/sync.rs b/tests/sync.rs index 7949928..8e7a7b9 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -33,6 +33,7 @@ mod common; mod sync { pub mod adaptive_pagination; pub mod catchup; + pub mod descendant_sync; pub mod discovery; pub mod historic_recovery; pub mod historic_sync; diff --git a/tests/sync/descendant_sync.rs b/tests/sync/descendant_sync.rs new file mode 100644 index 0000000..7df362c --- /dev/null +++ b/tests/sync/descendant_sync.rs @@ -0,0 +1,82 @@ +//! Scheduled repository-descendant synchronization scenarios. + +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::{sync_helpers::*, TestRelay}; + +/// Events which reference a direct repository-thread member, but not the +/// repository or its root event, are recovered by the scheduled historic +/// pass. The recovered events do not recursively extend the frontier. +#[tokio::test] +async fn historic_sync_recovers_one_generation_of_parent_only_descendants() { + let source = TestRelay::start().await; + let keys = Keys::generate(); + let repo_id = "scheduled-descendant-history"; + + // Seed the source before the syncing relay exists so this exercises the + // complete historic rotation, not ordinary live coverage. + let source_domains = [source.domain()]; + let source_refs = source_domains + .iter() + .map(String::as_str) + .collect::>(); + let (_announcement, _source_git) = + setup_announcement_on_relay(&source, &keys, &source_refs, repo_id).await; + let source_client = TestClient::new(source.url(), keys.clone()) + .await + .expect("connect to source relay"); + let issue = + build_layer2_issue_event(&keys, &repo_coord(&keys, repo_id), "Repository thread root") + .expect("build issue"); + let direct_reply = build_layer3_reply_with_e_tag(&keys, &issue.id, "Direct reply") + .expect("build direct reply"); + let parent_only = build_layer3_reply_with_e_tag( + &keys, + &direct_reply.id, + "Reply visible only through the scheduled descendant query", + ) + .expect("build parent-only descendant"); + let recursive = build_layer3_reply_with_e_tag( + &keys, + &parent_only.id, + "Deliberately out-of-scope recursive descendant", + ) + .expect("build recursive descendant"); + for event in [&issue, &direct_reply, &parent_only, &recursive] { + source_client + .send_event(event) + .await + .expect("seed source event"); + } + + let syncing = TestRelay::start_with_sync(None).await; + let domains = [source.domain(), syncing.domain()]; + let domain_refs = domains.iter().map(String::as_str).collect::>(); + let (_target_announcement, _target_git) = + setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await; + + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(parent_only.id), + Duration::from_secs(20), + ) + .await, + "scheduled historic rotation should recover a parent-only descendant" + ); + assert!( + !wait_for_event_on_relay( + syncing.url(), + Filter::new().id(recursive.id), + Duration::from_secs(7), + ) + .await, + "recovered descendants must not recursively extend the frontier" + ); + + source_client.disconnect().await; + syncing.stop().await; + source.stop().await; +} From 4dced403167346c531ef87e0bb33b6861246aa60 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 20:55:12 +0000 Subject: [PATCH 3/4] feat(sync): retain descendants when connection capacity permits Scheduled history is a safe completeness fallback but delays new parent-only collaboration events even on relays with ample subscription capacity. Descendant coverage should be immediate when it does not compromise the existing core product or recovery headroom. Admit the complete grouped descendant filter set only when the session ledger can still preserve its two control-plane slots and one usable transient slot, and when the learned retained-byte limit still leaves one maximum-sized transient request. Track the auxiliary subscription IDs separately so frontier changes, CLOSED, daily reset, and core consolidation can retire only descendant coverage. A failed CLOSE retains its ledger ownership and aborts core replacement rather than allowing local accounting to run ahead of the relay. Core filter grouping and rollback remain unchanged; auxiliary filters are intentionally not packed into partially full core groups. Constrained sessions continue using the independently complete historic rotation from the preceding commit. Recursive traversal, connection sharding, and minimum-churn tail compaction remain excluded. Validated with cargo check --lib, the descendant rotation unit test, and auxiliary_live_admission_preserves_one_historic_slot using a nostream-shaped advertised budget of 10. --- CHANGELOG.md | 11 +- docs/explanation/grasp-02-proactive-sync.md | 26 +-- src/sync/mod.rs | 169 +++++++++++++++++++- src/sync/relay_connection.rs | 124 ++++++++++++++ 4 files changed, 311 insertions(+), 19 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c13e671..ef99923 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,12 +10,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added - Recover repository-event descendants which reference a direct thread member - but omit the repository and root-event tags. One low-priority, EOSE-closing - historic filter starts globally every five seconds while rotating across - relays: a complete initial pass followed by rolling 24-hour passes. The - frontier is deliberately non-recursive and uses the existing subscription - ledger, pagination, and request pacing rather than adding retained - subscriptions or connection sharding. + but omit the repository and root-event tags. When the per-connection ledger + can retain complete descendant coverage while preserving control and + historic capacity, those filters stay live; otherwise bounded EOSE-closing + history provides eventual coverage. The frontier is deliberately + non-recursive and connection sharding remains unnecessary. ## [2.1.2] - 2026-08-08 diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 1a25699..7260a8e 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -788,24 +788,30 @@ live in - **Function**: `build_root_event_tag_filters(root_events, since)` - **Only for `SyncLevel::Full` repos** — purgatory announcements (`StateOnly`) skip this layer -### Scheduled Direct-Member Descendants +### Direct-Member Descendants Some collaboration events reference only their immediate parent. Once the ordinary Layer 3 filters have discovered an event which directly tags a repository root, each source relay is also queried for events whose `e`, `E`, -or `q` tags reference that direct member. These filters are not retained live: +or `q` tags reference that direct member. -- one filter starts globally on the existing five-second maintenance cadence; -- every request uses the existing historic REQ+EOSE path, including the - per-connection subscription ledger, proactive pacing, and pagination; -- the first stable frontier receives a complete historic pass, then subsequent - rotations use a rolling 24-hour `since` window; and +- complete descendant filters are retained live when they fit after core live + coverage while preserving the two control-plane slots and at least one + transient historic slot; +- auxiliary subscriptions are tracked separately and retired before core + consolidation or restoration, so they never enter the core rollback set; +- when the complete live set does not fit, one EOSE-closing filter starts on + the existing maintenance cadence through the ordinary historic queue, + pagination, shared ledger, and request pacing; +- the fallback receives a complete initial pass, then rolling 24-hour passes; + and - recovered descendants never enter the frontier, so this is deliberately one additional generation rather than recursive thread traversal. -The rotation restarts with a complete pass after a reconnect or daily -reconciliation. This makes an interrupted filter recoverable without adding a -durable cursor, another scheduler, retained subscription pressure, or +An unexpected auxiliary CLOSED retires the remaining descendant subscriptions +without rebuilding core coverage and falls back to history. Reconnect and daily +reconciliation reconstruct the mode from current session capacity. This keeps +the feature complete without durable cursors, another capacity ledger, or multi-connection sharding. ### Combined Layer 2+3 (SyncLevel-Aware) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 0965cca..aec2cef 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -707,6 +707,12 @@ struct DescendantSyncRotation { historic: bool, } +#[derive(Debug)] +struct DescendantLiveCoverage { + frontier: HashSet, + subscription_ids: Vec, +} + impl Default for DescendantSyncRotation { fn default() -> Self { Self { @@ -1469,6 +1475,9 @@ pub struct SyncManager { byte_limited_live_relays: HashMap, /// Per-relay snapshots for low-priority, EOSE-closing descendant queries. descendant_sync_rotations: HashMap, + /// Auxiliary persistent descendant coverage, kept separate from the core + /// desired set so it can be retired without replacing healthy core REQs. + descendant_live_coverage: HashMap, /// Round-robin cursor so the maintenance loop starts at most one extra /// query globally per tick. descendant_relay_cursor: usize, @@ -1573,6 +1582,7 @@ impl SyncManager { deferred_consolidations: DeferredConsolidations::default(), byte_limited_live_relays: HashMap::new(), descendant_sync_rotations: HashMap::new(), + descendant_live_coverage: HashMap::new(), descendant_relay_cursor: 0, disconnect_tx: None, eose_tx: None, @@ -2801,6 +2811,9 @@ impl SyncManager { async fn daily_sync(&mut self, relay_url: &str) { tracing::info!(relay = %relay_url, "Starting daily sync"); self.cancel_deferred_consolidation(relay_url, "daily sync reset"); + let _ = self + .close_descendant_live_coverage(relay_url, "daily sync reset") + .await; self.descendant_sync_rotations.remove(relay_url); // Get connection @@ -3651,6 +3664,39 @@ impl SyncManager { members } + async fn close_descendant_live_coverage( + &mut self, + relay_url: &str, + reason: &'static str, + ) -> bool { + let Some(coverage) = self.descendant_live_coverage.remove(relay_url) else { + return true; + }; + if let Some(connection) = self.connections.get(relay_url) { + if let Err(error) = connection + .close_live_subscriptions(&coverage.subscription_ids) + .await + { + tracing::warn!( + relay = %relay_url, + %error, + reason, + "Could not retire auxiliary descendant live coverage" + ); + self.descendant_live_coverage + .insert(relay_url.to_string(), coverage); + return false; + } + } + tracing::info!( + relay = %relay_url, + subscription_count = coverage.subscription_ids.len(), + reason, + "Retired auxiliary descendant live coverage" + ); + true + } + /// Start at most one low-priority descendant query globally per tick. /// /// Each relay first receives a complete historic pass over a stable @@ -3688,16 +3734,85 @@ impl SyncManager { self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); let root_events = targets[&relay_url].root_events.clone(); + let mut refreshed_members = None; + if let Some(existing) = self.descendant_live_coverage.get(&relay_url) { + let members = self.direct_thread_members(&root_events).await; + if members == existing.frontier { + return; + } + refreshed_members = Some(members); + if !self + .close_descendant_live_coverage(&relay_url, "frontier changed") + .await + { + return; + } + } + let needs_snapshot = self .descendant_sync_rotations .get(&relay_url) .is_none_or(|rotation| rotation.next_filter >= rotation.filters.len()); if needs_snapshot { - let members = self.direct_thread_members(&root_events).await; + let members = match refreshed_members { + Some(members) => members, + None => self.direct_thread_members(&root_events).await, + }; if members.is_empty() { return; } + let now = Timestamp::now(); + let live_since = Timestamp::from( + now.as_secs() + .saturating_sub(DESCENDANT_RECENT_WINDOW_SECS.min(15 * 60)), + ); + let live_filters = filters::tagged_one_of_our_root_event_filters( + &members, + Some(live_since), + ); + if let Some(connection) = self.connections.get(&relay_url).cloned() { + let groups = live_filter_groups(&live_filters, connection.max_filters_per_req()); + if connection.can_admit_auxiliary_live_groups(&groups) { + match connection + .subscribe_auxiliary_live_filter_groups(groups) + .await + { + Ok(subscription_ids) => { + let filter_count = live_filters.len(); + self.descendant_live_coverage.insert( + relay_url.clone(), + DescendantLiveCoverage { + frontier: members, + subscription_ids: subscription_ids.clone(), + }, + ); + self.descendant_sync_rotations.remove(&relay_url); + tracing::info!( + relay = %relay_url, + filter_count, + subscription_count = subscription_ids.len(), + "Installed auxiliary descendant live coverage" + ); + return; + } + Err(error) => { + tracing::warn!( + relay = %relay_url, + %error, + "Descendant live admission failed; using historic rotation" + ); + } + } + } else { + tracing::info!( + relay = %relay_url, + filter_count = live_filters.len(), + "Descendant live coverage does not fit while preserving historic capacity" + ); + } + } + let rotation = self .descendant_sync_rotations .entry(relay_url.clone()) @@ -4744,6 +4859,7 @@ impl SyncManager { // EOSE. Restart the bounded full rotation so that page is not treated // as complete on the next session. self.descendant_sync_rotations.remove(relay_url); + self.descendant_live_coverage.remove(relay_url); // Check if this was an intentional disconnect (Disconnecting status) let was_intentional = { @@ -5456,6 +5572,16 @@ impl SyncManager { "Starting consolidation" ); + // Descendant coverage is auxiliary and independently reconstructible. + // Retire it before capturing the core rollback set so it can neither + // displace newly required core coverage nor become part of that set. + if !self + .close_descendant_live_coverage(relay_url, "core consolidation") + .await + { + return false; + } + let now = Timestamp::now(); let since = Timestamp::from(now.as_secs().saturating_sub(QUICK_RECONNECT_WINDOW_SECS)); let complete_live = self.complete_live_filters(relay_url, Some(since)).await; @@ -5529,6 +5655,33 @@ impl SyncManager { live_generation: Option, live_filter_count: Option, ) { + let is_descendant_live = self + .descendant_live_coverage + .get(relay_url) + .is_some_and(|coverage| coverage.subscription_ids.contains(&subscription_id)); + if is_descendant_live { + let coverage = self + .descendant_live_coverage + .remove(relay_url) + .expect("descendant subscription belonged to tracked coverage"); + if let Some(connection) = self.connections.get(relay_url) { + let _ = connection + .close_live_subscriptions(&coverage.subscription_ids) + .await; + } + self.descendant_sync_rotations + .entry(relay_url.to_string()) + .or_default() + .refresh(coverage.frontier, Timestamp::now()); + tracing::warn!( + relay = %relay_url, + sub_id = %subscription_id, + reason, + "Descendant live subscription closed; using historic fallback" + ); + return; + } + if let Some(limit) = subscription_state_byte_limit(reason) { tracing::warn!( relay = %relay_url, @@ -5654,13 +5807,23 @@ impl SyncManager { } async fn restore_live_coverage_after_closed(&mut self, relay_url: &str, generation: u64) { - let Some(connection) = self.connections.get(relay_url) else { + let Some(current_generation) = self + .connections + .get(relay_url) + .map(RelayConnection::current_subscription_generation) + else { return; }; - if connection.current_subscription_generation() != generation { + if current_generation != generation { tracing::debug!(relay = %relay_url, generation, "Ignoring stale live CLOSED from retired session"); return; } + if !self + .close_descendant_live_coverage(relay_url, "core live restoration") + .await + { + return; + } let now = Timestamp::now(); let since = Timestamp::from(now.as_secs().saturating_sub(QUICK_RECONNECT_WINDOW_SECS)); let filters = self.complete_live_filters(relay_url, Some(since)).await; diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index a769db6..4993ec2 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -1170,6 +1170,100 @@ impl RelayConnection { <= usable } + /// Whether auxiliary persistent coverage fits while preserving one usable + /// ledger slot for transient history. The two control-plane slots have + /// already been removed from `usable` when the session ledger was built. + pub fn can_admit_auxiliary_live_groups(&self, filter_groups: &[Vec]) -> bool { + let usable = self + .subscription_usable_slots + .load(std::sync::atomic::Ordering::Relaxed); + let held = self + .live_req_permits_held + .lock() + .expect("live permit map poisoned"); + if held + .len() + .saturating_add(filter_groups.len()) + .saturating_add(1) + > usable + { + return false; + } + + let remote_limit = self.remote_subscription_byte_limit(); + if remote_limit.is_none() { + return true; + } + let existing_bytes: usize = held + .values() + .map(|held| { + ClientMessage::req(SubscriptionId::generate(), held.filters.clone()) + .as_json() + .len() + }) + .sum(); + let auxiliary_bytes: usize = filter_groups + .iter() + .map(|filters| { + ClientMessage::req(SubscriptionId::generate(), filters.clone()) + .as_json() + .len() + }) + .sum(); + existing_bytes + .checked_add(auxiliary_bytes) + .and_then(|used| used.checked_add(super::SUBSCRIPTION_BYTE_RESERVED_MARGIN)) + .is_some_and(|used| used <= remote_limit.unwrap()) + } + + /// Open separately tracked auxiliary live groups. The caller must first + /// use [`Self::can_admit_auxiliary_live_groups`] so one transient slot is + /// preserved; the shared ledger remains the final admission authority. + pub async fn subscribe_auxiliary_live_filter_groups( + &self, + filter_groups: Vec>, + ) -> Result, String> { + if !self.can_admit_auxiliary_live_groups(&filter_groups) { + return Err(format!( + "Auxiliary live coverage would consume historic capacity for {}", + self.url + )); + } + self.subscribe_live_filter_groups(filter_groups).await + } + + /// Close a caller-owned subset of persistent subscriptions and return + /// their ledger slots only after CLOSE has been accepted by the SDK. + pub async fn close_live_subscriptions( + &self, + subscription_ids: &[SubscriptionId], + ) -> Result<(), String> { + let mut failures = Vec::new(); + for subscription_id in subscription_ids { + if let Err(error) = self.client.unsubscribe(subscription_id).await { + tracing::debug!( + relay = %self.url, + sub_id = %subscription_id, + %error, + "Failed to close caller-owned live subscription" + ); + failures.push(format!("{subscription_id}: {error}")); + continue; + } + self.release_live_req_permit(subscription_id); + } + if failures.is_empty() { + Ok(()) + } else { + Err(format!( + "Failed to close {} live subscriptions on {}: {}", + failures.len(), + self.url, + failures.join("; ") + )) + } + } + pub fn complete_live_set_fits(&self, slots: usize) -> bool { slots <= self @@ -2902,6 +2996,36 @@ mod tests { assert_eq!(connection.subscription_budget().available_permits(), 8); } + #[tokio::test] + async fn auxiliary_live_admission_preserves_one_historic_slot() { + let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate()); + connection.reset_subscription_budget(Some(10)); + + for index in 0..6 { + let permit = connection.acquire_subscription_slots(1).await.unwrap(); + connection.live_req_permits_held.lock().unwrap().insert( + SubscriptionId::new(format!("core-{index}")), + HeldLiveSubscription { + generation: permit.generation, + _ledger_slot: permit.permit, + filters: vec![Filter::new().kind(Kind::Custom(20_100 + index))], + }, + ); + } + + let one_group = vec![vec![Filter::new().kind(Kind::Custom(20_200))]]; + assert!(connection.can_admit_auxiliary_live_groups(&one_group)); + + let two_groups = vec![ + vec![Filter::new().kind(Kind::Custom(20_201))], + vec![Filter::new().kind(Kind::Custom(20_202))], + ]; + assert!( + !connection.can_admit_auxiliary_live_groups(&two_groups), + "two auxiliary groups would consume the final historic slot" + ); + } + #[tokio::test] async fn failed_live_group_rolls_back_opened_groups_and_full_capacity() { let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate()); From 9f3c943943101b3d9e18ad7c6408a5d460bc0dd0 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 21:08:22 +0000 Subject: [PATCH 4/4] feat(sync): rotate constrained descendants per relay A global five-second cursor makes descendant recovery latency scale with every relay, while advancing at request admission can skip a filter whose queued or active subscription later fails. Permanent coverage also needs a historic baseline because limit:0 protects only events arriving after admission. Give each constrained relay one EOSE-confirmed fallback cursor. Every five-second maintenance pass may start one filter per constrained relay, but never a second while that relay has a queued or active descendant batch. Successful terminal completion advances the cursor; failure retains the same filter. Later requests begin at the last successful upper bound minus fifteen minutes. All work uses the existing transient permit queue without request-class priority. Pair every permanent live admission with complete descendant history, and augment ordinary core historic batches with direct-member filters already known locally. A test-only NIP-11 subscription limit forces fallback end to end; a normal source proves retained live delivery. Large frontiers are verified to split across several grouped subscriptions. Cursors remain session-local, descendants remain non-recursive, core and auxiliary filter packing remain separate, and multi-connection sharding is excluded. Correctness assumes failed or CLOSED batches are removed through the existing terminal paths, which now release the matching in-flight cursor without advancing it. Validated locally with cargo check --lib, cargo test --lib (673 passed), both descendant integration scenarios, and the large-frontier grouping test. Archive production at the pre-amend behavior-identical tip installed permanent coverage, split six filters across two subscriptions, and rotated constrained Shakespeare filters sequentially with successful terminal results. The amendment only corrects the terminal log label so it also describes the historic baseline used by permanent coverage. --- docs/explanation/grasp-02-proactive-sync.md | 21 +- docs/explanation/sync-scaling-constraints.md | 29 +- src/sync/mod.rs | 481 ++++++++++++------- tests/common/relay.rs | 19 + tests/sync/descendant_sync.rs | 86 +++- 5 files changed, 445 insertions(+), 191 deletions(-) diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 7260a8e..d146d7d 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -801,18 +801,23 @@ or `q` tags reference that direct member. - auxiliary subscriptions are tracked separately and retired before core consolidation or restoration, so they never enter the core rollback set; - when the complete live set does not fit, one EOSE-closing filter starts on - the existing maintenance cadence through the ordinary historic queue, - pagination, shared ledger, and request pacing; -- the fallback receives a complete initial pass, then rolling 24-hour passes; - and + each constrained relay on the existing five-second maintenance cadence, + through the ordinary historic queue, pagination, shared ledger, and request + pacing; +- ordinary historic batches include the currently known direct-member filters; +- fallback filters keep an in-memory cursor, advance it only after successful + EOSE, and query from the preceding successful upper bound with 15 minutes of + overlap; and - recovered descendants never enter the frontier, so this is deliberately one additional generation rather than recursive thread traversal. An unexpected auxiliary CLOSED retires the remaining descendant subscriptions -without rebuilding core coverage and falls back to history. Reconnect and daily -reconciliation reconstruct the mode from current session capacity. This keeps -the feature complete without durable cursors, another capacity ledger, or -multi-connection sharding. +without rebuilding core coverage and falls back to history. A filter already +queued or active blocks another fallback filter for that relay; failure leaves +the same filter and cursor at the head. Reconnect and daily reconciliation +reconstruct the mode from current session capacity. This keeps the feature +complete without durable cursors, another capacity ledger, request-class +priority, or multi-connection sharding. ### Combined Layer 2+3 (SyncLevel-Aware) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index dc84ca8..1f5d4c6 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -301,24 +301,31 @@ Derived from the tightest commonly observed values; all sizing below assumes: ## Our Approach: A Per-Connection Budget Ledger -Each relay connection owns one implemented budget ledger of B subscription slots. Three -consumers share it, in priority order: +Each relay connection owns one implemented budget ledger of B subscription +slots. Four consumers share it, in priority order: -1. **Live subscriptions** (persistent, `limit: 0`) — the product; sized first. +1. **Core live subscriptions** (persistent, `limit: 0`) — the product; sized + first. 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) — negentropy rounds, +3. **Historic sync and dependency recovery** (transient) — at least one usable + slot remains after live admission; negentropy rounds, 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. 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. Live filter groups -are packed first and admitted atomically against the advertised -subscription-count budget: if the complete live set cannot fit, +(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: if the complete core live set cannot fit, the existing subscriptions are consolidated into the byte- and filter-count -bounded REQ groups first. If the consolidated live set still cannot fit, +bounded REQ groups first. If the consolidated core live set still cannot fit, partial coverage is not opened, historic work is deferred, and a warning surfaces the condition. Reconnect and consolidation build Layer 1, 2, and 3 coverage together and reserve the whole grouped set before opening it. A @@ -346,6 +353,12 @@ live coverage. Each reconnect closes the retired ledger and creates a new generation; queued or late borrowers therefore fail before sending on the new SDK session and cannot inflate or bypass its capacity. +Descendant live subscriptions are kept outside the core rollback set. Core +consolidation or restoration first closes them, and aborts if CLOSE cannot be +sent, so auxiliary coverage cannot silently consume capacity needed by newly +required core filters. An auxiliary CLOSED retires its remaining group and +falls back to EOSE-closing history without rebuilding healthy core coverage. + Some relays additionally cap the cumulative serialized REQ state retained by one connection. NIP-11 has no field for this limit, so it cannot be negotiated before the first refusal. A CLOSED reason of the rust-nostr form `active diff --git a/src/sync/mod.rs b/src/sync/mod.rs index aec2cef..b16a595 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -65,7 +65,7 @@ const MAX_PURGATORY_FILTER_ACTIONS_PER_TICK: usize = 1; const MAX_PURGATORY_DEPENDENCY_IDS_PER_QUERY: usize = 100; const SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS: usize = 20; const SEMANTIC_FALLBACK_MAX_DELIVERED_PERCENT: usize = 10; -const DESCENDANT_RECENT_WINDOW_SECS: u64 = 24 * 60 * 60; +const DESCENDANT_FALLBACK_OVERLAP_SECS: u64 = 15 * 60; fn should_use_semantic_fallback(requested_count: usize, received_count: usize) -> bool { requested_count >= SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS @@ -702,9 +702,22 @@ pub enum PendingBatchPurpose { #[derive(Debug)] struct DescendantSyncRotation { frontier: HashSet, - filters: Vec, + filters: Vec, next_filter: usize, - historic: bool, + in_flight: Option, +} + +#[derive(Debug)] +struct DescendantFilterCursor { + filter: Filter, + last_successful_until: Option, +} + +#[derive(Debug, Clone, Copy)] +struct DescendantFilterInFlight { + batch_id: u64, + filter_index: usize, + until: Timestamp, } #[derive(Debug)] @@ -719,25 +732,65 @@ impl Default for DescendantSyncRotation { frontier: HashSet::new(), filters: Vec::new(), next_filter: 0, - historic: true, + in_flight: None, } } } impl DescendantSyncRotation { - fn refresh(&mut self, members: HashSet, now: Timestamp) { - if self.historic - && !self.filters.is_empty() - && self.next_filter >= self.filters.len() - && self.frontier == members - { - self.historic = false; - } + fn refresh(&mut self, members: HashSet) { + let previous: HashMap> = self + .filters + .drain(..) + .map(|cursor| (cursor.filter.as_json(), cursor.last_successful_until)) + .collect(); self.frontier = members.clone(); - let since = (!self.historic) - .then(|| Timestamp::from(now.as_secs().saturating_sub(DESCENDANT_RECENT_WINDOW_SECS))); - self.filters = filters::tagged_one_of_our_root_event_filters(&members, since); + self.filters = filters::tagged_one_of_our_root_event_filters(&members, None) + .into_iter() + .map(|filter| DescendantFilterCursor { + last_successful_until: previous.get(&filter.as_json()).copied().flatten(), + filter, + }) + .collect(); self.next_filter = 0; + self.in_flight = None; + } + + fn next_request(&self, now: Timestamp) -> Option<(usize, Filter, 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)) + } + + fn mark_started(&mut self, batch_id: u64, filter_index: usize, until: Timestamp) { + self.in_flight = Some(DescendantFilterInFlight { + batch_id, + filter_index, + until, + }); + } + + 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 { + return false; + }; + self.in_flight = None; + if succeeded { + self.filters[in_flight.filter_index].last_successful_until = Some(in_flight.until); + self.next_filter = (in_flight.filter_index + 1) % self.filters.len(); + } + true } } @@ -1226,9 +1279,9 @@ async fn run_daily_timer( /// (see [`missing_events`]). Recovery attempts are backed off per relay, so /// the tick itself stays cheap when nothing is due. /// -/// Finally, one relay at a time receives one low-priority repository-descendant -/// history query. Reusing this five-second maintenance cadence bounds the new -/// global query-start rate without another scheduler or configuration surface. +/// Finally, one relay's live descendant admission is reconciled and every +/// constrained relay may advance one queued descendant history query. Reusing +/// this five-second cadence avoids another scheduler or configuration surface. async fn run_purgatory_announcement_sync( sync_manager: Arc>, mut shutdown_rx: broadcast::Receiver<()>, @@ -1478,8 +1531,8 @@ pub struct SyncManager { /// Auxiliary persistent descendant coverage, kept separate from the core /// desired set so it can be retired without replacing healthy core REQs. descendant_live_coverage: HashMap, - /// Round-robin cursor so the maintenance loop starts at most one extra - /// query globally per tick. + /// Round-robin cursor so only one relay performs descendant live-mode + /// derivation and admission work per maintenance tick. descendant_relay_cursor: usize, /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, @@ -2584,11 +2637,17 @@ impl SyncManager { let sync_method = batch.sync_method; let is_generic_filter = batch.purpose == PendingBatchPurpose::Announcements; - if batch.failed && batch.purpose == PendingBatchPurpose::Descendants { - // Do not let an interrupted page turn the initial complete pass - // into a recent-only rotation. Starting this tiny auxiliary state - // afresh is simpler and safer than persisting per-filter cursors. - self.descendant_sync_rotations.remove(relay_url); + if batch.purpose == PendingBatchPurpose::Descendants { + let succeeded = !batch.failed; + if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) { + rotation.mark_completed(batch_id, succeeded); + } + tracing::info!( + relay = %relay_url, + batch_id, + succeeded, + "Descendant historic query reached terminal batch state" + ); } let mut relay_index = self.relay_sync_index.write().await; @@ -3697,12 +3756,142 @@ impl SyncManager { true } - /// Start at most one low-priority descendant query globally per tick. - /// - /// Each relay first receives a complete historic pass over a stable - /// snapshot of direct thread members. Later passes use a rolling 24-hour - /// window. Every request uses the ordinary historic REQ+EOSE pagination - /// path, so it draws from the existing ledger and closes at EOSE. + async fn reconcile_descendant_mode( + &mut self, + relay_url: &str, + root_events: &HashSet, + ) { + let members = self.direct_thread_members(root_events).await; + if members.is_empty() { + let _ = self + .close_descendant_live_coverage(relay_url, "frontier became empty") + .await; + self.descendant_sync_rotations.remove(relay_url); + return; + } + if self + .descendant_live_coverage + .get(relay_url) + .is_some_and(|coverage| coverage.frontier == members) + { + return; + } + if self.descendant_live_coverage.contains_key(relay_url) + && !self + .close_descendant_live_coverage(relay_url, "frontier changed") + .await + { + return; + } + + let live_since = Timestamp::from( + Timestamp::now() + .as_secs() + .saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS), + ); + let historic_filters = filters::tagged_one_of_our_root_event_filters(&members, None); + let live_filters = + filters::tagged_one_of_our_root_event_filters(&members, Some(live_since)); + if let Some(connection) = self.connections.get(relay_url).cloned() { + let groups = live_filter_groups(&live_filters, connection.max_filters_per_req()); + if connection.can_admit_auxiliary_live_groups(&groups) { + match connection + .subscribe_auxiliary_live_filter_groups(groups) + .await + { + Ok(subscription_ids) => { + let filter_count = live_filters.len(); + self.descendant_live_coverage.insert( + relay_url.to_string(), + DescendantLiveCoverage { + frontier: members, + subscription_ids: subscription_ids.clone(), + }, + ); + self.descendant_sync_rotations.remove(relay_url); + tracing::info!( + relay = %relay_url, + filter_count, + subscription_count = subscription_ids.len(), + "Installed auxiliary descendant live coverage" + ); + // `limit:0` protects the future only. Pair every new + // live frontier with a complete EOSE-closing baseline + // so events stored before admission are not skipped. + let _ = self + .historic_sync_with_options( + relay_url, + historic_filters, + PendingItems::default(), + None, + PendingBatchPurpose::Descendants, + true, + ) + .await; + return; + } + Err(error) => { + tracing::warn!( + relay = %relay_url, + %error, + "Descendant live admission failed; using historic rotation" + ); + } + } + } else { + tracing::info!( + relay = %relay_url, + filter_count = live_filters.len(), + "Descendant live coverage does not fit while preserving historic capacity" + ); + } + } + + let rotation = self + .descendant_sync_rotations + .entry(relay_url.to_string()) + .or_default(); + if rotation.frontier != members { + rotation.refresh(members); + } + } + + async fn start_descendant_fallback(&mut self, relay_url: &str) { + let Some((filter_index, filter, until)) = self + .descendant_sync_rotations + .get(relay_url) + .and_then(|rotation| rotation.next_request(Timestamp::now())) + else { + return; + }; + let filter_count = self.descendant_sync_rotations[relay_url].filters.len(); + if let Some(batch_id) = self + .historic_sync_with_options( + relay_url, + vec![filter], + PendingItems::default(), + None, + PendingBatchPurpose::Descendants, + true, + ) + .await + { + if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) { + rotation.mark_started(batch_id, filter_index, until); + } + tracing::info!( + relay = %relay_url, + batch_id, + filter_index, + filter_count, + until = until.as_secs(), + "Started queued descendant fallback query" + ); + } + } + + /// Reconcile one live-mode decision and advance every constrained relay by + /// at most one EOSE-closing request on each five-second maintenance tick. async fn tick_descendant_sync(&mut self) { let targets = { let index = self.repo_sync_index.read().await; @@ -3730,128 +3919,21 @@ impl SyncManager { return; } - let relay_url = relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); + let reconcile_relay = + relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); - let root_events = targets[&relay_url].root_events.clone(); + self.reconcile_descendant_mode( + &reconcile_relay, + &targets[&reconcile_relay].root_events, + ) + .await; - let mut refreshed_members = None; - if let Some(existing) = self.descendant_live_coverage.get(&relay_url) { - let members = self.direct_thread_members(&root_events).await; - if members == existing.frontier { - return; - } - refreshed_members = Some(members); - if !self - .close_descendant_live_coverage(&relay_url, "frontier changed") - .await - { - return; - } - } - - let needs_snapshot = self - .descendant_sync_rotations - .get(&relay_url) - .is_none_or(|rotation| rotation.next_filter >= rotation.filters.len()); - if needs_snapshot { - let members = match refreshed_members { - Some(members) => members, - None => self.direct_thread_members(&root_events).await, - }; - if members.is_empty() { - return; - } - - let now = Timestamp::now(); - let live_since = Timestamp::from( - now.as_secs() - .saturating_sub(DESCENDANT_RECENT_WINDOW_SECS.min(15 * 60)), - ); - let live_filters = filters::tagged_one_of_our_root_event_filters( - &members, - Some(live_since), - ); - if let Some(connection) = self.connections.get(&relay_url).cloned() { - let groups = live_filter_groups(&live_filters, connection.max_filters_per_req()); - if connection.can_admit_auxiliary_live_groups(&groups) { - match connection - .subscribe_auxiliary_live_filter_groups(groups) - .await - { - Ok(subscription_ids) => { - let filter_count = live_filters.len(); - self.descendant_live_coverage.insert( - relay_url.clone(), - DescendantLiveCoverage { - frontier: members, - subscription_ids: subscription_ids.clone(), - }, - ); - self.descendant_sync_rotations.remove(&relay_url); - tracing::info!( - relay = %relay_url, - filter_count, - subscription_count = subscription_ids.len(), - "Installed auxiliary descendant live coverage" - ); - return; - } - Err(error) => { - tracing::warn!( - relay = %relay_url, - %error, - "Descendant live admission failed; using historic rotation" - ); - } - } - } else { - tracing::info!( - relay = %relay_url, - filter_count = live_filters.len(), - "Descendant live coverage does not fit while preserving historic capacity" - ); - } - } - - let rotation = self - .descendant_sync_rotations - .entry(relay_url.clone()) - .or_default(); - rotation.refresh(members, Timestamp::now()); - } - - let (filter, historic, filter_index, filter_count) = { - let rotation = self.descendant_sync_rotations.get(&relay_url).unwrap(); - ( - rotation.filters[rotation.next_filter].clone(), - rotation.historic, - rotation.next_filter, - rotation.filters.len(), - ) - }; - if self - .historic_sync_with_options( - &relay_url, - vec![filter], - PendingItems::default(), - None, - PendingBatchPurpose::Descendants, - true, - ) - .await - .is_some() - { - self.descendant_sync_rotations - .get_mut(&relay_url) - .unwrap() - .next_filter += 1; - tracing::debug!( - relay = %relay_url, - historic, - filter_index, - filter_count, - "Started scheduled repository descendant query" - ); + let constrained: Vec = relay_urls + .into_iter() + .filter(|relay_url| self.descendant_sync_rotations.contains_key(relay_url)) + .collect(); + for relay_url in constrained { + self.start_descendant_fallback(&relay_url).await; } } @@ -5647,6 +5729,14 @@ impl SyncManager { true } + fn release_removed_descendant_batch(&mut self, relay_url: &str, batch: &PendingBatch) { + if batch.purpose == PendingBatchPurpose::Descendants { + if let Some(rotation) = self.descendant_sync_rotations.get_mut(relay_url) { + rotation.mark_completed(batch.batch_id, false); + } + } + } + async fn handle_subscription_closed( &mut self, relay_url: &str, @@ -5672,7 +5762,7 @@ impl SyncManager { self.descendant_sync_rotations .entry(relay_url.to_string()) .or_default() - .refresh(coverage.frontier, Timestamp::now()); + .refresh(coverage.frontier); tracing::warn!( relay = %relay_url, sub_id = %subscription_id, @@ -5690,9 +5780,14 @@ impl SyncManager { ); self.byte_limited_live_relays .insert(relay_url.to_string(), Instant::now()); - if live_generation.is_none() { + let removed_batch = if live_generation.is_none() { let mut pending = self.pending_sync_index.write().await; - take_batch_containing_subscription(&mut pending, relay_url, &subscription_id); + take_batch_containing_subscription(&mut pending, relay_url, &subscription_id) + } else { + None + }; + if let Some(batch) = &removed_batch { + self.release_removed_descendant_batch(relay_url, batch); } let has_pending = self.has_pending_batches(relay_url).await; if self @@ -5710,6 +5805,7 @@ impl SyncManager { }; if let Some(batch) = removed_batch { + self.release_removed_descendant_batch(relay_url, &batch); tracing::warn!( relay = %relay_url, sub_id = %subscription_id, @@ -5764,7 +5860,8 @@ impl SyncManager { &subscription_id, ) }; - if removed_batch.is_some() { + if let Some(batch) = &removed_batch { + self.release_removed_descendant_batch(relay_url, batch); self.recompute_new_sync_filters_for_relay(relay_url).await; } else if let Some(generation) = live_generation { self.restore_live_coverage_after_closed(relay_url, generation) @@ -5780,6 +5877,9 @@ impl SyncManager { let mut pending = self.pending_sync_index.write().await; take_batch_containing_subscription(&mut pending, relay_url, &subscription_id) }; + if let Some(batch) = &removed_batch { + self.release_removed_descendant_batch(relay_url, batch); + } self.health_tracker.record_policy_refusal(relay_url); if let Some(metrics) = &self.metrics { metrics.record_policy_refusal(relay_url, category.label()); @@ -6311,10 +6411,16 @@ impl SyncManager { async fn historic_sync( &mut self, relay_url: &str, - filters: Vec, + mut filters: Vec, items: PendingItems, since: Option, ) -> Option { + if !items.root_events.is_empty() { + let members = self.direct_thread_members(&items.root_events).await; + filters.extend(filters::tagged_one_of_our_root_event_filters( + &members, None, + )); + } self.historic_sync_with_options( relay_url, filters, @@ -6742,30 +6848,57 @@ mod tests { } #[test] - fn descendant_rotation_becomes_recent_only_after_a_stable_historic_pass() { + fn descendant_rotation_advances_cursor_only_after_successful_eose() { let member = EventId::from_byte_array([7; 32]); let members = HashSet::from([member]); let now = Timestamp::from_secs(200_000); let mut rotation = DescendantSyncRotation::default(); - rotation.refresh(members.clone(), now); - assert!(rotation.historic); - assert!(rotation - .filters - .iter() - .all(|filter| !serde_json::to_value(filter) - .unwrap() - .as_object() - .unwrap() - .contains_key("since"))); + rotation.refresh(members); + let (filter_index, first, until) = rotation.next_request(now).unwrap(); + assert!(serde_json::to_value(first).unwrap().get("since").is_none()); + rotation.mark_started(41, filter_index, until); + assert!(rotation.next_request(now).is_none()); - rotation.next_filter = rotation.filters.len(); - rotation.refresh(members, now); - assert!(!rotation.historic); - assert!(rotation.filters.iter().all(|filter| { - serde_json::to_value(filter).unwrap()["since"] - == serde_json::json!(now.as_secs() - DESCENDANT_RECENT_WINDOW_SECS) - })); + 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()); + 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(); + assert_eq!( + serde_json::to_value(recent).unwrap()["since"], + serde_json::json!(now.as_secs() - DESCENDANT_FALLBACK_OVERLAP_SECS) + ); + } + + #[test] + fn large_descendant_frontier_splits_across_live_subscriptions() { + let members: HashSet = (0..600u32) + .map(|index| { + let mut bytes = [0u8; 32]; + bytes[..4].copy_from_slice(&index.to_be_bytes()); + EventId::from_byte_array(bytes) + }) + .collect(); + let filters = filters::tagged_one_of_our_root_event_filters(&members, None); + let groups = live_filter_groups(&filters, MAX_FILTERS_PER_REQ); + + assert!(filters.len() > 3, "the frontier must span byte chunks"); + assert!( + groups.len() > 1, + "large complete descendant coverage must reserve several subscriptions" + ); + assert!(groups.iter().all(|group| group.len() <= MAX_FILTERS_PER_REQ)); } #[test] diff --git a/tests/common/relay.rs b/tests/common/relay.rs index dc0c7eb..a6cce5a 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -84,6 +84,7 @@ struct RelayOptions { relay_data_path: Option, deletion_lifecycle: Option, rejected_hot_cache_duration_secs: Option, + relay_max_subscriptions: Option, /// Run with the production outbound target policy (reject non-global /// event-directed sync targets). The fixture default is permissive /// because the entire test infrastructure lives on loopback. @@ -309,6 +310,20 @@ impl TestRelay { .await } + /// Start a source relay advertising and enforcing a small per-connection + /// subscription budget. Sync scenarios use this to exercise graceful + /// degradation paths selected from NIP-11. + pub async fn start_with_relay_max_subscriptions(limit: usize) -> Self { + Self::start_internal( + port::reserve_port(), + RelayOptions { + relay_max_subscriptions: Some(limit), + ..RelayOptions::default() + }, + ) + .await + } + /// Start a relay with LMDB backend (persistent side DBs on a temp dir). pub async fn start_with_lmdb() -> Self { Self::start_internal( @@ -581,6 +596,10 @@ impl TestRelay { ); } + if let Some(limit) = options.relay_max_subscriptions { + cmd.env("NGIT_RELAY_MAX_SUBSCRIPTIONS", limit.to_string()); + } + // Add negentropy disable flag if requested if options.disable_negentropy { cmd.env("NGIT_SYNC_DISABLE_NEGENTROPY", "true"); diff --git a/tests/sync/descendant_sync.rs b/tests/sync/descendant_sync.rs index 7df362c..f5c21fe 100644 --- a/tests/sync/descendant_sync.rs +++ b/tests/sync/descendant_sync.rs @@ -6,12 +6,28 @@ use nostr_sdk::prelude::*; use crate::common::{sync_helpers::*, TestRelay}; +async fn wait_for_log(path: &std::path::Path, needle: &str, timeout: Duration) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + loop { + if std::fs::read_to_string(path) + .unwrap_or_default() + .contains(needle) + { + return true; + } + if tokio::time::Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + /// Events which reference a direct repository-thread member, but not the /// repository or its root event, are recovered by the scheduled historic /// pass. The recovered events do not recursively extend the frontier. #[tokio::test] async fn historic_sync_recovers_one_generation_of_parent_only_descendants() { - let source = TestRelay::start().await; + let source = TestRelay::start_with_relay_max_subscriptions(5).await; let keys = Keys::generate(); let repo_id = "scheduled-descendant-history"; @@ -66,6 +82,15 @@ async fn historic_sync_recovers_one_generation_of_parent_only_descendants() { .await, "scheduled historic rotation should recover a parent-only descendant" ); + assert!( + wait_for_log( + &syncing.log_path(), + "Started queued descendant fallback query", + Duration::from_secs(5), + ) + .await, + "the constrained source should exercise EOSE-closing fallback" + ); assert!( !wait_for_event_on_relay( syncing.url(), @@ -80,3 +105,62 @@ async fn historic_sync_recovers_one_generation_of_parent_only_descendants() { syncing.stop().await; source.stop().await; } + +/// A relay with spare NIP-11 subscription capacity receives permanent +/// descendant coverage, so a later parent-only event arrives without waiting +/// for the fallback rotation. +#[tokio::test] +async fn descendant_live_coverage_is_preferred_when_capacity_remains() { + let source = TestRelay::start().await; + let keys = Keys::generate(); + let repo_id = "live-descendant-coverage"; + let source_domains = [source.domain()]; + let source_refs = source_domains + .iter() + .map(String::as_str) + .collect::>(); + let (_announcement, _source_git) = + setup_announcement_on_relay(&source, &keys, &source_refs, repo_id).await; + let source_client = TestClient::new(source.url(), keys.clone()) + .await + .expect("connect to source relay"); + let issue = build_layer2_issue_event(&keys, &repo_coord(&keys, repo_id), "Live root") + .expect("build issue"); + let direct_reply = build_layer3_reply_with_e_tag(&keys, &issue.id, "Direct reply") + .expect("build direct reply"); + source_client.send_event(&issue).await.unwrap(); + source_client.send_event(&direct_reply).await.unwrap(); + + let syncing = TestRelay::start_with_sync(None).await; + let domains = [source.domain(), syncing.domain()]; + let domain_refs = domains.iter().map(String::as_str).collect::>(); + let (_target_announcement, _target_git) = + setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await; + assert!( + wait_for_log( + &syncing.log_path(), + "Installed auxiliary descendant live coverage", + Duration::from_secs(20), + ) + .await, + "spare source capacity should retain descendant coverage" + ); + + let parent_only = + build_layer3_reply_with_e_tag(&keys, &direct_reply.id, "Arrives through live coverage") + .expect("build parent-only descendant"); + source_client.send_event(&parent_only).await.unwrap(); + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(parent_only.id), + Duration::from_secs(10), + ) + .await, + "parent-only descendant should arrive through retained live coverage" + ); + + source_client.disconnect().await; + syncing.stop().await; + source.stop().await; +}