feat(deletion): rollback deleted state to previous version

This commit is contained in:
DanConwayDev
2026-06-17 15:48:49 +00:00
parent d5c7dbfbf6
commit 545c52d0f8
5 changed files with 785 additions and 9 deletions
+1 -7
View File
@@ -378,9 +378,6 @@ impl Nip34WritePolicy {
"Accepted announcement to purgatory: {} (waiting for git data)",
event_id_str
);
self.capture_superseded_replaceable_history(event).await;
accepted_purgatory("won't be served until git data arrives")
}
Err(e) => {
@@ -880,10 +877,7 @@ impl Nip34WritePolicy {
}
fn event_persists_to_main_db(result: &WritePolicyResult) -> bool {
matches!(
result,
WritePolicyResult::Accept | WritePolicyResult::Reject { status: true, .. }
)
matches!(result, WritePolicyResult::Accept)
}
async fn capture_superseded_replaceable_history(&self, incoming: &Event) {
+84
View File
@@ -183,6 +183,36 @@ impl ReplaceableHistoryStore {
.collect()
}
/// Return the latest superseded record at or before `cutoff` for a coordinate.
///
/// Ordering is deterministic:
/// 1. `superseded_at` descending (missing treated as 0)
/// 2. `captured_at` descending (missing treated as 0)
/// 3. `metadata_event_id` descending
pub async fn latest_superseded_before(
&self,
coordinate: &str,
cutoff: Timestamp,
) -> Option<SupersededRecord> {
let mut records = self
.superseded_records_for_coordinate_before(coordinate, cutoff)
.await;
records.sort_by(|a, b| {
b.superseded_at
.unwrap_or_else(|| Timestamp::from_secs(0))
.cmp(&a.superseded_at.unwrap_or_else(|| Timestamp::from_secs(0)))
.then_with(|| {
b.captured_at
.unwrap_or_else(|| Timestamp::from_secs(0))
.cmp(&a.captured_at.unwrap_or_else(|| Timestamp::from_secs(0)))
})
.then_with(|| b.metadata_event_id.cmp(&a.metadata_event_id))
});
records.into_iter().next()
}
pub async fn event_by_id(&self, id: &EventId) -> anyhow::Result<Option<Event>> {
self.db
.event_by_id(id)
@@ -298,4 +328,58 @@ mod tests {
.expect("payload exists");
assert_eq!(payload.id, v1.id);
}
#[tokio::test]
async fn latest_superseded_before_prefers_newest_superseded_timestamp() {
let store = ReplaceableHistoryStore::in_memory();
let keys = Keys::generate();
let coord = format!("30617:{}:repo", keys.public_key().to_hex());
let v1 = announcement(&keys, "repo", 100, "v1");
let v2 = announcement(&keys, "repo", 120, "v2");
let v3 = announcement(&keys, "repo", 140, "v3");
store
.archive_superseded_event(&v1, &v2, &coord)
.await
.expect("archive v1");
store
.archive_superseded_event(&v2, &v3, &coord)
.await
.expect("archive v2");
let latest = store
.latest_superseded_before(&coord, Timestamp::from_secs(200))
.await
.expect("latest record");
assert_eq!(latest.superseded_event_id, Some(v2.id));
}
#[tokio::test]
async fn latest_superseded_before_respects_cutoff() {
let store = ReplaceableHistoryStore::in_memory();
let keys = Keys::generate();
let coord = format!("30617:{}:repo", keys.public_key().to_hex());
let v1 = announcement(&keys, "repo", 100, "v1");
let v2 = announcement(&keys, "repo", 120, "v2");
let v3 = announcement(&keys, "repo", 140, "v3");
store
.archive_superseded_event(&v1, &v2, &coord)
.await
.expect("archive v1");
store
.archive_superseded_event(&v2, &v3, &coord)
.await
.expect("archive v2");
let latest = store
.latest_superseded_before(&coord, Timestamp::from_secs(110))
.await
.expect("latest record before cutoff");
assert_eq!(latest.superseded_event_id, Some(v1.id));
}
}
+378
View File
@@ -43,6 +43,12 @@ use nostr_relay_builder::prelude::{
use tar::Builder as TarBuilder;
use super::{PolicyContext, RetentionReason};
use crate::git;
use crate::git::authorization::{
collect_authorized_maintainers, fetch_repository_data_excluding_purgatory,
};
use crate::git::process;
use crate::nostr::events::RepositoryState;
use crate::nostr::holding::{DeletionSource, GitArchiveMetadata, HoldingMetadata};
use crate::nostr::policy::reject_invalid;
@@ -76,6 +82,14 @@ pub struct DeletionPolicy {
ctx: PolicyContext,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct StateRollbackPlan {
coordinate: String,
identifier: String,
deleted_event_id: EventId,
cutoff: Timestamp,
}
impl DeletionPolicy {
pub fn new(ctx: PolicyContext) -> Self {
Self { ctx }
@@ -140,6 +154,8 @@ impl DeletionPolicy {
}
}
let rollback_plans = self.collect_state_rollback_plans(event).await;
// Process purgatory removals (synchronous, in-memory).
self.remove_purgatory_targets(event);
@@ -147,6 +163,16 @@ impl DeletionPolicy {
let mut moved_ids = HashSet::new();
self.delete_main_db_targets(event, &mut moved_ids).await;
let mut identifiers_to_realign = HashSet::new();
for plan in rollback_plans {
identifiers_to_realign.insert(plan.identifier.clone());
self.rollback_deleted_active_state(&plan).await;
}
for identifier in identifiers_to_realign {
self.realign_identifier_state(&identifier).await;
}
// Record the deletion so future re-submission is rejected.
if let Err(e) = self.ctx.tombstones.record_deletion(event).await {
tracing::error!(error = %e, "Failed to record deletion tombstone");
@@ -173,6 +199,358 @@ impl DeletionPolicy {
.collect()
}
async fn collect_state_rollback_plans(&self, deletion: &Event) -> Vec<StateRollbackPlan> {
let mut plans = Vec::new();
for target_id in Self::e_tag_ids(deletion) {
let Ok(Some(target)) = self.ctx.database.event_by_id(&target_id).await else {
continue;
};
if target.kind != Kind::RepoState || target.pubkey != deletion.pubkey {
continue;
}
let Some(identifier) = identifier_from_event(&target) else {
continue;
};
let Some(active) = self
.latest_state_for_coordinate(&target.pubkey, &identifier, None)
.await
else {
continue;
};
if active.id != target.id {
continue;
}
plans.push(StateRollbackPlan {
coordinate: format!("30618:{}:{}", target.pubkey.to_hex(), identifier),
identifier,
deleted_event_id: target.id,
cutoff: target.created_at,
});
}
for coordinate in Self::a_tag_coordinates(deletion) {
let Some((kind, coord_pubkey, identifier)) = Self::parse_coordinate(&coordinate) else {
continue;
};
if kind != Kind::RepoState || coord_pubkey != deletion.pubkey {
continue;
}
let Some(active_any_cutoff) = self
.latest_state_for_coordinate(&coord_pubkey, &identifier, None)
.await
else {
continue;
};
if active_any_cutoff.created_at > deletion.created_at {
continue;
}
plans.push(StateRollbackPlan {
coordinate,
identifier,
deleted_event_id: active_any_cutoff.id,
cutoff: deletion.created_at,
});
}
plans
.into_iter()
.collect::<HashSet<_>>()
.into_iter()
.collect()
}
fn a_tag_coordinates(event: &Event) -> Vec<String> {
event
.tags
.iter()
.filter_map(|tag| {
let v = tag.as_slice();
if v.len() >= 2 && v[0] == "a" {
Some(v[1].clone())
} else {
None
}
})
.collect()
}
fn parse_coordinate(coordinate: &str) -> Option<(Kind, PublicKey, String)> {
let parts: Vec<&str> = coordinate.splitn(3, ':').collect();
if parts.len() != 3 {
return None;
}
let kind_num = parts[0].parse::<u16>().ok()?;
let pubkey = PublicKey::from_hex(parts[1]).ok()?;
Some((Kind::from(kind_num), pubkey, parts[2].to_string()))
}
async fn latest_state_for_coordinate(
&self,
author: &PublicKey,
identifier: &str,
until: Option<Timestamp>,
) -> Option<Event> {
let mut filter = Filter::new()
.kind(Kind::RepoState)
.author(*author)
.custom_tag(
SingleLetterTag::lowercase(Alphabet::D),
identifier.to_string(),
);
if let Some(cutoff) = until {
filter = filter.until(cutoff);
}
let Ok(events) = self.ctx.database.query(filter).await else {
return None;
};
events.into_iter().max_by(|a, b| {
a.created_at
.cmp(&b.created_at)
.then_with(|| a.id.cmp(&b.id))
})
}
async fn rollback_deleted_active_state(&self, plan: &StateRollbackPlan) {
let Some(record) = self
.ctx
.history
.latest_superseded_before(&plan.coordinate, plan.cutoff)
.await
else {
tracing::info!(
coordinate = %plan.coordinate,
deleted_event_id = %plan.deleted_event_id.to_hex(),
cutoff = plan.cutoff.as_secs(),
"No rollback candidate in history; entering no-active-state path"
);
return;
};
let Some(candidate_id) = record.superseded_event_id else {
tracing::warn!(
coordinate = %plan.coordinate,
metadata_event_id = %record.metadata_event_id.to_hex(),
"History record missing superseded event id; cannot rollback"
);
return;
};
let candidate = match self.ctx.history.event_by_id(&candidate_id).await {
Ok(Some(event)) => event,
Ok(None) => {
tracing::warn!(
coordinate = %plan.coordinate,
candidate_id = %candidate_id.to_hex(),
"History payload missing for rollback candidate"
);
return;
}
Err(e) => {
tracing::warn!(
coordinate = %plan.coordinate,
candidate_id = %candidate_id.to_hex(),
error = %e,
"History lookup failed for rollback candidate"
);
return;
}
};
let Ok(candidate_state) = RepositoryState::from_event(candidate.clone()) else {
tracing::warn!(
coordinate = %plan.coordinate,
candidate_id = %candidate.id.to_hex(),
"Rollback candidate is not a valid state event"
);
return;
};
let repo_data = match fetch_repository_data_excluding_purgatory(
&self.ctx.database,
&candidate_state.identifier,
)
.await
{
Ok(data) => data,
Err(e) => {
tracing::warn!(
coordinate = %plan.coordinate,
error = %e,
"Failed to load repository data before rollback"
);
return;
}
};
if crate::git::authorization::pubkey_authorised_for_repo_owners(
&candidate.pubkey,
&repo_data,
)
.is_empty()
{
tracing::warn!(
coordinate = %plan.coordinate,
candidate_id = %candidate.id.to_hex(),
author = %candidate.pubkey.to_hex(),
"Rollback candidate author is no longer authorized; skipping restore"
);
return;
}
if let Err(e) = self.ctx.database.save_event(&candidate).await {
tracing::warn!(
coordinate = %plan.coordinate,
candidate_id = %candidate.id.to_hex(),
error = %e,
"Failed to restore rollback candidate into main DB"
);
return;
}
tracing::info!(
coordinate = %plan.coordinate,
deleted_event_id = %plan.deleted_event_id.to_hex(),
restored_event_id = %candidate.id.to_hex(),
"Restored deleted active state from replaceable history"
);
}
async fn realign_identifier_state(&self, identifier: &str) {
let repo_data =
match fetch_repository_data_excluding_purgatory(&self.ctx.database, identifier).await {
Ok(data) => data,
Err(e) => {
tracing::warn!(
identifier = %identifier,
error = %e,
"Failed to reload repository data for post-deletion state alignment"
);
return;
}
};
if repo_data.announcements.is_empty() {
return;
}
let by_owner = collect_authorized_maintainers(&repo_data.announcements);
for announcement in &repo_data.announcements {
let owner_hex = announcement.event.pubkey.to_hex();
let Some(maintainers) = by_owner.get(&owner_hex) else {
continue;
};
let target_repo_path = self.ctx.git_data_path.join(announcement.repo_path());
if !target_repo_path.is_dir() {
continue;
}
let latest_state = repo_data
.states
.iter()
.filter(|state| maintainers.contains(&state.event.pubkey.to_hex()))
.max_by(|a, b| {
a.event
.created_at
.cmp(&b.event.created_at)
.then_with(|| a.event.id.cmp(&b.event.id))
});
match latest_state {
Some(state) => {
let source_repo_path =
repo_data
.announcements
.iter()
.find_map(|candidate_announcement| {
let candidate_path = self
.ctx
.git_data_path
.join(candidate_announcement.repo_path());
if candidate_path.is_dir()
&& state.branches.iter().all(|branch| {
branch.commit.starts_with("ref: ")
|| git::oid_exists(&candidate_path, &branch.commit)
})
&& state
.tags
.iter()
.all(|tag| git::oid_exists(&candidate_path, &tag.commit))
{
Some(candidate_path)
} else {
None
}
});
if let Some(source_repo_path) = source_repo_path {
let result = process::process_state_with_git_data(
state,
&source_repo_path,
&repo_data,
&self.ctx.git_data_path,
);
tracing::info!(
identifier = %identifier,
owner = %owner_hex,
state_event_id = %state.event.id.to_hex(),
repos_synced = result.repos_synced,
refs_created = result.refs_created,
refs_updated = result.refs_updated,
refs_deleted = result.refs_deleted,
"Realigned git refs after state deletion/rollback"
);
}
}
None => {
let refs = match git::list_refs(&target_repo_path) {
Ok(refs) => refs,
Err(e) => {
tracing::warn!(
identifier = %identifier,
owner = %owner_hex,
repo = %target_repo_path.display(),
error = %e,
"Failed to list refs for no-active-state cleanup"
);
continue;
}
};
let mut deleted_refs = 0usize;
for (ref_name, _) in refs {
if (ref_name.starts_with("refs/heads/")
|| ref_name.starts_with("refs/tags/"))
&& git::delete_ref(&target_repo_path, &ref_name).is_ok()
{
deleted_refs += 1;
}
}
tracing::info!(
identifier = %identifier,
owner = %owner_hex,
repo = %target_repo_path.display(),
deleted_refs,
"No active authorized state remains; cleared managed refs"
);
}
}
}
}
/// Hard-delete the targeted events from the main database.
///
/// Handles both `e`-tag (by event id) and `a`-tag (by coordinate) targets.
+269 -2
View File
@@ -7,12 +7,70 @@ mod common;
use common::{
announcement_coordinate, announcement_served_by_coordinate, build_deletion,
publish_served_repo_with_state_event, TestRelay,
publish_served_announcement_with_state_for_identifier, publish_served_repo,
publish_served_repo_with_state_event, publish_served_repo_with_state_event_and_maintainers,
TestRelay,
};
use grasp_audit::{AuditClient, AuditConfig};
use grasp_audit::{AuditClient, AuditConfig, DETERMINISTIC_COMMIT_HASH};
use nostr_sdk::prelude::*;
use std::collections::HashSet;
use std::process::Command;
use std::time::Duration;
fn build_state_with_branches(
client: &AuditClient,
repo_id: &str,
created_at: Timestamp,
branches: &[&str],
) -> Event {
let mut tags = vec![Tag::identifier(repo_id)];
for branch in branches {
tags.push(Tag::custom(
format!("refs/heads/{branch}"),
vec![DETERMINISTIC_COMMIT_HASH.to_string()],
));
}
tags.push(Tag::custom(
"HEAD",
vec!["ref: refs/heads/main".to_string()],
));
EventBuilder::new(Kind::RepoState, "")
.tags(tags)
.custom_created_at(created_at)
.finalize(client.keys())
.expect("build state event")
}
fn repo_path(relay: &TestRelay, client: &AuditClient, repo_id: &str) -> std::path::PathBuf {
relay
.git_data_path()
.join(client.public_key().to_bech32().expect("npub"))
.join(format!("{repo_id}.git"))
}
fn list_repo_refs(repo_path: &std::path::Path) -> HashSet<String> {
let output = Command::new("git")
.args([
"--git-dir",
repo_path.to_str().expect("repo path"),
"show-ref",
])
.output()
.expect("run git show-ref");
if !output.status.success() {
return HashSet::new();
}
String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|line| line.split_whitespace().nth(1))
.map(|s| s.to_string())
.collect()
}
/// Deleting an announcement by coordinate must hard-delete the repository state
/// event (kind 30618) for the same repository from the main DB.
#[tokio::test]
@@ -153,3 +211,212 @@ async fn test_state_events_are_independent_per_repository() {
"unrelated repository coordinate must remain served"
);
}
#[tokio::test]
async fn test_e_delete_active_state_rolls_back_and_realigns_refs() {
let relay = TestRelay::start().await;
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
.await
.expect("create audit client");
let (_announcement, repo_id) = publish_served_repo(&client, "state-rollback-e").await;
let base = Timestamp::from_secs(Timestamp::now().as_secs() + 20);
let v1 = build_state_with_branches(&client, &repo_id, base, &["main"]);
let v2 = build_state_with_branches(
&client,
&repo_id,
Timestamp::from_secs(base.as_secs() + 1),
&["main", "feature"],
);
client.send_event(v1.clone()).await.expect("send v1 state");
client.send_event(v2.clone()).await.expect("send v2 state");
tokio::time::sleep(Duration::from_millis(500)).await;
let owner_repo = repo_path(&relay, &client, &repo_id);
let refs_before = list_repo_refs(&owner_repo);
assert!(
refs_before.contains("refs/heads/feature"),
"feature ref must exist before rollback deletion"
);
client
.send_event(build_deletion(&client, &[v2.id], &[]))
.await
.expect("send e-tag deletion for active state");
tokio::time::sleep(Duration::from_millis(700)).await;
assert!(
client
.is_event_on_relay(v1.id)
.await
.expect("query v1 after rollback"),
"previous state must be restored"
);
assert!(
!client
.is_event_on_relay(v2.id)
.await
.expect("query v2 after rollback"),
"deleted active state must stay deleted"
);
let refs_after = list_repo_refs(&owner_repo);
assert!(refs_after.contains("refs/heads/main"));
assert!(
!refs_after.contains("refs/heads/feature"),
"feature ref must be removed after rollback realignment"
);
relay.stop().await;
}
#[tokio::test]
async fn test_a_delete_coordinate_cutoff_rolls_back_to_latest_before_cutoff() {
let relay = TestRelay::start().await;
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
.await
.expect("create audit client");
let (_announcement, repo_id) = publish_served_repo(&client, "state-rollback-a").await;
let base = Timestamp::from_secs(Timestamp::now().as_secs() + 40);
let v1 = build_state_with_branches(&client, &repo_id, base, &["main"]);
let v2 = build_state_with_branches(
&client,
&repo_id,
Timestamp::from_secs(base.as_secs() + 1),
&["main", "feature"],
);
client.send_event(v1.clone()).await.expect("send v1 state");
client.send_event(v2.clone()).await.expect("send v2 state");
tokio::time::sleep(Duration::from_millis(500)).await;
let coordinate = format!("30618:{}:{}", client.public_key().to_hex(), repo_id);
let deletion = EventBuilder::new(Kind::EventDeletion, "")
.tags(vec![Tag::custom("a", vec![coordinate])])
.custom_created_at(v2.created_at)
.finalize(client.keys())
.expect("build coordinate deletion");
client
.send_event(deletion)
.await
.expect("send a-tag deletion with cutoff");
tokio::time::sleep(Duration::from_millis(700)).await;
assert!(
client
.is_event_on_relay(v1.id)
.await
.expect("query v1 after coordinate rollback"),
"cutoff deletion must rollback to latest superseded <= cutoff"
);
assert!(
!client
.is_event_on_relay(v2.id)
.await
.expect("query v2 after coordinate rollback"),
"state at cutoff must be deleted"
);
relay.stop().await;
}
#[tokio::test]
async fn test_delete_only_active_state_with_no_history_clears_managed_refs() {
let relay = TestRelay::start().await;
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
.await
.expect("create audit client");
let (_announcement, repo_id, original_state) =
publish_served_repo_with_state_event(&client, "state-rollback-none").await;
let owner_repo = repo_path(&relay, &client, &repo_id);
assert!(
list_repo_refs(&owner_repo).contains("refs/heads/main"),
"main ref must exist before deleting the only active state"
);
client
.send_event(build_deletion(&client, &[original_state.id], &[]))
.await
.expect("send deletion for only active state");
tokio::time::sleep(Duration::from_millis(700)).await;
let refs_after = list_repo_refs(&owner_repo);
assert!(
!refs_after.iter().any(|r| r.starts_with("refs/heads/")),
"no branch refs should remain when there is no active state"
);
assert!(
!refs_after.iter().any(|r| r.starts_with("refs/tags/")),
"no tag refs should remain when there is no active state"
);
relay.stop().await;
}
#[tokio::test]
async fn test_multi_maintainer_state_deletion_rolls_back_without_affecting_other_maintainer() {
let relay = TestRelay::start().await;
let client_a = AuditClient::new(relay.url(), AuditConfig::isolated())
.await
.expect("create maintainer A client");
let client_b =
AuditClient::new_with_keys(relay.url(), AuditConfig::isolated(), Keys::generate())
.await
.expect("create maintainer B client");
let (_announcement_a, repo_id, _state_a) =
publish_served_repo_with_state_event_and_maintainers(
&client_a,
"state-rollback-multi",
&[client_b.public_key().to_hex()],
)
.await;
let (_announcement_b, state_b) =
publish_served_announcement_with_state_for_identifier(&client_b, &repo_id).await;
let base = Timestamp::from_secs(Timestamp::now().as_secs() + 60);
let a_v1 = build_state_with_branches(&client_a, &repo_id, base, &["main"]);
let a_v2 = build_state_with_branches(
&client_a,
&repo_id,
Timestamp::from_secs(base.as_secs() + 1),
&["main", "feature"],
);
client_a
.send_event(a_v1.clone())
.await
.expect("send maintainer A v1");
client_a
.send_event(a_v2.clone())
.await
.expect("send maintainer A v2");
tokio::time::sleep(Duration::from_millis(700)).await;
client_a
.send_event(build_deletion(&client_a, &[a_v2.id], &[]))
.await
.expect("delete maintainer A active state");
tokio::time::sleep(Duration::from_millis(700)).await;
assert!(
client_a
.is_event_on_relay(a_v1.id)
.await
.expect("query restored maintainer A state"),
"maintainer A should rollback to previous state"
);
assert!(
client_b
.is_event_on_relay(state_b.id)
.await
.expect("query maintainer B state"),
"maintainer B state must remain served"
);
relay.stop().await;
}
+53
View File
@@ -279,3 +279,56 @@ async fn serving_behavior_unchanged_latest_version_is_still_served() {
relay.stop().await;
}
#[tokio::test]
async fn history_not_captured_for_purgatory_only_announcement() {
let relay = TestRelay::start_with_lmdb().await;
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
.await
.expect("create audit client");
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 npub = client.public_key().to_bech32().expect("pubkey to npub");
let repo_id = format!(
"history-purgatory-only-{}",
&Keys::generate().public_key().to_hex()[..8]
);
let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "")
.tags(vec![
Tag::identifier(&repo_id),
Tag::custom("name", vec![repo_id.clone()]),
Tag::custom(
"clone",
vec![format!("http://{}/{}/{}.git", relay_domain, npub, repo_id)],
),
Tag::custom("relays", vec![relay_url]),
])
.finalize(client.keys())
.expect("build purgatory-only announcement");
client
.send_event(announcement)
.await
.expect("send announcement to purgatory");
tokio::time::sleep(Duration::from_millis(300)).await;
let coordinate = format!("30617:{}:{}", client.public_key().to_hex(), repo_id);
let history = open_history(relay.relay_data_path()).await;
let latest = history
.latest_superseded_before(&coordinate, Timestamp::now())
.await;
assert!(
latest.is_none(),
"purgatory-only announcements must not capture superseded history"
);
relay.stop().await;
}