mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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.
This commit is contained in:
+66
-46
@@ -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<SubscriptionId>,
|
||||
HashMap<SubscriptionId, PaginationState>,
|
||||
) {
|
||||
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!(
|
||||
|
||||
Reference in New Issue
Block a user