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; + } +}