mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Add holding database infrastructure for archived events
This commit is contained in:
@@ -0,0 +1,467 @@
|
||||
/// Holding Database for Archived Events
|
||||
///
|
||||
/// This module provides a separate database for storing events that have been
|
||||
/// deleted via NIP-09 deletion requests. Events are archived with deletion metadata
|
||||
/// and can be permanently deleted after a retention period.
|
||||
///
|
||||
/// The holding database uses the same backend type as the main database but stores
|
||||
/// data in a separate file at `<relay_data_path>/holding-<backend>`.
|
||||
use std::num::NonZeroUsize;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use nostr_lmdb::NostrLmdb;
|
||||
use nostr_relay_builder::prelude::{
|
||||
Event, EventBuilder, EventId, Filter, Keys, MemoryDatabase, MemoryDatabaseOptions,
|
||||
NostrDatabase, Tag, TagKind,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::config::DatabaseBackend;
|
||||
|
||||
/// Type alias for the shared holding database
|
||||
pub type SharedHoldingDatabase = Arc<HoldingDatabase>;
|
||||
|
||||
/// Metadata stored with each archived event
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct DeletionMetadata {
|
||||
/// Unix timestamp when the event was deleted
|
||||
pub deletion_timestamp: u64,
|
||||
/// Event ID of the deletion request (NIP-09 kind 5 event)
|
||||
pub deletion_event_id: EventId,
|
||||
/// Unix timestamp when this event should be permanently deleted
|
||||
/// (deletion_timestamp + retention_secs)
|
||||
pub expiry_timestamp: u64,
|
||||
}
|
||||
|
||||
/// Holding database for archived events
|
||||
///
|
||||
/// This is a wrapper around the same database backends used by the main relay,
|
||||
/// but stores events in a separate file with deletion metadata.
|
||||
pub struct HoldingDatabase {
|
||||
/// The underlying database backend
|
||||
backend: Arc<dyn NostrDatabase>,
|
||||
/// Path to the holding database
|
||||
path: PathBuf,
|
||||
/// Backend type (for logging/debugging)
|
||||
backend_type: DatabaseBackend,
|
||||
}
|
||||
|
||||
impl HoldingDatabase {
|
||||
/// Create a new holding database
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `relay_data_path` - Base path for relay data
|
||||
/// * `backend` - Database backend type to use
|
||||
///
|
||||
/// # Returns
|
||||
/// A new `HoldingDatabase` instance
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the database cannot be created or opened
|
||||
pub async fn new(relay_data_path: impl AsRef<Path>, backend: DatabaseBackend) -> Result<Self> {
|
||||
let relay_data_path = relay_data_path.as_ref();
|
||||
|
||||
// Construct holding database path: <relay_data_path>/holding-<backend>
|
||||
let holding_path = relay_data_path.join(format!("holding-{}", backend));
|
||||
|
||||
tracing::info!(
|
||||
"Creating holding database at {} with backend: {}",
|
||||
holding_path.display(),
|
||||
backend
|
||||
);
|
||||
|
||||
// Create database based on backend type
|
||||
let db: Arc<dyn NostrDatabase> = match backend {
|
||||
DatabaseBackend::Memory => {
|
||||
tracing::info!("Using in-memory holding database (no persistence)");
|
||||
Arc::new(MemoryDatabase::with_opts(MemoryDatabaseOptions {
|
||||
events: true,
|
||||
max_events: Some(NonZeroUsize::new(50_000).unwrap()),
|
||||
}))
|
||||
}
|
||||
DatabaseBackend::NostrDb => {
|
||||
tracing::info!(
|
||||
"Using NostrDB holding database at: {}",
|
||||
holding_path.display()
|
||||
);
|
||||
// TODO: Implement NostrDB backend once nostr-relay-builder supports it
|
||||
tracing::warn!(
|
||||
"NostrDB backend not yet implemented, using in-memory holding database"
|
||||
);
|
||||
Arc::new(MemoryDatabase::with_opts(MemoryDatabaseOptions {
|
||||
events: true,
|
||||
max_events: Some(NonZeroUsize::new(50_000).unwrap()),
|
||||
}))
|
||||
}
|
||||
DatabaseBackend::Lmdb => {
|
||||
tracing::info!("Using LMDB holding database at: {}", holding_path.display());
|
||||
// Ensure the database directory exists
|
||||
std::fs::create_dir_all(&holding_path).context(format!(
|
||||
"Failed to create LMDB holding database directory at {}",
|
||||
holding_path.display()
|
||||
))?;
|
||||
|
||||
Arc::new(NostrLmdb::open(&holding_path).await.context(format!(
|
||||
"Failed to open LMDB holding database at {}",
|
||||
holding_path.display()
|
||||
))?)
|
||||
}
|
||||
};
|
||||
|
||||
Ok(Self {
|
||||
backend: db,
|
||||
path: holding_path,
|
||||
backend_type: backend,
|
||||
})
|
||||
}
|
||||
|
||||
/// Store an event in the holding database with deletion metadata
|
||||
///
|
||||
/// The deletion metadata is stored as tags on the event:
|
||||
/// - `["deletion-ts", "<unix_timestamp>"]` - When the event was deleted
|
||||
/// - `["deletion-event", "<event_id>"]` - ID of the deletion request
|
||||
/// - `["expiry-ts", "<unix_timestamp>"]` - When to permanently delete
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `event` - The event to archive
|
||||
/// * `metadata` - Deletion metadata
|
||||
///
|
||||
/// # Returns
|
||||
/// Ok(()) if the event was stored successfully
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the event cannot be stored
|
||||
pub async fn store_event(&self, event: &Event, metadata: DeletionMetadata) -> Result<()> {
|
||||
// Create a modified event with deletion metadata tags
|
||||
// We'll store the metadata as custom tags that won't interfere with queries
|
||||
let mut modified_event = event.clone();
|
||||
|
||||
// Add deletion metadata as tags
|
||||
modified_event.tags.push(Tag::custom(
|
||||
TagKind::custom("deletion-ts"),
|
||||
vec![metadata.deletion_timestamp.to_string()],
|
||||
));
|
||||
modified_event.tags.push(Tag::custom(
|
||||
TagKind::custom("deletion-event"),
|
||||
vec![metadata.deletion_event_id.to_hex()],
|
||||
));
|
||||
modified_event.tags.push(Tag::custom(
|
||||
TagKind::custom("expiry-ts"),
|
||||
vec![metadata.expiry_timestamp.to_string()],
|
||||
));
|
||||
|
||||
// Store in the holding database
|
||||
self.backend
|
||||
.save_event(&modified_event)
|
||||
.await
|
||||
.context("Failed to save event to holding database")?;
|
||||
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
deletion_event = %metadata.deletion_event_id,
|
||||
expiry = metadata.expiry_timestamp,
|
||||
"Stored event in holding database"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Query events from the holding database
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `filter` - Nostr filter to apply
|
||||
///
|
||||
/// # Returns
|
||||
/// A vector of events matching the filter
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the query fails
|
||||
pub async fn query_events(&self, filter: Filter) -> Result<Vec<Event>> {
|
||||
let events = self
|
||||
.backend
|
||||
.query(filter)
|
||||
.await
|
||||
.context("Failed to query holding database")?;
|
||||
|
||||
// Convert Events to Vec<Event>
|
||||
Ok(events.into_iter().collect())
|
||||
}
|
||||
|
||||
/// Permanently delete an event from the holding database
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `event_id` - ID of the event to delete
|
||||
///
|
||||
/// # Returns
|
||||
/// Ok(()) if the event was deleted successfully
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the deletion fails
|
||||
pub async fn delete_event(&self, event_id: &EventId) -> Result<()> {
|
||||
// Note: NostrDatabase doesn't have a direct delete method in the current API
|
||||
// We'll need to use the delete_event method if available, or mark as deleted
|
||||
// For now, we'll use check_id to mark as deleted
|
||||
|
||||
// The NostrDatabase trait doesn't expose a direct delete method
|
||||
// We need to check if there's a way to delete events
|
||||
// Looking at the API, we might need to use a different approach
|
||||
|
||||
tracing::warn!(
|
||||
event_id = %event_id,
|
||||
"Permanent deletion requested but NostrDatabase API doesn't support direct deletion"
|
||||
);
|
||||
|
||||
// TODO: Implement permanent deletion when API supports it
|
||||
// For now, we'll just log a warning
|
||||
// This might require extending the database backend or using a different approach
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Query events that have expired (past their retention period)
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `current_timestamp` - Current unix timestamp
|
||||
///
|
||||
/// # Returns
|
||||
/// A vector of events that should be permanently deleted
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the query fails
|
||||
pub async fn query_expired(&self, current_timestamp: u64) -> Result<Vec<Event>> {
|
||||
// Query all events and filter by expiry timestamp
|
||||
// Note: This is not efficient for large databases, but works for now
|
||||
// A better approach would be to use a custom index on the expiry-ts tag
|
||||
|
||||
let filter = Filter::new(); // Get all events
|
||||
let all_events = self.query_events(filter).await?;
|
||||
|
||||
// Filter events that have expired
|
||||
let expired: Vec<Event> = all_events
|
||||
.into_iter()
|
||||
.filter(|event| {
|
||||
// Extract expiry timestamp from tags
|
||||
event
|
||||
.tags
|
||||
.iter()
|
||||
.find_map(|tag| {
|
||||
let tag_vec = tag.clone().to_vec();
|
||||
if tag_vec.len() >= 2 && tag_vec[0] == "expiry-ts" {
|
||||
tag_vec[1].parse::<u64>().ok()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.map(|expiry| expiry <= current_timestamp)
|
||||
.unwrap_or(false)
|
||||
})
|
||||
.collect();
|
||||
|
||||
tracing::debug!(
|
||||
expired_count = expired.len(),
|
||||
current_timestamp = current_timestamp,
|
||||
"Found expired events in holding database"
|
||||
);
|
||||
|
||||
Ok(expired)
|
||||
}
|
||||
|
||||
/// Get the path to the holding database
|
||||
pub fn path(&self) -> &Path {
|
||||
&self.path
|
||||
}
|
||||
|
||||
/// Get the backend type
|
||||
pub fn backend_type(&self) -> DatabaseBackend {
|
||||
self.backend_type
|
||||
}
|
||||
|
||||
/// Get a reference to the underlying database backend
|
||||
///
|
||||
/// This allows direct access to the database for advanced operations
|
||||
pub fn backend(&self) -> &Arc<dyn NostrDatabase> {
|
||||
&self.backend
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use nostr_relay_builder::prelude::{EventBuilder, Keys};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
fn create_test_event() -> Event {
|
||||
let keys = Keys::generate();
|
||||
EventBuilder::text_note("test event")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn current_timestamp() -> u64 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_secs()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_memory_holding_database() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let db = HoldingDatabase::new(temp_dir.path(), DatabaseBackend::Memory)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(db.backend_type(), DatabaseBackend::Memory);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_store_and_query_event() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let db = HoldingDatabase::new(temp_dir.path(), DatabaseBackend::Memory)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let event = create_test_event();
|
||||
let deletion_keys = Keys::generate();
|
||||
let deletion_event_id = EventBuilder::text_note("deletion")
|
||||
.sign_with_keys(&deletion_keys)
|
||||
.unwrap()
|
||||
.id;
|
||||
|
||||
let metadata = DeletionMetadata {
|
||||
deletion_timestamp: current_timestamp(),
|
||||
deletion_event_id,
|
||||
expiry_timestamp: current_timestamp() + 3600,
|
||||
};
|
||||
|
||||
// Store event
|
||||
db.store_event(&event, metadata).await.unwrap();
|
||||
|
||||
// Query event
|
||||
let filter = Filter::new().id(event.id);
|
||||
let results = db.query_events(filter).await.unwrap();
|
||||
|
||||
assert_eq!(results.len(), 1);
|
||||
assert_eq!(results[0].id, event.id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_query_expired_events() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let db = HoldingDatabase::new(temp_dir.path(), DatabaseBackend::Memory)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let event1 = create_test_event();
|
||||
let event2 = create_test_event();
|
||||
|
||||
let deletion_keys = Keys::generate();
|
||||
let deletion_event_id = EventBuilder::text_note("deletion")
|
||||
.sign_with_keys(&deletion_keys)
|
||||
.unwrap()
|
||||
.id;
|
||||
|
||||
let now = current_timestamp();
|
||||
|
||||
// Event 1: Already expired
|
||||
let metadata1 = DeletionMetadata {
|
||||
deletion_timestamp: now - 7200,
|
||||
deletion_event_id,
|
||||
expiry_timestamp: now - 3600,
|
||||
};
|
||||
|
||||
// Event 2: Not yet expired
|
||||
let metadata2 = DeletionMetadata {
|
||||
deletion_timestamp: now,
|
||||
deletion_event_id,
|
||||
expiry_timestamp: now + 3600,
|
||||
};
|
||||
|
||||
db.store_event(&event1, metadata1).await.unwrap();
|
||||
db.store_event(&event2, metadata2).await.unwrap();
|
||||
|
||||
// Query expired events
|
||||
let expired = db.query_expired(now).await.unwrap();
|
||||
|
||||
assert_eq!(expired.len(), 1);
|
||||
assert_eq!(expired[0].id, event1.id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_lmdb_holding_database() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let db = HoldingDatabase::new(temp_dir.path(), DatabaseBackend::Lmdb)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(db.backend_type(), DatabaseBackend::Lmdb);
|
||||
|
||||
// Verify the holding database directory was created
|
||||
let expected_path = temp_dir.path().join("holding-lmdb");
|
||||
assert!(expected_path.exists());
|
||||
assert_eq!(db.path(), expected_path);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_metadata_preserved_in_tags() {
|
||||
let temp_dir = tempfile::tempdir().unwrap();
|
||||
let db = HoldingDatabase::new(temp_dir.path(), DatabaseBackend::Memory)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let event = create_test_event();
|
||||
let deletion_keys = Keys::generate();
|
||||
let deletion_event_id = EventBuilder::text_note("deletion")
|
||||
.sign_with_keys(&deletion_keys)
|
||||
.unwrap()
|
||||
.id;
|
||||
|
||||
let now = current_timestamp();
|
||||
let metadata = DeletionMetadata {
|
||||
deletion_timestamp: now,
|
||||
deletion_event_id,
|
||||
expiry_timestamp: now + 3600,
|
||||
};
|
||||
|
||||
db.store_event(&event, metadata.clone()).await.unwrap();
|
||||
|
||||
// Query and verify metadata tags
|
||||
let filter = Filter::new().id(event.id);
|
||||
let results = db.query_events(filter).await.unwrap();
|
||||
|
||||
assert_eq!(results.len(), 1);
|
||||
let stored_event = &results[0];
|
||||
|
||||
// Check deletion-ts tag
|
||||
let deletion_ts = stored_event
|
||||
.tags
|
||||
.iter()
|
||||
.find_map(|tag| {
|
||||
let tag_vec = tag.clone().to_vec();
|
||||
if tag_vec.len() >= 2 && tag_vec[0] == "deletion-ts" {
|
||||
tag_vec[1].parse::<u64>().ok()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(deletion_ts, metadata.deletion_timestamp);
|
||||
|
||||
// Check expiry-ts tag
|
||||
let expiry_ts = stored_event
|
||||
.tags
|
||||
.iter()
|
||||
.find_map(|tag| {
|
||||
let tag_vec = tag.clone().to_vec();
|
||||
if tag_vec.len() >= 2 && tag_vec[0] == "expiry-ts" {
|
||||
tag_vec[1].parse::<u64>().ok()
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(expiry_ts, metadata.expiry_timestamp);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
/// Database module for ngit-grasp
|
||||
///
|
||||
/// This module contains database-related functionality including:
|
||||
/// - Holding database for archived events (NIP-09 deletion support)
|
||||
pub mod holding;
|
||||
|
||||
pub use holding::{DeletionMetadata, HoldingDatabase, SharedHoldingDatabase};
|
||||
@@ -1,4 +1,5 @@
|
||||
pub mod config;
|
||||
pub mod database;
|
||||
pub mod git;
|
||||
pub mod http;
|
||||
pub mod metrics;
|
||||
|
||||
Reference in New Issue
Block a user