From 74c431a5e56458a9bfc06c09dfed252b3e7aa7cd Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 24 Jun 2026 12:06:40 +0100 Subject: [PATCH] fix(nostr): avoid duplicate superseded history capture Superseded replaceable/addressable history is now captured by the admission hooks for normal accepted announcement/state writes, and directly by the purgatory promotion before-save hook for promoted events. This keeps capture before the replacement is persisted, so the old active event is still queryable, while avoiding a second capture from the central save path used by sync and hot-cache reprocessing. Remove the now-empty accepted-event pre-persistence hook and the purgatory SaveContext variants that only existed to route through it. Add a regression test for the sync-style admit-then-save path that crosses a timestamp boundary and asserts only one history metadata record is created. --- src/nostr/builder.rs | 10 +-- src/nostr/persistence.rs | 35 ++------- src/purgatory/promotion_hooks.rs | 19 ++--- tests/replaceable_history.rs | 125 +++++++++++++++++++++++++++++++ 4 files changed, 139 insertions(+), 50 deletions(-) diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index a58bfba..1bba2f0 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -153,15 +153,7 @@ impl Nip34WritePolicy { } 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; + EventPersistence::new(&self.ctx.database) } /// Save an already-admitted event through the central accepted-event path. diff --git a/src/nostr/persistence.rs b/src/nostr/persistence.rs index 5fd6e72..7091027 100644 --- a/src/nostr/persistence.rs +++ b/src/nostr/persistence.rs @@ -1,13 +1,12 @@ //! 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. +//! `WritePolicyResult::Accept`, while sync/hot-cache paths save directly. This +//! module provides the direct accepted-event save facade for paths outside the +//! relay builder. use nostr_relay_builder::prelude::Event; -use crate::nostr::lifecycle::DeletionService; use crate::nostr::SharedDatabase; /// Where an accepted event is being prepared/saved from. @@ -20,38 +19,16 @@ pub enum SaveContext { 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; + pub fn new(database: &'a SharedDatabase) -> Self { + Self { database } } /// Save an already-accepted event to the main database via the central @@ -61,8 +38,6 @@ impl<'a> EventPersistence<'a> { event: &Event, context: SaveContext, ) -> anyhow::Result<()> { - self.before_accepted_event_persisted(event, context).await; - self.database .save_event(event) .await diff --git a/src/purgatory/promotion_hooks.rs b/src/purgatory/promotion_hooks.rs index 6674cec..9716408 100644 --- a/src/purgatory/promotion_hooks.rs +++ b/src/purgatory/promotion_hooks.rs @@ -57,18 +57,15 @@ 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; - }; + tracing::trace!( + event_id = %event.id, + kind = event.kind.as_u16(), + ?context, + "Preparing purgatory-promoted event for main database persistence" + ); - let context = match context { - PurgatorySaveContext::State => SaveContext::PurgatoryStatePromotion, - PurgatorySaveContext::Announcement => SaveContext::PurgatoryAnnouncementPromotion, - PurgatorySaveContext::Pr => SaveContext::PurgatoryPrPromotion, - }; - - write_policy - .prepare_accepted_event_persistence(event, context) + self.deletion + .capture_superseded_replaceable_history(event) .await; } diff --git a/tests/replaceable_history.rs b/tests/replaceable_history.rs index 117a1f0..2335641 100644 --- a/tests/replaceable_history.rs +++ b/tests/replaceable_history.rs @@ -2,10 +2,20 @@ mod common; +use clap::Parser; use common::{publish_served_repo, TestRelay}; use grasp_audit::{AuditClient, AuditConfig, DETERMINISTIC_COMMIT_HASH}; +use ngit_grasp::config::Config; +use ngit_grasp::grasp06::receive::new_repo_init_locks; +use ngit_grasp::nostr::builder::Nip34WritePolicy; use ngit_grasp::nostr::lifecycle::ReplaceableHistoryStore; +use ngit_grasp::nostr::lifecycle::{HoldingStore, RepositoryLifecycle, Tombstones}; +use ngit_grasp::nostr::persistence::SaveContext; +use ngit_grasp::nostr::SharedDatabase; +use ngit_grasp::purgatory::Purgatory; use nostr_sdk::prelude::*; +use std::net::{IpAddr, Ipv4Addr, SocketAddr}; +use std::sync::Arc; use std::time::Duration; async fn open_history(relay_data_path: &std::path::Path) -> ReplaceableHistoryStore { @@ -66,6 +76,57 @@ fn build_state_version( .expect("build state version") } +fn build_policy_for_history_regression( + database: SharedDatabase, + history: ReplaceableHistoryStore, + git_data_path: &std::path::Path, +) -> Nip34WritePolicy { + let config = Config::parse_from([ + "ngit-grasp-test", + "--domain", + "test.example.com", + "--archive-all", + "--archive-read-only", + "false", + ]); + let purgatory = Arc::new(Purgatory::new(git_data_path.to_path_buf())); + + Nip34WritePolicy::new( + database, + Tombstones::in_memory(), + HoldingStore::in_memory(), + RepositoryLifecycle::in_memory(), + history, + git_data_path.to_path_buf(), + purgatory, + config, + new_repo_init_locks(), + ) +} + +fn build_policy_announcement_version( + keys: &Keys, + repo_id: &str, + created_at: Timestamp, + content: &str, +) -> Event { + let npub = keys.public_key().to_bech32().expect("pubkey to npub"); + + EventBuilder::new(Kind::GitRepoAnnouncement, content) + .tags(vec![ + Tag::identifier(repo_id), + Tag::custom("name", vec![repo_id.to_string()]), + Tag::custom( + "clone", + vec![format!("http://test.example.com/{}/{}.git", npub, repo_id)], + ), + Tag::custom("relays", vec!["wss://test.example.com".to_string()]), + ]) + .custom_created_at(created_at) + .finalize(keys) + .expect("build policy announcement version") +} + #[tokio::test] async fn newer_30618_supersedes_older_and_preserves_old_in_history() { let relay = TestRelay::start_with_lmdb().await; @@ -176,6 +237,70 @@ async fn newer_30617_supersedes_older_and_preserves_old_in_history() { relay.stop().await; } +#[tokio::test] +async fn sync_style_admit_then_save_captures_superseded_history_once() { + use nostr_relay_builder::prelude::{WritePolicy, WritePolicyResult}; + + let git_dir = tempfile::tempdir().expect("git tempdir"); + let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let history = ReplaceableHistoryStore::in_memory(); + let policy = build_policy_for_history_regression(db.clone(), history.clone(), git_dir.path()); + let keys = Keys::generate(); + let repo_id = format!("history-sync-once-{}", &keys.public_key().to_hex()[..8]); + let old_announcement = build_policy_announcement_version( + &keys, + &repo_id, + Timestamp::from_secs(Timestamp::now().as_secs() + 1), + "old", + ); + let new_announcement = build_policy_announcement_version( + &keys, + &repo_id, + Timestamp::from_secs(old_announcement.created_at.as_secs() + 30), + "new", + ); + + db.save_event(&old_announcement) + .await + .expect("save old announcement"); + + let dummy_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0); + let result = policy.admit_event(&new_announcement, &dummy_addr).await; + assert!( + matches!(result, WritePolicyResult::Accept), + "new announcement must be accepted before central save" + ); + + // Cross a timestamp-second boundary before saving. The old persistence hook + // also captured history here, producing a distinct metadata event instead + // of collapsing onto the admission hook's metadata ID. + tokio::time::sleep(Duration::from_millis(1100)).await; + policy + .save_accepted_event(&new_announcement, SaveContext::RelaySync) + .await + .expect("save accepted announcement"); + + let coordinate = format!("30617:{}:{}", keys.public_key().to_hex(), repo_id); + let records = history + .superseded_records_for_coordinate_before( + &coordinate, + Timestamp::from_secs(new_announcement.created_at.as_secs() + 1), + ) + .await; + let matching_records = records + .iter() + .filter(|r| { + r.superseded_event_id == Some(old_announcement.id) + && r.replaced_by == Some(new_announcement.id) + }) + .count(); + + assert_eq!( + matching_records, 1, + "sync-style admit then save must not duplicate superseded history metadata" + ); +} + #[tokio::test] async fn replaceable_history_survives_restart_lmdb() { let relay = TestRelay::start_with_lmdb().await;