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());