mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
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.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user