refactor(nostr): share cascade planner with cleanup

This commit is contained in:
DanConwayDev
2026-06-23 16:10:15 +01:00
parent d9d795628c
commit c76067db14
7 changed files with 118 additions and 2052 deletions
+58 -68
View File
@@ -3,7 +3,7 @@
//! Provides the new `maintenance cleanup` pipeline and keeps
//! `cleanup-empty-repos` as a compatibility wrapper.
use std::collections::{HashMap, HashSet};
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
@@ -19,13 +19,15 @@ use crate::git::authorization::{
};
use crate::git::process;
use crate::nostr::events::RepositoryAnnouncement;
use crate::nostr::lifecycle::{build_event_graph, traverse_and_mark_deletions, RetentionReason};
use crate::nostr::lifecycle::plan_cascade_retention_for_deleted_announcements;
use crate::purgatory::can_apply_state;
const GRAPH_DEPENDENT_KINDS: &[u16] = &[1111, 1617, 1618, 1619, 1621, 1630, 1631, 1632, 1633];
#[derive(Debug, Args)]
pub struct MaintenanceCleanupArgs {
/// Public relay domain used for GRASP-06 PR-root retention checks.
#[arg(long, env = "NGIT_DOMAIN", default_value = "localhost:7334")]
pub domain: String,
/// Path to relay LMDB data.
#[arg(long, env = "NGIT_RELAY_DATA_PATH", default_value = "./data/relay")]
pub relay_data_path: String,
@@ -54,6 +56,10 @@ pub struct MaintenanceCleanupArgs {
/// Backward-compatible args for `cleanup-empty-repos`.
#[derive(Debug, Args)]
pub struct CleanupArgs {
/// Public relay domain used for GRASP-06 PR-root retention checks.
#[arg(long, env = "NGIT_DOMAIN", default_value = "localhost:7334")]
pub domain: String,
#[arg(long, env = "NGIT_RELAY_DATA_PATH", default_value = "./data/relay")]
pub relay_data_path: String,
@@ -75,6 +81,7 @@ enum PassProfile {
#[derive(Debug, Clone)]
struct MaintenanceOptions {
domain: String,
execute: bool,
prune_orphans: bool,
prune_missing_git: bool,
@@ -151,6 +158,7 @@ struct MaintenanceRunReport {
pub async fn run(args: &CleanupArgs) -> Result<()> {
let options = MaintenanceOptions {
domain: args.domain.clone(),
execute: args.execute,
prune_orphans: args.purge_orphans,
prune_missing_git: args.execute,
@@ -175,6 +183,7 @@ pub async fn run(args: &CleanupArgs) -> Result<()> {
pub async fn run_maintenance_cleanup(args: &MaintenanceCleanupArgs) -> Result<()> {
let options = MaintenanceOptions {
domain: args.domain.clone(),
execute: args.execute,
prune_orphans: args.prune_orphans,
prune_missing_git: args.prune_missing_git,
@@ -252,6 +261,7 @@ async fn run_engine_with_db(
report.cascade = run_cascade_reconcile(
&database,
&planned_prune_ids,
&options.domain,
options.execute,
options.prune_after,
)
@@ -346,6 +356,7 @@ async fn build_scan_plan(
async fn run_cascade_reconcile(
database: &Arc<dyn NostrDatabase>,
planned_prune_ids: &HashSet<EventId>,
domain: &str,
execute: bool,
prune_after: Option<Duration>,
) -> Result<CascadeReport> {
@@ -354,46 +365,30 @@ async fn run_cascade_reconcile(
..Default::default()
};
let candidates = query_cascade_candidates(database).await?;
report.scanned_candidates = candidates.len();
let announcements = query_cascade_announcements(database).await?;
report.scanned_candidates = announcements.len();
if candidates.is_empty() {
if announcements.is_empty() {
return Ok(report);
}
let graph = build_event_graph(&candidates);
let mut kept: HashMap<EventId, RetentionReason> = HashMap::new();
let planned_announcements: Vec<Event> = announcements
.iter()
.filter(|event| {
event.kind == Kind::GitRepoAnnouncement && planned_prune_ids.contains(&event.id)
})
.cloned()
.collect();
for event in &candidates {
if event.kind == Kind::GitRepoAnnouncement && !planned_prune_ids.contains(&event.id) {
kept.insert(
event.id,
RetentionReason {
event_id: event.id,
valid_announcements: vec![announcement_address_of(event)],
event_type: "Repository Announcement".to_string(),
},
);
}
}
let plan =
plan_cascade_retention_for_deleted_announcements(database, &planned_announcements, domain)
.await
.map_err(|e| anyhow!("Cascade reconcile planning failed: {e}"))?;
let traversal = traverse_and_mark_deletions(&graph, &kept, EventId::all_zeros());
let mut deletable_ids: Vec<EventId> = Vec::new();
for event in &candidates {
if !traversal.delete.contains(&event.id) {
continue;
}
let is_announcement = event.kind == Kind::GitRepoAnnouncement;
if is_announcement && !planned_prune_ids.contains(&event.id) {
continue;
}
if !is_announcement || planned_prune_ids.contains(&event.id) {
deletable_ids.push(event.id);
}
}
let mut deletable_ids: Vec<EventId> = plan.orphan_event_ids.into_iter().collect();
deletable_ids.extend(planned_announcements.iter().map(|event| event.id));
deletable_ids.sort();
deletable_ids.dedup();
report.sample_deleted_ids = deletable_ids
.iter()
@@ -808,30 +803,13 @@ async fn count_events(database: &Arc<dyn NostrDatabase>, filter: Filter) -> Resu
.len())
}
async fn query_cascade_candidates(database: &Arc<dyn NostrDatabase>) -> Result<Vec<Event>> {
let mut events = Vec::new();
events.extend(
database
.query(Filter::new().kind(Kind::GitRepoAnnouncement))
.await
.context("Failed querying announcements for cascade")?,
);
for kind_num in GRAPH_DEPENDENT_KINDS {
events.extend(
database
.query(Filter::new().kind(Kind::from(*kind_num)))
.await
.with_context(|| format!("Failed querying kind {} for cascade", kind_num))?,
);
}
Ok(events)
}
fn announcement_address_of(event: &Event) -> String {
let identifier = identifier_from_event(event).unwrap_or_default();
format!("30617:{}:{}", event.pubkey.to_hex(), identifier)
async fn query_cascade_announcements(database: &Arc<dyn NostrDatabase>) -> Result<Vec<Event>> {
Ok(database
.query(Filter::new().kind(Kind::GitRepoAnnouncement))
.await
.context("Failed querying announcements for cascade")?
.into_iter()
.collect())
}
fn identifier_from_event(event: &Event) -> Option<String> {
@@ -1044,6 +1022,7 @@ mod tests {
fn default_options() -> MaintenanceOptions {
MaintenanceOptions {
domain: "localhost:7334".to_string(),
execute: false,
prune_orphans: false,
prune_missing_git: false,
@@ -1157,26 +1136,37 @@ mod tests {
}
#[tokio::test]
async fn cascade_reconcile_cleans_legacy_orphan_dependents() {
async fn cascade_reconcile_deletes_dependents_of_planned_prune() {
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
let tmp = tempdir().unwrap();
let git_path = tmp.path();
let keys = Keys::generate();
let dangling = issue(
let ann = announcement(&keys, "repo-cascade-prune");
let st = state(
&keys,
"30617:ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff:ghost",
"repo-cascade-prune",
"dddddddddddddddddddddddddddddddddddddddd",
);
db.save_event(&dangling).await.unwrap();
let ann_addr = format!("30617:{}:{}", ann.pubkey.to_hex(), "repo-cascade-prune");
let dep = issue(&keys, &ann_addr);
db.save_event(&ann).await.unwrap();
db.save_event(&st).await.unwrap();
db.save_event(&dep).await.unwrap();
let mut options = default_options();
options.execute = true;
options.prune_missing_git = true;
let report = run_engine_with_db(db.clone(), tmp.path(), options)
let report = run_engine_with_db(db.clone(), git_path, options)
.await
.unwrap();
assert_eq!(count_kind(&db, Kind::GitRepoAnnouncement).await, 0);
assert_eq!(count_kind(&db, Kind::RepoState).await, 0);
assert_eq!(count_kind(&db, Kind::from(1621)).await, 0);
assert!(report.cascade.deleted_events >= 1);
assert!(report.cascade.deleted_events >= 2);
}
#[tokio::test]
@@ -1,577 +0,0 @@
/// Legacy Event Dependency Graph Builder
///
/// Builds a dependency graph for offline cleanup reconciliation scenarios.
/// This module is retained for `cleanup-empty-repos`; it is not used by the
/// production NIP-09 announcement cascade, which lives in `orchestration.rs`.
use std::collections::{HashMap, HashSet};
use nostr_relay_builder::prelude::{Event, EventId, Kind};
/// Maximum depth for graph traversal to prevent infinite loops
const DEFAULT_MAX_DEPTH: usize = 100;
/// Node in the event dependency graph
#[derive(Debug, Clone)]
pub struct EventNode {
/// Event ID
pub event_id: EventId,
/// Event kind
pub kind: Kind,
/// Event IDs this event references (via a/e/q tags)
pub references: HashSet<EventId>,
/// Announcement IDs that make this event acceptable (retention reasons)
/// For announcements themselves, this is empty (they're self-justifying)
/// For dependent events, this tracks which announcements they depend on
pub retention_reasons: HashSet<EventId>,
}
impl EventNode {
/// Create a new event node
pub fn new(event_id: EventId, kind: Kind) -> Self {
Self {
event_id,
kind,
references: HashSet::new(),
retention_reasons: HashSet::new(),
}
}
/// Add a reference to another event
pub fn add_reference(&mut self, referenced_id: EventId) {
self.references.insert(referenced_id);
}
/// Add a retention reason (announcement that justifies keeping this event)
pub fn add_retention_reason(&mut self, announcement_id: EventId) {
self.retention_reasons.insert(announcement_id);
}
}
/// Event dependency graph
///
/// Represents the dependency relationships between events in a repository.
/// Nodes are events, edges are references (a/e/q tags).
#[derive(Debug, Clone)]
pub struct EventGraph {
/// All nodes in the graph, indexed by event ID
nodes: HashMap<EventId, EventNode>,
/// Reverse index: event ID -> events that reference it
dependents: HashMap<EventId, HashSet<EventId>>,
/// Maximum depth for traversal
max_depth: usize,
}
impl EventGraph {
/// Create a new empty event graph
pub fn new() -> Self {
Self {
nodes: HashMap::new(),
dependents: HashMap::new(),
max_depth: DEFAULT_MAX_DEPTH,
}
}
/// Create a new event graph with custom max depth
pub fn with_max_depth(max_depth: usize) -> Self {
Self {
nodes: HashMap::new(),
dependents: HashMap::new(),
max_depth,
}
}
/// Add a node to the graph
pub fn add_node(&mut self, node: EventNode) {
let event_id = node.event_id;
// Update dependents index for all references
for referenced_id in &node.references {
self.dependents
.entry(*referenced_id)
.or_default()
.insert(event_id);
}
self.nodes.insert(event_id, node);
}
/// Add an edge (reference) between two events
///
/// # Arguments
/// * `from` - Event ID that contains the reference
/// * `to` - Event ID being referenced
pub fn add_edge(&mut self, from: EventId, to: EventId) {
// Add to node's references
if let Some(node) = self.nodes.get_mut(&from) {
node.add_reference(to);
}
// Update dependents index
self.dependents.entry(to).or_default().insert(from);
}
/// Get a node by event ID
pub fn get_node(&self, event_id: &EventId) -> Option<&EventNode> {
self.nodes.get(event_id)
}
/// Get all nodes in the graph
pub fn nodes(&self) -> &HashMap<EventId, EventNode> {
&self.nodes
}
/// Get events that this event depends on (references)
pub fn get_dependencies(&self, event_id: &EventId) -> HashSet<EventId> {
self.nodes
.get(event_id)
.map(|node| node.references.clone())
.unwrap_or_default()
}
/// Get events that depend on this event (reverse references)
pub fn get_dependents(&self, event_id: &EventId) -> HashSet<EventId> {
self.dependents.get(event_id).cloned().unwrap_or_default()
}
/// Get all transitive dependents of an event (up to max_depth)
///
/// Returns all events that transitively depend on the given event,
/// stopping at max_depth to prevent infinite loops.
pub fn get_transitive_dependents(&self, event_id: &EventId) -> HashSet<EventId> {
let mut result = HashSet::new();
let mut to_process = vec![(event_id, 0)];
let mut visited = HashSet::new();
while let Some((current_id, depth)) = to_process.pop() {
if depth >= self.max_depth {
continue;
}
if visited.contains(current_id) {
continue;
}
visited.insert(*current_id);
// Get direct dependents
if let Some(dependents) = self.dependents.get(current_id) {
for dependent_id in dependents {
result.insert(*dependent_id);
to_process.push((dependent_id, depth + 1));
}
}
}
result
}
/// Get the maximum traversal depth
pub fn max_depth(&self) -> usize {
self.max_depth
}
/// Get the number of nodes in the graph
pub fn node_count(&self) -> usize {
self.nodes.len()
}
/// Get the number of edges in the graph
pub fn edge_count(&self) -> usize {
self.nodes.values().map(|node| node.references.len()).sum()
}
}
impl Default for EventGraph {
fn default() -> Self {
Self::new()
}
}
/// Build an event dependency graph from a list of events
///
/// Creates a directed graph showing:
/// - Nodes: Events (with event IDs and kinds)
/// - Edges: References (a/e/q tags pointing to other events)
/// - Retention reasons: Which announcements make each event acceptable
///
/// # Arguments
/// * `events` - List of events to build graph from
///
/// # Returns
/// EventGraph with all dependencies mapped
pub fn build_event_graph(events: &[Event]) -> EventGraph {
let mut graph = EventGraph::new();
// First pass: Create nodes for all events
for event in events {
let mut node = EventNode::new(event.id, event.kind);
// Extract references from e/E/q event-id tags. Address refs are resolved
// after every node's address has been indexed.
for tag in event.tags.iter() {
let tag_vec = tag.as_slice();
if tag_vec.is_empty() {
continue;
}
match tag_vec[0].as_str() {
"e" | "E" if tag_vec.len() >= 2 => {
// Event reference
if let Ok(referenced_id) = EventId::from_hex(&tag_vec[1]) {
node.add_reference(referenced_id);
}
}
"q" if tag_vec.len() >= 2 && !tag_vec[1].contains(':') => {
if let Ok(referenced_id) = EventId::from_hex(&tag_vec[1]) {
node.add_reference(referenced_id);
}
}
_ => {}
}
}
graph.add_node(node);
}
// Second pass: Resolve address references and compute retention reasons
// Build a map of addresses to event IDs
let mut address_to_event: HashMap<String, EventId> = HashMap::new();
for event in events {
let kind_num = event.kind.as_u16();
if (30000..40000).contains(&kind_num) {
// Extract d-tag to build address
for tag in event.tags.iter() {
let tag_vec = tag.as_slice();
if tag_vec.len() >= 2 && tag_vec[0] == "d" {
let d_tag = &tag_vec[1];
let address = format!("{}:{}:{}", kind_num, event.pubkey.to_hex(), d_tag);
address_to_event.insert(address, event.id);
break;
}
}
} else if (10000..20000).contains(&kind_num) {
let address = format!("{}:{}", kind_num, event.pubkey.to_hex());
address_to_event.insert(address, event.id);
}
}
// Third pass: Resolve address references and add edges
for event in events {
for tag in event.tags.iter() {
let tag_vec = tag.as_slice();
let is_address_ref = tag_vec.len() >= 2 && matches!(tag_vec[0].as_str(), "a" | "A")
|| (tag_vec.len() >= 2 && tag_vec[0] == "q" && tag_vec[1].contains(':'));
if is_address_ref {
let address = &tag_vec[1];
if let Some(&referenced_id) = address_to_event.get(address) {
graph.add_edge(event.id, referenced_id);
}
}
}
}
// Fourth pass: Compute retention reasons
// An event is retained if it's an announcement (kind 30617) or if it references
// at least one announcement (directly or transitively)
let announcement_kind = Kind::from(30617);
// Find all announcements
let announcements: HashSet<EventId> = events
.iter()
.filter(|e| e.kind == announcement_kind)
.map(|e| e.id)
.collect();
// For each non-announcement event, find which announcements it depends on
for event in events {
if event.kind == announcement_kind {
// Announcements are self-justifying (no retention reasons needed)
continue;
}
// Find all announcements this event transitively depends on
let mut retention_reasons = HashSet::new();
let mut to_process = vec![event.id];
let mut visited = HashSet::new();
let mut depth = 0;
while let Some(current_id) = to_process.pop() {
if depth >= graph.max_depth {
break;
}
if visited.contains(&current_id) {
continue;
}
visited.insert(current_id);
// Check if this is an announcement
if announcements.contains(&current_id) {
retention_reasons.insert(current_id);
}
// Add dependencies to process
if let Some(node) = graph.nodes.get(&current_id) {
for referenced_id in &node.references {
to_process.push(*referenced_id);
}
}
depth += 1;
}
// Update the node with retention reasons
if let Some(node) = graph.nodes.get_mut(&event.id) {
node.retention_reasons = retention_reasons;
}
}
graph
}
#[cfg(test)]
mod tests {
use super::*;
use nostr::event::FinalizeEvent;
use nostr_relay_builder::prelude::{EventBuilder, Keys, Tag};
#[test]
fn test_event_node_creation() {
let event_id =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let kind = Kind::from(30617);
let node = EventNode::new(event_id, kind);
assert_eq!(node.event_id, event_id);
assert_eq!(node.kind, kind);
assert!(node.references.is_empty());
assert!(node.retention_reasons.is_empty());
}
#[test]
fn test_event_node_add_reference() {
let event_id =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let referenced_id =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000002")
.unwrap();
let mut node = EventNode::new(event_id, Kind::from(1617));
node.add_reference(referenced_id);
assert_eq!(node.references.len(), 1);
assert!(node.references.contains(&referenced_id));
}
#[test]
fn test_event_graph_empty() {
let graph = EventGraph::new();
assert_eq!(graph.node_count(), 0);
assert_eq!(graph.edge_count(), 0);
assert_eq!(graph.max_depth(), DEFAULT_MAX_DEPTH);
}
#[test]
fn test_event_graph_with_max_depth() {
let graph = EventGraph::with_max_depth(50);
assert_eq!(graph.max_depth(), 50);
}
#[test]
fn test_event_graph_add_node() {
let mut graph = EventGraph::new();
let event_id =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let node = EventNode::new(event_id, Kind::from(30617));
graph.add_node(node);
assert_eq!(graph.node_count(), 1);
assert!(graph.get_node(&event_id).is_some());
}
#[test]
fn test_event_graph_add_edge() {
let mut graph = EventGraph::new();
let event_id_1 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let event_id_2 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000002")
.unwrap();
graph.add_node(EventNode::new(event_id_1, Kind::from(1617)));
graph.add_node(EventNode::new(event_id_2, Kind::from(30617)));
graph.add_edge(event_id_1, event_id_2);
assert_eq!(graph.edge_count(), 1);
let deps = graph.get_dependencies(&event_id_1);
assert_eq!(deps.len(), 1);
assert!(deps.contains(&event_id_2));
let dependents = graph.get_dependents(&event_id_2);
assert_eq!(dependents.len(), 1);
assert!(dependents.contains(&event_id_1));
}
#[test]
fn test_build_event_graph_empty() {
let events: Vec<Event> = vec![];
let graph = build_event_graph(&events);
assert_eq!(graph.node_count(), 0);
assert_eq!(graph.edge_count(), 0);
}
#[test]
fn test_build_event_graph_single_event() {
let keys = Keys::generate();
let event = EventBuilder::new(Kind::from(30617), "test repo")
.tags(vec![Tag::custom("d", vec!["test-repo".to_string()])])
.finalize(&keys)
.unwrap();
let graph = build_event_graph(std::slice::from_ref(&event));
assert_eq!(graph.node_count(), 1);
assert_eq!(graph.edge_count(), 0);
let node = graph.get_node(&event.id).unwrap();
assert_eq!(node.event_id, event.id);
assert_eq!(node.kind, Kind::from(30617));
assert!(node.references.is_empty());
// Announcements are self-justifying
assert!(node.retention_reasons.is_empty());
}
#[test]
fn test_build_event_graph_with_reference() {
let keys = Keys::generate();
// Create announcement
let announcement = EventBuilder::new(Kind::from(30617), "test repo")
.tags(vec![Tag::custom("d", vec!["test-repo".to_string()])])
.finalize(&keys)
.unwrap();
// Create patch that references the announcement
let patch = EventBuilder::new(Kind::from(1617), "test patch")
.tags(vec![Tag::custom("e", vec![announcement.id.to_hex()])])
.finalize(&keys)
.unwrap();
let graph = build_event_graph(&[announcement.clone(), patch.clone()]);
assert_eq!(graph.node_count(), 2);
assert_eq!(graph.edge_count(), 1);
// Check patch references announcement
let deps = graph.get_dependencies(&patch.id);
assert_eq!(deps.len(), 1);
assert!(deps.contains(&announcement.id));
// Check announcement has patch as dependent
let dependents = graph.get_dependents(&announcement.id);
assert_eq!(dependents.len(), 1);
assert!(dependents.contains(&patch.id));
// Check retention reasons
let patch_node = graph.get_node(&patch.id).unwrap();
assert_eq!(patch_node.retention_reasons.len(), 1);
assert!(patch_node.retention_reasons.contains(&announcement.id));
}
#[test]
fn test_build_event_graph_circular_reference() {
let keys = Keys::generate();
// Create two events that reference each other (shouldn't happen in practice,
// but we need to handle it gracefully)
let event1 = EventBuilder::new(Kind::from(1617), "event 1")
.finalize(&keys)
.unwrap();
let event2 = EventBuilder::new(Kind::from(1617), "event 2")
.tags(vec![Tag::custom("e", vec![event1.id.to_hex()])])
.finalize(&keys)
.unwrap();
// Manually add circular reference (can't do this in real events)
// For this test, we'll just verify the graph handles it without infinite loops
let graph = build_event_graph(&[event1.clone(), event2.clone()]);
assert_eq!(graph.node_count(), 2);
// Get transitive dependents should not hang
let dependents = graph.get_transitive_dependents(&event1.id);
assert!(dependents.len() <= 2); // Should not infinite loop
}
#[test]
fn test_event_graph_transitive_dependents() {
let mut graph = EventGraph::new();
// Create a chain: event1 <- event2 <- event3
let event_id_1 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let event_id_2 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000002")
.unwrap();
let event_id_3 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000003")
.unwrap();
graph.add_node(EventNode::new(event_id_1, Kind::from(30617)));
graph.add_node(EventNode::new(event_id_2, Kind::from(1617)));
graph.add_node(EventNode::new(event_id_3, Kind::from(1619)));
graph.add_edge(event_id_2, event_id_1);
graph.add_edge(event_id_3, event_id_2);
// Get transitive dependents of event1
let dependents = graph.get_transitive_dependents(&event_id_1);
assert_eq!(dependents.len(), 2);
assert!(dependents.contains(&event_id_2));
assert!(dependents.contains(&event_id_3));
}
#[test]
fn test_event_graph_max_depth_limit() {
let mut graph = EventGraph::with_max_depth(2);
// Create a long chain: event1 <- event2 <- event3 <- event4
let event_id_1 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000001")
.unwrap();
let event_id_2 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000002")
.unwrap();
let event_id_3 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000003")
.unwrap();
let event_id_4 =
EventId::from_hex("0000000000000000000000000000000000000000000000000000000000000004")
.unwrap();
graph.add_node(EventNode::new(event_id_1, Kind::from(30617)));
graph.add_node(EventNode::new(event_id_2, Kind::from(1617)));
graph.add_node(EventNode::new(event_id_3, Kind::from(1619)));
graph.add_node(EventNode::new(event_id_4, Kind::from(1619)));
graph.add_edge(event_id_2, event_id_1);
graph.add_edge(event_id_3, event_id_2);
graph.add_edge(event_id_4, event_id_3);
// Get transitive dependents of event1 with max_depth=2
let dependents = graph.get_transitive_dependents(&event_id_1);
// Should only get event2 and event3 (depth 1 and 2), not event4 (depth 3)
assert!(dependents.len() <= 3); // Max depth limits traversal
}
}
+2 -11
View File
@@ -1,17 +1,8 @@
//! Announcement cascade deletion implementation.
//!
//! `orchestration` contains the production NIP-09/blacklist/whitelist cascade
//! algorithm. The `graph`, `reevaluation`, and `traversal` modules are legacy
//! reconciliation helpers retained for `cleanup-empty-repos`; their tests cover
//! that offline reconciliation behavior, not the live deletion cascade.
//! algorithm and the shared retention planner used by maintenance cleanup.
mod graph;
mod orchestration;
mod reevaluation;
mod traversal;
pub use graph::{build_event_graph, EventGraph, EventNode};
pub use reevaluation::{
reevaluate_events_without_announcement, ReevaluationResult, RetentionReason,
};
pub use traversal::{traverse_and_mark_deletions, CircularDependency, TraversalResult};
pub use orchestration::{plan_cascade_retention_for_deleted_announcements, CascadeRetentionPlan};
@@ -29,6 +29,14 @@ struct ParsedAddressRef {
identifier: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct CascadeRetentionPlan {
pub kept_event_ids: HashSet<EventId>,
pub orphan_event_ids: HashSet<EventId>,
pub component_nodes: usize,
pub component_edges: usize,
}
#[derive(Debug, Default)]
struct CascadeGraph {
nodes: HashMap<EventId, Event>,
@@ -147,19 +155,14 @@ impl DeletionPolicy {
return;
}
let deleted_announcement_ids: HashSet<EventId> = deletable_announcements
.iter()
.map(|event| event.id)
.collect();
let seed_event_ids = deleted_announcement_ids.clone();
let seed_addresses = HashSet::from([announcement_addr.to_string()]);
let graph = match self
.expand_cascade_graph(seed_event_ids, seed_addresses)
.await
let plan = match plan_cascade_retention_for_deleted_announcements(
&self.ctx.database,
&deletable_announcements,
&self.ctx.config.domain,
)
.await
{
Ok(graph) => graph,
Ok(plan) => plan,
Err(e) => {
tracing::warn!(error = %e, "Cascade deletion: graph expansion failed; falling back to simple coordinate deletion");
self.delete_announcement_coordinate_and_maybe_state(
@@ -175,22 +178,17 @@ impl DeletionPolicy {
}
};
let kept_events =
compute_retained_events(&graph, &deleted_announcement_ids, &self.ctx.config.domain);
let delete_events =
compute_deletable_events(&graph, &kept_events, &deleted_announcement_ids);
tracing::info!(
announcement = %announcement_addr,
component_nodes = graph.nodes.len(),
component_edges = graph.edge_count(),
kept = kept_events.len(),
deleted = delete_events.len(),
component_nodes = plan.component_nodes,
component_edges = plan.component_edges,
kept = plan.kept_event_ids.len(),
deleted = plan.orphan_event_ids.len(),
"Announcement cascade deletion: acceptance reachability complete"
);
// Hard-delete the orphaned dependent events by id.
let orphan_ids: Vec<EventId> = delete_events.into_iter().collect();
let orphan_ids: Vec<EventId> = plan.orphan_event_ids.into_iter().collect();
if orphan_ids.len() > MAX_CASCADE_ORPHAN_DELETES {
tracing::warn!(
announcement = %announcement_addr,
@@ -409,14 +407,6 @@ impl DeletionPolicy {
.await;
}
async fn expand_cascade_graph(
&self,
seed_event_ids: HashSet<EventId>,
seed_addresses: HashSet<String>,
) -> Result<CascadeGraph, CascadeExpansionError> {
expand_cascade_graph_from_db(&self.ctx.database, seed_event_ids, seed_addresses).await
}
async fn query_address_events_until(
&self,
address: &str,
@@ -426,6 +416,43 @@ impl DeletionPolicy {
}
}
pub async fn plan_cascade_retention_for_deleted_announcements(
database: &SharedDatabase,
deleted_announcements: &[Event],
domain: &str,
) -> Result<CascadeRetentionPlan, String> {
let deleted_announcement_ids: HashSet<EventId> = deleted_announcements
.iter()
.filter(|event| event.kind == Kind::GitRepoAnnouncement)
.map(|event| event.id)
.collect();
if deleted_announcement_ids.is_empty() {
return Ok(CascadeRetentionPlan::default());
}
let seed_addresses: HashSet<String> = deleted_announcements
.iter()
.filter(|event| event.kind == Kind::GitRepoAnnouncement)
.map(announcement_address_of)
.collect();
let graph =
expand_cascade_graph_from_db(database, deleted_announcement_ids.clone(), seed_addresses)
.await
.map_err(|e| e.to_string())?;
let kept_event_ids = compute_retained_events(&graph, &deleted_announcement_ids, domain);
let orphan_event_ids =
compute_deletable_events(&graph, &kept_event_ids, &deleted_announcement_ids);
Ok(CascadeRetentionPlan {
kept_event_ids,
orphan_event_ids,
component_nodes: graph.nodes.len(),
component_edges: graph.edge_count(),
})
}
async fn expand_cascade_graph_from_db(
database: &SharedDatabase,
seed_event_ids: HashSet<EventId>,
@@ -1,649 +0,0 @@
/// Legacy Re-evaluation Engine - Determines event retention after announcement deletion
///
/// When a repository announcement is deleted, this module re-evaluates all dependent
/// events to determine which should be kept (still acceptable through other announcements)
/// vs deleted (no longer acceptable).
///
/// This is critical for multi-maintainer scenarios where one maintainer deletes their
/// announcement but other maintainers' announcements should keep the events alive.
///
/// This module is retained for legacy/offline reconciliation callers. The
/// production NIP-09 announcement cascade algorithm lives in `orchestration.rs`.
use std::collections::{HashMap, HashSet};
use nostr_relay_builder::prelude::{Event, EventId, Filter, Kind};
use crate::nostr::SharedDatabase;
/// Result of re-evaluating events after an announcement deletion
#[derive(Debug, Clone)]
pub struct ReevaluationResult {
/// Events that should be KEPT (still acceptable through other announcements)
pub keep: HashMap<EventId, RetentionReason>,
/// Events that should be DELETED (no longer acceptable)
pub delete: HashSet<EventId>,
}
/// Reason why an event is being kept after re-evaluation
#[derive(Debug, Clone)]
pub struct RetentionReason {
/// Event ID being retained
pub event_id: EventId,
/// Announcement addresses that make this event acceptable
pub valid_announcements: Vec<String>,
/// Event type description for logging
pub event_type: String,
}
/// Re-evaluate events that were dependent on a deleted announcement
///
/// Legacy reconciliation helper. Do not use this as a description of the
/// production NIP-09 announcement cascade.
///
/// This function determines which events should be kept vs deleted when an
/// announcement is removed. It checks if events have alternative retention
/// reasons through other valid announcements.
///
/// # Arguments
/// * `database` - Database to query for announcements and events
/// * `deleted_address` - Address of the deleted announcement (format: `30617:<pubkey>:<identifier>`)
/// * `potentially_affected` - Event IDs that might be affected by the deletion
///
/// # Returns
/// `ReevaluationResult` with events to keep (with reasons) and events to delete
///
/// # Algorithm
/// For each potentially affected event:
/// 1. Query all active repository announcements (kind 30617) from database
/// 2. Extract all `a` tags from the event that reference announcements
/// 3. Check if any of those `a` tags reference announcements OTHER than the deleted one
/// 4. If YES: event is KEPT (has alternative retention reason)
/// 5. If NO: check if event is a repository announcement itself (multi-maintainer case)
/// 6. If repository announcement: check if other maintainers have announcements for same identifier
/// 7. If YES: event is KEPT (other maintainers exist)
/// 8. If NO: event is DELETED (no valid announcements remain)
pub async fn reevaluate_events_without_announcement(
database: &SharedDatabase,
deleted_address: &str,
potentially_affected: &HashSet<EventId>,
) -> Result<ReevaluationResult, String> {
let mut result = ReevaluationResult {
keep: HashMap::new(),
delete: HashSet::new(),
};
// Parse the deleted address to get identifier
let deleted_identifier = parse_identifier_from_address(deleted_address)?;
// Query all active repository announcements (kind 30617) for this identifier
let active_announcements = query_active_announcements(database, &deleted_identifier).await?;
// Build a set of valid announcement addresses (excluding the deleted one)
let valid_addresses: HashSet<String> = active_announcements
.iter()
.map(build_announcement_address)
.filter(|addr| addr != deleted_address)
.collect();
tracing::debug!(
deleted_address = %deleted_address,
identifier = %deleted_identifier,
active_count = active_announcements.len(),
valid_count = valid_addresses.len(),
"Re-evaluating {} events after announcement deletion",
potentially_affected.len()
);
// Re-evaluate each potentially affected event
for event_id in potentially_affected {
// Query the event from database
let filter = Filter::new().id(*event_id);
let events: Vec<Event> = database
.query(filter)
.await
.map_err(|e| format!("Database query failed: {}", e))?
.into_iter()
.collect();
if events.is_empty() {
tracing::debug!(
event_id = %event_id,
"Event not found in database during re-evaluation, marking for deletion (fail-secure)"
);
// If event is not in database, mark it for deletion (fail-secure)
result.delete.insert(*event_id);
continue;
}
let event = &events[0];
// Determine if this event should be kept or deleted
match evaluate_event_retention(event, &deleted_identifier, &valid_addresses).await {
Ok(Some(reason)) => {
tracing::debug!(
event_id = %event_id,
event_type = %reason.event_type,
valid_announcements = ?reason.valid_announcements,
"Event KEPT: has alternative retention reason"
);
result.keep.insert(*event_id, reason);
}
Ok(None) => {
tracing::debug!(
event_id = %event_id,
kind = %event.kind,
"Event DELETED: no valid announcements remain"
);
result.delete.insert(*event_id);
}
Err(e) => {
tracing::error!(
event_id = %event_id,
error = %e,
"Error evaluating event retention, marking for deletion (fail-secure)"
);
result.delete.insert(*event_id);
}
}
}
tracing::info!(
deleted_address = %deleted_address,
keep_count = result.keep.len(),
delete_count = result.delete.len(),
"Re-evaluation complete: {} kept, {} deleted",
result.keep.len(),
result.delete.len()
);
Ok(result)
}
/// Evaluate whether a single event should be retained after announcement deletion
///
/// Returns Some(RetentionReason) if event should be kept, None if it should be deleted.
async fn evaluate_event_retention(
event: &Event,
_deleted_identifier: &str,
valid_addresses: &HashSet<String>,
) -> Result<Option<RetentionReason>, String> {
let event_type = get_event_type_description(event);
// Special case: Repository announcements (kind 30617)
// These are kept if there are OTHER announcements for the same identifier
if event.kind == Kind::GitRepoAnnouncement {
let event_address = build_announcement_address(event);
// If this IS the deleted announcement, it should be deleted
if !valid_addresses.contains(&event_address) {
return Ok(None);
}
// This is a different announcement for the same identifier - keep it
return Ok(Some(RetentionReason {
event_id: event.id,
valid_announcements: vec![event_address],
event_type,
}));
}
// For all other events: check if they have `a` tags referencing valid announcements
let referenced_announcements = extract_announcement_references(event);
#[cfg(test)]
println!(
"Event {}: referenced {} announcements, {} valid addresses",
event.id,
referenced_announcements.len(),
valid_addresses.len()
);
// Filter to only valid announcements (excluding deleted one)
let valid_refs: Vec<String> = referenced_announcements
.into_iter()
.filter(|addr| valid_addresses.contains(addr))
.collect();
if valid_refs.is_empty() {
// No valid announcement references - event should be deleted
tracing::debug!(
event_id = %event.id,
"No valid announcement references found"
);
return Ok(None);
}
// Event has valid announcement references - keep it
tracing::debug!(
event_id = %event.id,
valid_refs_count = valid_refs.len(),
"Event has valid announcement references, keeping"
);
Ok(Some(RetentionReason {
event_id: event.id,
valid_announcements: valid_refs,
event_type,
}))
}
/// Query all active repository announcements for a given identifier
///
/// Returns all kind 30617 events with matching `d` tag (identifier).
async fn query_active_announcements(
database: &SharedDatabase,
identifier: &str,
) -> Result<Vec<Event>, String> {
// Query all repository announcements
let filter = Filter::new().kind(Kind::GitRepoAnnouncement);
let all_announcements: Vec<Event> = database
.query(filter)
.await
.map_err(|e| format!("Database query failed: {}", e))?
.into_iter()
.collect();
// Filter by identifier manually (some databases don't support custom_tag queries properly)
let events: Vec<Event> = all_announcements
.into_iter()
.filter(|event| {
event.tags.iter().any(|tag| {
let tag_vec = tag.as_slice();
tag_vec.len() >= 2 && tag_vec[0] == "d" && tag_vec[1] == identifier
})
})
.collect();
Ok(events)
}
/// Extract all announcement references (`a` tags) from an event
///
/// Returns addresses in format: `30617:<pubkey>:<identifier>`
fn extract_announcement_references(event: &Event) -> Vec<String> {
let mut refs = Vec::new();
for tag in event.tags.iter() {
let tag_vec = tag.as_slice();
if tag_vec.len() >= 2 && tag_vec[0] == "a" {
let address = &tag_vec[1];
// Only include kind 30617 (repository announcement) references
if address.starts_with("30617:") {
refs.push(address.to_string());
}
}
}
refs
}
/// Build announcement address from event
///
/// Format: `30617:<pubkey>:<identifier>`
fn build_announcement_address(event: &Event) -> String {
let identifier = event
.tags
.iter()
.find_map(|tag| {
let tag_vec = tag.as_slice();
if tag_vec.len() >= 2 && tag_vec[0] == "d" {
Some(tag_vec[1].to_string())
} else {
None
}
})
.unwrap_or_default();
format!("30617:{}:{}", event.pubkey.to_hex(), identifier)
}
/// Parse identifier from announcement address
///
/// Address format: `30617:<pubkey>:<identifier>`
fn parse_identifier_from_address(address: &str) -> Result<String, String> {
let parts: Vec<&str> = address.split(':').collect();
if parts.len() != 3 {
return Err(format!("Invalid address format: {}", address));
}
if parts[0] != "30617" {
return Err(format!(
"Not a repository announcement address: {}",
address
));
}
Ok(parts[2].to_string())
}
/// Get human-readable event type description for logging
fn get_event_type_description(event: &Event) -> String {
match event.kind {
Kind::GitRepoAnnouncement => "Repository Announcement".to_string(),
Kind::RepoState => "Repository State".to_string(),
Kind::GitPullRequest => "Pull Request".to_string(),
Kind::GitPullRequestUpdate => "PR Update".to_string(),
kind if kind == Kind::from(1621) => "Issue".to_string(),
kind if kind == Kind::from(1617) => "Patch".to_string(),
kind if kind == Kind::from(1630) => "Issue Status".to_string(),
kind if kind == Kind::from(1631) => "PR Status".to_string(),
kind if kind == Kind::from(1632) => "Patch Status".to_string(),
kind if kind == Kind::from(1633) => "Repository Status".to_string(),
_ => format!("Event (kind {})", event.kind.as_u16()),
}
}
#[cfg(test)]
mod tests {
use super::*;
use nostr::event::FinalizeEvent;
use nostr_memory::MemoryDatabase;
use nostr_relay_builder::prelude::{EventBuilder, Keys, NostrDatabase, Tag};
use std::sync::Arc;
/// Helper to create a test repository announcement
fn create_announcement(keys: &Keys, identifier: &str, domain: &str) -> Event {
EventBuilder::new(Kind::GitRepoAnnouncement, "Test repository")
.tags(vec![
Tag::custom("d", vec![identifier.to_string()]),
Tag::custom("clone", vec![format!("https://{}/{}", domain, identifier)]),
Tag::custom("relays", vec![format!("wss://{}", domain)]),
])
.finalize(keys)
.unwrap()
}
/// Helper to create a test issue event
fn create_issue(keys: &Keys, announcement_address: &str) -> Event {
EventBuilder::new(Kind::from(1621), "Test issue")
.tags(vec![Tag::custom(
"a",
vec![announcement_address.to_string()],
)])
.finalize(keys)
.unwrap()
}
/// Helper to create a test patch event
fn create_patch(keys: &Keys, announcement_address: &str) -> Event {
EventBuilder::new(Kind::from(1617), "Test patch")
.tags(vec![Tag::custom(
"a",
vec![announcement_address.to_string()],
)])
.finalize(keys)
.unwrap()
}
#[tokio::test]
async fn test_single_maintainer_all_deleted() {
// Setup: Single maintainer with announcement and dependent events
let keys = Keys::generate();
let identifier = "test-repo";
let domain = "relay.example.com";
let announcement = create_announcement(&keys, identifier, domain);
let announcement_addr = build_announcement_address(&announcement);
let issue = create_issue(&keys, &announcement_addr);
let patch = create_patch(&keys, &announcement_addr);
// Create in-memory database and store events
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
db.save_event(&announcement).await.unwrap();
db.save_event(&issue).await.unwrap();
db.save_event(&patch).await.unwrap();
// Re-evaluate after deleting the only announcement
let affected = HashSet::from([issue.id, patch.id]);
let result = reevaluate_events_without_announcement(&db, &announcement_addr, &affected)
.await
.unwrap();
// All events should be deleted (no other announcements)
assert_eq!(result.keep.len(), 0);
assert_eq!(result.delete.len(), 2);
assert!(result.delete.contains(&issue.id));
assert!(result.delete.contains(&patch.id));
}
#[tokio::test]
async fn test_two_maintainers_one_deletes() {
// Setup: Two maintainers with announcements for same repository
let maintainer1 = Keys::generate();
let maintainer2 = Keys::generate();
let identifier = "shared-repo";
let domain = "relay.example.com";
let announcement1 = create_announcement(&maintainer1, identifier, domain);
let announcement2 = create_announcement(&maintainer2, identifier, domain);
let addr1 = build_announcement_address(&announcement1);
let addr2 = build_announcement_address(&announcement2);
// Create issue referencing both announcements
let issue = EventBuilder::new(Kind::from(1621), "Test issue")
.tags(vec![
Tag::custom("a", vec![addr1.clone()]),
Tag::custom("a", vec![addr2.clone()]),
])
.finalize(&maintainer1)
.unwrap();
// Create patch referencing only first announcement
let patch = create_patch(&maintainer1, &addr1);
// Store in database
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
db.save_event(&announcement1).await.unwrap();
db.save_event(&announcement2).await.unwrap();
db.save_event(&issue).await.unwrap();
db.save_event(&patch).await.unwrap();
// Re-evaluate after deleting first announcement
let affected = HashSet::from([issue.id, patch.id]);
let result = reevaluate_events_without_announcement(&db, &addr1, &affected)
.await
.unwrap();
// Issue should be kept (references announcement2)
assert_eq!(result.keep.len(), 1);
assert!(result.keep.contains_key(&issue.id));
let issue_reason = result.keep.get(&issue.id).unwrap();
assert_eq!(issue_reason.valid_announcements.len(), 1);
assert_eq!(issue_reason.valid_announcements[0], addr2);
// Patch should be deleted (only referenced announcement1)
assert_eq!(result.delete.len(), 1);
assert!(result.delete.contains(&patch.id));
}
#[tokio::test]
async fn test_multiple_announcement_references() {
// Setup: Event with multiple announcement references
let maintainer1 = Keys::generate();
let maintainer2 = Keys::generate();
let maintainer3 = Keys::generate();
let identifier = "multi-repo";
let domain = "relay.example.com";
let announcement1 = create_announcement(&maintainer1, identifier, domain);
let announcement2 = create_announcement(&maintainer2, identifier, domain);
let announcement3 = create_announcement(&maintainer3, identifier, domain);
let addr1 = build_announcement_address(&announcement1);
let addr2 = build_announcement_address(&announcement2);
let addr3 = build_announcement_address(&announcement3);
// Create issue referencing all three announcements
let issue = EventBuilder::new(Kind::from(1621), "Test issue")
.tags(vec![
Tag::custom("a", vec![addr1.clone()]),
Tag::custom("a", vec![addr2.clone()]),
Tag::custom("a", vec![addr3.clone()]),
])
.finalize(&maintainer1)
.unwrap();
// Store in database
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
db.save_event(&announcement1).await.unwrap();
db.save_event(&announcement2).await.unwrap();
db.save_event(&announcement3).await.unwrap();
db.save_event(&issue).await.unwrap();
// Re-evaluate after deleting first announcement
let affected = HashSet::from([issue.id]);
let result = reevaluate_events_without_announcement(&db, &addr1, &affected)
.await
.unwrap();
// Issue should be kept (still has announcement2 and announcement3)
assert_eq!(result.keep.len(), 1);
assert!(result.keep.contains_key(&issue.id));
let reason = result.keep.get(&issue.id).unwrap();
assert_eq!(reason.valid_announcements.len(), 2);
assert!(reason.valid_announcements.contains(&addr2));
assert!(reason.valid_announcements.contains(&addr3));
assert!(!reason.valid_announcements.contains(&addr1));
assert_eq!(result.delete.len(), 0);
}
#[tokio::test]
async fn test_no_valid_announcements() {
// Setup: Event with no valid announcement references
let keys = Keys::generate();
let identifier = "orphan-repo";
let domain = "relay.example.com";
let announcement = create_announcement(&keys, identifier, domain);
let announcement_addr = build_announcement_address(&announcement);
// Create issue with no announcement references
let issue = EventBuilder::new(Kind::from(1621), "Orphan issue")
.finalize(&keys)
.unwrap();
// Store in database
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
db.save_event(&announcement).await.unwrap();
db.save_event(&issue).await.unwrap();
// Re-evaluate after deleting announcement
let affected = HashSet::from([issue.id]);
let result = reevaluate_events_without_announcement(&db, &announcement_addr, &affected)
.await
.unwrap();
// Issue should be deleted (no announcement references)
assert_eq!(result.keep.len(), 0);
assert_eq!(result.delete.len(), 1);
assert!(result.delete.contains(&issue.id));
}
#[tokio::test]
async fn test_repository_announcement_deletion() {
// Setup: Multiple announcements for same identifier
let maintainer1 = Keys::generate();
let maintainer2 = Keys::generate();
let identifier = "shared-repo";
let domain = "relay.example.com";
let announcement1 = create_announcement(&maintainer1, identifier, domain);
let announcement2 = create_announcement(&maintainer2, identifier, domain);
let addr1 = build_announcement_address(&announcement1);
let _addr2 = build_announcement_address(&announcement2);
// Store in database
let db: Arc<dyn NostrDatabase> = Arc::new(MemoryDatabase::unbounded());
db.save_event(&announcement1).await.unwrap();
db.save_event(&announcement2).await.unwrap();
// Re-evaluate announcements after deleting first one
let affected = HashSet::from([announcement1.id, announcement2.id]);
let result = reevaluate_events_without_announcement(&db, &addr1, &affected)
.await
.unwrap();
// announcement2 should be kept (different maintainer)
assert_eq!(result.keep.len(), 1);
assert!(result.keep.contains_key(&announcement2.id));
// announcement1 should be deleted (it's the deleted one)
assert_eq!(result.delete.len(), 1);
assert!(result.delete.contains(&announcement1.id));
}
#[test]
fn test_parse_identifier_from_address() {
let address = "30617:abc123:my-repo";
let identifier = parse_identifier_from_address(address).unwrap();
assert_eq!(identifier, "my-repo");
}
#[test]
fn test_parse_identifier_invalid_format() {
let address = "invalid";
let result = parse_identifier_from_address(address);
assert!(result.is_err());
}
#[test]
fn test_parse_identifier_wrong_kind() {
let address = "30618:abc123:my-repo";
let result = parse_identifier_from_address(address);
assert!(result.is_err());
}
#[test]
fn test_extract_announcement_references() {
let keys = Keys::generate();
let event = EventBuilder::new(Kind::from(1621), "Test")
.tags(vec![
Tag::custom("a", vec!["30617:abc:repo1".to_string()]),
Tag::custom("a", vec!["30617:def:repo2".to_string()]),
Tag::custom("a", vec!["30618:ghi:state".to_string()]), // Wrong kind
Tag::custom("e", vec!["event123".to_string()]), // Not an 'a' tag
])
.finalize(&keys)
.unwrap();
let refs = extract_announcement_references(&event);
assert_eq!(refs.len(), 2);
assert!(refs.contains(&"30617:abc:repo1".to_string()));
assert!(refs.contains(&"30617:def:repo2".to_string()));
}
#[test]
fn test_build_announcement_address() {
let keys = Keys::generate();
let event = create_announcement(&keys, "my-repo", "example.com");
let address = build_announcement_address(&event);
assert!(address.starts_with("30617:"));
assert!(address.ends_with(":my-repo"));
assert_eq!(address.split(':').count(), 3);
}
#[test]
fn test_get_event_type_description() {
let keys = Keys::generate();
let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "")
.finalize(&keys)
.unwrap();
assert_eq!(
get_event_type_description(&announcement),
"Repository Announcement"
);
let issue = EventBuilder::new(Kind::from(1621), "")
.finalize(&keys)
.unwrap();
assert_eq!(get_event_type_description(&issue), "Issue");
let pr = EventBuilder::new(Kind::GitPullRequest, "")
.finalize(&keys)
.unwrap();
assert_eq!(get_event_type_description(&pr), "Pull Request");
}
}
@@ -1,714 +0,0 @@
/// Legacy Graph Traversal and Circular Dependency Detection
///
/// Determines which events to keep vs delete for offline cleanup reconciliation
/// by:
/// 1. Starting from kept repository announcements (kind 30617)
/// 2. Traversing the dependency graph to find all reachable events
/// 3. Detecting circular dependencies and handling them appropriately
/// 4. Marking unreachable events for deletion
///
/// This module is retained for `cleanup-empty-repos`; it is not used by the
/// production NIP-09 announcement cascade, which lives in `orchestration.rs`.
use std::collections::{HashMap, HashSet, VecDeque};
use nostr_relay_builder::prelude::EventId;
use super::graph::EventGraph;
use super::reevaluation::RetentionReason;
/// Result of graph traversal and deletion marking
#[derive(Debug, Clone)]
pub struct TraversalResult {
/// Events that should be KEPT (reachable from kept announcements)
pub keep: HashSet<EventId>,
/// Events that should be DELETED (unreachable from kept announcements)
pub delete: HashSet<EventId>,
/// Circular dependencies detected (for logging/debugging)
pub circular_dependencies: Vec<CircularDependency>,
}
/// Information about a detected circular dependency
#[derive(Debug, Clone)]
pub struct CircularDependency {
/// Event IDs involved in the circular dependency
pub events: Vec<EventId>,
/// Whether this cycle is anchored (has a path to a kept announcement)
pub anchored: bool,
}
/// Traverse the event graph and mark events for deletion
///
/// Legacy reconciliation helper for `cleanup-empty-repos`. Do not use this as a
/// description of the production NIP-09 announcement cascade.
///
/// This function implements a BFS traversal starting from kept repository announcements,
/// marking all reachable events as KEEP and unreachable events as DELETE.
///
/// # Arguments
/// * `graph` - Event dependency graph built from all events
/// * `kept_events` - Map of events to keep (from re-evaluation), keyed by event ID
/// * `deleted_announcement_id` - ID of the deleted announcement (for logging)
///
/// # Returns
/// `TraversalResult` with events to keep, events to delete, and detected circular dependencies
///
/// # Algorithm
/// 1. Identify all kept announcements (kind 30617) from `kept_events`
/// 2. Use BFS to traverse from kept announcements, following dependency edges
/// 3. Mark all reachable events as KEEP
/// 4. Detect circular dependencies during traversal
/// 5. Handle circular dependencies:
/// - If cycle is anchored (reachable from kept announcement): KEEP all events in cycle
/// - If cycle is unanchored (isolated): DELETE all events in cycle
/// 6. Mark all unreachable events as DELETE
/// 7. Respect max_depth from graph to prevent infinite loops
pub fn traverse_and_mark_deletions(
graph: &EventGraph,
kept_events: &HashMap<EventId, RetentionReason>,
_deleted_announcement_id: EventId,
) -> TraversalResult {
let mut result = TraversalResult {
keep: HashSet::new(),
delete: HashSet::new(),
circular_dependencies: Vec::new(),
};
// Step 1: Identify kept announcements (kind 30617) as starting points
let kept_announcements: Vec<EventId> = kept_events
.iter()
.filter_map(|(event_id, reason)| {
// Check if this is a repository announcement
if reason.event_type == "Repository Announcement" {
Some(*event_id)
} else {
None
}
})
.collect();
tracing::debug!(
kept_announcements_count = kept_announcements.len(),
kept_events_count = kept_events.len(),
total_nodes = graph.node_count(),
"Starting graph traversal from kept announcements"
);
// Step 2: BFS traversal from kept announcements
let reachable = bfs_traverse_from_announcements(graph, &kept_announcements);
tracing::debug!(
reachable_count = reachable.len(),
"BFS traversal complete, found reachable events"
);
// Step 3: Detect circular dependencies
let circular_deps = detect_circular_dependencies(graph, &reachable);
tracing::debug!(
circular_deps_count = circular_deps.len(),
"Circular dependency detection complete"
);
// Step 4: Mark events as KEEP or DELETE
for event_id in graph.nodes().keys() {
if reachable.contains(event_id) {
// Event is reachable from a kept announcement - KEEP
result.keep.insert(*event_id);
} else {
// Event is not reachable - DELETE
result.delete.insert(*event_id);
}
}
result.circular_dependencies = circular_deps;
tracing::info!(
keep_count = result.keep.len(),
delete_count = result.delete.len(),
circular_deps_count = result.circular_dependencies.len(),
"Graph traversal complete: {} kept, {} deleted, {} circular dependencies",
result.keep.len(),
result.delete.len(),
result.circular_dependencies.len()
);
result
}
/// BFS traversal from kept announcements to find all reachable events
///
/// Uses breadth-first search to traverse the dependency graph, following both
/// forward references (dependencies) and backward references (dependents).
///
/// # Arguments
/// * `graph` - Event dependency graph
/// * `kept_announcements` - Starting points for traversal (kept announcements)
///
/// # Returns
/// Set of all event IDs reachable from kept announcements
fn bfs_traverse_from_announcements(
graph: &EventGraph,
kept_announcements: &[EventId],
) -> HashSet<EventId> {
let mut reachable = HashSet::new();
let mut queue = VecDeque::new();
let mut visited = HashSet::new();
// Initialize queue with kept announcements
for announcement_id in kept_announcements {
queue.push_back((*announcement_id, 0)); // (event_id, depth)
reachable.insert(*announcement_id);
}
let max_depth = graph.max_depth();
while let Some((current_id, depth)) = queue.pop_front() {
// Respect max depth to prevent infinite loops
if depth >= max_depth {
tracing::warn!(
event_id = %current_id,
depth = depth,
max_depth = max_depth,
"Reached max depth during BFS traversal"
);
continue;
}
// Skip if already visited
if visited.contains(&current_id) {
continue;
}
visited.insert(current_id);
// Traverse to dependents (events that reference this event)
// These are events that depend on the current event, so if the current
// event is kept, its dependents should also be kept
let dependents = graph.get_dependents(&current_id);
for dependent_id in dependents {
if !reachable.contains(&dependent_id) {
reachable.insert(dependent_id);
queue.push_back((dependent_id, depth + 1));
}
}
// Note: We do NOT traverse to dependencies (events this event references)
// because we're doing a forward traversal from announcements to dependents.
// If we traversed backwards, we'd be going from dependents to announcements,
// which is the opposite direction.
}
reachable
}
/// Detect circular dependencies in the graph
///
/// Finds cycles in the dependency graph and determines if they're anchored
/// (reachable from a kept announcement) or unanchored (isolated).
///
/// # Arguments
/// * `graph` - Event dependency graph
/// * `reachable` - Set of events reachable from kept announcements
///
/// # Returns
/// List of detected circular dependencies with anchoring information
fn detect_circular_dependencies(
graph: &EventGraph,
reachable: &HashSet<EventId>,
) -> Vec<CircularDependency> {
let mut circular_deps = Vec::new();
let mut visited = HashSet::new();
let mut recursion_stack = HashSet::new();
// Use DFS to detect cycles
for event_id in graph.nodes().keys() {
if !visited.contains(event_id) {
let mut path = Vec::new();
detect_cycles_dfs(
graph,
*event_id,
&mut visited,
&mut recursion_stack,
&mut path,
&mut circular_deps,
reachable,
);
}
}
circular_deps
}
/// DFS helper for cycle detection
///
/// Uses depth-first search with a recursion stack to detect cycles.
/// When a cycle is found, checks if it's anchored (any event in cycle is reachable).
#[allow(clippy::too_many_arguments)]
fn detect_cycles_dfs(
graph: &EventGraph,
current_id: EventId,
visited: &mut HashSet<EventId>,
recursion_stack: &mut HashSet<EventId>,
path: &mut Vec<EventId>,
circular_deps: &mut Vec<CircularDependency>,
reachable: &HashSet<EventId>,
) {
visited.insert(current_id);
recursion_stack.insert(current_id);
path.push(current_id);
// Get dependencies (events this event references)
let dependencies = graph.get_dependencies(&current_id);
for dep_id in dependencies {
if !visited.contains(&dep_id) {
// Continue DFS
detect_cycles_dfs(
graph,
dep_id,
visited,
recursion_stack,
path,
circular_deps,
reachable,
);
} else if recursion_stack.contains(&dep_id) {
// Found a cycle! Extract the cycle from the path
if let Some(cycle_start_idx) = path.iter().position(|&id| id == dep_id) {
let cycle_events: Vec<EventId> = path[cycle_start_idx..].to_vec();
// Check if cycle is anchored (any event in cycle is reachable)
let anchored = cycle_events.iter().any(|id| reachable.contains(id));
// Only add if we haven't seen this cycle before
// (cycles can be detected multiple times from different entry points)
let cycle_set: HashSet<EventId> = cycle_events.iter().copied().collect();
let is_duplicate = circular_deps.iter().any(|cd| {
let cd_set: HashSet<EventId> = cd.events.iter().copied().collect();
cd_set == cycle_set
});
if !is_duplicate {
tracing::debug!(
cycle_events = ?cycle_events,
anchored = anchored,
"Detected circular dependency"
);
circular_deps.push(CircularDependency {
events: cycle_events,
anchored,
});
}
}
}
}
path.pop();
recursion_stack.remove(&current_id);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::nostr::lifecycle::deletion::{build_event_graph, EventNode, RetentionReason};
use nostr::event::FinalizeEvent;
use nostr_relay_builder::prelude::{EventBuilder, EventId, Keys, Kind, Tag};
/// Helper to create a test event ID from a number
fn test_event_id(n: u8) -> EventId {
let hex = format!("{:064x}", n);
EventId::from_hex(&hex).unwrap()
}
#[test]
fn test_simple_linear_dependencies() {
// Create a simple linear dependency chain:
// announcement -> patch -> issue
let mut graph = EventGraph::new();
let announcement_id = test_event_id(1);
let patch_id = test_event_id(2);
let issue_id = test_event_id(3);
// Add nodes
graph.add_node(EventNode::new(announcement_id, Kind::from(30617)));
graph.add_node(EventNode::new(patch_id, Kind::from(1617)));
graph.add_node(EventNode::new(issue_id, Kind::from(1621)));
// Add edges (patch references announcement, issue references patch)
graph.add_edge(patch_id, announcement_id);
graph.add_edge(issue_id, patch_id);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement_id,
RetentionReason {
event_id: announcement_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement_id);
// All events should be kept (reachable from announcement)
assert_eq!(result.keep.len(), 3);
assert!(result.keep.contains(&announcement_id));
assert!(result.keep.contains(&patch_id));
assert!(result.keep.contains(&issue_id));
assert_eq!(result.delete.len(), 0);
assert_eq!(result.circular_dependencies.len(), 0);
}
#[test]
fn test_isolated_subgraph() {
// Create two subgraphs:
// 1. announcement1 -> patch1 (KEPT)
// 2. patch2 -> issue2 (isolated, should be DELETED)
let mut graph = EventGraph::new();
let announcement1_id = test_event_id(1);
let patch1_id = test_event_id(2);
let patch2_id = test_event_id(3);
let issue2_id = test_event_id(4);
// Add nodes
graph.add_node(EventNode::new(announcement1_id, Kind::from(30617)));
graph.add_node(EventNode::new(patch1_id, Kind::from(1617)));
graph.add_node(EventNode::new(patch2_id, Kind::from(1617)));
graph.add_node(EventNode::new(issue2_id, Kind::from(1621)));
// Add edges
graph.add_edge(patch1_id, announcement1_id);
graph.add_edge(issue2_id, patch2_id);
// Create kept_events with only announcement1
let mut kept_events = HashMap::new();
kept_events.insert(
announcement1_id,
RetentionReason {
event_id: announcement1_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement1_id);
// announcement1 and patch1 should be kept
assert_eq!(result.keep.len(), 2);
assert!(result.keep.contains(&announcement1_id));
assert!(result.keep.contains(&patch1_id));
// patch2 and issue2 should be deleted (isolated)
assert_eq!(result.delete.len(), 2);
assert!(result.delete.contains(&patch2_id));
assert!(result.delete.contains(&issue2_id));
}
#[test]
fn test_circular_dependency_anchored() {
// Create a circular dependency that's anchored to a kept announcement:
// announcement -> patch1 <-> patch2 (circular)
let mut graph = EventGraph::new();
let announcement_id = test_event_id(1);
let patch1_id = test_event_id(2);
let patch2_id = test_event_id(3);
// Add nodes
graph.add_node(EventNode::new(announcement_id, Kind::from(30617)));
graph.add_node(EventNode::new(patch1_id, Kind::from(1617)));
graph.add_node(EventNode::new(patch2_id, Kind::from(1617)));
// Add edges (circular: patch1 -> announcement, patch1 -> patch2, patch2 -> patch1)
graph.add_edge(patch1_id, announcement_id);
graph.add_edge(patch1_id, patch2_id);
graph.add_edge(patch2_id, patch1_id);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement_id,
RetentionReason {
event_id: announcement_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement_id);
// All events should be kept (circular dependency is anchored)
assert_eq!(result.keep.len(), 3);
assert!(result.keep.contains(&announcement_id));
assert!(result.keep.contains(&patch1_id));
assert!(result.keep.contains(&patch2_id));
assert_eq!(result.delete.len(), 0);
// Should detect the circular dependency
assert_eq!(result.circular_dependencies.len(), 1);
assert!(result.circular_dependencies[0].anchored);
}
#[test]
fn test_circular_dependency_unanchored() {
// Create an isolated circular dependency (no announcement):
// patch1 <-> patch2 (circular, isolated)
let mut graph = EventGraph::new();
let patch1_id = test_event_id(1);
let patch2_id = test_event_id(2);
// Add nodes
graph.add_node(EventNode::new(patch1_id, Kind::from(1617)));
graph.add_node(EventNode::new(patch2_id, Kind::from(1617)));
// Add edges (circular: patch1 -> patch2, patch2 -> patch1)
graph.add_edge(patch1_id, patch2_id);
graph.add_edge(patch2_id, patch1_id);
// No kept events (no announcements)
let kept_events = HashMap::new();
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, test_event_id(99));
// Both events should be deleted (unanchored circular dependency)
assert_eq!(result.keep.len(), 0);
assert_eq!(result.delete.len(), 2);
assert!(result.delete.contains(&patch1_id));
assert!(result.delete.contains(&patch2_id));
// Should detect the circular dependency as unanchored
assert_eq!(result.circular_dependencies.len(), 1);
assert!(!result.circular_dependencies[0].anchored);
}
#[test]
fn test_deep_dependency_chain() {
// Create a deep dependency chain (5 levels):
// announcement -> e1 -> e2 -> e3 -> e4
let mut graph = EventGraph::new();
let announcement_id = test_event_id(1);
let e1_id = test_event_id(2);
let e2_id = test_event_id(3);
let e3_id = test_event_id(4);
let e4_id = test_event_id(5);
// Add nodes
graph.add_node(EventNode::new(announcement_id, Kind::from(30617)));
graph.add_node(EventNode::new(e1_id, Kind::from(1617)));
graph.add_node(EventNode::new(e2_id, Kind::from(1617)));
graph.add_node(EventNode::new(e3_id, Kind::from(1617)));
graph.add_node(EventNode::new(e4_id, Kind::from(1617)));
// Add edges (linear chain)
graph.add_edge(e1_id, announcement_id);
graph.add_edge(e2_id, e1_id);
graph.add_edge(e3_id, e2_id);
graph.add_edge(e4_id, e3_id);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement_id,
RetentionReason {
event_id: announcement_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement_id);
// All events should be kept
assert_eq!(result.keep.len(), 5);
assert!(result.keep.contains(&announcement_id));
assert!(result.keep.contains(&e1_id));
assert!(result.keep.contains(&e2_id));
assert!(result.keep.contains(&e3_id));
assert!(result.keep.contains(&e4_id));
assert_eq!(result.delete.len(), 0);
}
#[test]
fn test_max_depth_limit() {
// Create a deep chain that exceeds max depth
let mut graph = EventGraph::with_max_depth(3);
let announcement_id = test_event_id(1);
let e1_id = test_event_id(2);
let e2_id = test_event_id(3);
let e3_id = test_event_id(4);
let e4_id = test_event_id(5);
let e5_id = test_event_id(6);
// Add nodes
graph.add_node(EventNode::new(announcement_id, Kind::from(30617)));
graph.add_node(EventNode::new(e1_id, Kind::from(1617)));
graph.add_node(EventNode::new(e2_id, Kind::from(1617)));
graph.add_node(EventNode::new(e3_id, Kind::from(1617)));
graph.add_node(EventNode::new(e4_id, Kind::from(1617)));
graph.add_node(EventNode::new(e5_id, Kind::from(1617)));
// Add edges (linear chain)
graph.add_edge(e1_id, announcement_id);
graph.add_edge(e2_id, e1_id);
graph.add_edge(e3_id, e2_id);
graph.add_edge(e4_id, e3_id);
graph.add_edge(e5_id, e4_id);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement_id,
RetentionReason {
event_id: announcement_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement_id);
// Should keep events within max_depth (announcement, e1, e2, e3)
// e4 and e5 might not be reached due to depth limit
assert!(result.keep.contains(&announcement_id));
assert!(result.keep.contains(&e1_id));
assert!(result.keep.contains(&e2_id));
assert!(result.keep.contains(&e3_id));
// Note: Depending on BFS implementation, e4 and e5 might be deleted
// or kept. The important thing is that max_depth prevents infinite loops.
}
#[test]
fn test_complex_reference_graph() {
// Create a complex graph with multiple paths:
// announcement
// / | \
// e1 e2 e3
// | | |
// e4 e5 e6
// \ | /
// e7
let mut graph = EventGraph::new();
let announcement_id = test_event_id(1);
let e1_id = test_event_id(2);
let e2_id = test_event_id(3);
let e3_id = test_event_id(4);
let e4_id = test_event_id(5);
let e5_id = test_event_id(6);
let e6_id = test_event_id(7);
let e7_id = test_event_id(8);
// Add nodes
graph.add_node(EventNode::new(announcement_id, Kind::from(30617)));
graph.add_node(EventNode::new(e1_id, Kind::from(1617)));
graph.add_node(EventNode::new(e2_id, Kind::from(1617)));
graph.add_node(EventNode::new(e3_id, Kind::from(1617)));
graph.add_node(EventNode::new(e4_id, Kind::from(1617)));
graph.add_node(EventNode::new(e5_id, Kind::from(1617)));
graph.add_node(EventNode::new(e6_id, Kind::from(1617)));
graph.add_node(EventNode::new(e7_id, Kind::from(1617)));
// Add edges
graph.add_edge(e1_id, announcement_id);
graph.add_edge(e2_id, announcement_id);
graph.add_edge(e3_id, announcement_id);
graph.add_edge(e4_id, e1_id);
graph.add_edge(e5_id, e2_id);
graph.add_edge(e6_id, e3_id);
graph.add_edge(e7_id, e4_id);
graph.add_edge(e7_id, e5_id);
graph.add_edge(e7_id, e6_id);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement_id,
RetentionReason {
event_id: announcement_id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement_id);
// All events should be kept (all reachable from announcement)
assert_eq!(result.keep.len(), 8);
assert!(result.keep.contains(&announcement_id));
assert!(result.keep.contains(&e1_id));
assert!(result.keep.contains(&e2_id));
assert!(result.keep.contains(&e3_id));
assert!(result.keep.contains(&e4_id));
assert!(result.keep.contains(&e5_id));
assert!(result.keep.contains(&e6_id));
assert!(result.keep.contains(&e7_id));
assert_eq!(result.delete.len(), 0);
}
#[test]
fn test_integration_with_build_event_graph() {
// Integration test using build_event_graph from Phase 3A
let keys = Keys::generate();
// Create announcement
let announcement = EventBuilder::new(Kind::from(30617), "test repo")
.tags(vec![Tag::custom("d", vec!["test-repo".to_string()])])
.finalize(&keys)
.unwrap();
// Create patch that references the announcement
let patch = EventBuilder::new(Kind::from(1617), "test patch")
.tags(vec![Tag::custom("e", vec![announcement.id.to_hex()])])
.finalize(&keys)
.unwrap();
// Create issue that references the patch
let issue = EventBuilder::new(Kind::from(1621), "test issue")
.tags(vec![Tag::custom("e", vec![patch.id.to_hex()])])
.finalize(&keys)
.unwrap();
// Build graph
let graph = build_event_graph(&[announcement.clone(), patch.clone(), issue.clone()]);
// Create kept_events with announcement
let mut kept_events = HashMap::new();
kept_events.insert(
announcement.id,
RetentionReason {
event_id: announcement.id,
valid_announcements: vec![],
event_type: "Repository Announcement".to_string(),
},
);
// Traverse
let result = traverse_and_mark_deletions(&graph, &kept_events, announcement.id);
// All events should be kept
assert_eq!(result.keep.len(), 3);
assert!(result.keep.contains(&announcement.id));
assert!(result.keep.contains(&patch.id));
assert!(result.keep.contains(&issue.id));
assert_eq!(result.delete.len(), 0);
}
}
+1 -3
View File
@@ -19,6 +19,4 @@ pub use startup::{
WhitelistRestoreStats,
};
pub use cascade::{build_event_graph, EventGraph, EventNode};
pub use cascade::{reevaluate_events_without_announcement, ReevaluationResult, RetentionReason};
pub use cascade::{traverse_and_mark_deletions, CircularDependency, TraversalResult};
pub use cascade::{plan_cascade_retention_for_deleted_announcements, CascadeRetentionPlan};