From d5c7dbfbf60ae9e1f0103b2624c6228a508d1d64 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 17 Jun 2026 15:33:15 +0000 Subject: [PATCH] feat(deletion): persist replaceable history for rollback --- src/nostr/builder.rs | 113 ++++++++++++- src/nostr/history.rs | 301 +++++++++++++++++++++++++++++++++++ src/nostr/mod.rs | 1 + src/nostr/policy/mod.rs | 6 + tests/nip09_blacklist_ops.rs | 2 + tests/replaceable_history.rs | 281 ++++++++++++++++++++++++++++++++ 6 files changed, 701 insertions(+), 3 deletions(-) create mode 100644 src/nostr/history.rs create mode 100644 tests/replaceable_history.rs diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index c782296..13c0718 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -18,6 +18,7 @@ use tar::Archive as TarArchive; use crate::config::{Config, DatabaseBackend}; use crate::nostr::events::RepositoryAnnouncement; +use crate::nostr::history::ReplaceableHistoryStore; use crate::nostr::holding::RecoveryMetadataRecord; use crate::nostr::policy::{ accepted_purgatory, duplicate, reject_error, reject_invalid, reject_restricted, @@ -72,6 +73,7 @@ impl Nip34WritePolicy { database: SharedDatabase, tombstones: crate::nostr::tombstones::Tombstones, holding: crate::nostr::holding::HoldingStore, + history: crate::nostr::history::ReplaceableHistoryStore, git_data_path: impl Into, purgatory: std::sync::Arc, config: crate::config::Config, @@ -82,6 +84,7 @@ impl Nip34WritePolicy { database, tombstones, holding, + history, git_data_path, purgatory, config.clone(), @@ -120,6 +123,11 @@ impl Nip34WritePolicy { &self.ctx.holding } + /// Get a reference to the replaceable-history store. + pub fn history(&self) -> &crate::nostr::history::ReplaceableHistoryStore { + &self.ctx.history + } + /// Startup-only blacklist parity pass. /// /// Scans already-stored kind-30617 announcements, identifies entries @@ -305,6 +313,8 @@ impl Nip34WritePolicy { self.check_purgatory_state_events_for_identifier(&announcement.identifier) .await; + self.capture_superseded_replaceable_history(event).await; + WritePolicyResult::Accept } Err(e) => { @@ -357,6 +367,7 @@ impl Nip34WritePolicy { e ); } + self.capture_superseded_replaceable_history(event).await; return WritePolicyResult::Accept; } @@ -368,6 +379,8 @@ impl Nip34WritePolicy { event_id_str ); + self.capture_superseded_replaceable_history(event).await; + accepted_purgatory("won't be served until git data arrives") } Err(e) => { @@ -396,6 +409,8 @@ impl Nip34WritePolicy { self.check_purgatory_state_events_for_identifier(&announcement.identifier) .await; + self.capture_superseded_replaceable_history(event).await; + WritePolicyResult::Accept } Err(e) => { @@ -444,7 +459,12 @@ impl Nip34WritePolicy { .process_state_event(event, is_synced) .await { - Ok(poilicy_result) => poilicy_result, + Ok(policy_result) => { + if Self::event_persists_to_main_db(&policy_result) { + self.capture_superseded_replaceable_history(event).await; + } + policy_result + } Err(e) => { let npub = event .pubkey @@ -844,6 +864,86 @@ impl Nip34WritePolicy { )) } + fn extract_addressable_identifier(event: &Event) -> Option { + event.tags.iter().find_map(|tag| { + let v = tag.as_slice(); + if v.len() >= 2 && v[0] == "d" { + Some(v[1].clone()) + } else { + None + } + }) + } + + fn should_capture_replaceable_history(kind: Kind) -> bool { + kind == Kind::GitRepoAnnouncement || kind == Kind::RepoState + } + + fn event_persists_to_main_db(result: &WritePolicyResult) -> bool { + matches!( + result, + WritePolicyResult::Accept | WritePolicyResult::Reject { status: true, .. } + ) + } + + async fn capture_superseded_replaceable_history(&self, incoming: &Event) { + if !Self::should_capture_replaceable_history(incoming.kind) { + return; + } + + let Some(coordinate) = Self::event_coordinate(incoming) else { + return; + }; + + let Some(identifier) = Self::extract_addressable_identifier(incoming) else { + tracing::warn!( + event_id = %incoming.id.to_hex(), + kind = incoming.kind.as_u16(), + "Skipping history capture for addressable event without d tag" + ); + return; + }; + + let filter = Filter::new() + .kind(incoming.kind) + .author(incoming.pubkey) + .identifier(identifier); + + let existing = match self.ctx.database.query(filter).await { + Ok(events) => events, + Err(e) => { + tracing::warn!( + error = %e, + event_id = %incoming.id.to_hex(), + coordinate = %coordinate, + "Failed querying existing events for history capture" + ); + return; + } + }; + + for superseded in existing { + if superseded.id == incoming.id || superseded.created_at >= incoming.created_at { + continue; + } + + if let Err(e) = self + .ctx + .history + .archive_superseded_event(&superseded, incoming, &coordinate) + .await + { + tracing::warn!( + error = %e, + superseded_event_id = %superseded.id.to_hex(), + replaced_by = %incoming.id.to_hex(), + coordinate = %coordinate, + "Failed to archive superseded replaceable/addressable event" + ); + } + } + } + /// Handle a NIP-62 request-to-vanish (kind 62). /// /// Reproduces the LMDB backend's vanish behaviour now that we run with @@ -1007,6 +1107,8 @@ pub struct RelayWithDatabase { pub write_policy: Nip34WritePolicy, /// Holding store used for deletion archival + expiry cleanup pub holding: crate::nostr::holding::HoldingStore, + /// Replaceable-history store used for rollback history capture. + pub history: crate::nostr::history::ReplaceableHistoryStore, } /// Create a configured LocalRelay with full GRASP-01 validation @@ -1033,10 +1135,11 @@ pub async fn create_relay( // would silently hard-delete events and block re-submission via its own // internal tables that we cannot inspect, conflicting with our purgatory / // bare-repo bookkeeping. - let (database, tombstones, holding): ( + let (database, tombstones, holding, history): ( SharedDatabase, crate::nostr::tombstones::Tombstones, crate::nostr::holding::HoldingStore, + ReplaceableHistoryStore, ) = match config.database_backend { DatabaseBackend::Memory => { tracing::info!("Using in-memory database (no persistence)"); @@ -1056,6 +1159,7 @@ pub async fn create_relay( Arc::new(db), crate::nostr::tombstones::Tombstones::in_memory(), crate::nostr::holding::HoldingStore::in_memory(), + ReplaceableHistoryStore::in_memory(), ) } DatabaseBackend::Lmdb => { @@ -1086,7 +1190,8 @@ pub async fn create_relay( Path::new(&config.effective_git_data_path()), ) .await?; - (Arc::new(db), tombstones, holding) + let history = ReplaceableHistoryStore::open_lmdb(db_path).await?; + (Arc::new(db), tombstones, holding, history) } }; @@ -1119,6 +1224,7 @@ pub async fn create_relay( database.clone(), tombstones, holding.clone(), + history.clone(), &git_data_path, purgatory, config.clone(), @@ -1151,6 +1257,7 @@ pub async fn create_relay( database, write_policy, holding, + history, }) } diff --git a/src/nostr/history.rs b/src/nostr/history.rs new file mode 100644 index 0000000..8f09544 --- /dev/null +++ b/src/nostr/history.rs @@ -0,0 +1,301 @@ +//! Durable history storage for superseded replaceable/addressable events. +//! +//! This store keeps prior versions of selected coordinates (currently kinds +//! 30617 + 30618) so rollback can be implemented later without depending on +//! backend-internal replaceable retention behavior. + +use std::path::Path; +use std::sync::Arc; + +use nostr_lmdb::NostrLmdb; +use nostr_memory::MemoryDatabase; +use nostr_relay_builder::prelude::{ + Alphabet, Event, EventBuilder, EventId, Filter, FinalizeEvent, Keys, Kind, NostrDatabase, + SingleLetterTag, Tag, Timestamp, +}; + +/// Directory name (under `relay_data_path`) for LMDB replaceable history. +pub const HISTORY_DIR: &str = "replaceable-history"; + +/// Internal metadata event kind for history captures. +pub const HISTORY_METADATA_KIND: u16 = 9906; +pub const HISTORY_SUPERSEDED_AT_TAG: &str = "history-superseded-at"; +pub const HISTORY_CAPTURED_AT_TAG: &str = "history-captured-at"; +pub const HISTORY_REPLACED_BY_TAG: &str = "history-replaced-by"; +pub const HISTORY_KIND_TAG: &str = "history-kind"; +pub const HISTORY_AUTHOR_PUBKEY_TAG: &str = "history-author-pubkey"; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SupersededRecord { + pub metadata_event_id: EventId, + pub superseded_event_id: Option, + pub coordinate: Option, + pub superseded_at: Option, + pub captured_at: Option, + pub replaced_by: Option, +} + +#[derive(Clone)] +pub struct ReplaceableHistoryStore { + db: Arc, + metadata_signer: Keys, +} + +impl std::fmt::Debug for ReplaceableHistoryStore { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ReplaceableHistoryStore") + .finish_non_exhaustive() + } +} + +impl ReplaceableHistoryStore { + /// Open a persistent LMDB-backed history store. + pub async fn open_lmdb(relay_data_path: &Path) -> anyhow::Result { + let path = relay_data_path.join(HISTORY_DIR); + std::fs::create_dir_all(&path).map_err(|e| { + anyhow::anyhow!( + "Failed to create replaceable-history directory {}: {}", + path.display(), + e + ) + })?; + + let db = NostrLmdb::builder(&path) + .process_nip09(false) + .process_nip62(false) + .build() + .await + .map_err(|e| { + anyhow::anyhow!( + "Failed to open replaceable-history LMDB at {}: {}", + path.display(), + e + ) + })?; + + Ok(Self { + db: Arc::new(db), + metadata_signer: Keys::generate(), + }) + } + + /// In-memory history store (non-persistent). + pub fn in_memory() -> Self { + Self { + db: Arc::new(MemoryDatabase::unbounded()), + metadata_signer: Keys::generate(), + } + } + + /// Persist one superseded event payload + metadata row for later rollback. + pub async fn archive_superseded_event( + &self, + superseded: &Event, + replaced_by: &Event, + coordinate: &str, + ) -> anyhow::Result<()> { + self.db.save_event(superseded).await.map_err(|e| { + anyhow::anyhow!( + "Failed to store superseded event {} in history DB: {e}", + superseded.id + ) + })?; + + let captured_at = Timestamp::now(); + let mut tags = vec![ + Tag::event(superseded.id), + Tag::custom("a", vec![coordinate.to_string()]), + Tag::custom( + HISTORY_SUPERSEDED_AT_TAG, + vec![superseded.created_at.as_secs().to_string()], + ), + Tag::custom( + HISTORY_CAPTURED_AT_TAG, + vec![captured_at.as_secs().to_string()], + ), + Tag::custom(HISTORY_REPLACED_BY_TAG, vec![replaced_by.id.to_hex()]), + Tag::custom(HISTORY_KIND_TAG, vec![superseded.kind.as_u16().to_string()]), + Tag::custom(HISTORY_AUTHOR_PUBKEY_TAG, vec![superseded.pubkey.to_hex()]), + ]; + + if let Some(identifier) = extract_identifier(superseded) { + tags.push(Tag::custom("d", vec![identifier])); + } + + let metadata_event = EventBuilder::new(Kind::from(HISTORY_METADATA_KIND), "") + .tags(tags) + .custom_created_at(captured_at) + .finalize(&self.metadata_signer) + .map_err(|e| anyhow::anyhow!("Failed to build replaceable-history metadata: {e}"))?; + + self.db.save_event(&metadata_event).await.map_err(|e| { + anyhow::anyhow!( + "Failed to store replaceable-history metadata for {}: {e}", + superseded.id + ) + })?; + + Ok(()) + } + + pub async fn superseded_records_for_coordinate_before( + &self, + coordinate: &str, + before_or_at: Timestamp, + ) -> Vec { + let filter = Filter::new() + .kind(Kind::from(HISTORY_METADATA_KIND)) + .custom_tag( + SingleLetterTag::lowercase(Alphabet::A), + coordinate.to_string(), + ); + + let events = match self.db.query(filter).await { + Ok(events) => events, + Err(e) => { + tracing::error!( + error = %e, + coordinate, + "History DB query failed" + ); + return Vec::new(); + } + }; + + events + .into_iter() + .filter_map(|metadata| { + let superseded_at = parse_timestamp_tag(&metadata, HISTORY_SUPERSEDED_AT_TAG) + .or_else(|| Some(metadata.created_at)); + if superseded_at.is_some_and(|ts| ts > before_or_at) { + return None; + } + + Some(SupersededRecord { + metadata_event_id: metadata.id, + superseded_event_id: parse_event_id_tag(&metadata, "e"), + coordinate: parse_string_tag(&metadata, "a"), + superseded_at, + captured_at: parse_timestamp_tag(&metadata, HISTORY_CAPTURED_AT_TAG), + replaced_by: parse_event_id_tag(&metadata, HISTORY_REPLACED_BY_TAG), + }) + }) + .collect() + } + + pub async fn event_by_id(&self, id: &EventId) -> anyhow::Result> { + self.db + .event_by_id(id) + .await + .map_err(|e| anyhow::anyhow!("Failed history lookup for event {}: {e}", id)) + } +} + +fn parse_string_tag(metadata: &Event, tag_name: &str) -> Option { + metadata.tags.iter().find_map(|tag| { + let v = tag.as_slice(); + if v.len() >= 2 && v[0] == tag_name { + Some(v[1].clone()) + } else { + None + } + }) +} + +fn parse_timestamp_tag(metadata: &Event, tag_name: &str) -> Option { + parse_string_tag(metadata, tag_name) + .and_then(|v| v.parse::().ok()) + .map(Timestamp::from_secs) +} + +fn parse_event_id_tag(metadata: &Event, tag_name: &str) -> Option { + parse_string_tag(metadata, tag_name).and_then(|v| EventId::from_hex(&v).ok()) +} + +fn extract_identifier(event: &Event) -> Option { + event.tags.iter().find_map(|tag| { + let v = tag.as_slice(); + if v.len() >= 2 && v[0] == "d" { + Some(v[1].clone()) + } else { + None + } + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn announcement(keys: &Keys, identifier: &str, created_at: u64, content: &str) -> Event { + EventBuilder::new(Kind::GitRepoAnnouncement, content) + .tags(vec![Tag::identifier(identifier)]) + .custom_created_at(Timestamp::from_secs(created_at)) + .finalize(keys) + .unwrap() + } + + #[tokio::test] + async fn archives_and_queries_superseded_records_by_coordinate_and_time() { + let store = ReplaceableHistoryStore::in_memory(); + let keys = Keys::generate(); + + let v1 = announcement(&keys, "repo", 100, "v1"); + let v2 = announcement(&keys, "repo", 200, "v2"); + let coord = format!("30617:{}:repo", keys.public_key().to_hex()); + + store + .archive_superseded_event(&v1, &v2, &coord) + .await + .expect("archive superseded event"); + + let records = store + .superseded_records_for_coordinate_before(&coord, Timestamp::from_secs(150)) + .await; + assert_eq!(records.len(), 1); + assert_eq!(records[0].superseded_event_id, Some(v1.id)); + assert_eq!(records[0].replaced_by, Some(v2.id)); + + let payload = store + .event_by_id(&v1.id) + .await + .expect("history payload lookup") + .expect("payload exists"); + assert_eq!(payload.content, "v1"); + } + + #[tokio::test] + async fn lmdb_history_survives_restart() { + let temp = tempfile::tempdir().expect("tempdir"); + let keys = Keys::generate(); + let coord = format!("30617:{}:repo", keys.public_key().to_hex()); + let v1 = announcement(&keys, "repo", 10, "old"); + let v2 = announcement(&keys, "repo", 20, "new"); + + { + let store = ReplaceableHistoryStore::open_lmdb(temp.path()) + .await + .expect("open store"); + store + .archive_superseded_event(&v1, &v2, &coord) + .await + .expect("archive"); + } + + let reopened = ReplaceableHistoryStore::open_lmdb(temp.path()) + .await + .expect("reopen store"); + let records = reopened + .superseded_records_for_coordinate_before(&coord, Timestamp::from_secs(999)) + .await; + assert_eq!(records.len(), 1); + assert_eq!(records[0].superseded_event_id, Some(v1.id)); + + let payload = reopened + .event_by_id(&v1.id) + .await + .expect("lookup") + .expect("payload exists"); + assert_eq!(payload.id, v1.id); + } +} diff --git a/src/nostr/mod.rs b/src/nostr/mod.rs index 287d0b8..a55aa7e 100644 --- a/src/nostr/mod.rs +++ b/src/nostr/mod.rs @@ -1,5 +1,6 @@ pub mod builder; pub mod events; +pub mod history; pub mod holding; pub mod policy; pub mod tombstones; diff --git a/src/nostr/policy/mod.rs b/src/nostr/policy/mod.rs index 0c4877b..bc71b53 100644 --- a/src/nostr/policy/mod.rs +++ b/src/nostr/policy/mod.rs @@ -47,6 +47,7 @@ use super::SharedDatabase; #[cfg(test)] use crate::grasp06::receive::new_repo_init_locks; use crate::grasp06::receive::RepoInitLocks; +use crate::nostr::history::ReplaceableHistoryStore; use crate::nostr::holding::HoldingStore; use crate::nostr::tombstones::Tombstones; use crate::purgatory::Purgatory; @@ -68,6 +69,8 @@ pub struct PolicyContext { pub tombstones: Tombstones, /// Holding DB for events moved out of main DB by deletion flows. pub holding: HoldingStore, + /// Durable rollback history for superseded replaceable/addressable events. + pub history: ReplaceableHistoryStore, pub git_data_path: std::path::PathBuf, pub purgatory: Arc, /// Local relay for notifying WebSocket subscribers (set after relay creation) @@ -86,6 +89,7 @@ impl PolicyContext { database: SharedDatabase, tombstones: Tombstones, holding: HoldingStore, + history: ReplaceableHistoryStore, git_data_path: impl Into, purgatory: Arc, config: crate::config::Config, @@ -96,6 +100,7 @@ impl PolicyContext { database, tombstones, holding, + history, git_data_path: git_data_path.into(), purgatory, local_relay: Arc::new(std::sync::RwLock::new(None)), @@ -120,6 +125,7 @@ impl PolicyContext { database, Tombstones::in_memory(), HoldingStore::in_memory(), + ReplaceableHistoryStore::in_memory(), git_data_path, purgatory, config, diff --git a/tests/nip09_blacklist_ops.rs b/tests/nip09_blacklist_ops.rs index fc74082..8a70130 100644 --- a/tests/nip09_blacklist_ops.rs +++ b/tests/nip09_blacklist_ops.rs @@ -9,6 +9,7 @@ use ngit_grasp::config::Config; use ngit_grasp::grasp06::receive::new_repo_init_locks; use ngit_grasp::metrics::{self, Metrics}; use ngit_grasp::nostr::builder::{Nip34WritePolicy, SharedDatabase}; +use ngit_grasp::nostr::history::ReplaceableHistoryStore; use ngit_grasp::nostr::holding::{ DeletionSource, GitArchiveMetadata, HoldingMetadata, HoldingStore, HOLDING_ARCHIVE_PATH_TAG, }; @@ -68,6 +69,7 @@ fn make_policy( database, Tombstones::in_memory(), holding, + ReplaceableHistoryStore::in_memory(), git_data_path.to_path_buf(), purgatory, config, diff --git a/tests/replaceable_history.rs b/tests/replaceable_history.rs new file mode 100644 index 0000000..ca1c598 --- /dev/null +++ b/tests/replaceable_history.rs @@ -0,0 +1,281 @@ +//! Integration tests for durable replaceable/addressable history capture. + +mod common; + +use common::{publish_served_repo, TestRelay}; +use grasp_audit::{AuditClient, AuditConfig, DETERMINISTIC_COMMIT_HASH}; +use ngit_grasp::nostr::history::ReplaceableHistoryStore; +use nostr_sdk::prelude::*; +use std::time::Duration; + +async fn open_history(relay_data_path: &std::path::Path) -> ReplaceableHistoryStore { + ReplaceableHistoryStore::open_lmdb(relay_data_path) + .await + .expect("open history db") +} + +async fn build_announcement_version( + client: &AuditClient, + repo_id: &str, + created_at: Timestamp, + content: &str, +) -> Event { + let relay_url = client + .relay_url() + .await + .expect("client should have relay url"); + let relay_domain = relay_url + .trim_start_matches("ws://") + .trim_start_matches("wss://") + .to_string(); + let http_url = format!("http://{}", relay_domain); + let npub = client.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!("{}/{}/{}.git", http_url, npub, repo_id)], + ), + Tag::custom("relays", vec![relay_url]), + ]) + .custom_created_at(created_at) + .finalize(client.keys()) + .expect("build announcement version") +} + +fn build_state_version( + client: &AuditClient, + repo_id: &str, + created_at: Timestamp, + content: &str, +) -> Event { + EventBuilder::new(Kind::RepoState, content) + .tags(vec![ + Tag::identifier(repo_id), + Tag::custom( + "refs/heads/main", + vec![DETERMINISTIC_COMMIT_HASH.to_string()], + ), + Tag::custom("HEAD", vec!["ref: refs/heads/main".to_string()]), + ]) + .custom_created_at(created_at) + .finalize(client.keys()) + .expect("build state version") +} + +#[tokio::test] +async fn newer_30618_supersedes_older_and_preserves_old_in_history() { + let relay = TestRelay::start_with_lmdb().await; + let client = AuditClient::new(relay.url(), AuditConfig::isolated()) + .await + .expect("create audit client"); + + let (_announcement, repo_id) = publish_served_repo(&client, "history-30618").await; + let base_ts = Timestamp::now(); + let old_state = build_state_version(&client, &repo_id, base_ts, "state-v1"); + let new_state = build_state_version( + &client, + &repo_id, + Timestamp::from_secs(base_ts.as_secs() + 30), + "state-v2", + ); + + client + .send_event(old_state.clone()) + .await + .expect("send old state"); + client + .send_event(new_state.clone()) + .await + .expect("send new state"); + tokio::time::sleep(Duration::from_millis(400)).await; + + let coordinate = format!("30618:{}:{}", client.public_key().to_hex(), repo_id); + let history = open_history(relay.relay_data_path()).await; + let records = history + .superseded_records_for_coordinate_before(&coordinate, Timestamp::now()) + .await; + + assert!( + records + .iter() + .any(|r| r.superseded_event_id == Some(old_state.id) + && r.replaced_by == Some(new_state.id)), + "history must contain old 30618 as superseded by the newer 30618" + ); + assert!( + history + .event_by_id(&old_state.id) + .await + .expect("lookup old state in history") + .is_some(), + "history payload must persist the superseded 30618 event" + ); + + relay.stop().await; +} + +#[tokio::test] +async fn newer_30617_supersedes_older_and_preserves_old_in_history() { + let relay = TestRelay::start_with_lmdb().await; + let client = AuditClient::new(relay.url(), AuditConfig::isolated()) + .await + .expect("create audit client"); + + let (first_announcement, repo_id) = publish_served_repo(&client, "history-30617").await; + let newer_announcement = build_announcement_version( + &client, + &repo_id, + Timestamp::from_secs(first_announcement.created_at.as_secs() + 30), + "re-announcement", + ) + .await; + + client + .send_event(newer_announcement.clone()) + .await + .expect("send newer announcement"); + tokio::time::sleep(Duration::from_millis(400)).await; + + let coordinate = format!("30617:{}:{}", client.public_key().to_hex(), repo_id); + let history = open_history(relay.relay_data_path()).await; + let records = history + .superseded_records_for_coordinate_before(&coordinate, Timestamp::now()) + .await; + + assert!( + records.iter().any(|r| { + r.superseded_event_id == Some(first_announcement.id) + && r.replaced_by == Some(newer_announcement.id) + }), + "history must contain old 30617 as superseded by the newer 30617" + ); + assert!( + history + .event_by_id(&first_announcement.id) + .await + .expect("lookup old announcement in history") + .is_some(), + "history payload must persist the superseded 30617 event" + ); + + relay.stop().await; +} + +#[tokio::test] +async fn replaceable_history_survives_restart_lmdb() { + let relay = TestRelay::start_with_lmdb().await; + let client = AuditClient::new(relay.url(), AuditConfig::isolated()) + .await + .expect("create audit client"); + + let (first_announcement, repo_id) = publish_served_repo(&client, "history-restart").await; + let newer_announcement = build_announcement_version( + &client, + &repo_id, + Timestamp::from_secs(first_announcement.created_at.as_secs() + 45), + "restart-reannouncement", + ) + .await; + client + .send_event(newer_announcement.clone()) + .await + .expect("send newer announcement"); + tokio::time::sleep(Duration::from_millis(400)).await; + + let coordinate = format!("30617:{}:{}", client.public_key().to_hex(), repo_id); + let history_before_restart = open_history(relay.relay_data_path()).await; + let pre_restart_records = history_before_restart + .superseded_records_for_coordinate_before(&coordinate, Timestamp::now()) + .await; + assert!( + pre_restart_records.iter().any(|r| { + r.superseded_event_id == Some(first_announcement.id) + && r.replaced_by == Some(newer_announcement.id) + }), + "history metadata must exist before restart" + ); + + let relay_data_path = relay.relay_data_path().clone(); + let git_data_path = relay.git_data_path().clone(); + let owner_hex = client.public_key().to_hex(); + relay.stop().await; + + let restarted = TestRelay::start_with_existing_lmdb_paths( + git_data_path, + relay_data_path.clone(), + None, + false, + ) + .await; + restarted.stop().await; + + let history = open_history(&relay_data_path).await; + let coordinate = format!("30617:{}:{}", owner_hex, repo_id); + let records = history + .superseded_records_for_coordinate_before(&coordinate, Timestamp::now()) + .await; + + assert!( + records.iter().any(|r| { + r.superseded_event_id == Some(first_announcement.id) + && r.replaced_by == Some(newer_announcement.id) + }), + "history metadata must survive LMDB restart" + ); + assert!( + history + .event_by_id(&first_announcement.id) + .await + .expect("lookup archived announcement after restart") + .is_some(), + "history payload must survive LMDB restart" + ); +} + +#[tokio::test] +async fn serving_behavior_unchanged_latest_version_is_still_served() { + let relay = TestRelay::start_with_lmdb().await; + let client = AuditClient::new(relay.url(), AuditConfig::isolated()) + .await + .expect("create audit client"); + + let (_announcement, repo_id) = publish_served_repo(&client, "history-serving").await; + let base_ts = Timestamp::now(); + let old_state = build_state_version(&client, &repo_id, base_ts, "state-old"); + let new_state = build_state_version( + &client, + &repo_id, + Timestamp::from_secs(base_ts.as_secs() + 30), + "state-new", + ); + + client + .send_event(old_state.clone()) + .await + .expect("send old state"); + client + .send_event(new_state.clone()) + .await + .expect("send new state"); + tokio::time::sleep(Duration::from_millis(400)).await; + + let by_id_old = client + .is_event_on_relay(old_state.id) + .await + .expect("query old state by id"); + let by_id_new = client + .is_event_on_relay(new_state.id) + .await + .expect("query new state by id"); + assert!( + !by_id_old, + "superseded state must remain non-served in the main relay behavior" + ); + assert!(by_id_new, "latest state must still be served"); + + relay.stop().await; +}