From ee95b112ba041f7bea8212e8d7f06337a1b8947c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 19 Jun 2026 22:06:13 +0100 Subject: [PATCH] refactor(nostr): centralize accepted event persistence --- src/git/sync.rs | 59 +++++++++++++++++++++++-- src/grasp06/receive.rs | 1 - src/http/mod.rs | 1 - src/nostr/builder.rs | 60 ++++++++++++++++++++------ src/nostr/mod.rs | 1 + src/nostr/persistence.rs | 81 +++++++++++++++++++++++++++++++++++ src/nostr/promotion_hooks.rs | 40 +++++++++++------ src/purgatory/sync/context.rs | 2 +- src/sync/mod.rs | 5 ++- 9 files changed, 215 insertions(+), 35 deletions(-) create mode 100644 src/nostr/persistence.rs diff --git a/src/git/sync.rs b/src/git/sync.rs index 50deba7..1733a6a 100644 --- a/src/git/sync.rs +++ b/src/git/sync.rs @@ -833,6 +833,10 @@ pub fn extract_identifier_from_pr_event(event: &Event) -> Option { /// git/purgatory lifecycle events. #[async_trait] pub trait PurgatoryPromotionHooks: Send + Sync { + /// Called immediately before an accepted purgatory event is saved to the + /// main database. + async fn before_event_saved(&self, _event: &Event, _context: PurgatorySaveContext) {} + /// Called immediately before a purgatory announcement is promoted. async fn before_announcement_promote(&self, _event: &Event, _identifier: &str) {} @@ -846,6 +850,31 @@ pub trait PurgatoryPromotionHooks: Send + Sync { async fn after_pr_saved(&self, _event: &Event) {} } +/// Neutral save context for purgatory events promoted by git sync. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PurgatorySaveContext { + State, + Announcement, + Pr, +} + +async fn save_promoted_event( + database: &SharedDatabase, + event: &Event, + context: PurgatorySaveContext, + promotion_hooks: Option<&dyn PurgatoryPromotionHooks>, +) -> anyhow::Result<()> { + if let Some(hooks) = promotion_hooks { + hooks.before_event_saved(event, context).await; + } + + database + .save_event(event) + .await + .map(|_| ()) + .map_err(|e| anyhow::anyhow!("Failed to save promoted event {}: {e}", event.id)) +} + #[allow(clippy::too_many_arguments)] pub async fn process_newly_available_git_data( source_repo_path: &Path, @@ -1152,7 +1181,14 @@ async fn process_purgatory_state_events( ); } else { // Save to database - match database.save_event(&entry.event).await { + match save_promoted_event( + database, + &entry.event, + PurgatorySaveContext::State, + promotion_hooks, + ) + .await + { Ok(_) => { info!( identifier = %identifier, @@ -1220,7 +1256,14 @@ async fn process_purgatory_state_events( } if let Some(ann_event) = purgatory.promote_announcement(owner, identifier) { - match database.save_event(&ann_event).await { + match save_promoted_event( + database, + &ann_event, + PurgatorySaveContext::Announcement, + promotion_hooks, + ) + .await + { Ok(_) => { info!( identifier = %identifier, @@ -1428,7 +1471,8 @@ async fn process_purgatory_pr_events( result.errors.extend(process_result.errors); // Save event to database - match database.save_event(event).await { + match save_promoted_event(database, event, PurgatorySaveContext::Pr, promotion_hooks).await + { Ok(_) => { info!( identifier = %identifier, @@ -1540,7 +1584,14 @@ async fn process_purgatory_announcements( if let Some(event) = announcement_event { // Save to database - match database.save_event(&event).await { + match save_promoted_event( + database, + &event, + PurgatorySaveContext::Announcement, + promotion_hooks, + ) + .await + { Ok(_) => { info!( identifier = %identifier, diff --git a/src/grasp06/receive.rs b/src/grasp06/receive.rs index d39011a..52aeba4 100644 --- a/src/grasp06/receive.rs +++ b/src/grasp06/receive.rs @@ -353,7 +353,6 @@ pub async fn handle_prs_receive_pack( let promotion_hooks = NostrPurgatoryPromotionHooks::git_push( &write_policy, &rejected_events_index, - &database, Some(&relay), ); diff --git a/src/http/mod.rs b/src/http/mod.rs index 1a51c04..99da0ae 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -530,7 +530,6 @@ impl Service> for HttpService { let promotion_hooks = NostrPurgatoryPromotionHooks::git_push( &write_policy, &rejected_events_index, - &database, Some(&relay), ); diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index 00928d5..225378b 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -19,6 +19,7 @@ use crate::nostr::deletion::{DeletionContext, DeletionService}; use crate::nostr::events::RepositoryAnnouncement; use crate::nostr::history::ReplaceableHistoryStore; use crate::nostr::lifecycle::RepositoryLifecycle; +use crate::nostr::persistence::{EventPersistence, SaveContext}; use crate::nostr::policy::{ accepted_purgatory, duplicate, reject_error, reject_invalid, reject_restricted, AnnouncementPolicy, AnnouncementResult, PolicyContext, PrEventPolicy, ReferenceResult, @@ -150,6 +151,29 @@ impl Nip34WritePolicy { &self.deletion } + fn event_persistence(&self) -> EventPersistence<'_> { + EventPersistence::new(&self.ctx.database, &self.deletion) + } + + /// Run central persistence hooks for an accepted event whose actual save is + /// performed by the relay builder after policy acceptance. + pub async fn prepare_accepted_event_persistence(&self, event: &Event, context: SaveContext) { + self.event_persistence() + .before_accepted_event_persisted(event, context) + .await; + } + + /// Save an already-admitted event through the central accepted-event path. + pub async fn save_accepted_event( + &self, + event: &Event, + context: SaveContext, + ) -> anyhow::Result<()> { + self.event_persistence() + .save_accepted_event(event, context) + .await + } + /// Extract repository identifier from event's 'd' tag. /// /// Used for structured logging when parsing fails - we try to extract @@ -243,9 +267,11 @@ impl Nip34WritePolicy { self.check_purgatory_state_events_for_identifier(&announcement.identifier) .await; - self.deletion - .capture_superseded_replaceable_history(event) - .await; + self.prepare_accepted_event_persistence( + event, + SaveContext::DirectWritePolicyAcceptance, + ) + .await; WritePolicyResult::Accept } @@ -300,9 +326,11 @@ impl Nip34WritePolicy { e ); } - self.deletion - .capture_superseded_replaceable_history(event) - .await; + self.prepare_accepted_event_persistence( + event, + SaveContext::DirectWritePolicyAcceptance, + ) + .await; return WritePolicyResult::Accept; } @@ -352,9 +380,11 @@ impl Nip34WritePolicy { self.check_purgatory_state_events_for_identifier(&announcement.identifier) .await; - self.deletion - .capture_superseded_replaceable_history(event) - .await; + self.prepare_accepted_event_persistence( + event, + SaveContext::DirectWritePolicyAcceptance, + ) + .await; WritePolicyResult::Accept } @@ -401,7 +431,7 @@ impl Nip34WritePolicy { async fn handle_state(&self, event: &Event, is_synced: bool) -> WritePolicyResult { match self.state_policy.validate(event) { StateResult::Accept => { - let promotion_hooks = NostrPurgatoryPromotionHooks::recovery_only(self.deletion()); + let promotion_hooks = NostrPurgatoryPromotionHooks::recovery_only(self); // Process state alignment asynchronously match self @@ -411,9 +441,11 @@ impl Nip34WritePolicy { { Ok(policy_result) => { if Self::event_persists_to_main_db(&policy_result) { - self.deletion - .capture_superseded_replaceable_history(event) - .await; + self.prepare_accepted_event_persistence( + event, + SaveContext::DirectWritePolicyAcceptance, + ) + .await; } policy_result } @@ -622,7 +654,7 @@ impl Nip34WritePolicy { ); for entry in state_events { - let promotion_hooks = NostrPurgatoryPromotionHooks::recovery_only(self.deletion()); + let promotion_hooks = NostrPurgatoryPromotionHooks::recovery_only(self); // Re-evaluate authorization with the new announcement match self diff --git a/src/nostr/mod.rs b/src/nostr/mod.rs index 729464b..87c86bf 100644 --- a/src/nostr/mod.rs +++ b/src/nostr/mod.rs @@ -4,6 +4,7 @@ pub mod events; pub mod history; pub mod holding; pub mod lifecycle; +pub mod persistence; pub mod policy; pub mod promotion_hooks; pub mod tombstones; diff --git a/src/nostr/persistence.rs b/src/nostr/persistence.rs new file mode 100644 index 0000000..18db1ef --- /dev/null +++ b/src/nostr/persistence.rs @@ -0,0 +1,81 @@ +//! Central persistence hooks for accepted GRASP events. +//! +//! The normal websocket write path is persisted by `nostr-relay-builder` after +//! `WritePolicyResult::Accept`, while sync/purgatory paths save directly. This +//! module keeps the pre-save side effects that must happen for every accepted +//! main-DB persistence path in one Nostr-owned facade. + +use nostr_relay_builder::prelude::Event; + +use crate::nostr::deletion::DeletionService; +use crate::nostr::SharedDatabase; + +/// Where an accepted event is being prepared/saved from. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SaveContext { + /// Websocket/local relay path. The database save happens outside this crate + /// after the write policy returns `Accept`. + DirectWritePolicyAcceptance, + /// Event accepted while consuming a remote relay subscription. + RelaySync, + /// Event accepted during rejected hot-cache reprocessing. + HotCacheReprocess, + /// State event released from git purgatory into the main DB. + PurgatoryStatePromotion, + /// Announcement event released from git purgatory into the main DB. + PurgatoryAnnouncementPromotion, + /// PR/PR-update event released from git purgatory into the main DB. + PurgatoryPrPromotion, +} + +/// Nostr-owned facade for accepted-event persistence side effects. +pub struct EventPersistence<'a> { + database: &'a SharedDatabase, + deletion: &'a DeletionService, +} + +impl<'a> EventPersistence<'a> { + pub fn new(database: &'a SharedDatabase, deletion: &'a DeletionService) -> Self { + Self { database, deletion } + } + + /// Run side effects that must occur immediately before an accepted event is + /// persisted to the main database. + pub async fn before_accepted_event_persisted(&self, event: &Event, context: SaveContext) { + tracing::trace!( + event_id = %event.id, + kind = event.kind.as_u16(), + ?context, + "Preparing accepted event for main database persistence" + ); + + self.deletion + .capture_superseded_replaceable_history(event) + .await; + } + + /// Save an already-accepted event to the main database via the central + /// persistence path. + pub async fn save_accepted_event( + &self, + event: &Event, + context: SaveContext, + ) -> anyhow::Result<()> { + self.before_accepted_event_persisted(event, context).await; + + self.database + .save_event(event) + .await + .map(|_| ()) + .map_err(|e| anyhow::anyhow!("Failed to save accepted event {}: {e}", event.id))?; + + tracing::trace!( + event_id = %event.id, + kind = event.kind.as_u16(), + ?context, + "Saved accepted event to main database" + ); + + Ok(()) + } +} diff --git a/src/nostr/promotion_hooks.rs b/src/nostr/promotion_hooks.rs index 0870fd2..1c27386 100644 --- a/src/nostr/promotion_hooks.rs +++ b/src/nostr/promotion_hooks.rs @@ -14,29 +14,27 @@ use nostr_relay_builder::LocalRelay; use nostr_sdk::prelude::PublicKey; use tracing::{debug, info, warn}; -use crate::git::sync::PurgatoryPromotionHooks; +use crate::git::sync::{PurgatoryPromotionHooks, PurgatorySaveContext}; use crate::nostr::builder::Nip34WritePolicy; use crate::nostr::deletion::DeletionService; use crate::nostr::events::RepositoryAnnouncement; -use crate::nostr::SharedDatabase; +use crate::nostr::persistence::SaveContext; use crate::sync::rejected_index::{EventType, RejectedEventsIndex}; pub struct NostrPurgatoryPromotionHooks<'a> { deletion: &'a DeletionService, write_policy: Option<&'a Nip34WritePolicy>, rejected_events_index: Option<&'a Arc>, - database: Option<&'a SharedDatabase>, local_relay: Option<&'a LocalRelay>, } impl<'a> NostrPurgatoryPromotionHooks<'a> { - /// Run deletion recovery hooks only. - pub fn recovery_only(deletion: &'a DeletionService) -> Self { + /// Run deletion recovery and accepted-event persistence hooks only. + pub fn recovery_only(write_policy: &'a Nip34WritePolicy) -> Self { Self { - deletion, - write_policy: None, + deletion: write_policy.deletion(), + write_policy: Some(write_policy), rejected_events_index: None, - database: None, local_relay: None, } } @@ -45,14 +43,12 @@ impl<'a> NostrPurgatoryPromotionHooks<'a> { pub fn git_push( write_policy: &'a Nip34WritePolicy, rejected_events_index: &'a Arc, - database: &'a SharedDatabase, local_relay: Option<&'a LocalRelay>, ) -> Self { Self { deletion: write_policy.deletion(), write_policy: Some(write_policy), rejected_events_index: Some(rejected_events_index), - database: Some(database), local_relay, } } @@ -60,6 +56,22 @@ impl<'a> NostrPurgatoryPromotionHooks<'a> { #[async_trait] impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks<'_> { + async fn before_event_saved(&self, event: &Event, context: PurgatorySaveContext) { + let Some(write_policy) = self.write_policy else { + return; + }; + + let context = match context { + PurgatorySaveContext::State => SaveContext::PurgatoryStatePromotion, + PurgatorySaveContext::Announcement => SaveContext::PurgatoryAnnouncementPromotion, + PurgatorySaveContext::Pr => SaveContext::PurgatoryPrPromotion, + }; + + write_policy + .prepare_accepted_event_persistence(event, context) + .await; + } + async fn before_announcement_promote(&self, event: &Event, identifier: &str) { self.deletion .maybe_recover_deleted_repository(event, identifier) @@ -67,10 +79,9 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks<'_> { } async fn after_announcement_saved(&self, event: &Event) { - let (Some(write_policy), Some(rejected_events_index), Some(database), Some(relay)) = ( + let (Some(write_policy), Some(rejected_events_index), Some(relay)) = ( self.write_policy, self.rejected_events_index, - self.database, self.local_relay, ) else { return; @@ -121,7 +132,10 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks<'_> { match write_policy.admit_event(&hot_event, &dummy_addr).await { WritePolicyResult::Accept => { - match database.save_event(&hot_event).await { + match write_policy + .save_accepted_event(&hot_event, SaveContext::HotCacheReprocess) + .await + { Ok(_) => { relay.notify_event(hot_event.clone()); info!( diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index a1a435c..85014e0 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -523,7 +523,7 @@ impl SyncContext for RealSyncContext { let promotion_hooks = self .write_policy .as_ref() - .map(|wp| NostrPurgatoryPromotionHooks::recovery_only(wp.deletion())); + .map(NostrPurgatoryPromotionHooks::recovery_only); let promotion_hooks = promotion_hooks .as_ref() .map(|hooks| hooks as &dyn PurgatoryPromotionHooks); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 2bb9c22..83ddc06 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1788,6 +1788,7 @@ impl SyncManager { &write_policy, &local_relay, &rejected_events_index, + crate::nostr::persistence::SaveContext::RelaySync, ) .await; // Only record metric when event is actually saved @@ -2637,6 +2638,7 @@ impl SyncManager { write_policy, local_relay, rejected_events_index, + crate::nostr::persistence::SaveContext::HotCacheReprocess, )) .await; @@ -2699,6 +2701,7 @@ impl SyncManager { write_policy: &Nip34WritePolicy, local_relay: &LocalRelay, rejected_events_index: &Arc, + save_context: crate::nostr::persistence::SaveContext, ) -> ProcessResult { use nostr_relay_builder::prelude::{WritePolicy, WritePolicyResult}; use std::net::{IpAddr, Ipv4Addr, SocketAddr}; @@ -2723,7 +2726,7 @@ impl SyncManager { WritePolicyResult::Accept => { // Save event to database - if let Err(e) = database.save_event(event).await { + if let Err(e) = write_policy.save_accepted_event(event, save_context).await { tracing::error!( event_id = %event.id, relay = %relay_url,