diff --git a/CHANGELOG.md b/CHANGELOG.md index e349e73..ef99923 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,15 @@ 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. 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 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..d146d7d 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,37 @@ live in - **Function**: `build_root_event_tag_filters(root_events, since)` - **Only for `SyncLevel::Full` repos** — purgatory announcements (`StateOnly`) skip this layer +### 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. + +- 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 + 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. 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) The `build_sync_level_aware_filters()` function combines both layers, partitioning repos by `SyncLevel`: 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/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 6a9319b..b16a595 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_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 @@ -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,108 @@ 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, + 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)] +struct DescendantLiveCoverage { + frontier: HashSet, + subscription_ids: Vec, +} + +impl Default for DescendantSyncRotation { + fn default() -> Self { + Self { + frontier: HashSet::new(), + filters: Vec::new(), + next_filter: 0, + in_flight: None, + } + } +} + +impl DescendantSyncRotation { + 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(); + 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 + } +} + /// Items included in a pending batch #[derive(Debug, Clone, Default)] pub struct PendingItems { @@ -1172,6 +1278,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'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<()>, @@ -1187,6 +1297,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 +1526,14 @@ 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, + /// 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 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>, /// Channel for EOSE notifications (set during run) @@ -1515,6 +1634,9 @@ 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_live_coverage: HashMap::new(), + descendant_relay_cursor: 0, disconnect_tx: None, eose_tx: None, subscription_closed_tx: None, @@ -2507,14 +2629,26 @@ 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.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; @@ -2567,7 +2701,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 +2719,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 +2870,10 @@ 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 let connection = match self.connections.get(relay_url) { @@ -3528,16 +3667,276 @@ 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 + } + + 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 + } + + 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; + 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 reconcile_relay = + relay_urls[self.descendant_relay_cursor % relay_urls.len()].clone(); + self.descendant_relay_cursor = self.descendant_relay_cursor.wrapping_add(1); + self.reconcile_descendant_mode( + &reconcile_relay, + &targets[&reconcile_relay].root_events, + ) + .await; + + 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; + } + } + /// 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 +4937,11 @@ 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); + self.descendant_live_coverage.remove(relay_url); // Check if this was an intentional disconnect (Disconnecting status) let was_intentional = { @@ -5250,6 +5654,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; @@ -5315,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, @@ -5323,6 +5745,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); + 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, @@ -5331,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 @@ -5351,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, @@ -5405,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) @@ -5421,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()); @@ -5448,13 +5907,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; @@ -5856,6 +6325,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, @@ -5877,11 +6409,37 @@ impl SyncManager { /// * `Some(batch_id)` - Batch was created and sync initiated /// * `None` - No connection or sync failed to start async fn historic_sync( + &mut self, + relay_url: &str, + 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, + 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!( @@ -5922,8 +6480,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(); @@ -5945,6 +6504,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, @@ -6162,53 +6722,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!( @@ -6221,6 +6738,7 @@ impl SyncManager { // Create PendingBatch for REQ+EOSE let batch = PendingBatch { batch_id, + purpose, items, outstanding_subs: subscription_ids, sync_method: SyncMethod::ReqEose, @@ -6329,6 +6847,60 @@ mod tests { assert!(!should_use_semantic_fallback(19, 0)); } + #[test] + 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); + 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()); + + 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] fn group_filters_for_req_respects_count_and_byte_budgets() { // Many small filters group by the count cap. @@ -6843,6 +7415,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, @@ -6945,6 +7518,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, @@ -6980,6 +7554,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, @@ -7298,6 +7873,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, @@ -7321,6 +7897,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/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()); 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.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..f5c21fe --- /dev/null +++ b/tests/sync/descendant_sync.rs @@ -0,0 +1,166 @@ +//! Scheduled repository-descendant synchronization scenarios. + +use std::time::Duration; + +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_with_relay_max_subscriptions(5).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_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(), + 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; +} + +/// 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; +}