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.
This commit is contained in:
DanConwayDev
2026-08-19 20:57:55 +00:00
parent c268ea3085
commit ed8bdae231
4 changed files with 144 additions and 21 deletions
+2 -2
View File
@@ -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<AcceptedEventSaveOutcome> {
self.event_persistence()
.save_accepted_event(event, context)
.await
+29 -6
View File
@@ -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<AcceptedEventSaveOutcome> {
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)
}
}
+25 -3
View File
@@ -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,
+88 -10
View File
@@ -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();