mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
Merge #d6e07abe: Fix invitation recovery during busy relay startup
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsddcr6h6cafdqp6qmcxh9cd8fxc6tr7nc6asx60krjxeuqx34zkjc9p0928 PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 PR description: Fixes maintainer invitation acceptance that remains in purgatory while a production GRASP server starts and ingests a large bootstrap history. The live reproduction showed two independent blockers: startup waited minutes for retention cleanup across more than 52,000 deletion requests, then a rate-limit notice and paginated EOSE formed a circular wait on the SyncManager lock. Fresh invitation announcements were accepted into purgatory, but the five-second recovery pass could not run. This PR starts retention catch-up after the relay and sync workers are live, resumes the exact historic page outside the actor lock after rate-limit cooldown, and replaces repeated full-index action construction with lightweight dirty-relay recomputation. Coverage verifies deferred retention cleanup, preserves generic pagination while the actor remains available, limits relay recomputation to the changed batch, and keeps the existing maintainer invitation/state integration scenarios passing.
This commit is contained in:
@@ -13,6 +13,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fix maintainer invitation recovery after a production restart by starting the
|
||||
relay and SyncManager before scanning retained deletion requests. The
|
||||
potentially long retention catch-up now begins immediately in the existing
|
||||
background maintenance task instead of blocking purgatory processing.
|
||||
- Fix maintainer invitation recovery during a rate-limited historic sync by
|
||||
resuming pagination outside the SyncManager lock. Fresh relay batches now
|
||||
mark only their changed relays for recomputation, avoiding repeated
|
||||
full-index filter construction while a large bootstrap history is arriving.
|
||||
- Fix maintainer invitation acceptance by prioritizing fresh purgatory
|
||||
dependencies within a bounded reconciliation pass and retaining the source
|
||||
relays needed to recover inviter events after the short-lived hot cache
|
||||
|
||||
@@ -63,14 +63,6 @@ impl DeletionRuntime {
|
||||
/// Run deletion-owned startup tasks before the relay begins serving traffic.
|
||||
pub async fn run_startup_tasks(&self) -> Result<()> {
|
||||
log_startup_reconciliation(self.service.run_startup_reconciliation().await?);
|
||||
// Enumeration failure is fatal before serving traffic; individual record
|
||||
// failures are conservatively retained and retried by the timer.
|
||||
let request_stats = self
|
||||
.service
|
||||
.cleanup_expired_requests(Timestamp::now())
|
||||
.await?;
|
||||
log_request_cleanup("startup catch-up", request_stats);
|
||||
self.run_holding_startup_cleanup().await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -84,16 +76,22 @@ impl DeletionRuntime {
|
||||
let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
|
||||
|
||||
let handle = tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval_at(
|
||||
tokio::time::Instant::now() + interval_duration,
|
||||
interval_duration,
|
||||
);
|
||||
let mut first_pass = true;
|
||||
let mut interval =
|
||||
tokio::time::interval_at(tokio::time::Instant::now(), interval_duration);
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = interval.tick() => {
|
||||
// The first pass is the startup catch-up. It deliberately
|
||||
// runs here, after the relay and SyncManager are live:
|
||||
// large retained request sets can take minutes to scan,
|
||||
// but must not prevent fresh purgatory work from starting.
|
||||
match service.cleanup_expired_requests(Timestamp::now()).await {
|
||||
Ok(stats) => log_request_cleanup("periodic pass", stats),
|
||||
Ok(stats) => log_request_cleanup(
|
||||
if first_pass { "startup catch-up" } else { "periodic pass" },
|
||||
stats,
|
||||
),
|
||||
Err(error) => tracing::warn!(error = %error, "Deletion-request cleanup periodic pass failed"),
|
||||
}
|
||||
match holding.cleanup_expired_with_lifecycle(&lifecycle, Timestamp::now(), retention).await {
|
||||
@@ -112,6 +110,7 @@ impl DeletionRuntime {
|
||||
tracing::warn!(error = %e, "Holding cleanup periodic pass failed");
|
||||
}
|
||||
}
|
||||
first_pass = false;
|
||||
}
|
||||
changed = shutdown_rx.changed() => {
|
||||
if changed.is_ok() && *shutdown_rx.borrow() {
|
||||
@@ -134,33 +133,6 @@ impl DeletionRuntime {
|
||||
handle,
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_holding_startup_cleanup(&self) {
|
||||
match self
|
||||
.holding
|
||||
.cleanup_expired_with_lifecycle(
|
||||
&self.lifecycle,
|
||||
Timestamp::now(),
|
||||
self.holding_retention,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(stats) => {
|
||||
if stats.expired_records > 0 {
|
||||
tracing::info!(
|
||||
examined = stats.metadata_examined,
|
||||
expired = stats.expired_records,
|
||||
metadata_deleted = stats.metadata_deleted,
|
||||
archived_events_deleted = stats.archived_events_deleted,
|
||||
"Holding cleanup startup catch-up completed"
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "Holding cleanup startup catch-up failed");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct DeletionCleanupTask {
|
||||
@@ -284,3 +256,101 @@ fn log_whitelist_restore(stats: WhitelistRestoreStats) {
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use nostr_memory::MemoryDatabase;
|
||||
use nostr_relay_builder::prelude::{
|
||||
EventBuilder, EventId, FinalizeEvent, Keys, Kind, NostrDatabase, Tag,
|
||||
};
|
||||
|
||||
use super::*;
|
||||
use crate::grasp06::receive::new_repo_init_locks;
|
||||
use crate::nostr::lifecycle::{
|
||||
DeletionContext, HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle,
|
||||
RequestClassification, Tombstones,
|
||||
};
|
||||
use crate::purgatory::Purgatory;
|
||||
|
||||
#[tokio::test]
|
||||
async fn retained_request_cleanup_starts_after_blocking_startup_finishes() {
|
||||
let database = Arc::new(MemoryDatabase::unbounded());
|
||||
let tombstones = Tombstones::in_memory();
|
||||
let holding = HoldingStore::in_memory();
|
||||
let lifecycle = RepositoryLifecycle::in_memory();
|
||||
let config = crate::config::Config {
|
||||
deletion_request_retention_unused_served_secs: 1,
|
||||
deletion_request_retention_unused_unserved_gating_additional_secs: 1,
|
||||
..crate::config::Config::for_testing()
|
||||
};
|
||||
let service = DeletionService::new(DeletionContext::new(
|
||||
"test.example.com",
|
||||
database.clone(),
|
||||
tombstones.clone(),
|
||||
holding.clone(),
|
||||
lifecycle.clone(),
|
||||
ReplaceableHistoryStore::in_memory(),
|
||||
PathBuf::new(),
|
||||
Arc::new(Purgatory::new(PathBuf::new())),
|
||||
config,
|
||||
new_repo_init_locks(),
|
||||
));
|
||||
let runtime = DeletionRuntime::new(
|
||||
service,
|
||||
holding,
|
||||
Duration::from_secs(60),
|
||||
Duration::from_secs(60),
|
||||
);
|
||||
let request = EventBuilder::new(Kind::EventDeletion, "")
|
||||
.tags(vec![Tag::event(EventId::all_zeros())])
|
||||
.finalize(&Keys::generate())
|
||||
.expect("deletion request should sign");
|
||||
|
||||
tombstones
|
||||
.record_request(
|
||||
&request,
|
||||
Timestamp::from_secs(1),
|
||||
RequestClassification::LocallyActionable,
|
||||
)
|
||||
.await
|
||||
.expect("request lifecycle should be retained");
|
||||
database
|
||||
.save_event(&request)
|
||||
.await
|
||||
.expect("request payload should be served");
|
||||
|
||||
runtime
|
||||
.run_startup_tasks()
|
||||
.await
|
||||
.expect("blocking startup reconciliation should complete");
|
||||
assert!(
|
||||
tombstones
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.expect("lifecycle lookup should succeed")
|
||||
.is_some(),
|
||||
"retention cleanup must not delay relay and SyncManager startup"
|
||||
);
|
||||
|
||||
let cleanup = runtime.spawn_cleanup_task();
|
||||
tokio::time::timeout(Duration::from_secs(2), async {
|
||||
loop {
|
||||
if tombstones
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.expect("lifecycle lookup should succeed")
|
||||
.is_none()
|
||||
{
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("background startup catch-up should run immediately");
|
||||
cleanup.shutdown().await;
|
||||
}
|
||||
}
|
||||
|
||||
+171
-48
@@ -40,8 +40,6 @@ pub use self_subscriber::SelfSubscriber;
|
||||
|
||||
// Re-export health tracking types
|
||||
pub use health::RelayHealthTracker;
|
||||
use tokio::time::sleep;
|
||||
|
||||
use std::collections::{HashMap, HashSet, VecDeque};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
@@ -439,6 +437,18 @@ fn take_drained_batch_as_failed(
|
||||
Some(batch)
|
||||
}
|
||||
|
||||
fn mark_deferred_pagination(
|
||||
batch: &mut PendingBatch,
|
||||
completed_sub_id: &SubscriptionId,
|
||||
) -> SubscriptionId {
|
||||
let deferred_sub_id = SubscriptionId::new(format!(
|
||||
"deferred-pagination-{}-{completed_sub_id}",
|
||||
batch.batch_id
|
||||
));
|
||||
batch.outstanding_subs.insert(deferred_sub_id.clone());
|
||||
deferred_sub_id
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// SyncManager - Main Entry Point
|
||||
// =============================================================================
|
||||
@@ -1118,48 +1128,43 @@ impl SyncManager {
|
||||
let relay_url_for_pagination = relay_url.to_string();
|
||||
let batch_id = batch.batch_id;
|
||||
|
||||
// Drop the lock before async operations
|
||||
drop(pending);
|
||||
|
||||
// Wait for rate limiting to clear before pagination continues
|
||||
// A NOTICE can arrive immediately before this page's EOSE.
|
||||
// Never wait for the cooldown here: the caller owns the
|
||||
// global SyncManager mutex, while the health checker that
|
||||
// clears the cooldown needs that same mutex. Keep a
|
||||
// sentinel in the batch and resume this exact page from a
|
||||
// detached worker so generic history is not lost or
|
||||
// restarted from page one.
|
||||
if self.health_tracker.is_rate_limited(relay_url) {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
batch_id = batch_id,
|
||||
"Relay is rate limited, waiting before pagination"
|
||||
);
|
||||
let deferred_sub_id = mark_deferred_pagination(batch, &sub_id);
|
||||
drop(pending);
|
||||
|
||||
// Loop until rate limit clears, sleeping with jitter between checks
|
||||
while self.health_tracker.is_rate_limited(relay_url) {
|
||||
let jitter_secs = 1 + (rand::random::<u64>() % 5); // 1-5 seconds
|
||||
sleep(Duration::from_secs(jitter_secs)).await;
|
||||
}
|
||||
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
batch_id = batch_id,
|
||||
"Rate limit cleared, continuing pagination"
|
||||
);
|
||||
let batch_exists = {
|
||||
let pending = self.pending_sync_index.read().await;
|
||||
pending
|
||||
.get(&relay_url_for_pagination)
|
||||
.map(|batches| batches.iter().any(|b| b.batch_id == batch_id))
|
||||
.unwrap_or(false)
|
||||
};
|
||||
|
||||
// If we were rate limited, verify batch still exists after waiting
|
||||
// (batches are wiped during disconnect, so avoid orphaned pagination)
|
||||
if !batch_exists {
|
||||
tracing::debug!(
|
||||
let Some(connection) =
|
||||
self.connections.get(&relay_url_for_pagination).cloned()
|
||||
else {
|
||||
tracing::error!(
|
||||
relay = %relay_url_for_pagination,
|
||||
batch_id = batch_id,
|
||||
"Batch no longer exists after rate limit wait, skipping pagination"
|
||||
batch_id,
|
||||
"Cannot defer rate-limited pagination without a relay connection"
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
Self::spawn_deferred_pagination(
|
||||
connection,
|
||||
self.health_tracker.clone(),
|
||||
self.pending_sync_index.clone(),
|
||||
relay_url_for_pagination,
|
||||
batch_id,
|
||||
deferred_sub_id,
|
||||
next_filter,
|
||||
until_timestamp,
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// Drop the lock before async operations
|
||||
drop(pending);
|
||||
|
||||
// Subscribe to next page and add to outstanding_subs
|
||||
let mut next_page_started = false;
|
||||
if let Some(conn) = self.connections.get(&relay_url_for_pagination) {
|
||||
@@ -1515,6 +1520,86 @@ impl SyncManager {
|
||||
self.confirm_batch(relay_url, completed_batch).await;
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn spawn_deferred_pagination(
|
||||
connection: RelayConnection,
|
||||
health_tracker: Arc<RelayHealthTracker>,
|
||||
pending_sync_index: PendingSyncIndex,
|
||||
relay_url: String,
|
||||
batch_id: u64,
|
||||
deferred_sub_id: SubscriptionId,
|
||||
next_filter: Filter,
|
||||
until_timestamp: Timestamp,
|
||||
) {
|
||||
tokio::spawn(async move {
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
batch_id,
|
||||
until = %until_timestamp,
|
||||
"Rate limited during historic pagination; deferring the exact next page without blocking the sync actor"
|
||||
);
|
||||
|
||||
loop {
|
||||
if let Some(remaining) = health_tracker.get_remaining_backoff(&relay_url) {
|
||||
tokio::time::sleep(remaining + Duration::from_millis(10)).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
// Hold only the pending-index lock while subscribing. This
|
||||
// prevents a very fast EOSE from reaching the actor before its
|
||||
// new subscription ID is registered, without blocking the
|
||||
// SyncManager mutex or the purgatory timer.
|
||||
let mut pending = pending_sync_index.write().await;
|
||||
let Some(batch) = pending.get_mut(&relay_url).and_then(|batches| {
|
||||
batches.iter_mut().find(|batch| batch.batch_id == batch_id)
|
||||
}) else {
|
||||
tracing::debug!(
|
||||
relay = %relay_url,
|
||||
batch_id,
|
||||
"Deferred pagination batch no longer exists"
|
||||
);
|
||||
return;
|
||||
};
|
||||
if !batch.outstanding_subs.contains(&deferred_sub_id) {
|
||||
return;
|
||||
}
|
||||
|
||||
match connection.subscribe_filter(next_filter.clone(), true).await {
|
||||
Ok(new_sub_id) => {
|
||||
batch.outstanding_subs.remove(&deferred_sub_id);
|
||||
batch.outstanding_subs.insert(new_sub_id.clone());
|
||||
batch.pagination_state.insert(
|
||||
new_sub_id.clone(),
|
||||
PaginationState {
|
||||
event_count: 0,
|
||||
min_created_at: None,
|
||||
original_filter: next_filter,
|
||||
},
|
||||
);
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
new_sub_id = %new_sub_id,
|
||||
batch_id,
|
||||
until = %until_timestamp,
|
||||
"Deferred pagination resumed after rate-limit cooldown"
|
||||
);
|
||||
return;
|
||||
}
|
||||
Err(error) => {
|
||||
drop(pending);
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
batch_id,
|
||||
error = %error,
|
||||
"Deferred pagination subscription failed; retrying"
|
||||
);
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// Confirm a completed batch by moving items to RelayState
|
||||
///
|
||||
/// This method is used by both sync paths (REQ+EOSE and Negentropy) to
|
||||
@@ -1898,9 +1983,14 @@ impl SyncManager {
|
||||
action = action_rx.recv() => {
|
||||
match action {
|
||||
Some(add_filters) => {
|
||||
// Process AddFilters action directly
|
||||
// SelfSubscriber actions are dirty-relay signals.
|
||||
// Recompute against current pending and confirmed
|
||||
// state instead of trusting a queued full-index
|
||||
// snapshot that may already be stale.
|
||||
let mut manager = sync_manager.lock().await;
|
||||
manager.handle_new_sync_filters(add_filters).await;
|
||||
manager
|
||||
.recompute_new_sync_filters_for_relay(&add_filters.relay_url)
|
||||
.await;
|
||||
}
|
||||
None => break,
|
||||
}
|
||||
@@ -2769,6 +2859,18 @@ impl SyncManager {
|
||||
async fn recompute_new_sync_filters_for_relay(&mut self, relay_url: &str) {
|
||||
use crate::sync::algorithms::{compute_actions, derive_relay_targets};
|
||||
|
||||
let relay_url = match canonical_relay_key(relay_url) {
|
||||
Ok(relay_url) => relay_url,
|
||||
Err(error) => {
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
error = %error,
|
||||
"Ignoring invalid dirty relay signal"
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
// Get current state from indexes (need to collect to avoid holding locks)
|
||||
let all_targets = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
@@ -2776,7 +2878,7 @@ impl SyncManager {
|
||||
};
|
||||
|
||||
// Filter to only targets for this specific relay
|
||||
let relay_target = match all_targets.get(relay_url) {
|
||||
let relay_target = match all_targets.get(&relay_url) {
|
||||
Some(target) => target.clone(),
|
||||
None => {
|
||||
tracing::debug!(
|
||||
@@ -2789,7 +2891,7 @@ impl SyncManager {
|
||||
|
||||
// Build single-relay targets map for compute_actions
|
||||
let mut single_relay_targets = std::collections::HashMap::new();
|
||||
single_relay_targets.insert(relay_url.to_string(), relay_target);
|
||||
single_relay_targets.insert(relay_url.clone(), relay_target);
|
||||
|
||||
// Compute actions for new items
|
||||
let actions = {
|
||||
@@ -4413,8 +4515,7 @@ impl SyncManager {
|
||||
filter_idx = idx,
|
||||
remote_count = remote_count,
|
||||
local_count = reconciliation.local.len(),
|
||||
remote_ids = ?remote_excluding_ids,
|
||||
"[DIAG TRACE] ✓ Negentropy diff results for filter {}", idx
|
||||
"Negentropy diff completed for filter"
|
||||
);
|
||||
if remote_count > 0 {
|
||||
all_remote_ids.extend(remote_excluding_ids);
|
||||
@@ -4502,16 +4603,12 @@ impl SyncManager {
|
||||
.map(|c| Filter::new().ids(c.iter().copied()))
|
||||
.collect();
|
||||
|
||||
// DEBUG TRACING: Log that we're requesting events by ID
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
batch_id = batch_id,
|
||||
total_event_ids = all_remote_ids.len(),
|
||||
filter_chunks = ids_filters.len(),
|
||||
event_ids = ?all_remote_ids,
|
||||
"[DIAG TRACE] ✓ Creating {} subscription(s) to fetch {} missing event(s) by ID",
|
||||
ids_filters.len(),
|
||||
all_remote_ids.len()
|
||||
"Creating subscriptions to fetch missing events by ID"
|
||||
);
|
||||
|
||||
let mut subscription_ids = HashSet::new();
|
||||
@@ -4994,6 +5091,32 @@ mod tests {
|
||||
assert!(!pending[relay_url][0].failed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rate_limited_pagination_keeps_batch_pending_without_actor_wait() {
|
||||
let relay_url = "wss://pagination.example";
|
||||
let completed_sub = SubscriptionId::new("completed-page");
|
||||
let mut batch = PendingBatch {
|
||||
batch_id: 73,
|
||||
items: PendingItems::default(),
|
||||
outstanding_subs: HashSet::new(),
|
||||
sync_method: SyncMethod::ReqEose,
|
||||
pagination_state: HashMap::new(),
|
||||
requested_event_ids: None,
|
||||
received_event_ids: None,
|
||||
retry_count: 0,
|
||||
failed: false,
|
||||
};
|
||||
|
||||
let deferred = mark_deferred_pagination(&mut batch, &completed_sub);
|
||||
assert!(batch.outstanding_subs.contains(&deferred));
|
||||
|
||||
let mut pending = HashMap::from([(relay_url.to_string(), vec![batch])]);
|
||||
assert!(
|
||||
take_drained_batch_as_failed(&mut pending, relay_url, 73).is_none(),
|
||||
"the deferred exact page must prevent a generic historic batch from being confirmed early"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn relay_disconnect_waits_for_pending_and_historic_sync_work() {
|
||||
let mut source = RelayState {
|
||||
|
||||
+44
-34
@@ -77,6 +77,13 @@ impl PendingUpdates {
|
||||
}
|
||||
}
|
||||
|
||||
fn dirty_relays(updates: &HashMap<String, RepoSyncNeeds>) -> HashSet<String> {
|
||||
updates
|
||||
.values()
|
||||
.flat_map(|needs| needs.relays.iter().cloned())
|
||||
.collect()
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// SelfSubscriber - Main Component
|
||||
// =============================================================================
|
||||
@@ -549,11 +556,9 @@ impl SelfSubscriber {
|
||||
|
||||
/// Process accumulated batch
|
||||
///
|
||||
/// Updates the RepoSyncIndex with discovered repos, then derives per-relay
|
||||
/// targets and sends RelayAction messages to the SyncManager.
|
||||
/// Updates the RepoSyncIndex with discovered repos, then marks only the
|
||||
/// relays changed by this batch for SyncManager recomputation.
|
||||
async fn process_batch(&self, pending: &mut PendingUpdates) {
|
||||
use crate::sync::algorithms::derive_relay_targets;
|
||||
|
||||
let updates = pending.take();
|
||||
|
||||
if updates.is_empty() {
|
||||
@@ -574,6 +579,8 @@ impl SelfSubscriber {
|
||||
);
|
||||
}
|
||||
|
||||
let dirty_relays = dirty_relays(&updates);
|
||||
|
||||
// Update RepoSyncIndex
|
||||
let mut index = self.repo_sync_index.write().await;
|
||||
|
||||
@@ -601,41 +608,22 @@ impl SelfSubscriber {
|
||||
);
|
||||
}
|
||||
|
||||
// Derive per-relay targets from the updated index
|
||||
let targets = derive_relay_targets(&index);
|
||||
drop(index); // Release lock before async operations
|
||||
drop(index);
|
||||
|
||||
// For each relay, send AddFilters action directly
|
||||
// SyncManager's handle_new_sync_filters auto-spawns connection for unknown relays
|
||||
for (relay_url, needs) in targets {
|
||||
// Send only the relay keys dirtied by this batch. The SyncManager
|
||||
// recomputes their current desired work against pending and confirmed
|
||||
// state, avoiding an O(all repositories × all relays) rebuild for
|
||||
// every historic batch.
|
||||
for relay_url in dirty_relays {
|
||||
// Skip our own relay URL (we're subscribed to ourselves via self-subscription)
|
||||
if relay_url.contains(&self.relay_domain) {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Build filters for these repos (sync-level-aware)
|
||||
let filters = crate::sync::filters::build_sync_level_aware_filters(
|
||||
&needs.repos,
|
||||
&needs.state_only_repos,
|
||||
&needs.root_events,
|
||||
None,
|
||||
);
|
||||
|
||||
// Log before moving values
|
||||
let repo_count = needs.repos.len() + needs.state_only_repos.len();
|
||||
let event_count = needs.root_events.len();
|
||||
|
||||
// Combine all repos into pending items
|
||||
let mut all_repos = needs.repos;
|
||||
all_repos.extend(needs.state_only_repos);
|
||||
|
||||
let action = AddFilters {
|
||||
relay_url: relay_url.clone(),
|
||||
items: crate::sync::PendingItems {
|
||||
repos: all_repos,
|
||||
root_events: needs.root_events,
|
||||
},
|
||||
filters,
|
||||
items: crate::sync::PendingItems::default(),
|
||||
filters: Vec::new(),
|
||||
};
|
||||
|
||||
if let Err(e) = self.action_tx.send(action).await {
|
||||
@@ -647,9 +635,7 @@ impl SelfSubscriber {
|
||||
} else {
|
||||
tracing::info!(
|
||||
relay = %relay_url,
|
||||
repo_count = repo_count,
|
||||
event_count = event_count,
|
||||
"Sent AddFilters action to SyncManager"
|
||||
"Marked relay dirty for SyncManager recomputation"
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -716,4 +702,28 @@ mod tests {
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dirty_relay_signals_are_limited_to_the_current_batch() {
|
||||
let updates = HashMap::from([(
|
||||
"30617:owner:new-repo".to_string(),
|
||||
RepoSyncNeeds {
|
||||
relays: HashSet::from([
|
||||
"wss://owner.example".to_string(),
|
||||
"wss://shared.example".to_string(),
|
||||
]),
|
||||
root_events: HashSet::new(),
|
||||
sync_level: SyncLevel::Full,
|
||||
},
|
||||
)]);
|
||||
|
||||
assert_eq!(
|
||||
dirty_relays(&updates),
|
||||
HashSet::from([
|
||||
"wss://owner.example".to_string(),
|
||||
"wss://shared.example".to_string(),
|
||||
]),
|
||||
"a fresh historic batch must not rebuild work for unrelated indexed relays"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user