diff --git a/CHANGELOG.md b/CHANGELOG.md index d866b56..f247a7d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/src/nostr/lifecycle/deletion/runtime.rs b/src/nostr/lifecycle/deletion/runtime.rs index 6d162b5..4d6783a 100644 --- a/src/nostr/lifecycle/deletion/runtime.rs +++ b/src/nostr/lifecycle/deletion/runtime.rs @@ -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; + } +} diff --git a/src/sync/mod.rs b/src/sync/mod.rs index e5ad8ae..5f6c9b0 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -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::() % 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, + 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 { diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index fffec30..0153b03 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -77,6 +77,13 @@ impl PendingUpdates { } } +fn dirty_relays(updates: &HashMap) -> HashSet { + 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" + ); + } }