From e780185babe4dce2adefcae5d5fab09cf5af90a4 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 16 Jul 2026 20:08:42 +0100 Subject: [PATCH] feat: reconcile deletion requests at startup --- docs/explanation/repository-lifecycle.md | 23 ++ src/nostr/lifecycle/deletion/mod.rs | 4 +- src/nostr/lifecycle/deletion/policy.rs | 2 +- src/nostr/lifecycle/deletion/runtime.rs | 5 +- src/nostr/lifecycle/deletion/service.rs | 13 +- src/nostr/lifecycle/deletion/startup.rs | 401 ++++++++++++++++++++++- src/nostr/lifecycle/tombstones.rs | 167 ++++++++-- src/server.rs | 5 +- 8 files changed, 576 insertions(+), 44 deletions(-) diff --git a/docs/explanation/repository-lifecycle.md b/docs/explanation/repository-lifecycle.md index 4d822fc..93c48e1 100644 --- a/docs/explanation/repository-lifecycle.md +++ b/docs/explanation/repository-lifecycle.md @@ -244,6 +244,25 @@ request a fresh 30-day probation window. This deliberately favors preservation over immediate production cleanup; normal cleanup moves requests that remain unused out of the served main database after that window. +Startup performs this migration before blacklist/whitelist reconciliation and +before the relay starts serving traffic. It discovers signed kind-5 and kind-62 +payloads from both the served and tombstone databases, including tombstone +payloads whose metadata is missing or malformed, and deduplicates by signed +event ID. Valid existing metadata is preserved; missing metadata is written +with one startup timestamp and the original payload is promoted to the served +database for its fresh probation. The pass is restart-safe: replacement +metadata preserves the earliest `first_seen_at` and greatest `last_used_at`. + +The same startup pass reclassifies retained requests using current relay +configuration and reconciles served copies at the exact lifecycle deadline. +In disrespector mode it read-only evaluates unused targeting requests against +main and purgatory data, marking only requests that would currently have an +effect as used; it never changes target data during that evaluation. Used +targeting requests remain served in that mode, while non-targeting NIP-62 and +unused no-op requests retain their ordinary bounded served schedule. Critical +query, metadata, promotion, or removal failures fail startup rather than +allowing traffic to begin with an incomplete lifecycle reconciliation. + The relay's current deletion-disrespector configuration governs retained requests; receipt-time mode is not permanent metadata. Switching into disrespector mode removes their local gates and keeps requests that are used or @@ -252,6 +271,10 @@ 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. + ## Core consistency invariant ngit-grasp enforces the following invariant across admission, serving, deletion, diff --git a/src/nostr/lifecycle/deletion/mod.rs b/src/nostr/lifecycle/deletion/mod.rs index 84b4a79..447beea 100644 --- a/src/nostr/lifecycle/deletion/mod.rs +++ b/src/nostr/lifecycle/deletion/mod.rs @@ -47,8 +47,8 @@ pub use policy::DeletionPolicy; pub use runtime::{run_holding_eject, DeletionCleanupTask, DeletionRuntime, HoldingEjectArgs}; pub use service::DeletionService; pub use startup::{ - BlacklistParityStats, BlacklistRestoreStats, StartupReconciliationStats, WhitelistParityStats, - WhitelistRestoreStats, + BlacklistParityStats, BlacklistRestoreStats, RequestLifecycleStartupStats, + StartupReconciliationStats, WhitelistParityStats, WhitelistRestoreStats, }; #[cfg(test)] diff --git a/src/nostr/lifecycle/deletion/policy.rs b/src/nostr/lifecycle/deletion/policy.rs index cda533e..f1c759e 100644 --- a/src/nostr/lifecycle/deletion/policy.rs +++ b/src/nostr/lifecycle/deletion/policy.rs @@ -253,7 +253,7 @@ impl DeletionPolicy { /// Read-only normal-policy target evaluation for disrespector mode. /// This intentionally only queries the main database and purgatory snapshots; /// it must not call deletion, holding, archive, rollback, or cascade helpers. - async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result { + pub(super) async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result { for id in Self::e_tag_ids(event) { if self.ctx.database.event_by_id(&id).await?.is_some() { return Ok(true); diff --git a/src/nostr/lifecycle/deletion/runtime.rs b/src/nostr/lifecycle/deletion/runtime.rs index 9701bbb..b6b07dd 100644 --- a/src/nostr/lifecycle/deletion/runtime.rs +++ b/src/nostr/lifecycle/deletion/runtime.rs @@ -61,9 +61,10 @@ impl DeletionRuntime { } /// Run deletion-owned startup tasks before the relay begins serving traffic. - pub async fn run_startup_tasks(&self) { - log_startup_reconciliation(self.service.run_startup_reconciliation().await); + pub async fn run_startup_tasks(&self) -> Result<()> { + log_startup_reconciliation(self.service.run_startup_reconciliation().await?); self.run_holding_startup_cleanup().await; + Ok(()) } /// Spawn deletion-owned background maintenance tasks. diff --git a/src/nostr/lifecycle/deletion/service.rs b/src/nostr/lifecycle/deletion/service.rs index 7c49c2f..b16a5f1 100644 --- a/src/nostr/lifecycle/deletion/service.rs +++ b/src/nostr/lifecycle/deletion/service.rs @@ -408,7 +408,7 @@ impl DeletionService { DeletionAdmissionHooks { deletion: self } } - fn vanish_targets_this_relay(&self, event: &Event) -> bool { + pub(crate) fn vanish_targets_this_relay(&self, event: &Event) -> bool { let relay_urls = Self::relay_url_candidates(self.ctx.domain()); if relay_urls.is_empty() { @@ -422,6 +422,17 @@ impl DeletionService { }) } + pub(crate) async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result { + self.policy.would_delete_stored_target(event).await + } + + pub(crate) async fn would_nip62_vanish_stored_data( + &self, + event: &Event, + ) -> anyhow::Result { + self.policy.would_nip62_vanish_stored_data(event).await + } + fn relay_url_candidates(domain: &str) -> Vec { let domain = domain.trim().trim_end_matches('/'); if domain.is_empty() { diff --git a/src/nostr/lifecycle/deletion/startup.rs b/src/nostr/lifecycle/deletion/startup.rs index 2c99caa..04bca8b 100644 --- a/src/nostr/lifecycle/deletion/startup.rs +++ b/src/nostr/lifecycle/deletion/startup.rs @@ -1,8 +1,11 @@ use nostr::nips::nip19::ToBech32; -use nostr_relay_builder::prelude::{Filter, Kind, PublicKey, Timestamp}; +use std::collections::HashMap; + +use anyhow::Result; +use nostr_relay_builder::prelude::{Event, Filter, Kind, PublicKey, Timestamp}; use crate::nostr::events::RepositoryAnnouncement; -use crate::nostr::lifecycle::DeletionSource; +use crate::nostr::lifecycle::{DeletionSource, RequestClassification, RequestLifecycleRecord}; use super::DeletionService; @@ -48,12 +51,30 @@ pub struct WhitelistRestoreStats { #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct StartupReconciliationStats { + pub request_lifecycle: RequestLifecycleStartupStats, pub blacklist_parity: BlacklistParityStats, pub blacklist_restore: BlacklistRestoreStats, pub whitelist_parity: WhitelistParityStats, pub whitelist_restore: WhitelistRestoreStats, } +/// Startup migration and policy reconciliation counters for retained deletion +/// requests. A failure is returned to the caller rather than represented as a +/// successful zero-count pass. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct RequestLifecycleStartupStats { + pub main_payloads_scanned: usize, + pub tombstone_payloads_scanned: usize, + pub unique_requests_discovered: usize, + pub requests_migrated: usize, + pub valid_metadata_existing: usize, + pub requests_reclassified: usize, + pub disrespector_requests_newly_used: usize, + pub main_requests_promoted: usize, + pub main_requests_removed: usize, + pub failures: usize, +} + // --------------------------------------------------------------------------- // DeletionService startup pass implementations // --------------------------------------------------------------------------- @@ -63,22 +84,229 @@ impl DeletionService { /// /// Ordering is intentionally owned here so relay startup only depends on a /// single facade while deletion internals decide sequencing: - /// 1) blacklist parity delete - /// 2) blacklist auto-restore - /// 3) whitelist parity delete - /// 4) whitelist restore - pub async fn run_startup_reconciliation(&self) -> StartupReconciliationStats { + /// 1) deletion request lifecycle migration and reconciliation + /// 2) blacklist parity delete + /// 3) blacklist auto-restore + /// 4) whitelist parity delete + /// 5) whitelist restore + pub async fn run_startup_reconciliation(&self) -> Result { + let request_lifecycle = self + .run_request_lifecycle_startup_reconciliation(Timestamp::now()) + .await?; let blacklist_parity = self.run_startup_blacklist_parity_pass().await; let blacklist_restore = self.run_startup_blacklist_restore_pass().await; let whitelist_parity = self.run_startup_whitelist_parity_pass().await; let whitelist_restore = self.run_startup_whitelist_restore_pass().await; - StartupReconciliationStats { + Ok(StartupReconciliationStats { + request_lifecycle, blacklist_parity, blacklist_restore, whitelist_parity, whitelist_restore, + }) + } + + /// Migrate historical request payloads and reconcile their current-policy + /// serving state. The timestamp is supplied for deterministic tests and is + /// captured once by the production startup pass. + pub(crate) async fn run_request_lifecycle_startup_reconciliation( + &self, + now: Timestamp, + ) -> Result { + let main = self.request_payloads_from_main().await?; + let tombstone_payloads = self.ctx.tombstones().request_payloads().await?; + let mut stats = RequestLifecycleStartupStats { + main_payloads_scanned: main.len(), + tombstone_payloads_scanned: tombstone_payloads.len(), + ..Default::default() + }; + let mut requests: HashMap<_, Event> = HashMap::new(); + for request in main.into_iter().chain(tombstone_payloads) { + requests.entry(request.id).or_insert(request); } + stats.unique_requests_discovered = requests.len(); + + for request in requests.into_values() { + let existing = self + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await?; + let mut record = match existing { + Some(record) => { + stats.valid_metadata_existing += 1; + record + } + None => { + let classification = self.current_classification(&request); + let record = self + .ctx + .tombstones() + .record_request(&request, now, classification) + .await?; + self.promote_request(&request, &mut stats).await?; + stats.requests_migrated += 1; + record + } + }; + + let desired = self.current_classification(&record.request); + if record.classification != desired { + record = self + .ctx + .tombstones() + .reclassify_request(&record.request.id, desired) + .await? + .ok_or_else(|| { + anyhow::anyhow!( + "lifecycle disappeared during reclassification: {}", + record.request.id + ) + })?; + stats.requests_reclassified += 1; + } + + let targeting = record.request.kind == Kind::EventDeletion + || self.vanish_targets_this_relay(&record.request); + if self.ctx.config.deletion_request_disrespector + && targeting + && record.last_used_at.is_none() + { + let would_use = if record.request.kind == Kind::EventDeletion { + self.would_delete_stored_target(&record.request).await? + } else { + self.would_nip62_vanish_stored_data(&record.request).await? + }; + if would_use { + record = self + .ctx + .tombstones() + .mark_request_used(&record.request.id, now) + .await? + .ok_or_else(|| { + anyhow::anyhow!( + "lifecycle disappeared while marking used: {}", + record.request.id + ) + })?; + stats.disrespector_requests_newly_used += 1; + } + } + + let serve_indefinitely = self.ctx.config.deletion_request_disrespector + && targeting + && record.last_used_at.is_some(); + if serve_indefinitely || self.request_is_served(&record, now)? { + self.promote_request(&record.request, &mut stats).await?; + } else { + self.remove_request_from_main(&record.request, &mut stats) + .await?; + } + } + tracing::info!( + main_payloads_scanned = stats.main_payloads_scanned, + tombstone_payloads_scanned = stats.tombstone_payloads_scanned, + unique_requests = stats.unique_requests_discovered, + migrated = stats.requests_migrated, + existing_metadata = stats.valid_metadata_existing, + reclassified = stats.requests_reclassified, + disrespector_newly_used = stats.disrespector_requests_newly_used, + promoted = stats.main_requests_promoted, + removed = stats.main_requests_removed, + failures = stats.failures, + "Deletion-request startup lifecycle reconciliation completed" + ); + Ok(stats) + } + + async fn request_payloads_from_main(&self) -> Result> { + let mut requests = Vec::new(); + for kind in [Kind::EventDeletion, Kind::RequestToVanish] { + requests.extend( + self.ctx + .database() + .query(Filter::new().kind(kind)) + .await + .map_err(|e| { + anyhow::anyhow!( + "Failed to list main database deletion-request payloads: {e}" + ) + })?, + ); + } + Ok(requests) + } + + fn current_classification(&self, request: &Event) -> RequestClassification { + if request.kind == Kind::RequestToVanish && !self.vanish_targets_this_relay(request) { + RequestClassification::NonTargetingNip62 + } else if self.ctx.config.deletion_request_disrespector { + RequestClassification::Disrespector + } else { + RequestClassification::LocallyActionable + } + } + + 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, + stats: &mut RequestLifecycleStartupStats, + ) -> Result<()> { + let exists = self + .ctx + .database() + .event_by_id(&request.id) + .await? + .is_some(); + if !exists { + self.ctx.database().save_event(request).await?; + stats.main_requests_promoted += 1; + } + Ok(()) + } + + async fn remove_request_from_main( + &self, + request: &Event, + stats: &mut RequestLifecycleStartupStats, + ) -> Result<()> { + let exists = self + .ctx + .database() + .event_by_id(&request.id) + .await? + .is_some(); + if exists { + self.ctx + .database() + .delete(Filter::new().ids(vec![request.id])) + .await?; + stats.main_requests_removed += 1; + } + Ok(()) } /// Startup-only blacklist parity pass. @@ -475,3 +703,160 @@ impl DeletionService { stats } } + +#[cfg(test)] +mod tests { + use std::path::PathBuf; + use std::sync::Arc; + + use nostr_relay_builder::prelude::{EventBuilder, FinalizeEvent, Keys}; + + use super::super::DeletionContext; + use super::*; + use crate::grasp06::receive::new_repo_init_locks; + use crate::nostr::lifecycle::{ + HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle, 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_used_served_after_last_used_secs: 20, + ..crate::config::Config::for_testing() + }; + DeletionService::new(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(keys: &Keys) -> Event { + EventBuilder::new(Kind::EventDeletion, "") + .finalize(keys) + .unwrap() + } + + #[tokio::test] + async fn migrates_main_only_request_with_fresh_probation_idempotently() { + let service = service(false); + let request = deletion(&Keys::generate()); + service.ctx.database.save_event(&request).await.unwrap(); + + let stats = service + .run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(100)) + .await + .unwrap(); + assert_eq!(stats.requests_migrated, 1); + assert_eq!(stats.main_requests_promoted, 0); + let record = service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .unwrap(); + assert_eq!(record.first_seen_at, Timestamp::from_secs(100)); + assert_eq!( + record.classification, + RequestClassification::LocallyActionable + ); + + let rerun = service + .run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(200)) + .await + .unwrap(); + assert_eq!(rerun.requests_migrated, 0); + assert_eq!( + service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .unwrap() + .first_seen_at, + Timestamp::from_secs(100) + ); + } + + #[tokio::test] + async fn tombstone_only_request_is_migrated_and_promoted() { + let service = service(false); + let request = deletion(&Keys::generate()); + service + .ctx + .tombstones() + .save_request_payload_without_metadata(&request) + .await + .unwrap(); + + let stats = service + .run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(100)) + .await + .unwrap(); + assert_eq!(stats.requests_migrated, 1); + assert_eq!(stats.main_requests_promoted, 1); + assert!(service + .ctx + .database + .event_by_id(&request.id) + .await + .unwrap() + .is_some()); + } + + #[tokio::test] + async fn reclassifies_used_request_without_erasing_timestamps() { + let service = service(true); + let request = deletion(&Keys::generate()); + service + .ctx + .tombstones() + .record_request( + &request, + Timestamp::from_secs(10), + RequestClassification::LocallyActionable, + ) + .await + .unwrap(); + service + .ctx + .tombstones() + .mark_request_used(&request.id, Timestamp::from_secs(20)) + .await + .unwrap(); + + service + .run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(30)) + .await + .unwrap(); + let record = service + .ctx + .tombstones() + .lifecycle_for_request_result(&request.id) + .await + .unwrap() + .unwrap(); + assert_eq!(record.first_seen_at, Timestamp::from_secs(10)); + assert_eq!(record.last_used_at, Some(Timestamp::from_secs(20))); + assert_eq!(record.classification, RequestClassification::Disrespector); + assert!(service + .ctx + .database + .event_by_id(&request.id) + .await + .unwrap() + .is_some()); + } +} diff --git a/src/nostr/lifecycle/tombstones.rs b/src/nostr/lifecycle/tombstones.rs index f8b4651..73544e9 100644 --- a/src/nostr/lifecycle/tombstones.rs +++ b/src/nostr/lifecycle/tombstones.rs @@ -196,6 +196,18 @@ impl Tombstones { } } + #[cfg(test)] + pub(crate) async fn save_request_payload_without_metadata( + &self, + event: &Event, + ) -> anyhow::Result<()> { + self.db + .save_event(event) + .await + .map(|_| ()) + .map_err(|error| anyhow::anyhow!("Failed to save test request payload: {error}")) + } + /// Persist a signed deletion/vanish request and relay-generated lifecycle /// metadata. Exact replays retain their original `first_seen_at`. pub async fn record_request( @@ -208,7 +220,7 @@ impl Tombstones { anyhow::bail!("Tombstone lifecycle accepts only kind-5 or kind-62 requests"); } - if let Some(record) = self.lifecycle_for_request(&event.id).await { + if let Some(record) = self.lifecycle_for_request_result(&event.id).await? { return Ok(record); } @@ -328,7 +340,7 @@ impl Tombstones { request_id: &EventId, used_at: Timestamp, ) -> anyhow::Result> { - let Some(existing) = self.lifecycle_for_request(request_id).await else { + let Some(existing) = self.lifecycle_for_request_result(request_id).await? else { return Ok(None); }; let last_used_at = existing.last_used_at.max(Some(used_at)); @@ -336,37 +348,36 @@ impl Tombstones { return Ok(Some(existing)); } - let replacement = self.build_metadata( - *request_id, + self.replace_lifecycle_metadata( + request_id, existing.first_seen_at, last_used_at, existing.classification, - )?; - self.db.save_event(&replacement).await.map_err(|e| { - anyhow::anyhow!("Failed to update lifecycle metadata for {request_id}: {e}") - })?; - let updated = self - .lifecycle_for_request(request_id) - .await - .ok_or_else(|| { - anyhow::anyhow!("Updated lifecycle metadata for {request_id} was not readable") - })?; - let stale_ids: Vec<_> = self - .metadata_events_for_request(request_id) - .await - .into_iter() - .map(|event| event.id) - .filter(|id| *id != updated.metadata_event_id) - .collect(); - if !stale_ids.is_empty() { - self.db - .delete(Filter::new().ids(stale_ids)) - .await - .map_err(|e| { - anyhow::anyhow!("Failed to compact lifecycle metadata for {request_id}: {e}") - })?; + ) + .await + } + + /// Update a request's policy classification without changing its lifecycle + /// timestamps. The replacement is durable before stale rows are compacted, + /// so an interrupted update always leaves a readable lifecycle record. + pub async fn reclassify_request( + &self, + request_id: &EventId, + classification: RequestClassification, + ) -> anyhow::Result> { + let Some(existing) = self.lifecycle_for_request_result(request_id).await? else { + return Ok(None); + }; + if existing.classification == classification { + return Ok(Some(existing)); } - Ok(self.lifecycle_for_request(request_id).await) + self.replace_lifecycle_metadata( + request_id, + existing.first_seen_at, + existing.last_used_at, + classification, + ) + .await } /// List canonical request lifecycles for migration and retention work. @@ -391,10 +402,89 @@ impl Tombstones { records } + /// 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. + pub async fn request_payloads(&self) -> anyhow::Result> { + let mut requests = Vec::new(); + for kind in [Kind::EventDeletion, Kind::RequestToVanish] { + let events = self + .db + .query(Filter::new().kind(kind)) + .await + .map_err(|error| { + anyhow::anyhow!("Failed to list tombstone request payloads: {error}") + })?; + requests.extend(events); + } + requests.sort_by_key(|event| event.id.to_hex()); + requests.dedup_by_key(|event| event.id); + Ok(requests) + } + fn is_request_kind(kind: Kind) -> bool { kind == Kind::EventDeletion || kind == Kind::RequestToVanish } + async fn replace_lifecycle_metadata( + &self, + request_id: &EventId, + first_seen_at: Timestamp, + last_used_at: Option, + classification: RequestClassification, + ) -> anyhow::Result> { + // Metadata timestamps only order relay-private replacement rows. Ensure + // this row is newer even when multiple updates happen in one wall-clock + // second, otherwise canonical selection could retain the old class. + let newest_existing = self + .metadata_events_for_request_result(request_id) + .await? + .iter() + .map(|event| event.created_at) + .max(); + let replacement_created_at = match newest_existing { + Some(existing) if existing >= Timestamp::now() => Timestamp::from_secs( + existing + .as_secs() + .checked_add(1) + .ok_or_else(|| anyhow::anyhow!("lifecycle metadata timestamp overflow"))?, + ), + _ => Timestamp::now(), + }; + let replacement = self.build_metadata_at( + *request_id, + first_seen_at, + last_used_at, + classification, + replacement_created_at, + )?; + self.db.save_event(&replacement).await.map_err(|e| { + anyhow::anyhow!("Failed to update lifecycle metadata for {request_id}: {e}") + })?; + let updated = self + .lifecycle_for_request_result(request_id) + .await? + .ok_or_else(|| { + anyhow::anyhow!("Updated lifecycle metadata for {request_id} was not readable") + })?; + let stale_ids: Vec<_> = self + .metadata_events_for_request_result(request_id) + .await? + .into_iter() + .map(|event| event.id) + .filter(|id| *id != updated.metadata_event_id) + .collect(); + if !stale_ids.is_empty() { + self.db + .delete(Filter::new().ids(stale_ids)) + .await + .map_err(|e| { + anyhow::anyhow!("Failed to compact lifecycle metadata for {request_id}: {e}") + })?; + } + self.lifecycle_for_request_result(request_id).await + } + async fn request_payload(&self, request_id: &EventId) -> anyhow::Result> { for kind in [Kind::EventDeletion, Kind::RequestToVanish] { let events = self.db.query(Filter::new().kind(kind)).await.map_err(|e| { @@ -413,6 +503,23 @@ impl Tombstones { first_seen_at: Timestamp, last_used_at: Option, classification: RequestClassification, + ) -> anyhow::Result { + self.build_metadata_at( + request_id, + first_seen_at, + last_used_at, + classification, + Timestamp::now(), + ) + } + + fn build_metadata_at( + &self, + request_id: EventId, + first_seen_at: Timestamp, + last_used_at: Option, + classification: RequestClassification, + created_at: Timestamp, ) -> anyhow::Result { let mut tags = vec![ Tag::event(request_id), @@ -433,10 +540,12 @@ impl Tombstones { } EventBuilder::new(Kind::from(TOMBSTONE_REQUEST_METADATA_KIND), "") .tags(tags) + .custom_created_at(created_at) .finalize(&self.metadata_signer) .map_err(|e| anyhow::anyhow!("Failed to build tombstone lifecycle metadata: {e}")) } + #[cfg(test)] async fn metadata_events_for_request(&self, request_id: &EventId) -> Vec { match self.metadata_events_for_request_result(request_id).await { Ok(events) => events, diff --git a/src/server.rs b/src/server.rs index d47a831..ff6cd24 100644 --- a/src/server.rs +++ b/src/server.rs @@ -188,7 +188,10 @@ impl RelayServer { .set_local_relay(relay_runtime.relay.clone()); let deletion_runtime = relay_runtime.deletion_runtime.clone(); - deletion_runtime.run_startup_tasks().await; + deletion_runtime + .run_startup_tasks() + .await + .context("failed deletion lifecycle startup reconciliation")?; // Wire the GRASP-06 `/prs/` filesystem cleanup context into // purgatory so the standard expiry sweep can delete dangling