From d203f9de5683b48ee8483681ae138956bc62a604 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 20:43:08 +0000 Subject: [PATCH] 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!(