mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
fix(nostr): expand cascade deletion by reference graph
This commit is contained in:
@@ -205,7 +205,8 @@ pub fn build_event_graph(events: &[Event]) -> EventGraph {
|
||||
for event in events {
|
||||
let mut node = EventNode::new(event.id, event.kind);
|
||||
|
||||
// Extract references from a/e/q tags
|
||||
// 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() {
|
||||
@@ -213,16 +214,16 @@ pub fn build_event_graph(events: &[Event]) -> EventGraph {
|
||||
}
|
||||
|
||||
match tag_vec[0].as_str() {
|
||||
"e" | "q" if tag_vec.len() >= 2 => {
|
||||
"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);
|
||||
}
|
||||
}
|
||||
"a" if tag_vec.len() >= 2 => {
|
||||
// Address reference - need to find the event with this address
|
||||
// For now, we'll handle this in a second pass after all nodes are added
|
||||
// Store the address temporarily (we'll resolve it later)
|
||||
"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);
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
@@ -235,9 +236,8 @@ pub fn build_event_graph(events: &[Event]) -> EventGraph {
|
||||
// Build a map of addresses to event IDs
|
||||
let mut address_to_event: HashMap<String, EventId> = HashMap::new();
|
||||
for event in events {
|
||||
// Check if this is a replaceable/addressable event (kind 30000-39999)
|
||||
let kind_num = event.kind.as_u16();
|
||||
if (30000..=39999).contains(&kind_num) {
|
||||
if (30000..40000).contains(&kind_num) {
|
||||
// Extract d-tag to build address
|
||||
for tag in event.tags.iter() {
|
||||
let tag_vec = tag.as_slice();
|
||||
@@ -248,6 +248,9 @@ pub fn build_event_graph(events: &[Event]) -> EventGraph {
|
||||
break;
|
||||
}
|
||||
}
|
||||
} else if (10000..20000).contains(&kind_num) {
|
||||
let address = format!("{}:{}", kind_num, event.pubkey.to_hex());
|
||||
address_to_event.insert(address, event.id);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -255,7 +258,9 @@ pub fn build_event_graph(events: &[Event]) -> EventGraph {
|
||||
for event in events {
|
||||
for tag in event.tags.iter() {
|
||||
let tag_vec = tag.as_slice();
|
||||
if tag_vec.len() >= 2 && tag_vec[0] == "a" {
|
||||
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);
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::collections::{HashMap, HashSet, VecDeque};
|
||||
|
||||
use nostr_relay_builder::prelude::{
|
||||
Alphabet, Event, EventId, Filter, Kind, PublicKey, SingleLetterTag, Timestamp,
|
||||
@@ -7,50 +7,90 @@ use nostr_relay_builder::prelude::{
|
||||
use super::super::policy::{
|
||||
announcement_address_of, identifier_from_event, owner_directory_component, DeletionPolicy,
|
||||
};
|
||||
use super::{
|
||||
build_event_graph, reevaluate_events_without_announcement, traverse_and_mark_deletions,
|
||||
RetentionReason,
|
||||
};
|
||||
use crate::nostr::lifecycle::{DeletionSource, HoldingMetadata};
|
||||
use crate::nostr::lifecycle::{DeletionSource, HoldingMetadata, HOLDING_METADATA_KIND};
|
||||
use crate::nostr::SharedDatabase;
|
||||
|
||||
/// NIP-34 event kinds (other than the kind-30617 announcement and kind-30618
|
||||
/// state event) that participate in the repository dependency graph. These are
|
||||
/// the events that can become orphaned when an announcement is deleted.
|
||||
const GRAPH_DEPENDENT_KINDS: &[u16] = &[
|
||||
1111, // NIP-22 comment
|
||||
1617, // patch
|
||||
1618, // pull request
|
||||
1619, // PR update
|
||||
1621, // issue
|
||||
1630, // issue status
|
||||
1631, // PR status
|
||||
1632, // patch status
|
||||
1633, // repository status
|
||||
];
|
||||
use crate::nostr::lifecycle::history::HISTORY_METADATA_KIND;
|
||||
|
||||
const MAX_CASCADE_CANDIDATE_EVENTS: usize = 50_000;
|
||||
const MAX_CASCADE_GRAPH_EDGES: usize = 250_000;
|
||||
const MAX_CASCADE_ORPHAN_DELETES: usize = 20_000;
|
||||
const CASCADE_QUERY_CHUNK_SIZE: usize = 100;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||
enum RefKey {
|
||||
Event(EventId),
|
||||
Address(String),
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
struct CascadeGraph {
|
||||
nodes: HashMap<EventId, Event>,
|
||||
event_to_refs: HashMap<EventId, HashSet<RefKey>>,
|
||||
ref_to_events: HashMap<RefKey, HashSet<EventId>>,
|
||||
}
|
||||
|
||||
impl CascadeGraph {
|
||||
fn add_event(&mut self, event: Event) {
|
||||
let event_id = event.id;
|
||||
let refs = extract_accepted_reference_keys(&event);
|
||||
|
||||
self.event_to_refs.insert(event_id, refs.clone());
|
||||
for ref_key in refs {
|
||||
self.ref_to_events
|
||||
.entry(ref_key)
|
||||
.or_default()
|
||||
.insert(event_id);
|
||||
}
|
||||
|
||||
if let Some(address) = event_address(&event) {
|
||||
self.ref_to_events
|
||||
.entry(RefKey::Address(address))
|
||||
.or_default()
|
||||
.insert(event_id);
|
||||
}
|
||||
self.ref_to_events
|
||||
.entry(RefKey::Event(event_id))
|
||||
.or_default()
|
||||
.insert(event_id);
|
||||
|
||||
self.nodes.insert(event_id, event);
|
||||
}
|
||||
|
||||
fn edge_count(&self) -> usize {
|
||||
self.event_to_refs.values().map(HashSet::len).sum()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
enum CascadeExpansionError {
|
||||
Query(String),
|
||||
NodeCap { max: usize },
|
||||
EdgeCap { max: usize },
|
||||
}
|
||||
|
||||
impl std::fmt::Display for CascadeExpansionError {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
match self {
|
||||
Self::Query(e) => write!(f, "{e}"),
|
||||
Self::NodeCap { max } => write!(f, "cascade graph node cap exceeded ({max})"),
|
||||
Self::EdgeCap { max } => write!(f, "cascade graph edge cap exceeded ({max})"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl DeletionPolicy {
|
||||
/// Graph-based cascade deletion of a repository announcement.
|
||||
///
|
||||
/// Mirrors the algorithm order of the multi-maintainer deletion design:
|
||||
///
|
||||
/// 1. Query the candidate event set (the announcement's siblings + all
|
||||
/// NIP-34 dependent events) and build the dependency graph via
|
||||
/// [`crate::nostr::lifecycle::deletion::build_event_graph`] (a/e/q tag edges).
|
||||
/// 2. Re-evaluate the candidate events WITHOUT the deleted announcement but
|
||||
/// WITH every OTHER maintainer's announcement for the same identifier via
|
||||
/// [`crate::nostr::lifecycle::deletion::reevaluate_events_without_announcement`], yielding the
|
||||
/// kept/deleted split (multi-maintainer retention).
|
||||
/// 3. Add every surviving announcement (a different identifier, or the same
|
||||
/// identifier authored by a different maintainer) to the kept set so the
|
||||
/// traversal keeps everything still anchored to a live announcement.
|
||||
/// 4. Traverse the graph from the kept announcements via
|
||||
/// [`crate::nostr::lifecycle::deletion::traverse_and_mark_deletions`] to compute the final orphan
|
||||
/// set (events unreachable from any surviving announcement).
|
||||
/// 5. Hard-delete the orphaned event ids plus the announcement coordinate
|
||||
/// itself from the main DB.
|
||||
/// Mirrors [`crate::nostr::policy::related::RelatedEventPolicy`] by loading
|
||||
/// the affected main-DB reference component from the deleted announcement's
|
||||
/// event/address, following both outgoing accepted-reference tags and
|
||||
/// incoming tag-index matches. The retained set is then recomputed as a
|
||||
/// fixed point after excluding the deleted announcement: surviving
|
||||
/// independently accepted anchors are kept, any node referencing a kept node
|
||||
/// is kept, and any node referenced by a kept node is kept. Nodes in the
|
||||
/// loaded component without a surviving acceptance path are moved to holding
|
||||
/// and hard-deleted from the main DB.
|
||||
///
|
||||
/// Holding orchestration is active in this cascade path via
|
||||
/// [`Self::archive_and_delete_filter`]: every deleted orphan/coordinate
|
||||
@@ -74,127 +114,67 @@ impl DeletionPolicy {
|
||||
}
|
||||
let identifier = parts[2];
|
||||
|
||||
// Step 1: candidate event set — all announcements (to know which other
|
||||
// maintainers survive) plus all NIP-34 dependent events.
|
||||
let candidates = match self.query_cascade_candidates().await {
|
||||
Ok(events) => events,
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "Cascade deletion: failed to query candidate events; falling back to simple coordinate deletion");
|
||||
self.delete_coordinate_from_main_db(
|
||||
author,
|
||||
announcement_addr,
|
||||
deletion_created_at,
|
||||
&mut moved_ids,
|
||||
source,
|
||||
)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
if candidates.len() > MAX_CASCADE_CANDIDATE_EVENTS {
|
||||
tracing::warn!(
|
||||
announcement = %announcement_addr,
|
||||
candidates = candidates.len(),
|
||||
max = MAX_CASCADE_CANDIDATE_EVENTS,
|
||||
"Cascade deletion candidate set exceeds limit; falling back to coordinate-only deletion"
|
||||
);
|
||||
self.delete_coordinate_from_main_db(
|
||||
author,
|
||||
announcement_addr,
|
||||
deletion_created_at,
|
||||
&mut moved_ids,
|
||||
source,
|
||||
)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
|
||||
let graph = build_event_graph(&candidates);
|
||||
|
||||
// The set of candidate event ids we feed into re-evaluation. The
|
||||
// re-evaluation engine queries each one back from the DB; passing the
|
||||
// full candidate set lets it decide retention for every dependent.
|
||||
let candidate_ids: HashSet<EventId> = candidates.iter().map(|e| e.id).collect();
|
||||
|
||||
// Step 2: re-evaluate dependents WITHOUT the deleted announcement but
|
||||
// WITH any other surviving maintainers' announcements.
|
||||
let reeval = match reevaluate_events_without_announcement(
|
||||
&self.ctx.database,
|
||||
announcement_addr,
|
||||
&candidate_ids,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(result) => result,
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "Cascade deletion: re-evaluation failed; falling back to simple coordinate deletion");
|
||||
self.delete_coordinate_from_main_db(
|
||||
author,
|
||||
announcement_addr,
|
||||
deletion_created_at,
|
||||
&mut moved_ids,
|
||||
source,
|
||||
)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let mut kept_events: HashMap<EventId, RetentionReason> = reeval.keep;
|
||||
|
||||
// Step 3: any announcement that is NOT the one being deleted is a
|
||||
// surviving anchor. Add it to the kept set (with a Repository
|
||||
// Announcement retention reason) so the traversal starts from it. This
|
||||
// covers both "different identifier" and "same identifier, different
|
||||
// maintainer" survivors.
|
||||
for event in &candidates {
|
||||
if event.kind != Kind::GitRepoAnnouncement {
|
||||
continue;
|
||||
}
|
||||
let event_addr = announcement_address_of(event);
|
||||
if event_addr == announcement_addr {
|
||||
// This is the announcement being deleted — do not keep it.
|
||||
continue;
|
||||
}
|
||||
kept_events
|
||||
.entry(event.id)
|
||||
.or_insert_with(|| RetentionReason {
|
||||
event_id: event.id,
|
||||
valid_announcements: vec![event_addr],
|
||||
event_type: "Repository Announcement".to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
// Step 4: traverse from kept announcements to find the final orphans.
|
||||
let deleted_announcement_id = candidates
|
||||
.iter()
|
||||
.find(|e| {
|
||||
e.kind == Kind::GitRepoAnnouncement
|
||||
&& announcement_address_of(e) == announcement_addr
|
||||
let deleted_announcement_id = self
|
||||
.query_address_events(announcement_addr)
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|events| {
|
||||
events
|
||||
.into_iter()
|
||||
.find(|e| e.kind == Kind::GitRepoAnnouncement)
|
||||
})
|
||||
.map(|e| e.id)
|
||||
.map(|event| event.id)
|
||||
.unwrap_or_else(EventId::all_zeros);
|
||||
|
||||
let traversal = traverse_and_mark_deletions(&graph, &kept_events, deleted_announcement_id);
|
||||
let mut seed_event_ids = HashSet::new();
|
||||
if deleted_announcement_id != EventId::all_zeros() {
|
||||
seed_event_ids.insert(deleted_announcement_id);
|
||||
}
|
||||
let seed_addresses = HashSet::from([announcement_addr.to_string()]);
|
||||
|
||||
let graph = match self
|
||||
.expand_cascade_graph(seed_event_ids, seed_addresses)
|
||||
.await
|
||||
{
|
||||
Ok(graph) => graph,
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "Cascade deletion: graph expansion failed; falling back to simple coordinate deletion");
|
||||
self.delete_coordinate_from_main_db(
|
||||
author,
|
||||
announcement_addr,
|
||||
deletion_created_at,
|
||||
&mut moved_ids,
|
||||
source,
|
||||
)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let kept_events = compute_retained_events(
|
||||
&graph,
|
||||
announcement_addr,
|
||||
deleted_announcement_id,
|
||||
&self.ctx.config.domain,
|
||||
);
|
||||
let delete_events = compute_deletable_events(
|
||||
&graph,
|
||||
&kept_events,
|
||||
announcement_addr,
|
||||
deleted_announcement_id,
|
||||
);
|
||||
|
||||
tracing::info!(
|
||||
announcement = %announcement_addr,
|
||||
kept = traversal.keep.len(),
|
||||
deleted = traversal.delete.len(),
|
||||
circular = traversal.circular_dependencies.len(),
|
||||
"Announcement cascade deletion: graph traversal complete"
|
||||
component_nodes = graph.nodes.len(),
|
||||
component_edges = graph.edge_count(),
|
||||
kept = kept_events.len(),
|
||||
deleted = delete_events.len(),
|
||||
"Announcement cascade deletion: acceptance reachability complete"
|
||||
);
|
||||
for circular in &traversal.circular_dependencies {
|
||||
tracing::debug!(
|
||||
events = ?circular.events,
|
||||
anchored = circular.anchored,
|
||||
"Cascade deletion: circular dependency"
|
||||
);
|
||||
}
|
||||
|
||||
// Step 5a: hard-delete the orphaned dependent events by id.
|
||||
let orphan_ids: Vec<EventId> = traversal.delete.into_iter().collect();
|
||||
// Hard-delete the orphaned dependent events by id.
|
||||
let orphan_ids: Vec<EventId> = delete_events.into_iter().collect();
|
||||
if orphan_ids.len() > MAX_CASCADE_ORPHAN_DELETES {
|
||||
tracing::warn!(
|
||||
announcement = %announcement_addr,
|
||||
@@ -391,69 +371,470 @@ impl DeletionPolicy {
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Query the candidate event set for the cascade graph: every repository
|
||||
/// announcement (so surviving maintainers can be identified) plus every
|
||||
/// NIP-34 dependent event kind.
|
||||
async fn query_cascade_candidates(&self) -> Result<Vec<Event>, String> {
|
||||
let mut events = Vec::new();
|
||||
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
|
||||
}
|
||||
|
||||
// All announcements.
|
||||
let ann_filter = Filter::new().kind(Kind::GitRepoAnnouncement);
|
||||
events.extend(
|
||||
self.ctx
|
||||
.database
|
||||
.query(ann_filter)
|
||||
.await
|
||||
.map_err(|e| format!("query announcements failed: {e}"))?,
|
||||
);
|
||||
|
||||
// All dependent kinds.
|
||||
for kind_num in GRAPH_DEPENDENT_KINDS {
|
||||
let filter = Filter::new().kind(Kind::from(*kind_num));
|
||||
events.extend(
|
||||
self.ctx
|
||||
.database
|
||||
.query(filter)
|
||||
.await
|
||||
.map_err(|e| format!("query kind {kind_num} failed: {e}"))?,
|
||||
);
|
||||
}
|
||||
|
||||
// Exclude GRASP-06 PR / PR-update subgraphs from cascade deletion.
|
||||
//
|
||||
// These events can intentionally target this relay's `/prs/` endpoint
|
||||
// via clone tags and should not be treated as regular announcement-
|
||||
// anchored dependents for announcement cascade deletion.
|
||||
let excluded_roots: HashSet<EventId> = events
|
||||
.iter()
|
||||
.filter(|event| {
|
||||
crate::grasp06::policy::event_names_relays_prs_endpoint(
|
||||
event,
|
||||
&self.ctx.config.domain,
|
||||
)
|
||||
})
|
||||
.map(|event| event.id)
|
||||
.collect();
|
||||
|
||||
if !excluded_roots.is_empty() {
|
||||
let graph = build_event_graph(&events);
|
||||
let mut excluded_ids = excluded_roots.clone();
|
||||
for root in &excluded_roots {
|
||||
excluded_ids.extend(graph.get_transitive_dependents(root));
|
||||
}
|
||||
|
||||
let total_before = events.len();
|
||||
events.retain(|event| !excluded_ids.contains(&event.id));
|
||||
|
||||
tracing::debug!(
|
||||
total_before,
|
||||
excluded_roots = excluded_roots.len(),
|
||||
excluded_related = excluded_ids.len().saturating_sub(excluded_roots.len()),
|
||||
total_after = events.len(),
|
||||
"Cascade deletion: excluded GRASP-06 PR candidate subgraphs"
|
||||
);
|
||||
}
|
||||
|
||||
Ok(events)
|
||||
async fn query_address_events(&self, address: &str) -> Result<Vec<Event>, String> {
|
||||
query_address_events(&self.ctx.database, address).await
|
||||
}
|
||||
}
|
||||
|
||||
async fn expand_cascade_graph_from_db(
|
||||
database: &SharedDatabase,
|
||||
seed_event_ids: HashSet<EventId>,
|
||||
seed_addresses: HashSet<String>,
|
||||
) -> Result<CascadeGraph, CascadeExpansionError> {
|
||||
let mut graph = CascadeGraph::default();
|
||||
let mut event_queue: VecDeque<EventId> = seed_event_ids.into_iter().collect();
|
||||
let mut address_queue: VecDeque<String> = seed_addresses.into_iter().collect();
|
||||
let mut visited_event_ids = HashSet::new();
|
||||
let mut visited_addresses = HashSet::new();
|
||||
let mut queried_incoming_event_refs = HashSet::new();
|
||||
let mut queried_incoming_address_refs = HashSet::new();
|
||||
|
||||
while !event_queue.is_empty() || !address_queue.is_empty() {
|
||||
let mut loaded = Vec::new();
|
||||
|
||||
let mut id_batch = Vec::new();
|
||||
while id_batch.len() < CASCADE_QUERY_CHUNK_SIZE {
|
||||
let Some(event_id) = event_queue.pop_front() else {
|
||||
break;
|
||||
};
|
||||
if visited_event_ids.insert(event_id) {
|
||||
id_batch.push(event_id);
|
||||
}
|
||||
}
|
||||
if !id_batch.is_empty() {
|
||||
loaded.extend(query_events_by_ids(database, &id_batch).await?);
|
||||
for event_id in id_batch {
|
||||
if queried_incoming_event_refs.insert(event_id) {
|
||||
let incoming = query_incoming_event_refs(database, event_id).await?;
|
||||
enqueue_loaded_events(incoming, &mut loaded, &mut event_queue);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut address_batch = Vec::new();
|
||||
while address_batch.len() < CASCADE_QUERY_CHUNK_SIZE {
|
||||
let Some(address) = address_queue.pop_front() else {
|
||||
break;
|
||||
};
|
||||
if visited_addresses.insert(address.clone()) {
|
||||
address_batch.push(address);
|
||||
}
|
||||
}
|
||||
for address in address_batch {
|
||||
loaded.extend(
|
||||
query_address_events(database, &address)
|
||||
.await
|
||||
.map_err(CascadeExpansionError::Query)?,
|
||||
);
|
||||
if queried_incoming_address_refs.insert(address.clone()) {
|
||||
let incoming = query_incoming_address_refs(database, &address).await?;
|
||||
enqueue_loaded_events(incoming, &mut loaded, &mut event_queue);
|
||||
}
|
||||
}
|
||||
|
||||
for event in loaded {
|
||||
if should_exclude_from_generic_graph(&event) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let is_new = !graph.nodes.contains_key(&event.id);
|
||||
graph.add_event(event.clone());
|
||||
if graph.nodes.len() > MAX_CASCADE_CANDIDATE_EVENTS {
|
||||
return Err(CascadeExpansionError::NodeCap {
|
||||
max: MAX_CASCADE_CANDIDATE_EVENTS,
|
||||
});
|
||||
}
|
||||
if graph.edge_count() > MAX_CASCADE_GRAPH_EDGES {
|
||||
return Err(CascadeExpansionError::EdgeCap {
|
||||
max: MAX_CASCADE_GRAPH_EDGES,
|
||||
});
|
||||
}
|
||||
|
||||
if !is_new || should_not_recurse_through(&event) {
|
||||
continue;
|
||||
}
|
||||
|
||||
for ref_key in extract_accepted_reference_keys(&event) {
|
||||
match ref_key {
|
||||
RefKey::Event(event_id) => {
|
||||
if !visited_event_ids.contains(&event_id) {
|
||||
event_queue.push_back(event_id);
|
||||
}
|
||||
}
|
||||
RefKey::Address(address) => {
|
||||
if !visited_addresses.contains(&address) {
|
||||
address_queue.push_back(address);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if queried_incoming_event_refs.insert(event.id) {
|
||||
let incoming = query_incoming_event_refs(database, event.id).await?;
|
||||
enqueue_loaded_events(incoming, &mut Vec::new(), &mut event_queue);
|
||||
}
|
||||
if let Some(address) = event_address(&event) {
|
||||
if queried_incoming_address_refs.insert(address.clone()) {
|
||||
let incoming = query_incoming_address_refs(database, &address).await?;
|
||||
enqueue_loaded_events(incoming, &mut Vec::new(), &mut event_queue);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(graph)
|
||||
}
|
||||
|
||||
fn enqueue_loaded_events(
|
||||
events: Vec<Event>,
|
||||
loaded: &mut Vec<Event>,
|
||||
event_queue: &mut VecDeque<EventId>,
|
||||
) {
|
||||
for event in events {
|
||||
event_queue.push_back(event.id);
|
||||
loaded.push(event);
|
||||
}
|
||||
}
|
||||
|
||||
async fn query_events_by_ids(
|
||||
database: &SharedDatabase,
|
||||
ids: &[EventId],
|
||||
) -> Result<Vec<Event>, CascadeExpansionError> {
|
||||
if ids.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
database
|
||||
.query(Filter::new().ids(ids.iter().copied()))
|
||||
.await
|
||||
.map(|events| events.into_iter().collect())
|
||||
.map_err(|e| CascadeExpansionError::Query(format!("query ids failed: {e}")))
|
||||
}
|
||||
|
||||
async fn query_address_events(
|
||||
database: &SharedDatabase,
|
||||
address: &str,
|
||||
) -> Result<Vec<Event>, String> {
|
||||
let parts: Vec<&str> = address.splitn(3, ':').collect();
|
||||
if parts.len() < 2 {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
let kind_num = parts[0]
|
||||
.parse::<u16>()
|
||||
.map_err(|e| format!("invalid address kind {address}: {e}"))?;
|
||||
let pubkey = PublicKey::from_hex(parts[1])
|
||||
.map_err(|e| format!("invalid address pubkey {address}: {e}"))?;
|
||||
|
||||
let kind = Kind::from(kind_num);
|
||||
let mut filter = Filter::new().kind(kind).author(pubkey);
|
||||
if parts.len() == 3 && kind.is_addressable() {
|
||||
filter = filter.custom_tag(
|
||||
SingleLetterTag::lowercase(Alphabet::D),
|
||||
parts[2].to_string(),
|
||||
);
|
||||
}
|
||||
|
||||
database
|
||||
.query(filter)
|
||||
.await
|
||||
.map(|events| events.into_iter().collect())
|
||||
.map_err(|e| format!("query address {address} failed: {e}"))
|
||||
}
|
||||
|
||||
async fn query_incoming_event_refs(
|
||||
database: &SharedDatabase,
|
||||
event_id: EventId,
|
||||
) -> Result<Vec<Event>, CascadeExpansionError> {
|
||||
let hex = event_id.to_hex();
|
||||
query_tag_variants(
|
||||
database,
|
||||
&[
|
||||
SingleLetterTag::lowercase(Alphabet::E),
|
||||
SingleLetterTag::uppercase(Alphabet::E),
|
||||
SingleLetterTag::lowercase(Alphabet::Q),
|
||||
],
|
||||
&hex,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn query_incoming_address_refs(
|
||||
database: &SharedDatabase,
|
||||
address: &str,
|
||||
) -> Result<Vec<Event>, CascadeExpansionError> {
|
||||
query_tag_variants(
|
||||
database,
|
||||
&[
|
||||
SingleLetterTag::lowercase(Alphabet::A),
|
||||
SingleLetterTag::uppercase(Alphabet::A),
|
||||
SingleLetterTag::lowercase(Alphabet::Q),
|
||||
],
|
||||
address,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn query_tag_variants(
|
||||
database: &SharedDatabase,
|
||||
tags: &[SingleLetterTag],
|
||||
value: &str,
|
||||
) -> Result<Vec<Event>, CascadeExpansionError> {
|
||||
let mut by_id = HashMap::new();
|
||||
for tag in tags {
|
||||
let events = database
|
||||
.query(Filter::new().custom_tag(*tag, value.to_string()))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
CascadeExpansionError::Query(format!("query tag #{tag}={value} failed: {e}"))
|
||||
})?;
|
||||
for event in events {
|
||||
by_id.insert(event.id, event);
|
||||
}
|
||||
}
|
||||
Ok(by_id.into_values().collect())
|
||||
}
|
||||
|
||||
fn compute_retained_events(
|
||||
graph: &CascadeGraph,
|
||||
deleted_address: &str,
|
||||
deleted_announcement_id: EventId,
|
||||
domain: &str,
|
||||
) -> HashSet<EventId> {
|
||||
let mut kept = HashSet::new();
|
||||
let mut vetoed_git_nodes = HashSet::new();
|
||||
|
||||
for event in graph.nodes.values() {
|
||||
if event.id == deleted_announcement_id
|
||||
|| event_address(event).as_deref() == Some(deleted_address)
|
||||
|| vetoed_git_nodes.contains(&event.id)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if is_independent_anchor(event, deleted_address, domain) {
|
||||
kept.insert(event.id);
|
||||
}
|
||||
}
|
||||
|
||||
loop {
|
||||
let before = kept.len();
|
||||
let kept_refs = refs_for_kept_events(graph, &kept);
|
||||
|
||||
for event in graph.nodes.values() {
|
||||
if kept.contains(&event.id)
|
||||
|| vetoed_git_nodes.contains(&event.id)
|
||||
|| event.id == deleted_announcement_id
|
||||
|| event_address(event).as_deref() == Some(deleted_address)
|
||||
|| should_exclude_from_generic_graph(event)
|
||||
|| should_not_recurse_through(event)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
let own_refs = graph
|
||||
.event_to_refs
|
||||
.get(&event.id)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let own_keys = event_identity_keys(event);
|
||||
|
||||
if own_refs
|
||||
.iter()
|
||||
.any(|ref_key| kept_refs.identities.contains(ref_key))
|
||||
|| own_keys
|
||||
.iter()
|
||||
.any(|ref_key| kept_refs.outgoing.contains(ref_key))
|
||||
{
|
||||
kept.insert(event.id);
|
||||
}
|
||||
}
|
||||
|
||||
let vetoed = veto_invalid_git_invariant_nodes(graph, &kept, domain);
|
||||
if !vetoed.is_empty() {
|
||||
for id in vetoed {
|
||||
kept.remove(&id);
|
||||
vetoed_git_nodes.insert(id);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
if kept.len() == before {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
kept
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
struct KeptRefSets {
|
||||
identities: HashSet<RefKey>,
|
||||
outgoing: HashSet<RefKey>,
|
||||
}
|
||||
|
||||
fn refs_for_kept_events(graph: &CascadeGraph, kept: &HashSet<EventId>) -> KeptRefSets {
|
||||
let mut refs = KeptRefSets::default();
|
||||
for event_id in kept {
|
||||
if let Some(event) = graph.nodes.get(event_id) {
|
||||
refs.identities.extend(event_identity_keys(event));
|
||||
}
|
||||
if let Some(event_refs) = graph.event_to_refs.get(event_id) {
|
||||
refs.outgoing.extend(event_refs.iter().cloned());
|
||||
}
|
||||
}
|
||||
refs
|
||||
}
|
||||
|
||||
fn compute_deletable_events(
|
||||
graph: &CascadeGraph,
|
||||
kept: &HashSet<EventId>,
|
||||
deleted_address: &str,
|
||||
deleted_announcement_id: EventId,
|
||||
) -> HashSet<EventId> {
|
||||
graph
|
||||
.nodes
|
||||
.values()
|
||||
.filter(|event| {
|
||||
!kept.contains(&event.id)
|
||||
&& event.id != deleted_announcement_id
|
||||
&& event_address(event).as_deref() != Some(deleted_address)
|
||||
&& !should_exclude_from_generic_graph(event)
|
||||
&& !should_not_recurse_through(event)
|
||||
})
|
||||
.map(|event| event.id)
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn is_independent_anchor(event: &Event, deleted_address: &str, domain: &str) -> bool {
|
||||
if event.kind == Kind::GitRepoAnnouncement {
|
||||
return event_address(event).as_deref() != Some(deleted_address);
|
||||
}
|
||||
|
||||
crate::grasp06::policy::event_names_relays_prs_endpoint(event, domain)
|
||||
}
|
||||
|
||||
fn veto_invalid_git_invariant_nodes(
|
||||
graph: &CascadeGraph,
|
||||
kept: &HashSet<EventId>,
|
||||
domain: &str,
|
||||
) -> HashSet<EventId> {
|
||||
graph
|
||||
.nodes
|
||||
.values()
|
||||
.filter(|event| kept.contains(&event.id))
|
||||
.filter(|event| {
|
||||
matches!(
|
||||
event.kind,
|
||||
Kind::GitPullRequest | Kind::GitPullRequestUpdate
|
||||
)
|
||||
})
|
||||
.filter(|event| !crate::grasp06::policy::event_names_relays_prs_endpoint(event, domain))
|
||||
.filter(|event| !event_references_kept_repo_announcement(event, graph, kept))
|
||||
.map(|event| event.id)
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn event_references_kept_repo_announcement(
|
||||
event: &Event,
|
||||
graph: &CascadeGraph,
|
||||
kept: &HashSet<EventId>,
|
||||
) -> bool {
|
||||
graph
|
||||
.event_to_refs
|
||||
.get(&event.id)
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.any(|ref_key| {
|
||||
graph.ref_to_events.get(ref_key).is_some_and(|ids| {
|
||||
ids.iter().any(|id| {
|
||||
kept.contains(id)
|
||||
&& graph
|
||||
.nodes
|
||||
.get(id)
|
||||
.is_some_and(|node| node.kind == Kind::GitRepoAnnouncement)
|
||||
})
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
fn extract_accepted_reference_keys(event: &Event) -> HashSet<RefKey> {
|
||||
let mut refs = HashSet::new();
|
||||
for tag in event.tags.iter() {
|
||||
let tag_vec = tag.as_slice();
|
||||
if tag_vec.len() < 2 {
|
||||
continue;
|
||||
}
|
||||
match tag_vec[0].as_str() {
|
||||
"a" | "A" => {
|
||||
refs.insert(RefKey::Address(tag_vec[1].to_string()));
|
||||
}
|
||||
"e" | "E" => {
|
||||
if let Ok(event_id) = EventId::from_hex(&tag_vec[1]) {
|
||||
refs.insert(RefKey::Event(event_id));
|
||||
}
|
||||
}
|
||||
"q" if tag_vec[1].contains(':') => {
|
||||
refs.insert(RefKey::Address(tag_vec[1].to_string()));
|
||||
}
|
||||
"q" => {
|
||||
if let Ok(event_id) = EventId::from_hex(&tag_vec[1]) {
|
||||
refs.insert(RefKey::Event(event_id));
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
refs
|
||||
}
|
||||
|
||||
fn event_identity_keys(event: &Event) -> HashSet<RefKey> {
|
||||
let mut keys = HashSet::from([RefKey::Event(event.id)]);
|
||||
if let Some(address) = event_address(event) {
|
||||
keys.insert(RefKey::Address(address));
|
||||
}
|
||||
keys
|
||||
}
|
||||
|
||||
fn event_address(event: &Event) -> Option<String> {
|
||||
if event.kind == Kind::GitRepoAnnouncement {
|
||||
return Some(announcement_address_of(event));
|
||||
}
|
||||
|
||||
if event.kind.is_addressable() {
|
||||
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].clone())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})?;
|
||||
Some(format!(
|
||||
"{}:{}:{}",
|
||||
event.kind.as_u16(),
|
||||
event.pubkey.to_hex(),
|
||||
identifier
|
||||
))
|
||||
} else if event.kind.is_replaceable() {
|
||||
Some(format!("{}:{}", event.kind.as_u16(), event.pubkey.to_hex()))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
fn should_exclude_from_generic_graph(event: &Event) -> bool {
|
||||
event.kind == Kind::RepoState
|
||||
}
|
||||
|
||||
fn should_not_recurse_through(event: &Event) -> bool {
|
||||
event.kind == Kind::EventDeletion
|
||||
|| event.kind == Kind::from(62)
|
||||
|| event.kind == Kind::from(HOLDING_METADATA_KIND)
|
||||
|| event.kind == Kind::from(HISTORY_METADATA_KIND)
|
||||
}
|
||||
|
||||
@@ -573,7 +573,11 @@ async fn test_cascade_recurses_to_unknown_kind_via_comment_only_edge() {
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||
|
||||
for (label, id) in [("issue", issue.id), ("comment", comment.id), ("kind1", kind1.id)] {
|
||||
for (label, id) in [
|
||||
("issue", issue.id),
|
||||
("comment", comment.id),
|
||||
("kind1", kind1.id),
|
||||
] {
|
||||
assert!(
|
||||
!client
|
||||
.is_event_on_relay(id)
|
||||
@@ -598,10 +602,18 @@ async fn test_unknown_kind_leaf_referenced_only_by_deleted_chain_is_removed() {
|
||||
let coordinate = announcement_coordinate(&announcement, &repo_id);
|
||||
|
||||
let issue = client
|
||||
.create_issue(&announcement, "Issue for deleted chain", "issue body", vec![])
|
||||
.create_issue(
|
||||
&announcement,
|
||||
"Issue for deleted chain",
|
||||
"issue body",
|
||||
vec![],
|
||||
)
|
||||
.expect("build issue");
|
||||
let kind1 = client
|
||||
.event_builder(Kind::from(1), "kind 1 note accepted via deleted-chain comment")
|
||||
.event_builder(
|
||||
Kind::from(1),
|
||||
"kind 1 note accepted via deleted-chain comment",
|
||||
)
|
||||
.tag(Tag::custom("e", vec![issue.id.to_hex()]))
|
||||
.build(&keys)
|
||||
.expect("build kind 1 note");
|
||||
@@ -647,7 +659,11 @@ async fn test_unknown_kind_leaf_referenced_only_by_deleted_chain_is_removed() {
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||
|
||||
for (label, id) in [("issue", issue.id), ("comment", comment.id), ("kind1", kind1.id)] {
|
||||
for (label, id) in [
|
||||
("issue", issue.id),
|
||||
("comment", comment.id),
|
||||
("kind1", kind1.id),
|
||||
] {
|
||||
assert!(
|
||||
!client
|
||||
.is_event_on_relay(id)
|
||||
@@ -661,7 +677,7 @@ async fn test_unknown_kind_leaf_referenced_only_by_deleted_chain_is_removed() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_unknown_kind_shared_with_retained_repo_is_kept_while_deleted_repo_issue_is_removed() {
|
||||
async fn test_unknown_kind_shared_with_retained_repo_keeps_deleted_repo_issue_reference() {
|
||||
let relay = TestRelay::start().await;
|
||||
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
|
||||
.await
|
||||
@@ -673,12 +689,7 @@ async fn test_unknown_kind_shared_with_retained_repo_is_kept_while_deleted_repo_
|
||||
let (announcement_b, _repo_id_b) = publish_served_repo(&client, "cascade-shared-repo-b").await;
|
||||
|
||||
let issue_a = client
|
||||
.create_issue(
|
||||
&announcement_a,
|
||||
"Repo A issue",
|
||||
"issue body",
|
||||
vec![],
|
||||
)
|
||||
.create_issue(&announcement_a, "Repo A issue", "issue body", vec![])
|
||||
.expect("build repo A issue");
|
||||
client
|
||||
.send_event(issue_a.clone())
|
||||
@@ -750,11 +761,11 @@ async fn test_unknown_kind_shared_with_retained_repo_is_kept_while_deleted_repo_
|
||||
tokio::time::sleep(Duration::from_millis(700)).await;
|
||||
|
||||
assert!(
|
||||
!client
|
||||
client
|
||||
.is_event_on_relay(issue_a.id)
|
||||
.await
|
||||
.expect("query repo A issue after deletion"),
|
||||
"issue_A ({}) should be deleted with repo A",
|
||||
"issue_A ({}) should remain because retained kind1 references it",
|
||||
issue_a.id
|
||||
);
|
||||
assert!(
|
||||
@@ -785,9 +796,11 @@ async fn test_deleted_repo_issue_retained_when_its_comment_is_quoted_in_retained
|
||||
.expect("create audit client");
|
||||
let keys = client.keys().clone();
|
||||
|
||||
let (announcement_a, repo_id_a) = publish_served_repo(&client, "cascade-quoted-comment-repo-a").await;
|
||||
let (announcement_a, repo_id_a) =
|
||||
publish_served_repo(&client, "cascade-quoted-comment-repo-a").await;
|
||||
let coordinate_a = announcement_coordinate(&announcement_a, &repo_id_a);
|
||||
let (announcement_b, _repo_id_b) = publish_served_repo(&client, "cascade-quoted-comment-repo-b").await;
|
||||
let (announcement_b, _repo_id_b) =
|
||||
publish_served_repo(&client, "cascade-quoted-comment-repo-b").await;
|
||||
|
||||
let issue_a = client
|
||||
.create_issue(&announcement_a, "Repo A issue", "issue body", vec![])
|
||||
|
||||
@@ -570,6 +570,8 @@ async fn test_shared_subgraph_survives_when_one_anchor_path_removed() {
|
||||
.expect("build top unique A event");
|
||||
|
||||
// Shared descendant references both tops and should survive through top_shared.
|
||||
// Under acceptance-reachability cascade, retaining shared_descendant also
|
||||
// retains top_unique_a because a kept node references it.
|
||||
let shared_descendant = client_a
|
||||
.event_builder(Kind::from(1111), "shared descendant")
|
||||
.tag(Tag::custom("e", vec![top_shared.id.to_hex()]))
|
||||
@@ -577,7 +579,8 @@ async fn test_shared_subgraph_survives_when_one_anchor_path_removed() {
|
||||
.build(client_a.keys())
|
||||
.expect("build shared descendant");
|
||||
|
||||
// Unique branch descendant should be deleted with top_unique_a.
|
||||
// Once top_unique_a is retained by the shared descendant, this unique
|
||||
// descendant remains acceptable because it references a kept node.
|
||||
let unique_descendant = client_a
|
||||
.event_builder(Kind::from(1111), "unique descendant")
|
||||
.tag(Tag::custom("e", vec![top_unique_a.id.to_hex()]))
|
||||
@@ -660,15 +663,15 @@ async fn test_shared_subgraph_survives_when_one_anchor_path_removed() {
|
||||
"top shared event must survive because anchor B remains"
|
||||
);
|
||||
assert!(
|
||||
!top_unique_a_survived,
|
||||
"top unique A event must be deleted after A deletion"
|
||||
top_unique_a_survived,
|
||||
"top unique A event must survive because retained shared_descendant references it"
|
||||
);
|
||||
assert!(
|
||||
shared_descendant_survived,
|
||||
"shared descendant must survive through still-anchored top_shared path"
|
||||
);
|
||||
assert!(
|
||||
!unique_descendant_survived,
|
||||
"unique descendant must be deleted with the A-only branch"
|
||||
unique_descendant_survived,
|
||||
"unique descendant must survive because retained top_unique_A remains acceptable"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user