From 337e56ff12980ce39fc55d39f4897a817947bdec Mon Sep 17 00:00:00 2001 From: Awiteb Date: Thu, 3 Sep 2026 20:57:34 +0000 Subject: [PATCH 1/4] database: add garbage collection function `NostrDatabase::collect_garbage` Added the `collect_garbage` function to the `NostrDatabase` trait. This gives the user (client or relay) direct control over garbage collection timing, removing the need for an automatic runtime on each database. Signed-off-by: Awiteb --- database/nostr-database/CHANGELOG.md | 6 ++++++ database/nostr-database/src/lib.rs | 5 +++++ 2 files changed, 11 insertions(+) diff --git a/database/nostr-database/CHANGELOG.md b/database/nostr-database/CHANGELOG.md index 650a191c4..e7482c3ae 100644 --- a/database/nostr-database/CHANGELOG.md +++ b/database/nostr-database/CHANGELOG.md @@ -27,6 +27,12 @@ --> +## Unreleased + +### Added + +- New function to collect garbage `NostrDatabase::collect_garbage` (https://github.com/nostrdevkit/nostr/pull/1474) + ## v0.45.1 - 2026/08/07 ### Fixed diff --git a/database/nostr-database/src/lib.rs b/database/nostr-database/src/lib.rs index 8d7b868f5..d95e0b3b2 100644 --- a/database/nostr-database/src/lib.rs +++ b/database/nostr-database/src/lib.rs @@ -198,4 +198,9 @@ pub trait NostrDatabase: Any + Debug + Send + Sync { /// Wipe all data fn wipe(&self) -> Pin> + Send + '_>>; + + /// Cleans up expired events and other unused resources. + fn collect_garbage(&self) -> Pin> + Send + '_>> { + Box::pin(async { Ok(()) }) + } } From 3c06dec20dc8f71b9af3082c941d7f0045e926a0 Mon Sep 17 00:00:00 2001 From: Awiteb Date: Thu, 3 Sep 2026 21:02:31 +0000 Subject: [PATCH 2/4] ndb: implment `NostrDatabase::collect_garbage` to return an error Signed-off-by: Awiteb --- database/nostr-ndb/src/lib.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/database/nostr-ndb/src/lib.rs b/database/nostr-ndb/src/lib.rs index 10659f268..8c5da0ead 100644 --- a/database/nostr-ndb/src/lib.rs +++ b/database/nostr-ndb/src/lib.rs @@ -191,6 +191,11 @@ impl NostrDatabase for NdbDatabase { fn wipe(&self) -> Pin> + Send + '_>> { Box::pin(async move { Err(Error::unsupported("wiping is not supported by nostrdb")) }) } + + #[inline] + fn collect_garbage(&self) -> Pin> + Send + '_>> { + Box::pin(async move { Err(Error::unsupported("delete is not supported by nostrdb")) }) + } } fn ndb_query<'a>( From 42c9afaa6ea849db84b49e6ca793b7951aabe1ff Mon Sep 17 00:00:00 2001 From: Awiteb Date: Tue, 8 Sep 2026 23:56:24 +0000 Subject: [PATCH 3/4] database-test-suite: `test_expire_events` test Ensure expired events are correctly removed when calling `NostrDatabase::collect_garbage`. The test saves an expired event and verifies it gets deleted after garbage collection. Signed-off-by: Awiteb --- database/nostr-database-test-suite/src/lib.rs | 37 +++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/database/nostr-database-test-suite/src/lib.rs b/database/nostr-database-test-suite/src/lib.rs index c63c9ffc2..f778d8662 100644 --- a/database/nostr-database-test-suite/src/lib.rs +++ b/database/nostr-database-test-suite/src/lib.rs @@ -1055,5 +1055,42 @@ macro_rules! database_unit_tests { let status = store.save_event(&new_event).await.unwrap(); assert_eq!(status, SaveEventStatus::Rejected(RejectedReason::Vanished)); } + + #[tokio::test] + async fn test_expire_events() { + let store: $store_type = $setup_fn().await; + let features = store.features(); + + if !features.event_expiration { + println!("Skipping event expiration tests as the database doesn't support it!"); + return; + } + + let keys = Keys::generate(); + let event = EventBuilder::new(Kind::TextNote, "Nothing in my mind") + .tag(Tag::expiration(Timestamp::now() - 10)) + .finalize(&keys) + .unwrap(); + + assert!( + store.save_event(&event).await.unwrap().is_success(), + "Database should not check if the event is expired" + ); + + let count = store + .count(Filter::new().author(keys.public_key())) + .await + .unwrap(); + assert_eq!(count, 1, "We stored a single event"); + + // Remove expired events + store.collect_garbage().await.unwrap(); + + let count = store + .count(Filter::new().author(keys.public_key())) + .await + .unwrap(); + assert_eq!(count, 0, "Garbage collected"); + } }; } From 62128717041996513a908262ac31a2a52733106d Mon Sep 17 00:00:00 2001 From: Awiteb Date: Wed, 9 Sep 2026 01:07:18 +0000 Subject: [PATCH 4/4] lmdb: implement `NostrDatabase::collect_garbage` Adds garbage collection to delete expired events. This implements the `event_expiration` feature, making the database track and manage event expiration. - Creates a new `expirations` database to track events with expiration timestamps - Adds an `is_indexable` field to `TagIndexKeySet` because the `expiration` tag is non-indexable - Bumps the database version to 3 to migrate existing events and populate the expirations tracker Signed-off-by: Awiteb --- database/nostr-lmdb/CHANGELOG.md | 6 + database/nostr-lmdb/src/lib.rs | 9 +- database/nostr-lmdb/src/store/event.rs | 10 ++ database/nostr-lmdb/src/store/ingester.rs | 37 ++++++ database/nostr-lmdb/src/store/lmdb/index.rs | 23 +++- database/nostr-lmdb/src/store/lmdb/mod.rs | 121 ++++++++++++++++++-- database/nostr-lmdb/src/store/mod.rs | 17 ++- 7 files changed, 212 insertions(+), 11 deletions(-) diff --git a/database/nostr-lmdb/CHANGELOG.md b/database/nostr-lmdb/CHANGELOG.md index 8665734af..8fd19d22c 100644 --- a/database/nostr-lmdb/CHANGELOG.md +++ b/database/nostr-lmdb/CHANGELOG.md @@ -27,6 +27,12 @@ --> +## Unreleased + +### Added + +- Support event expiration (https://github.com/nostrdevkit/nostr/pull/1474) + ## v0.45.2 - 2026/08/19 ### Fixed diff --git a/database/nostr-lmdb/src/lib.rs b/database/nostr-lmdb/src/lib.rs index 5fbeeff81..d95cccf61 100644 --- a/database/nostr-lmdb/src/lib.rs +++ b/database/nostr-lmdb/src/lib.rs @@ -180,7 +180,7 @@ impl NostrDatabase for NostrLmdb { fn features(&self) -> Features { Features { persistent: true, - event_expiration: false, + event_expiration: true, full_text_search: true, request_to_vanish: true, } @@ -239,6 +239,13 @@ impl NostrDatabase for NostrLmdb { fn wipe(&self) -> Pin> + Send + '_>> { Box::pin(async move { Ok(self.db.wipe().await?) }) } + + fn collect_garbage(&self) -> Pin> + Send + '_>> { + Box::pin(async move { + self.db.delete_expired().await?; + Ok(()) + }) + } } #[cfg(test)] diff --git a/database/nostr-lmdb/src/store/event.rs b/database/nostr-lmdb/src/store/event.rs index 1f7c427a5..3a348b80e 100644 --- a/database/nostr-lmdb/src/store/event.rs +++ b/database/nostr-lmdb/src/store/event.rs @@ -59,6 +59,16 @@ impl<'a> DatabaseTag<'a> { } } + /// Return the expiration timestamp if it's `expiration` tag + #[inline] + pub(super) fn expiration(&self) -> Option { + if self.kind() == "expiration" { + return self.content().and_then(|t| u64::from_str(t).ok()); + } + + None + } + /// Into owned tag pub(super) fn into_owned(self) -> Tag { let buf: Vec = self.buf.into_iter().map(|t| t.into_owned()).collect(); diff --git a/database/nostr-lmdb/src/store/ingester.rs b/database/nostr-lmdb/src/store/ingester.rs index 592123f80..6f0278c87 100644 --- a/database/nostr-lmdb/src/store/ingester.rs +++ b/database/nostr-lmdb/src/store/ingester.rs @@ -42,6 +42,10 @@ enum OperationResult { result: Result<(), StoreError>, tx: Option>>, }, + DeleteExpired { + result: Result<(), StoreError>, + tx: Option>>, + }, Wipe { result: Result<(), StoreError>, tx: Option>>, @@ -79,6 +83,15 @@ impl OperationResult { tracing::error!(error = %e, "Delete operation failed in batch"); } } + Self::DeleteExpired { result, tx } => { + if let Some(tx) = tx { + if tx.send(result).is_err() { + tracing::debug!("Failed to send delete expired result: receiver dropped"); + } + } else if let Err(e) = result { + tracing::error!(error = %e, "delete expired operation failed in batch"); + } + } Self::Wipe { result, tx } => { if let Some(tx) = tx { if tx.send(result).is_err() { @@ -104,6 +117,9 @@ enum IngesterOperation { filter: Filter, tx: Option>>, }, + DeleteExpired { + tx: Option>>, + }, Wipe { tx: Option>>, }, @@ -125,6 +141,10 @@ impl IngesterOperation { result: Err(error), tx, }, + Self::DeleteExpired { tx } => OperationResult::DeleteExpired { + result: Err(error), + tx, + }, Self::Wipe { tx } => OperationResult::Wipe { result: Err(error), tx, @@ -175,6 +195,16 @@ impl IngesterItem { (item, rx) } + #[must_use] + pub(super) fn delete_expired_with_feedback() -> (Self, oneshot::Receiver>) + { + let (tx, rx) = oneshot::channel(); + let item: Self = Self { + operation: IngesterOperation::DeleteExpired { tx: Some(tx) }, + }; + (item, rx) + } + #[must_use] pub(super) fn wipe_with_feedback() -> (Self, oneshot::Receiver>) { let (tx, rx) = oneshot::channel(); @@ -342,6 +372,10 @@ impl Ingester { let result = self.db.delete(txn, filter); OperationResult::Delete { result, tx } } + IngesterOperation::DeleteExpired { tx } => { + let result = self.db.delete_expired(txn); + OperationResult::DeleteExpired { result, tx } + } IngesterOperation::Wipe { tx } => { let result = self.db.wipe(txn); OperationResult::Wipe { result, tx } @@ -364,6 +398,9 @@ fn mark_all_as_failed(results: &mut [OperationResult]) { OperationResult::Delete { result: res, .. } => { *res = Err(StoreError::BatchTransactionFailed) } + OperationResult::DeleteExpired { result: res, .. } => { + *res = Err(StoreError::BatchTransactionFailed) + } OperationResult::Wipe { result: res, .. } => { *res = Err(StoreError::BatchTransactionFailed) } diff --git a/database/nostr-lmdb/src/store/lmdb/index.rs b/database/nostr-lmdb/src/store/lmdb/index.rs index 69be63552..fb79fd459 100644 --- a/database/nostr-lmdb/src/store/lmdb/index.rs +++ b/database/nostr-lmdb/src/store/lmdb/index.rs @@ -18,10 +18,14 @@ const KIND_BE: usize = 2; const TAG_VALUE_PAD_LEN: usize = 182; // TODO: use fixed-size arrays instead of vectors +// TODO: convert `TagIndexKeySet` to an enum instead of `is_indexable`? pub(super) struct TagIndexKeySet { + /// Whether the tag is indexable + pub(super) is_indexable: bool, pub(super) atc_index: Vec, pub(super) ktc_index: Vec, pub(super) tc_index: Vec, + pub(super) expiration: Option, } // TODO: use fixed-size arrays instead of vectors @@ -49,7 +53,7 @@ impl EventIndexKeys { // Index by kind (with created_at and id) let kc_index: Vec = make_kc_index_key(event.kind, event.created_at, event.id); - let tags = event + let mut tags: Vec = event .tags .iter() .filter_map(|t| t.extract()) @@ -77,13 +81,30 @@ impl EventIndexKeys { make_tc_index_key(&tag_name, tag_value, event.created_at, event.id); TagIndexKeySet { + is_indexable: true, atc_index, ktc_index, tc_index, + expiration: None, } }) .collect(); + // Collect expiration tags, they can't be with `tags` because it's not a single letter tag + let expiration_tags = event + .tags + .iter() + .filter_map(|t| t.expiration()) + .map(|expire_at| TagIndexKeySet { + is_indexable: false, + atc_index: Vec::new(), + ktc_index: Vec::new(), + tc_index: Vec::new(), + expiration: Some(expire_at), + }); + // Extend the tags with expiration tags + tags.extend(expiration_tags); + Self { id: *event.id, ci_index, diff --git a/database/nostr-lmdb/src/store/lmdb/mod.rs b/database/nostr-lmdb/src/store/lmdb/mod.rs index 87d68f9e1..ded1040b9 100644 --- a/database/nostr-lmdb/src/store/lmdb/mod.rs +++ b/database/nostr-lmdb/src/store/lmdb/mod.rs @@ -22,12 +22,13 @@ use super::error::{MigrationError, StoreError}; use super::event::DatabaseEvent; use super::filter::DatabaseFilter; use crate::NostrLmdbBuilder; +use crate::store::event::DatabaseTag; const EVENT_ID_ALL_ZEROS: [u8; 32] = [0; 32]; const EVENT_ID_ALL_255: [u8; 32] = [255; 32]; /// Current database schema version -const DB_VERSION: u64 = 2; +const DB_VERSION: u64 = 3; const DB_VERSION_KEY: &[u8] = b"db_version"; #[derive(Debug)] @@ -109,6 +110,8 @@ pub(crate) struct Lmdb { deleted_coordinates: Database>, // Coordinate, UNIX timestamp /// Vanished public keys vanished_public_keys: Database, // Public key + /// Event expiration tracker + expirations: Database>, // Event ID, Expiration timestamp /// Database metadata (version, etc) metadata: Database>, // Key, Value } @@ -119,7 +122,7 @@ impl Lmdb { let env: Env = unsafe { EnvOpenOptions::new() .flags(EnvFlags::NO_TLS) - .max_dbs(12 + builder.additional_dbs) + .max_dbs(13 + builder.additional_dbs) .max_readers(builder.max_readers) .map_size(builder.map_size) .open(builder.path)? @@ -183,6 +186,11 @@ impl Lmdb { .types::() .name("vanished-public-keys") .create(&mut txn)?; + let expirations = env + .database_options() + .types::>() + .name("expirations") + .create(&mut txn)?; let metadata = env .database_options() .types::>() @@ -212,6 +220,7 @@ impl Lmdb { deleted_ids, deleted_coordinates, vanished_public_keys, + expirations, metadata, }; @@ -241,6 +250,10 @@ impl Lmdb { self.migrate_v1_to_v2(&mut txn)?; } + if current_version < 3 { + self.migrate_v2_to_v3(&mut txn)?; + } + // Update version self.metadata.put(&mut txn, DB_VERSION_KEY, &DB_VERSION)?; txn.commit()?; @@ -296,6 +309,51 @@ impl Lmdb { Ok(()) } + /// Migrates the database from schema version 2 to version 3 by building + /// an expiration tracker for events with `expiration` tags. + fn migrate_v2_to_v3(&self, txn: &mut RwTxn) -> Result<(), StoreError> { + tracing::info!("Starting migration to schema v3: Building expiration tracker"); + tracing::debug!( + "Scanning {} total events for expiration tags", + self.events.len(txn)? + ); + + // Collect all events containing expiration timestamps + let expiring_events: Vec<([u8; 32], u64)> = { + let mut events_with_expiration = Vec::new(); + for result in self.events.iter(txn)? { + let (_, event_bytes) = result?; + + // Decode event + if let Ok(event) = DatabaseEvent::from_flatbuf(event_bytes) { + // Extract expiration timestamp from event tags + if let Some(timestamp) = event.tags.iter().find_map(DatabaseTag::expiration) { + events_with_expiration.push((*event.id, timestamp)); + } + } + } + events_with_expiration + }; + + tracing::info!( + "Found {} events with expiration tags", + expiring_events.len(), + ); + + // Build expiration index + tracing::debug!("Building expiration index entries"); + for (event_id, expire_at) in expiring_events.iter() { + self.expirations.put(txn, event_id, expire_at)?; + } + + tracing::info!( + "Successfully migrated to schema v3: Created expiration tracker with {} entries", + expiring_events.len() + ); + + Ok(()) + } + /// Get a read transaction /// /// This should never block the current thread @@ -334,9 +392,17 @@ impl Lmdb { self.kc_index.put(txn, &index.kc_index, &index.id)?; for tag in index.tags.into_iter() { - self.atc_index.put(txn, &tag.atc_index, &index.id)?; - self.ktc_index.put(txn, &tag.ktc_index, &index.id)?; - self.tc_index.put(txn, &tag.tc_index, &index.id)?; + // If the tag is indexable means it is a single letter tag that we index + if tag.is_indexable { + self.atc_index.put(txn, &tag.atc_index, &index.id)?; + self.ktc_index.put(txn, &tag.ktc_index, &index.id)?; + self.tc_index.put(txn, &tag.tc_index, &index.id)?; + } + + // The expiration tag is not indexable, it's just in the `expirations` tracker + if let Some(ref timestamp) = tag.expiration { + self.expirations.put(txn, &index.id, timestamp)?; + } } Ok(()) @@ -368,9 +434,12 @@ impl Lmdb { // Delete tag indexes for tag in &index.tags { - self.atc_index.delete(txn, &tag.atc_index)?; - self.ktc_index.delete(txn, &tag.ktc_index)?; - self.tc_index.delete(txn, &tag.tc_index)?; + if tag.is_indexable { + self.atc_index.delete(txn, &tag.atc_index)?; + self.ktc_index.delete(txn, &tag.ktc_index)?; + self.tc_index.delete(txn, &tag.tc_index)?; + } + self.expirations.delete(txn, &index.id)?; } Ok(()) @@ -397,6 +466,7 @@ impl Lmdb { self.deleted_ids.clear(txn)?; self.deleted_coordinates.clear(txn)?; self.vanished_public_keys.clear(txn)?; + self.expirations.clear(txn)?; Ok(()) } @@ -551,6 +621,41 @@ impl Lmdb { Ok(()) } + /// Delete expired events + pub fn delete_expired(&self, txn: &mut RwTxn) -> Result<(), StoreError> { + tracing::info!( + "Processing {} events that can expire", + self.expirations.len(txn)? + ); + + // Collect all expired events + let expired_indexes = { + let now = Timestamp::now().as_secs(); + let mut expired = Vec::new(); + for result in self.expirations.iter(txn)? { + let (event_id, expire_at) = result?; + if now >= expire_at { + if let Some(event_index_keys) = self + .get_event_by_id(txn, event_id)? + .map(EventIndexKeys::new) + { + expired.push(event_index_keys); + } + } + } + expired + }; + + tracing::info!("Found {} expired events", expired_indexes.len()); + + // Remove them, we can safely mutate the transaction + for index in expired_indexes { + self.remove(txn, &index)?; + } + + Ok(()) + } + pub fn count(&self, txn: &RoTxn, filter: Filter) -> Result { // Check if we can use fast counting let can_fast_count: bool = filter.ids.is_none() diff --git a/database/nostr-lmdb/src/store/mod.rs b/database/nostr-lmdb/src/store/mod.rs index 7d35813a4..79521a447 100644 --- a/database/nostr-lmdb/src/store/mod.rs +++ b/database/nostr-lmdb/src/store/mod.rs @@ -122,7 +122,14 @@ impl Store { let txn: RoTxn = db.read_txn()?; let output = db.query(&txn, filter)?; - events.extend(output.into_iter().map(|e| e.into_owned())); + let now = Timestamp::now(); + + events.extend( + output + .into_iter() + .filter(|e| !e.is_expired_at(now)) + .map(|e| e.into_owned()), + ); txn.commit()?; Ok(events) @@ -157,6 +164,14 @@ impl Store { rx.await? } + pub(super) async fn delete_expired(&self) -> Result<(), StoreError> { + let (item, rx) = IngesterItem::delete_expired_with_feedback(); + self.ingester + .send(item) + .map_err(|_| StoreError::FlumeSend)?; + rx.await? + } + pub(super) async fn wipe(&self) -> Result<(), StoreError> { let (item, rx) = IngesterItem::wipe_with_feedback(); self.ingester