diff --git a/docs/explanation/repository-lifecycle.md b/docs/explanation/repository-lifecycle.md index 93c48e1..c2ae8b9 100644 --- a/docs/explanation/repository-lifecycle.md +++ b/docs/explanation/repository-lifecycle.md @@ -12,9 +12,10 @@ and purgatory transitions. > **Status:** Lifecycle metadata, NIP-09 and NIP-62 lifecycle admission > (including disrespector read-only would-have-deleted classification), -> deterministic admission-use attribution, and measured destructive outcomes are -> implemented. Request cleanup/expiry, migration, target-set deduplication -> removal, and final policy-mode reconciliation remain future stages. +> deterministic admission-use attribution, measured destructive outcomes, +> startup reconciliation, and periodic request cleanup/permanent expiry are +> implemented. Target-set deduplication removal, repair-command retirement, and +> final rollout coverage remain future work. ### Production motivation @@ -271,9 +272,16 @@ mode reapplies normal gating and the 9-month served plus 3-month unserved expiry schedule using retained lifecycle timestamps. Startup reconciliation must apply these transitions before the relay begins serving traffic. -Periodic request cleanup, permanent tombstone expiry, removal of target-set -deduplication and the `repair-deletion-requests` command, and final rollout -coverage remain future work. +Startup catch-up and periodic cleanup use the holding-cleanup cadence. They +apply the same timestamp-derived half-open lifecycle boundaries, remove expired +served copies from Main before Tombstones, and permanently remove payload plus +lifecycle metadata only after the additional gating period. A shared lifecycle +transition lock prevents cleanup from acting on a stale record while admission +is promoting and marking it used. + +Removal of target-set deduplication, retirement of the +`repair-deletion-requests` command, and final rollout coverage remain future +work. ## Core consistency invariant diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 16df4bc..578b041 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -113,6 +113,30 @@ lazy_static! { .expect("register holding cleanup last-run metric"); metric }; + static ref DELETION_REQUEST_CLEANUP_RUNS_TOTAL: Counter = { + let metric = Counter::with_opts(Opts::new( + "ngit_deletion_request_cleanup_runs_total", + "Number of deletion-request cleanup passes run", + )).expect("build deletion request cleanup runs metric"); + REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup runs metric"); + metric + }; + static ref DELETION_REQUEST_CLEANUP_REMOVED_TOTAL: CounterVec = { + let metric = CounterVec::new(Opts::new( + "ngit_deletion_request_cleanup_removed_total", + "Deletion-request cleanup removals by storage type", + ), &["type"]).expect("build deletion request cleanup removed metric"); + REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup removed metric"); + metric + }; + static ref DELETION_REQUEST_CLEANUP_OUTCOMES_TOTAL: CounterVec = { + let metric = CounterVec::new(Opts::new( + "ngit_deletion_request_cleanup_outcomes_total", + "Deletion-request cleanup failures and concurrent skips", + ), &["outcome"]).expect("build deletion request cleanup outcomes metric"); + REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup outcomes metric"); + metric + }; static ref RECOVERY_TOTAL: CounterVec = { let metric = CounterVec::new( Opts::new( @@ -247,6 +271,30 @@ pub fn record_holding_cleanup_run( .set(archive_files_deleted as f64); } +pub fn record_deletion_request_cleanup_run( + main_removed: usize, + tombstone_removed: usize, + metadata_removed: usize, + failures: usize, + skipped: usize, +) { + DELETION_REQUEST_CLEANUP_RUNS_TOTAL.inc(); + for (kind, value) in [ + ("main", main_removed), + ("tombstone", tombstone_removed), + ("metadata", metadata_removed), + ] { + DELETION_REQUEST_CLEANUP_REMOVED_TOTAL + .with_label_values(&[kind]) + .inc_by(value as f64); + } + for (outcome, value) in [("failure", failures), ("stale_or_concurrent_skip", skipped)] { + DELETION_REQUEST_CLEANUP_OUTCOMES_TOTAL + .with_label_values(&[outcome]) + .inc_by(value as f64); + } +} + pub fn record_recovery_attempt() { RECOVERY_TOTAL.with_label_values(&["attempted"]).inc(); } diff --git a/src/nostr/lifecycle/deletion/cleanup.rs b/src/nostr/lifecycle/deletion/cleanup.rs new file mode 100644 index 0000000..7c2eeb5 --- /dev/null +++ b/src/nostr/lifecycle/deletion/cleanup.rs @@ -0,0 +1,457 @@ +use anyhow::Result; +use nostr_relay_builder::prelude::{Filter, Kind, Timestamp}; + +use crate::nostr::lifecycle::RequestLifecycleRecord; + +use super::DeletionService; + +/// Result counters for one bounded deletion-request retention pass. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct RequestCleanupStats { + pub canonical_records_examined: usize, + pub main_payloads_removed: usize, + pub tombstone_payloads_removed: usize, + pub metadata_rows_removed: usize, + pub indefinitely_retained_disrespector_requests: usize, + pub stale_or_concurrent_records_skipped: usize, + pub failures: usize, +} + +impl DeletionService { + /// Reconcile served copies and permanently expire elapsed request lifecycles. + /// The supplied time makes every retention boundary deterministic in tests. + pub(crate) async fn cleanup_expired_requests( + &self, + now: Timestamp, + ) -> Result { + let records = self.ctx.tombstones().lifecycle_records_result().await?; + let mut stats = RequestCleanupStats { + canonical_records_examined: records.len(), + ..Default::default() + }; + + for selected in records { + if let Err(error) = self.cleanup_one_request(&selected, now, &mut stats).await { + stats.failures += 1; + tracing::warn!(request_id = %selected.request.id, error = %error, "Deletion-request cleanup failed for record"); + } + } + + match self.ctx.tombstones().remove_orphan_metadata().await { + Ok(removed) => stats.metadata_rows_removed += removed, + Err(error) => { + stats.failures += 1; + tracing::warn!(error = %error, "Deletion-request cleanup failed to remove orphan metadata"); + } + } + Ok(stats) + } + + async fn cleanup_one_request( + &self, + selected: &RequestLifecycleRecord, + now: Timestamp, + stats: &mut RequestCleanupStats, + ) -> Result<()> { + let tombstones = self.ctx.tombstones(); + let _guard = tombstones.lock_lifecycle().await; + let Some(current) = tombstones + .lifecycle_for_request_result(&selected.request.id) + .await? + else { + stats.stale_or_concurrent_records_skipped += 1; + return Ok(()); + }; + // A mark-used/reclassification replacement between enumeration and this + // lock acquisition means this pass must not act on its stale decision. + if current != *selected { + stats.stale_or_concurrent_records_skipped += 1; + return Ok(()); + } + + let targeting = current.request.kind == Kind::EventDeletion + || self.vanish_targets_this_relay(¤t.request); + if self.ctx.config.deletion_request_disrespector + && targeting + && current.last_used_at.is_some() + { + stats.indefinitely_retained_disrespector_requests += 1; + return Ok(()); + } + + if self.request_is_served(¤t, now)? { + return Ok(()); + } + + // Main must be removed before Tombstones. Holding the same lifecycle + // lock as gate attribution prevents stale cleanup from unserving a + // request that an admission just promoted and marked used. + if self + .ctx + .database() + .event_by_id(¤t.request.id) + .await? + .is_some() + { + self.ctx + .database() + .delete(Filter::new().ids(vec![current.request.id])) + .await?; + stats.main_payloads_removed += 1; + } + + if !self.request_is_expired(¤t, now)? { + return Ok(()); + } + let removed = tombstones + .permanently_delete_request_locked(¤t.request.id) + .await?; + stats.tombstone_payloads_removed += removed.payloads_deleted; + stats.metadata_rows_removed += removed.metadata_deleted; + Ok(()) + } + + pub(crate) fn request_is_served( + &self, + record: &RequestLifecycleRecord, + now: Timestamp, + ) -> Result { + let (anchor, duration) = match record.last_used_at { + Some(used) => ( + used, + self.ctx + .config + .deletion_request_retention_used_served_after_last_used(), + ), + None => ( + record.first_seen_at, + self.ctx.config.deletion_request_retention_unused_served(), + ), + }; + let deadline = anchor + .as_secs() + .checked_add(duration.as_secs()) + .ok_or_else(|| { + anyhow::anyhow!("deletion-request served retention deadline overflow") + })?; + Ok(now.as_secs() < deadline) + } +} + +#[cfg(test)] +mod tests { + use std::path::PathBuf; + use std::sync::Arc; + + use nostr_relay_builder::prelude::{Event, EventBuilder, EventId, FinalizeEvent, Keys, Tag}; + + use super::*; + use crate::grasp06::receive::new_repo_init_locks; + use crate::nostr::lifecycle::{ + HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle, RequestClassification, + Tombstones, + }; + use crate::purgatory::Purgatory; + + fn service(disrespector: bool) -> DeletionService { + let db = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let config = crate::config::Config { + deletion_request_disrespector: disrespector, + deletion_request_retention_unused_served_secs: 10, + deletion_request_retention_unused_unserved_gating_additional_secs: 5, + deletion_request_retention_used_served_after_last_used_secs: 20, + deletion_request_retention_used_unserved_gating_additional_secs: 5, + ..crate::config::Config::for_testing() + }; + DeletionService::new(super::super::DeletionContext::new( + "test.example.com", + db, + Tombstones::in_memory(), + HoldingStore::in_memory(), + RepositoryLifecycle::in_memory(), + ReplaceableHistoryStore::in_memory(), + PathBuf::new(), + Arc::new(Purgatory::new(PathBuf::new())), + config, + new_repo_init_locks(), + )) + } + + fn deletion() -> Event { + EventBuilder::new(Kind::EventDeletion, "") + .tags(vec![Tag::event(EventId::all_zeros())]) + .finalize(&Keys::generate()) + .unwrap() + } + + fn non_targeting_vanish() -> Event { + EventBuilder::new(Kind::RequestToVanish, "") + .tags(vec![Tag::custom( + "relay", + vec!["wss://elsewhere.example".to_owned()], + )]) + .finalize(&Keys::generate()) + .unwrap() + } + + async fn retained(service: &DeletionService, request: &Event, first_seen: u64) { + service + .ctx + .tombstones() + .record_request( + request, + Timestamp::from_secs(first_seen), + RequestClassification::LocallyActionable, + ) + .await + .unwrap(); + service.ctx.database().save_event(request).await.unwrap(); + } + + async fn main_has(service: &DeletionService, request: &Event) -> bool { + service + .ctx + .database() + .event_by_id(&request.id) + .await + .unwrap() + .is_some() + } + + #[tokio::test] + async fn unused_boundaries_keep_gate_then_permanently_expire() { + let service = service(false); + let request = deletion(); + retained(&service, &request, 0).await; + assert_eq!( + service + .cleanup_expired_requests(Timestamp::from_secs(9)) + .await + .unwrap() + .main_payloads_removed, + 0 + ); + assert!(main_has(&service, &request).await); + + let stats = service + .cleanup_expired_requests(Timestamp::from_secs(10)) + .await + .unwrap(); + assert_eq!(stats.main_payloads_removed, 1); + assert!(!main_has(&service, &request).await); + assert!(service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .is_some()); + assert_eq!( + service + .ctx + .tombstones() + .event_deletion_candidates(&EventId::all_zeros(), &request.pubkey) + .await + .unwrap() + .len(), + 1 + ); + assert_eq!( + service + .ctx + .tombstones() + .lifecycle_records_result() + .await + .unwrap() + .len(), + 1 + ); + + let stats = service + .cleanup_expired_requests(Timestamp::from_secs(15)) + .await + .unwrap(); + assert_eq!(stats.tombstone_payloads_removed, 1); + assert!(service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .is_none()); + } + + #[tokio::test] + async fn used_boundaries_and_reuse_reset_the_lifecycle() { + let service = service(false); + let request = deletion(); + retained(&service, &request, 0).await; + service + .ctx + .tombstones() + .mark_request_used(&request.id, Timestamp::from_secs(20)) + .await + .unwrap(); + assert!( + service + .cleanup_expired_requests(Timestamp::from_secs(39)) + .await + .unwrap() + .main_payloads_removed + == 0 + ); + service + .ctx + .tombstones() + .mark_request_used(&request.id, Timestamp::from_secs(39)) + .await + .unwrap(); + assert!( + service + .cleanup_expired_requests(Timestamp::from_secs(40)) + .await + .unwrap() + .main_payloads_removed + == 0 + ); + let stats = service + .cleanup_expired_requests(Timestamp::from_secs(59)) + .await + .unwrap(); + assert_eq!(stats.main_payloads_removed, 1); + assert!(service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .is_some()); + let stats = service + .cleanup_expired_requests(Timestamp::from_secs(64)) + .await + .unwrap(); + assert_eq!(stats.tombstone_payloads_removed, 1); + } + + #[tokio::test] + async fn disrespector_only_retains_used_targeting_requests_indefinitely() { + let service = service(true); + let used = deletion(); + retained(&service, &used, 0).await; + service + .ctx + .tombstones() + .mark_request_used(&used.id, Timestamp::from_secs(1)) + .await + .unwrap(); + let stats = service + .cleanup_expired_requests(Timestamp::from_secs(10_000)) + .await + .unwrap(); + assert_eq!(stats.indefinitely_retained_disrespector_requests, 1); + assert!(main_has(&service, &used).await); + + let unused = deletion(); + retained(&service, &unused, 0).await; + service + .cleanup_expired_requests(Timestamp::from_secs(15)) + .await + .unwrap(); + assert!(service + .ctx + .tombstones() + .lifecycle_for_request_result(&unused.id) + .await + .unwrap() + .is_none()); + } + + #[tokio::test] + async fn non_targeting_nip62_expires_and_cleanup_is_idempotent() { + let service = service(false); + let request = non_targeting_vanish(); + retained(&service, &request, 0).await; + assert!( + service + .cleanup_expired_requests(Timestamp::from_secs(15)) + .await + .unwrap() + .tombstone_payloads_removed + == 1 + ); + let rerun = service + .cleanup_expired_requests(Timestamp::from_secs(15)) + .await + .unwrap(); + assert_eq!( + rerun.main_payloads_removed + + rerun.tombstone_payloads_removed + + rerun.metadata_rows_removed, + 0 + ); + } + + #[tokio::test] + async fn stale_cleanup_snapshot_cannot_unserve_a_marked_used_request() { + let service = service(false); + let request = deletion(); + retained(&service, &request, 0).await; + let selected = service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .unwrap(); + // This models admission's metadata update and Main promotion happening + // after enumeration but before cleanup acquires the shared lock. + service + .ctx + .tombstones() + .mark_request_used(&request.id, Timestamp::from_secs(14)) + .await + .unwrap(); + let mut stats = RequestCleanupStats::default(); + service + .cleanup_one_request(&selected, Timestamp::from_secs(15), &mut stats) + .await + .unwrap(); + assert_eq!(stats.stale_or_concurrent_records_skipped, 1); + assert!(main_has(&service, &request).await); + } + + #[tokio::test] + async fn startup_reconciliation_and_periodic_cleanup_share_serving_boundary() { + let startup = service(false); + let periodic = service(false); + let startup_request = deletion(); + let periodic_request = deletion(); + retained(&startup, &startup_request, 0).await; + retained(&periodic, &periodic_request, 0).await; + + startup + .run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(10)) + .await + .unwrap(); + periodic + .cleanup_expired_requests(Timestamp::from_secs(10)) + .await + .unwrap(); + assert!(!main_has(&startup, &startup_request).await); + assert!(!main_has(&periodic, &periodic_request).await); + assert!(startup + .ctx + .tombstones() + .lifecycle_for_request_result(&startup_request.id) + .await + .unwrap() + .is_some()); + assert!(periodic + .ctx + .tombstones() + .lifecycle_for_request_result(&periodic_request.id) + .await + .unwrap() + .is_some()); + } +} diff --git a/src/nostr/lifecycle/deletion/mod.rs b/src/nostr/lifecycle/deletion/mod.rs index 447beea..7d6076d 100644 --- a/src/nostr/lifecycle/deletion/mod.rs +++ b/src/nostr/lifecycle/deletion/mod.rs @@ -1,5 +1,6 @@ mod archival; mod cascade; +mod cleanup; mod context; mod policy; mod pr_refs; @@ -42,6 +43,7 @@ impl DeletionOutcome { } } +pub use cleanup::RequestCleanupStats; pub use context::DeletionContext; pub use policy::DeletionPolicy; pub use runtime::{run_holding_eject, DeletionCleanupTask, DeletionRuntime, HoldingEjectArgs}; diff --git a/src/nostr/lifecycle/deletion/policy.rs b/src/nostr/lifecycle/deletion/policy.rs index f1c759e..5466e8e 100644 --- a/src/nostr/lifecycle/deletion/policy.rs +++ b/src/nostr/lifecycle/deletion/policy.rs @@ -334,16 +334,15 @@ impl DeletionPolicy { return true; } } - Kind::RepoState => { + Kind::RepoState if self .ctx .purgatory .find_state(&identifier) .into_iter() - .any(|entry| entry.author == owner && covered(entry.event.created_at)) - { - return true; - } + .any(|entry| entry.author == owner && covered(entry.event.created_at)) => + { + return true; } _ => {} } diff --git a/src/nostr/lifecycle/deletion/runtime.rs b/src/nostr/lifecycle/deletion/runtime.rs index b6b07dd..6d162b5 100644 --- a/src/nostr/lifecycle/deletion/runtime.rs +++ b/src/nostr/lifecycle/deletion/runtime.rs @@ -14,7 +14,7 @@ use super::startup::{ BlacklistParityStats, BlacklistRestoreStats, StartupReconciliationStats, WhitelistParityStats, WhitelistRestoreStats, }; -use super::DeletionService; +use super::{DeletionService, RequestCleanupStats}; #[derive(Debug, Args)] pub struct HoldingEjectArgs { @@ -63,6 +63,13 @@ 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(()) } @@ -71,6 +78,7 @@ impl DeletionRuntime { pub fn spawn_cleanup_task(&self) -> DeletionCleanupTask { let holding = self.holding.clone(); let lifecycle = self.lifecycle.clone(); + let service = self.service.clone(); let retention = self.holding_retention; let interval_duration = self.holding_cleanup_interval; let (shutdown_tx, mut shutdown_rx) = watch::channel(false); @@ -84,6 +92,10 @@ impl DeletionRuntime { loop { tokio::select! { _ = interval.tick() => { + match service.cleanup_expired_requests(Timestamp::now()).await { + Ok(stats) => log_request_cleanup("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 { Ok(stats) => { if stats.expired_records > 0 { @@ -103,7 +115,7 @@ impl DeletionRuntime { } changed = shutdown_rx.changed() => { if changed.is_ok() && *shutdown_rx.borrow() { - tracing::info!("Holding cleanup task received shutdown signal"); + tracing::info!("Deletion lifecycle maintenance task received shutdown signal"); break; } } @@ -114,7 +126,7 @@ impl DeletionRuntime { tracing::info!( retention_secs = retention.as_secs(), interval_secs = interval_duration.as_secs(), - "Holding cleanup task started" + "Deletion lifecycle maintenance task started" ); DeletionCleanupTask { @@ -160,7 +172,7 @@ impl DeletionCleanupTask { pub async fn shutdown(self) { let _ = self.shutdown_tx.send(true); if let Err(e) = self.handle.await { - tracing::warn!(error = %e, "Holding cleanup task join failed during shutdown"); + tracing::warn!(error = %e, "Deletion lifecycle maintenance task join failed during shutdown"); } } } @@ -200,6 +212,27 @@ fn log_startup_reconciliation(stats: StartupReconciliationStats) { log_whitelist_restore(stats.whitelist_restore); } +fn log_request_cleanup(phase: &str, stats: RequestCleanupStats) { + crate::metrics::record_deletion_request_cleanup_run( + stats.main_payloads_removed, + stats.tombstone_payloads_removed, + stats.metadata_rows_removed, + stats.failures, + stats.stale_or_concurrent_records_skipped, + ); + tracing::info!( + phase, + examined = stats.canonical_records_examined, + main_removed = stats.main_payloads_removed, + tombstone_removed = stats.tombstone_payloads_removed, + metadata_removed = stats.metadata_rows_removed, + indefinitely_retained = stats.indefinitely_retained_disrespector_requests, + skipped = stats.stale_or_concurrent_records_skipped, + failures = stats.failures, + "Deletion-request cleanup completed" + ); +} + fn log_blacklist_parity(stats: BlacklistParityStats) { if stats.scanned_announcements > 0 || stats.matched_announcements > 0 { tracing::info!( diff --git a/src/nostr/lifecycle/deletion/service.rs b/src/nostr/lifecycle/deletion/service.rs index b16a5f1..10b05a5 100644 --- a/src/nostr/lifecycle/deletion/service.rs +++ b/src/nostr/lifecycle/deletion/service.rs @@ -93,13 +93,44 @@ impl DeletionService { .into_values() .min_by(|(left, _), (right, _)| Self::winner_order(left, right))?; + // Serialize promotion plus last-used attribution with request cleanup. + // Cleanup always re-reads metadata under this same lock, so it cannot + // remove the Main copy selected by an admission using an old snapshot. + let _lifecycle_guard = tombstones.lock_lifecycle().await; + let Some(winner) = (match tombstones + .lifecycle_for_request_result(&winner.request.id) + .await + { + Ok(record) => record, + Err(error) => { + return Some(Self::gate_error( + "re-reading deletion-request lifecycle", + error, + )) + } + }) else { + return Some(Self::gate_error( + "re-reading deletion-request lifecycle", + anyhow::anyhow!("request {} disappeared", winner.request.id), + )); + }; + match self.request_is_expired(&winner, Timestamp::now()) { + Ok(true) => return None, + Ok(false) => {} + Err(error) => { + return Some(Self::gate_error( + "evaluating deletion-request retention", + error, + )) + } + } let previous_last_used_at = winner.last_used_at; if let Err(error) = self.ctx.database().save_event(&winner.request).await { return Some(Self::gate_error("promoting deletion request", error)); } let used_at = Timestamp::now(); let updated = match tombstones - .mark_request_used(&winner.request.id, used_at) + .mark_request_used_locked(&winner.request.id, used_at) .await { Ok(Some(record)) if record.last_used_at >= Some(used_at) => record, @@ -149,7 +180,11 @@ impl DeletionService { }) } - fn request_is_expired(&self, record: &RequestLifecycleRecord, now: Timestamp) -> Result { + pub(crate) fn request_is_expired( + &self, + record: &RequestLifecycleRecord, + now: Timestamp, + ) -> Result { let (anchor, served, additional) = match record.last_used_at { Some(last_used_at) => ( last_used_at, diff --git a/src/nostr/lifecycle/deletion/startup.rs b/src/nostr/lifecycle/deletion/startup.rs index 04bca8b..f2159c9 100644 --- a/src/nostr/lifecycle/deletion/startup.rs +++ b/src/nostr/lifecycle/deletion/startup.rs @@ -5,7 +5,7 @@ use anyhow::Result; use nostr_relay_builder::prelude::{Event, Filter, Kind, PublicKey, Timestamp}; use crate::nostr::events::RepositoryAnnouncement; -use crate::nostr::lifecycle::{DeletionSource, RequestClassification, RequestLifecycleRecord}; +use crate::nostr::lifecycle::{DeletionSource, RequestClassification}; use super::DeletionService; @@ -248,28 +248,6 @@ impl DeletionService { } } - fn request_is_served(&self, record: &RequestLifecycleRecord, now: Timestamp) -> Result { - let (anchor, duration) = match record.last_used_at { - Some(used) => ( - used, - self.ctx - .config - .deletion_request_retention_used_served_after_last_used(), - ), - None => ( - record.first_seen_at, - self.ctx.config.deletion_request_retention_unused_served(), - ), - }; - let deadline = anchor - .as_secs() - .checked_add(duration.as_secs()) - .ok_or_else(|| { - anyhow::anyhow!("deletion-request served retention deadline overflow") - })?; - Ok(now.as_secs() < deadline) - } - async fn promote_request( &self, request: &Event, diff --git a/src/nostr/lifecycle/tombstones.rs b/src/nostr/lifecycle/tombstones.rs index 73544e9..9d18e31 100644 --- a/src/nostr/lifecycle/tombstones.rs +++ b/src/nostr/lifecycle/tombstones.rs @@ -45,6 +45,8 @@ use std::collections::HashSet; use std::path::Path; use std::sync::Arc; +use tokio::sync::{Mutex, MutexGuard}; + use nostr_lmdb::NostrLmdb; use nostr_memory::MemoryDatabase; use nostr_relay_builder::prelude::{ @@ -101,6 +103,13 @@ pub struct RequestLifecycleRecord { pub classification: RequestClassification, } +/// Counts returned after permanently removing one request lifecycle. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct RequestLifecycleDeletionStats { + pub payloads_deleted: usize, + pub metadata_deleted: usize, +} + #[derive(Debug, Default)] struct DeletionTargets { e_ids: HashSet, @@ -141,6 +150,11 @@ impl DeletionTargets { pub struct Tombstones { db: Arc, metadata_signer: Keys, + /// Serializes lifecycle metadata changes with Main-database transitions + /// performed by `DeletionService`. This is intentionally one simple lock: + /// request lifecycles are a small control plane and correctness matters more + /// than parallel metadata writes. + lifecycle_lock: Arc>, } impl std::fmt::Debug for Tombstones { @@ -179,6 +193,7 @@ impl Tombstones { Ok(Self { db: Arc::new(db), metadata_signer: Keys::generate(), + lifecycle_lock: Arc::new(Mutex::new(())), }) } @@ -193,9 +208,17 @@ impl Tombstones { .build(), ), metadata_signer: Keys::generate(), + lifecycle_lock: Arc::new(Mutex::new(())), } } + /// Acquire the lifecycle transition lock. Callers that also transition the + /// Main database must hold this across the Main operation and the matching + /// metadata update/removal. + pub(crate) async fn lock_lifecycle(&self) -> MutexGuard<'_, ()> { + self.lifecycle_lock.lock().await + } + #[cfg(test)] pub(crate) async fn save_request_payload_without_metadata( &self, @@ -215,6 +238,17 @@ impl Tombstones { event: &Event, first_seen_at: Timestamp, classification: RequestClassification, + ) -> anyhow::Result { + let _guard = self.lock_lifecycle().await; + self.record_request_locked(event, first_seen_at, classification) + .await + } + + pub(crate) async fn record_request_locked( + &self, + event: &Event, + first_seen_at: Timestamp, + classification: RequestClassification, ) -> anyhow::Result { if !Self::is_request_kind(event.kind) { anyhow::bail!("Tombstone lifecycle accepts only kind-5 or kind-62 requests"); @@ -339,6 +373,15 @@ impl Tombstones { &self, request_id: &EventId, used_at: Timestamp, + ) -> anyhow::Result> { + let _guard = self.lock_lifecycle().await; + self.mark_request_used_locked(request_id, used_at).await + } + + pub(crate) async fn mark_request_used_locked( + &self, + request_id: &EventId, + used_at: Timestamp, ) -> anyhow::Result> { let Some(existing) = self.lifecycle_for_request_result(request_id).await? else { return Ok(None); @@ -364,6 +407,16 @@ impl Tombstones { &self, request_id: &EventId, classification: RequestClassification, + ) -> anyhow::Result> { + let _guard = self.lock_lifecycle().await; + self.reclassify_request_locked(request_id, classification) + .await + } + + pub(crate) async fn reclassify_request_locked( + &self, + request_id: &EventId, + classification: RequestClassification, ) -> anyhow::Result> { let Some(existing) = self.lifecycle_for_request_result(request_id).await? else { return Ok(None); @@ -402,6 +455,86 @@ impl Tombstones { records } + /// Fallible enumeration for destructive maintenance. Unlike + /// `lifecycle_records`, this never turns a database failure into an + /// apparently empty result. Private metadata is not returned. + pub async fn lifecycle_records_result(&self) -> anyhow::Result> { + let requests = self.request_payloads().await?; + let mut records = Vec::new(); + for request in requests { + if let Some(record) = self.lifecycle_for_request_result(&request.id).await? { + records.push(record); + } + } + records.sort_by_key(|record| record.request.id.to_hex()); + Ok(records) + } + + /// Remove one payload and every metadata row linked to it. This is + /// idempotent when either side was already removed. The caller must hold + /// `lock_lifecycle` and has already completed any required Main removal. + pub(crate) async fn permanently_delete_request_locked( + &self, + request_id: &EventId, + ) -> anyhow::Result { + let payloads_deleted = usize::from(self.request_payload(request_id).await?.is_some()); + let metadata = self.metadata_events_for_request_result(request_id).await?; + if payloads_deleted > 0 { + self.db + .delete(Filter::new().ids(vec![*request_id])) + .await + .map_err(|error| { + anyhow::anyhow!("Failed to remove tombstone request {request_id}: {error}") + })?; + } + let metadata_deleted = metadata.len(); + if metadata_deleted > 0 { + self.db + .delete(Filter::new().ids(metadata.into_iter().map(|event| event.id))) + .await + .map_err(|error| { + anyhow::anyhow!("Failed to remove lifecycle metadata for {request_id}: {error}") + })?; + } + Ok(RequestLifecycleDeletionStats { + payloads_deleted, + metadata_deleted, + }) + } + + /// Remove metadata whose unambiguous `e` link has no retained request. + /// This is deliberately separate from canonical enumeration: malformed + /// metadata remains invisible to queries but is not guessed at. + pub async fn remove_orphan_metadata(&self) -> anyhow::Result { + let _guard = self.lock_lifecycle().await; + let metadata = self + .db + .query(Filter::new().kind(Kind::from(TOMBSTONE_REQUEST_METADATA_KIND))) + .await + .map_err(|error| anyhow::anyhow!("Failed to enumerate tombstone metadata: {error}"))?; + let mut stale = Vec::new(); + for event in metadata { + let Some(request_id) = + Self::tag_value(&event, "e").and_then(|value| EventId::from_hex(value).ok()) + else { + continue; + }; + if self.request_payload(&request_id).await?.is_none() { + stale.push(event.id); + } + } + let removed = stale.len(); + if removed > 0 { + self.db + .delete(Filter::new().ids(stale)) + .await + .map_err(|error| { + anyhow::anyhow!("Failed to remove orphan lifecycle metadata: {error}") + })?; + } + Ok(removed) + } + /// List the original kind-5 and kind-62 payloads retained by this store. /// This deliberately does not require lifecycle metadata, and never exposes /// relay-private metadata events to callers.