Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 37 additions & 0 deletions database/nostr-database-test-suite/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
};
}
6 changes: 6 additions & 0 deletions database/nostr-database/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions database/nostr-database/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -198,4 +198,9 @@ pub trait NostrDatabase: Any + Debug + Send + Sync {

/// Wipe all data
fn wipe(&self) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>>;

/// Cleans up expired events and other unused resources.
fn collect_garbage(&self) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>> {
Box::pin(async { Ok(()) })
}
}
6 changes: 6 additions & 0 deletions database/nostr-lmdb/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@

-->

## Unreleased

### Added

- Support event expiration (https://github.com/nostrdevkit/nostr/pull/1474)

## v0.45.2 - 2026/08/19

### Fixed
Expand Down
9 changes: 8 additions & 1 deletion database/nostr-lmdb/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
Expand Down Expand Up @@ -239,6 +239,13 @@ impl NostrDatabase for NostrLmdb {
fn wipe(&self) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>> {
Box::pin(async move { Ok(self.db.wipe().await?) })
}

fn collect_garbage(&self) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + '_>> {
Box::pin(async move {
self.db.delete_expired().await?;
Ok(())
})
}
}

#[cfg(test)]
Expand Down
10 changes: 10 additions & 0 deletions database/nostr-lmdb/src/store/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,16 @@ impl<'a> DatabaseTag<'a> {
}
}

/// Return the expiration timestamp if it's `expiration` tag
#[inline]
pub(super) fn expiration(&self) -> Option<u64> {
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<String> = self.buf.into_iter().map(|t| t.into_owned()).collect();
Expand Down
37 changes: 37 additions & 0 deletions database/nostr-lmdb/src/store/ingester.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,10 @@ enum OperationResult {
result: Result<(), StoreError>,
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
},
DeleteExpired {
result: Result<(), StoreError>,
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
},
Wipe {
result: Result<(), StoreError>,
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
Expand Down Expand Up @@ -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() {
Expand All @@ -104,6 +117,9 @@ enum IngesterOperation {
filter: Filter,
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
},
DeleteExpired {
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
},
Wipe {
tx: Option<oneshot::Sender<Result<(), StoreError>>>,
},
Expand All @@ -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,
Expand Down Expand Up @@ -175,6 +195,16 @@ impl IngesterItem {
(item, rx)
}

#[must_use]
pub(super) fn delete_expired_with_feedback() -> (Self, oneshot::Receiver<Result<(), StoreError>>)
{
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<Result<(), StoreError>>) {
let (tx, rx) = oneshot::channel();
Expand Down Expand Up @@ -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 }
Expand All @@ -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)
}
Expand Down
23 changes: 22 additions & 1 deletion database/nostr-lmdb/src/store/lmdb/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u8>,
pub(super) ktc_index: Vec<u8>,
pub(super) tc_index: Vec<u8>,
pub(super) expiration: Option<u64>,
}

// TODO: use fixed-size arrays instead of vectors
Expand Down Expand Up @@ -49,7 +53,7 @@ impl EventIndexKeys {
// Index by kind (with created_at and id)
let kc_index: Vec<u8> = make_kc_index_key(event.kind, event.created_at, event.id);

let tags = event
let mut tags: Vec<TagIndexKeySet> = event
.tags
.iter()
.filter_map(|t| t.extract())
Expand Down Expand Up @@ -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,
Expand Down
Loading
Loading