feat(deletion): add holding DB expiry cleanup engine

This commit is contained in:
DanConwayDev
2026-06-17 13:02:54 +00:00
parent 920c784162
commit 412cd8cbd4
7 changed files with 535 additions and 19 deletions
+27 -10
View File
@@ -1,9 +1,9 @@
# Deletion Request Support (NIP-09)
**Status:** ✅ **PARTIALLY IMPLEMENTED (holding DB foundation live)**
**Status:** ✅ **PARTIALLY IMPLEMENTED (holding DB cleanup engine live)**
This document now reflects both the implemented single-node deletion path
(main DB + tombstones + holding DB foundation) and the still-planned
(main DB + tombstones + holding DB cleanup engine) and the still-planned
expiry/archive/recovery architecture.
---
@@ -15,7 +15,7 @@ expiry/archive/recovery architecture.
> What follows in this section is what is **actually built today** and is the
> foundation that work builds on. The `deletion-request-disrespector` archival
> mode and the holding-DB move semantics are now implemented (see below);
> cleanup expiry, git archival, and recovery remain planned.
> cleanup expiry is now implemented; git archival and recovery remain planned.
Up to the rust-nostr 0.45 bump, ngit-grasp relied on the LMDB backend's
*automatic* NIP-09 / NIP-62 processing (`NostrLmdb` defaults
@@ -562,12 +562,17 @@ When implementation is complete, the following documentation will be updated:
### Completed in the current implementation
The following are implemented now (expiry/archive/recovery still pending):
The following are implemented now (archive/recovery still pending):
- ngit-grasp-owned NIP-09/NIP-62 handling (backend auto-processing disabled)
- tombstone persistence and deletion re-submission gate
- holding-DB foundation (same backend family as main DB) with move-before-delete
semantics for NIP-09 deletion paths
- holding-DB retention cleanup engine:
- startup catch-up pass,
- periodic expiry cleanup task,
- graceful shutdown signal handling for the cleanup task,
- idempotent handling of partial/missing holding data during cleanup
- announcement cascade deletion with multi-maintainer retention semantics
- extended dependent-kind coverage, including PR chain kinds
(`1618`/`1619`/`1631`/`1632`)
@@ -600,12 +605,12 @@ using served PR/PR-update fixtures (event + git-ref promotion).
[`tests/nip09_holding_db.rs`](../../tests/nip09_holding_db.rs) for holding-DB
assertions; existing cascade suites remain focused on cascade semantics.
**Phase 3: Expiry/cleanup engine (the complex part)** 🔄
- Implement retention expiry policy for holding DB (and associated metadata).
- Add cleanup scheduler + startup catch-up behavior.
- Handle partial/corrupt/missing cleanup inputs safely and idempotently.
- Add focused tests for short-retention timing, deterministic expiry windows, and
restart/offline cleanup catch-up.
**Phase 3: Expiry/cleanup engine (the complex part)** ✅
- Implemented retention expiry policy for holding DB metadata + archived events.
- Added startup catch-up cleanup and periodic background cleanup orchestration.
- Cleanup is idempotent and safe with partial/missing holding data.
- Added focused cleanup tests in
[`tests/nip09_holding_cleanup.rs`](../../tests/nip09_holding_cleanup.rs).
**Phase 4: Git archive + metadata lifecycle** 🔄
- Add git archive creation (`.tar.gz`) on deletion and archive metadata linkage.
@@ -624,6 +629,18 @@ using served PR/PR-update fixtures (event + git-ref promotion).
- Observability: metrics for holding DB size/count, cleanup, recoveries, ejections.
- Concurrency/race analysis, max-depth/scale limits, lock strategy finalization.
### Deferred items after Phase 3
The next phases still intentionally deferred are:
1. **Git archive lifecycle (Phase 4):** create and clean up `.tar.gz` git
archives linked to holding entries.
2. **Recovery workflow (Phase 5):** restore events/git data from holding/archive
on explicit/operator-approved flows.
3. **Blacklist runtime deletion parity (Phase 6):** there is currently no active
runtime blacklist-triggered deletion path in the codebase; when introduced,
it must route through holding with `DeletionSource::Blacklist`.
> ✅ Resolved prerequisite from earlier plan: rust-nostr backend auto-processing
> behavior is understood; ngit-grasp runs with `process_nip09(false)` and
> `process_nip62(false)` and owns deletion handling itself.
+3 -3
View File
@@ -487,9 +487,9 @@ pub struct Config {
/// relay an archival server, preventing "left-pad" scenarios by ensuring at least
/// some relays preserve deleted content.
///
/// This setting ONLY affects NIP-09 user-initiated deletions. It does NOT prevent
/// blacklist-triggered deletions, which are an operator moderation mechanism that
/// archival relays still need.
/// This setting ONLY affects NIP-09 user-initiated deletions. Runtime
/// blacklist-triggered deletion orchestration is deferred, so there is no
/// active blacklist deletion path to bypass today.
///
/// When `true`, NIP-09 (`"deletion"`) is NOT advertised in the NIP-11 supported
/// NIPs list so clients can discover that the relay does not honour deletions.
+83
View File
@@ -4,6 +4,7 @@ use std::{path::PathBuf, sync::Arc};
use anyhow::Result;
use clap::Parser;
use tokio::signal;
use tokio::sync::watch;
use tracing::{error, info, warn};
use tracing_subscriber::{EnvFilter, FmtSubscriber};
@@ -309,6 +310,82 @@ async fn run_relay(config: Config) -> Result<()> {
// Start HTTP server with integrated relay and database
info!("Starting HTTP server on {}", config.bind_address);
// Holding DB expiry engine:
// 1) startup catch-up pass (runs once)
// 2) periodic background cleanup pass
let holding_cleanup_store = relay_with_db.holding.clone();
match holding_cleanup_store
.cleanup_expired(
nostr_sdk::prelude::Timestamp::now(),
nostr::holding::DEFAULT_RETENTION,
)
.await
{
Ok(stats) => {
if stats.expired_records > 0 {
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) => {
warn!(error = %e, "Holding cleanup startup catch-up failed");
}
}
let (holding_cleanup_shutdown_tx, mut holding_cleanup_shutdown_rx) = watch::channel(false);
let holding_cleanup_interval = nostr::holding::DEFAULT_CLEANUP_INTERVAL;
let holding_cleanup_task = tokio::spawn(async move {
let mut interval = tokio::time::interval_at(
tokio::time::Instant::now() + holding_cleanup_interval,
holding_cleanup_interval,
);
loop {
tokio::select! {
_ = interval.tick() => {
match holding_cleanup_store
.cleanup_expired(
nostr_sdk::prelude::Timestamp::now(),
nostr::holding::DEFAULT_RETENTION,
)
.await
{
Ok(stats) => {
if stats.expired_records > 0 {
info!(
examined = stats.metadata_examined,
expired = stats.expired_records,
metadata_deleted = stats.metadata_deleted,
archived_events_deleted = stats.archived_events_deleted,
"Holding cleanup periodic pass completed"
);
}
}
Err(e) => {
warn!(error = %e, "Holding cleanup periodic pass failed");
}
}
}
changed = holding_cleanup_shutdown_rx.changed() => {
if changed.is_ok() && *holding_cleanup_shutdown_rx.borrow() {
info!("Holding cleanup task received shutdown signal");
break;
}
}
}
}
});
info!(
retention_secs = nostr::holding::DEFAULT_RETENTION.as_secs(),
interval_secs = nostr::holding::DEFAULT_CLEANUP_INTERVAL.as_secs(),
"Holding cleanup task started"
);
// Wrap write_policy in Arc for sharing between HTTP server connections
let http_write_policy = Arc::new(relay_with_db.write_policy.clone());
@@ -361,6 +438,12 @@ async fn run_relay(config: Config) -> Result<()> {
}
}
// Stop holding cleanup task cleanly.
let _ = holding_cleanup_shutdown_tx.send(true);
if let Err(e) = holding_cleanup_task.await {
warn!(error = %e, "Holding cleanup task join failed during shutdown");
}
// Save purgatory state to disk
let purgatory_save_path = PathBuf::from(&git_data_path).join("purgatory-state.json");
if let Err(e) = shutdown_purgatory.save_to_disk(&purgatory_save_path) {
+9 -1
View File
@@ -101,6 +101,11 @@ impl Nip34WritePolicy {
&self.ctx.purgatory
}
/// Get a reference to the holding store for retention cleanup orchestration.
pub fn holding(&self) -> &crate::nostr::holding::HoldingStore {
&self.ctx.holding
}
/// Set the local relay for purgatory notifications.
///
/// This must be called after the relay is created since the relay depends
@@ -861,6 +866,8 @@ pub struct RelayWithDatabase {
pub database: SharedDatabase,
/// The write policy used for event validation
pub write_policy: Nip34WritePolicy,
/// Holding store used for deletion archival + expiry cleanup
pub holding: crate::nostr::holding::HoldingStore,
}
/// Create a configured LocalRelay with full GRASP-01 validation
@@ -968,7 +975,7 @@ pub async fn create_relay(
let write_policy = Nip34WritePolicy::new(
database.clone(),
tombstones,
holding,
holding.clone(),
&git_data_path,
purgatory,
config.clone(),
@@ -1000,5 +1007,6 @@ pub async fn create_relay(
relay,
database,
write_policy,
holding,
})
}
+243
View File
@@ -6,6 +6,7 @@
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use nostr_lmdb::NostrLmdb;
use nostr_memory::MemoryDatabase;
@@ -20,6 +21,12 @@ pub const HOLDING_DIR: &str = "holding";
/// Internal metadata event kind stored alongside archived events.
pub const HOLDING_METADATA_KIND: u16 = 9905;
/// Default retention window for deleted events kept in the holding DB.
pub const DEFAULT_RETENTION: Duration = Duration::from_secs(90 * 24 * 60 * 60);
/// Default interval between background holding cleanup passes.
pub const DEFAULT_CLEANUP_INTERVAL: Duration = Duration::from_secs(24 * 60 * 60);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DeletionSource {
Nip09,
@@ -49,6 +56,20 @@ pub struct HoldingStore {
metadata_signer: Keys,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ExpiredMetadataRecord {
pub metadata_event_id: EventId,
pub archived_event_id: Option<EventId>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct CleanupStats {
pub metadata_examined: usize,
pub expired_records: usize,
pub metadata_deleted: usize,
pub archived_events_deleted: usize,
}
impl std::fmt::Debug for HoldingStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HoldingStore").finish_non_exhaustive()
@@ -150,4 +171,226 @@ impl HoldingStore {
}
}
}
fn parse_holding_deleted_at(metadata_event: &Event) -> Option<Timestamp> {
metadata_event.tags.iter().find_map(|tag| {
let v = tag.as_slice();
if v.len() >= 2 && v[0] == "holding-deleted-at" {
v[1].parse::<u64>().ok().map(Timestamp::from_secs)
} else {
None
}
})
}
fn parse_archived_event_id(metadata_event: &Event) -> Option<EventId> {
metadata_event.tags.iter().find_map(|tag| {
let v = tag.as_slice();
if v.len() >= 2 && v[0] == "e" {
EventId::from_hex(&v[1]).ok()
} else {
None
}
})
}
/// Find expired metadata records in holding based on `holding-deleted-at`
/// and a retention window.
///
/// Metadata entries that are missing or have an invalid `holding-deleted-at`
/// value are skipped (treated as non-expired) to keep cleanup fail-safe.
pub async fn expired_records(
&self,
now: Timestamp,
retention: Duration,
) -> Vec<ExpiredMetadataRecord> {
let cutoff = now.as_secs().saturating_sub(retention.as_secs());
let filter = Filter::new().kind(Kind::from(HOLDING_METADATA_KIND));
let metadata_events = match self.db.query(filter).await {
Ok(events) => events,
Err(e) => {
tracing::error!(error = %e, "Holding cleanup failed to query metadata events");
return Vec::new();
}
};
metadata_events
.into_iter()
.filter_map(|metadata_event| {
let deleted_at = Self::parse_holding_deleted_at(&metadata_event)?;
if deleted_at.as_secs() <= cutoff {
Some(ExpiredMetadataRecord {
metadata_event_id: metadata_event.id,
archived_event_id: Self::parse_archived_event_id(&metadata_event),
})
} else {
None
}
})
.collect()
}
/// Delete one expired metadata record plus its archived event (when present).
///
/// This is idempotent: missing archived events and missing metadata rows are
/// tolerated and do not error.
pub async fn delete_expired_record(
&self,
record: ExpiredMetadataRecord,
) -> anyhow::Result<(bool, bool)> {
let mut metadata_deleted = false;
let mut event_deleted = false;
if let Some(event_id) = record.archived_event_id {
let exists = self.db.event_by_id(&event_id).await.map_err(|e| {
anyhow::anyhow!("Failed to look up archived holding event {}: {e}", event_id)
})?;
if exists.is_some() {
self.db
.delete(Filter::new().ids(vec![event_id]))
.await
.map_err(|e| {
anyhow::anyhow!("Failed to delete archived holding event {}: {e}", event_id)
})?;
event_deleted = true;
}
}
let metadata_exists = self
.db
.event_by_id(&record.metadata_event_id)
.await
.map_err(|e| {
anyhow::anyhow!(
"Failed to look up holding metadata event {}: {e}",
record.metadata_event_id
)
})?;
if metadata_exists.is_some() {
self.db
.delete(Filter::new().ids(vec![record.metadata_event_id]))
.await
.map_err(|e| {
anyhow::anyhow!(
"Failed to delete holding metadata event {}: {e}",
record.metadata_event_id
)
})?;
metadata_deleted = true;
}
Ok((metadata_deleted, event_deleted))
}
/// Cleanup pass for expired holding records.
///
/// - Computes expiry from `holding-deleted-at` metadata + `retention`.
/// - Removes expired metadata rows and best-effort removes their archived
/// event payload.
/// - Idempotent with missing metadata/events and safe to run repeatedly.
pub async fn cleanup_expired(
&self,
now: Timestamp,
retention: Duration,
) -> anyhow::Result<CleanupStats> {
let filter = Filter::new().kind(Kind::from(HOLDING_METADATA_KIND));
let metadata_examined = self
.db
.query(filter)
.await
.map_err(|e| anyhow::anyhow!("Failed to query holding metadata for cleanup: {e}"))?
.len();
let expired = self.expired_records(now, retention).await;
let mut stats = CleanupStats {
metadata_examined,
expired_records: expired.len(),
..CleanupStats::default()
};
for record in expired {
let (metadata_deleted, event_deleted) = self.delete_expired_record(record).await?;
if metadata_deleted {
stats.metadata_deleted += 1;
}
if event_deleted {
stats.archived_events_deleted += 1;
}
}
Ok(stats)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_event(content: &str, created_at: Timestamp) -> Event {
EventBuilder::new(Kind::TextNote, content)
.custom_created_at(created_at)
.finalize(&Keys::generate())
.expect("build test event")
}
#[tokio::test]
async fn cleanup_ignores_corrupt_metadata_deleted_at() {
let store = HoldingStore::in_memory();
let metadata = EventBuilder::new(Kind::from(HOLDING_METADATA_KIND), "")
.tags(vec![
Tag::custom("e", vec![EventId::all_zeros().to_hex()]),
Tag::custom("holding-deleted-at", vec!["not-a-number".to_string()]),
])
.finalize(&Keys::generate())
.expect("build metadata event");
store.db.save_event(&metadata).await.expect("save metadata");
let stats = store
.cleanup_expired(Timestamp::now(), Duration::from_secs(1))
.await
.expect("cleanup should succeed");
assert_eq!(stats.expired_records, 0);
assert_eq!(stats.metadata_deleted, 0);
}
#[tokio::test]
async fn cleanup_deletes_metadata_even_if_archived_event_is_missing() {
let store = HoldingStore::in_memory();
let event = test_event("to archive", Timestamp::from_secs(10));
let deleted_at = Timestamp::from_secs(11);
store
.archive_event(
&event,
&HoldingMetadata {
deleted_at,
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
},
)
.await
.expect("archive event");
// Simulate partial/corrupt holding state: archived payload missing while
// metadata still exists.
store
.db
.delete(Filter::new().ids(vec![event.id]))
.await
.expect("delete archived payload directly");
let stats = store
.cleanup_expired(Timestamp::from_secs(10_000), Duration::from_secs(1))
.await
.expect("cleanup should succeed");
assert_eq!(stats.expired_records, 1);
assert_eq!(stats.archived_events_deleted, 0);
assert_eq!(stats.metadata_deleted, 1);
assert!(store.metadata_for_event(&event.id).await.is_empty());
}
}
+6 -5
View File
@@ -97,8 +97,8 @@ impl DeletionPolicy {
/// main database) so clients see an OK and the request is preserved, but it
/// is NOT acted upon — no purgatory eviction, no main-DB deletion, and no
/// tombstone recording. The targeted events therefore remain fully
/// accessible. This only affects NIP-09 user-initiated deletions; it does
/// not influence blacklist-triggered deletions.
/// accessible. This only affects NIP-09 user-initiated deletions. Runtime
/// blacklist-triggered deletion flow is currently deferred.
pub async fn handle(&self, event: &Event) -> WritePolicyResult {
// Archival mode: store the deletion request but do not process it.
if self.ctx.config.deletion_request_disrespector {
@@ -248,9 +248,10 @@ impl DeletionPolicy {
/// 5. Hard-delete the orphaned event ids plus the announcement coordinate
/// itself from the main DB.
///
/// All holding-DB / archival / recovery orchestration from the original
/// multi-maintainer prototype is deliberately omitted — this is a pure
/// hard-delete cascade.
/// Phase-2 holding orchestration is active in this cascade path via
/// [`Self::archive_and_delete_filter`]: every deleted orphan/coordinate
/// target is moved into holding before main-DB deletion. Archive-file
/// lifecycle and recovery orchestration remain future phases.
async fn cascade_delete_announcement(
&self,
author: &PublicKey,
+164
View File
@@ -0,0 +1,164 @@
//! Focused integration tests for Phase 3 holding-DB expiry/cleanup behavior.
use ngit_grasp::nostr::holding::{
DeletionSource, HoldingMetadata, HoldingStore, DEFAULT_CLEANUP_INTERVAL, DEFAULT_RETENTION,
};
use nostr_sdk::prelude::*;
use std::time::Duration;
fn build_test_event(content: &str, created_at: Timestamp) -> Event {
EventBuilder::new(Kind::TextNote, content)
.custom_created_at(created_at)
.finalize(&Keys::generate())
.expect("build test event")
}
#[tokio::test]
async fn cleanup_removes_expired_holding_events_and_metadata() {
let store = HoldingStore::in_memory();
let event = build_test_event("expired", Timestamp::from_secs(100));
store
.archive_event(
&event,
&HoldingMetadata {
deleted_at: Timestamp::from_secs(120),
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
},
)
.await
.expect("archive event");
let stats = store
.cleanup_expired(Timestamp::from_secs(200), Duration::from_secs(30))
.await
.expect("cleanup should succeed");
assert_eq!(stats.expired_records, 1);
assert_eq!(stats.metadata_deleted, 1);
assert_eq!(stats.archived_events_deleted, 1);
assert!(!store.has_event(&event.id).await);
assert!(store.metadata_for_event(&event.id).await.is_empty());
}
#[tokio::test]
async fn cleanup_preserves_non_expired_entries() {
let store = HoldingStore::in_memory();
let event = build_test_event("fresh", Timestamp::from_secs(100));
store
.archive_event(
&event,
&HoldingMetadata {
deleted_at: Timestamp::from_secs(190),
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
},
)
.await
.expect("archive event");
let stats = store
.cleanup_expired(Timestamp::from_secs(200), Duration::from_secs(30))
.await
.expect("cleanup should succeed");
assert_eq!(stats.expired_records, 0);
assert!(store.has_event(&event.id).await);
assert_eq!(store.metadata_for_event(&event.id).await.len(), 1);
}
#[tokio::test]
async fn startup_catchup_cleanup_removes_expired_entries_after_restart() {
let temp = tempfile::tempdir().expect("create temp dir");
let event_id = {
let store = HoldingStore::open_lmdb(temp.path())
.await
.expect("open holding");
let event = build_test_event("restart-expired", Timestamp::from_secs(10));
let id = event.id;
store
.archive_event(
&event,
&HoldingMetadata {
deleted_at: Timestamp::from_secs(20),
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
},
)
.await
.expect("archive event");
id
};
// Simulate process restart: reopen store and run startup catch-up pass.
let reopened = HoldingStore::open_lmdb(temp.path())
.await
.expect("reopen holding");
let stats = reopened
.cleanup_expired(Timestamp::from_secs(300), Duration::from_secs(60))
.await
.expect("startup catch-up cleanup should succeed");
assert_eq!(stats.expired_records, 1);
assert!(!reopened.has_event(&event_id).await);
assert!(reopened.metadata_for_event(&event_id).await.is_empty());
}
#[tokio::test]
async fn cleanup_is_idempotent_with_partial_missing_event_data() {
let store = HoldingStore::in_memory();
let event = build_test_event("partial", Timestamp::from_secs(50));
let metadata = HoldingMetadata {
deleted_at: Timestamp::from_secs(55),
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
};
// Two metadata rows referencing the same archived event: first delete removes
// the event, second row becomes metadata-only cleanup.
store
.archive_event(&event, &metadata)
.await
.expect("archive event first time");
let metadata_with_coordinate = HoldingMetadata {
coordinate: Some("30617:deadbeef:partial".to_string()),
..metadata.clone()
};
store
.archive_event(&event, &metadata_with_coordinate)
.await
.expect("archive event second time");
assert_eq!(store.metadata_for_event(&event.id).await.len(), 2);
let first = store
.cleanup_expired(Timestamp::from_secs(500), Duration::from_secs(1))
.await
.expect("first cleanup should succeed");
assert_eq!(first.expired_records, 2);
assert_eq!(first.metadata_deleted, 2);
assert_eq!(first.archived_events_deleted, 1);
let second = store
.cleanup_expired(Timestamp::from_secs(500), Duration::from_secs(1))
.await
.expect("second cleanup should still succeed");
assert_eq!(second.expired_records, 0);
assert_eq!(second.metadata_deleted, 0);
assert_eq!(second.archived_events_deleted, 0);
assert!(!store.has_event(&event.id).await);
assert!(store.metadata_for_event(&event.id).await.is_empty());
}
#[test]
fn cleanup_defaults_are_non_zero() {
assert!(DEFAULT_RETENTION.as_secs() > 0);
assert!(DEFAULT_CLEANUP_INTERVAL.as_secs() > 0);
}