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; +}