mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat(deletion): persist replaceable history for rollback
This commit is contained in:
+110
-3
@@ -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<std::path::PathBuf>,
|
||||
purgatory: std::sync::Arc<crate::purgatory::Purgatory>,
|
||||
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<String> {
|
||||
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,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -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<EventId>,
|
||||
pub coordinate: Option<String>,
|
||||
pub superseded_at: Option<Timestamp>,
|
||||
pub captured_at: Option<Timestamp>,
|
||||
pub replaced_by: Option<EventId>,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct ReplaceableHistoryStore {
|
||||
db: Arc<dyn NostrDatabase>,
|
||||
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<Self> {
|
||||
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<SupersededRecord> {
|
||||
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<Option<Event>> {
|
||||
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<String> {
|
||||
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<Timestamp> {
|
||||
parse_string_tag(metadata, tag_name)
|
||||
.and_then(|v| v.parse::<u64>().ok())
|
||||
.map(Timestamp::from_secs)
|
||||
}
|
||||
|
||||
fn parse_event_id_tag(metadata: &Event, tag_name: &str) -> Option<EventId> {
|
||||
parse_string_tag(metadata, tag_name).and_then(|v| EventId::from_hex(&v).ok())
|
||||
}
|
||||
|
||||
fn extract_identifier(event: &Event) -> Option<String> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
pub mod builder;
|
||||
pub mod events;
|
||||
pub mod history;
|
||||
pub mod holding;
|
||||
pub mod policy;
|
||||
pub mod tombstones;
|
||||
|
||||
@@ -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<Purgatory>,
|
||||
/// 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<std::path::PathBuf>,
|
||||
purgatory: Arc<Purgatory>,
|
||||
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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
Reference in New Issue
Block a user