From ed8bdae231408a58e0962c80eb15dfe92fde462c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 19 Aug 2026 20:57:55 +0000 Subject: [PATCH] fix(sync): count only inserted events as saved Repeated cold syncs reported nearly identical saved totals because the accepted-event persistence facade discarded SaveEventStatus. A stale replaceable event rejected as Replaced was consequently counted as a new insert, broadcast, and allowed to trigger downstream dependency work.\n\nReturn a bounded saved-or-duplicate outcome from the central facade. Sync now reports Duplicate for database Duplicate/Replaced results and broadcasts or expands dependencies only after SaveEventStatus::Success. Hot-cache promotion likewise avoids broadcasting superseded events while retaining the dependency retry needed for an already-stored maintainer announcement.\n\nThis deliberately does not change write-policy admission, database replacement semantics, or the historical mailbox filters. A focused regression preloads a newer replaceable event and verifies that syncing its predecessor is classified as duplicate without displacing the stored event.\n\nValidation: git diff --check. The focused test and broader suite were not executed in this session per operator instruction. --- src/nostr/builder.rs | 4 +- src/nostr/persistence.rs | 35 ++++++++++-- src/purgatory/promotion_hooks.rs | 28 ++++++++- src/sync/mod.rs | 98 ++++++++++++++++++++++++++++---- 4 files changed, 144 insertions(+), 21 deletions(-) diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index 3586b2e..5b8e23f 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -24,7 +24,7 @@ use crate::nostr::lifecycle::{ DeletionContext, DeletionRuntime, DeletionService, HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle, Tombstones, }; -use crate::nostr::persistence::{EventPersistence, SaveContext}; +use crate::nostr::persistence::{AcceptedEventSaveOutcome, EventPersistence, SaveContext}; use crate::nostr::policy::{ accepted_purgatory, duplicate, reject_error, reject_invalid, reject_restricted, AnnouncementPolicy, AnnouncementResult, IdentityAdmission, PolicyContext, PrEventPolicy, @@ -257,7 +257,7 @@ impl Nip34WritePolicy { &self, event: &Event, context: SaveContext, - ) -> anyhow::Result<()> { + ) -> anyhow::Result { self.event_persistence() .save_accepted_event(event, context) .await diff --git a/src/nostr/persistence.rs b/src/nostr/persistence.rs index 1e3f1df..da360e5 100644 --- a/src/nostr/persistence.rs +++ b/src/nostr/persistence.rs @@ -5,7 +5,7 @@ //! module provides the direct accepted-event save facade for paths outside the //! relay builder. -use nostr_sdk::prelude::Event; +use nostr_sdk::prelude::{Event, RejectedReason, SaveEventStatus}; use crate::nostr::SharedDatabase; @@ -21,6 +21,15 @@ pub enum SaveContext { HotCacheReprocess, } +/// Durable outcome of saving an event that already passed write policy. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum AcceptedEventSaveOutcome { + /// The database inserted the event. + Saved, + /// The database already had this event or retained a newer replaceable event. + Duplicate, +} + /// Nostr-owned facade for accepted-event persistence side effects. pub struct EventPersistence<'a> { database: &'a SharedDatabase, @@ -37,20 +46,34 @@ impl<'a> EventPersistence<'a> { &self, event: &Event, context: SaveContext, - ) -> anyhow::Result<()> { - self.database + ) -> anyhow::Result { + let status = self + .database .save_event(event) .await - .map(|_| ()) .map_err(|e| anyhow::anyhow!("Failed to save accepted event {}: {e}", event.id))?; + let outcome = match status { + SaveEventStatus::Success => AcceptedEventSaveOutcome::Saved, + SaveEventStatus::Rejected( + RejectedReason::Duplicate | RejectedReason::Replaced, + ) => AcceptedEventSaveOutcome::Duplicate, + SaveEventStatus::Rejected(reason) => { + return Err(anyhow::anyhow!( + "Database rejected accepted event {}: {reason:?}", + event.id + )); + } + }; + tracing::trace!( event_id = %event.id, kind = event.kind.as_u16(), ?context, - "Saved accepted event to main database" + ?outcome, + "Applied accepted event to main database" ); - Ok(()) + Ok(outcome) } } diff --git a/src/purgatory/promotion_hooks.rs b/src/purgatory/promotion_hooks.rs index 9437017..111ac5c 100644 --- a/src/purgatory/promotion_hooks.rs +++ b/src/purgatory/promotion_hooks.rs @@ -18,7 +18,7 @@ use crate::git::sync::{PurgatoryPromotionHooks, PurgatorySaveContext}; use crate::nostr::builder::Nip34WritePolicy; use crate::nostr::events::RepositoryAnnouncement; use crate::nostr::lifecycle::DeletionService; -use crate::nostr::persistence::SaveContext; +use crate::nostr::persistence::{AcceptedEventSaveOutcome, SaveContext}; use crate::sync::rejected_index::{EventType, RejectedEventsIndex}; pub struct NostrPurgatoryPromotionHooks { @@ -88,10 +88,17 @@ impl NostrPurgatoryPromotionHooks { .save_accepted_event(&state, SaveContext::HotCacheReprocess) .await { - Ok(_) => { + Ok(AcceptedEventSaveOutcome::Saved) => { rejected_events_index.remove(&state.id); relay.notify_event(state.clone()); } + Ok(AcceptedEventSaveOutcome::Duplicate) => { + rejected_events_index.remove(&state.id); + debug!( + event_id = %state.id, + "Re-processed state event was already superseded or stored" + ); + } Err(e) => { warn!( event_id = %state.id, @@ -194,7 +201,7 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks { .save_accepted_event(&hot_event, SaveContext::HotCacheReprocess) .await { - Ok(_) => { + Ok(AcceptedEventSaveOutcome::Saved) => { rejected_events_index.remove(&hot_event.id); relay.notify_event(hot_event.clone()); self.reprocess_state_dependencies( @@ -210,6 +217,21 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks { "Maintainer announcement accepted and saved on re-processing" ); } + Ok(AcceptedEventSaveOutcome::Duplicate) => { + rejected_events_index.remove(&hot_event.id); + self.reprocess_state_dependencies( + write_policy, + rejected_events_index, + relay, + &hot_event.pubkey, + &announcement.identifier, + ) + .await; + debug!( + event_id = %hot_event.id, + "Re-processed maintainer announcement was already superseded or stored" + ); + } Err(e) => { warn!( event_id = %hot_event.id, diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 1582ccb..c86b2ac 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -588,7 +588,7 @@ pub enum SyncMethod { pub enum ProcessResult { /// Event was new and saved to database Saved, - /// Event already existed in database + /// Event already existed or a newer replaceable event was retained Duplicate, /// Event added to Purgatory Purgatory, @@ -7567,15 +7567,25 @@ impl SyncManager { match result { WritePolicyResult::Accept => { // Save event to database - - if let Err(e) = write_policy.save_accepted_event(event, save_context).await { - tracing::error!( - event_id = %event.id, - relay = %relay_url, - error = %e, - "Failed to save synced event" - ); - return ProcessResult::PersistenceError; + match write_policy.save_accepted_event(event, save_context).await { + Ok(crate::nostr::persistence::AcceptedEventSaveOutcome::Saved) => {} + Ok(crate::nostr::persistence::AcceptedEventSaveOutcome::Duplicate) => { + tracing::trace!( + event_id = %event.id, + relay = %relay_url, + "Database retained an existing event during sync" + ); + return ProcessResult::Duplicate; + } + Err(e) => { + tracing::error!( + event_id = %event.id, + relay = %relay_url, + error = %e, + "Failed to save synced event" + ); + return ProcessResult::PersistenceError; + } } // Broadcast to WebSocket subscribers (enables recursive relay discovery) @@ -9308,6 +9318,74 @@ mod tests { assert!(!rejected.contains(&child.id)); } + #[tokio::test] + async fn replaced_sync_event_is_reported_as_duplicate() { + let directory = tempfile::tempdir().expect("create test directory"); + let git_data_path = directory.path().join("git"); + let mut config = Config::for_testing(); + config.git_data_path = git_data_path.to_string_lossy().into_owned(); + config.relay_data_path = directory + .path() + .join("relay") + .to_string_lossy() + .into_owned(); + let purgatory = Arc::new(crate::purgatory::Purgatory::new(git_data_path)); + let runtime = crate::nostr::builder::create_relay( + &config, + purgatory, + crate::grasp06::receive::RepoInitLocks::default(), + None, + ) + .await + .expect("create test relay runtime"); + let rejected = Arc::new(RejectedEventsIndex::new( + Duration::from_secs(120), + Duration::from_secs(604800), + )); + let keys = Keys::generate(); + let older = EventBuilder::new(Kind::GitUserGraspList, "older") + .custom_created_at(Timestamp::from_secs(100)) + .finalize(&keys) + .expect("build older replaceable event"); + let newer = EventBuilder::new(Kind::GitUserGraspList, "newer") + .custom_created_at(Timestamp::from_secs(200)) + .finalize(&keys) + .expect("build newer replaceable event"); + runtime + .stores + .database + .save_event(&newer) + .await + .expect("save newer replaceable event"); + + let result = SyncManager::process_event_static( + &older, + "wss://source.example", + &runtime.stores.database, + &runtime.write_policy, + &runtime.relay, + &rejected, + crate::nostr::persistence::SaveContext::RelaySync, + ) + .await; + + assert_eq!(result, ProcessResult::Duplicate); + assert!(runtime + .stores + .database + .event_by_id(&older.id) + .await + .unwrap() + .is_none()); + assert!(runtime + .stores + .database + .event_by_id(&newer.id) + .await + .unwrap() + .is_some()); + } + #[test] fn private_members_include_only_owners_of_accepted_relays() { let configured = Keys::generate().public_key();