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!(