diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index d0180b1..245fc94 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -596,12 +596,13 @@ The rejected events index solves two critical problems during sync: 1. **Negentropy sync efficiency**: Prevents repeatedly downloading events that will be rejected again 2. **Race condition resolution**: Enables immediate hot-cache re-processing or an exact-ID fetch after the full event expires -**Two-Tier Architecture:** +**Bounded Architecture:** | Tier | Duration | Storage | Purpose | |------|----------|---------|---------| | Hot Cache | 2 minutes | Full events | Immediate re-processing when dependencies arrive | | Cold Index | 7 days | Metadata only | Prevent re-fetch during negentropy sync | +| Related Dependency Index | 7 days | Full events + reference keys + relay hints | Retry policy orphans when either side of their relationship is accepted | **Event Flow:** @@ -627,6 +628,13 @@ Exact-ID Recovery Negentropy Sync │ └──▶ Exclude Cold Index IDs from "missing events" calculation + +Related Event Rejected as an Orphan + │ + ├──▶ Retain the full event in the crash-safe checkpoint + ├──▶ Match later accepted IDs and addresses in both directions + ├──▶ Re-process an iterative closure of newly unblocked events + └──▶ Keep retryable results pending; remove terminal outcomes ``` Exact-ID recovery runs as bounded background work, so a slow relay does not hold @@ -639,6 +647,9 @@ timeout. - Repository announcements (kind 30617) rejected for not listing this service or a dependency-resolvable maintainer validation failure - State events (kind 30618) rejected for missing announcements or dependency-resolvable authorization failures - Permanently invalid events classified as `Other` remain excluded and are not exact-refetched +- Related comments, reactions, zap requests, and other policy orphans are + retained under fixed limits: 1,024 entries, 8 MiB total, 128 KiB per event, + and 512 attempts per triggered closure **Source Code:** [`src/sync/rejected_index.rs`](../../src/sync/rejected_index.rs) diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index ca536f1..dc32fd1 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -960,6 +960,23 @@ This prevents both failure modes: broad synchronization does not repeatedly download known-invalid events, while a dependency-resolvable event cannot become permanently suppressed merely because its full hot-cache copy expired. +Related events have a second dependency race: a comment, reaction, zap request, +or other repository-thread event can arrive before the event or repository +address that makes it admissible. Sync-originated `restricted` orphan results +are therefore retained as full events in the same crash-safe checkpoint for up +to seven days. Acceptance of either side of the relationship re-evaluates a +bounded iterative closure, covering both backward references (the orphan points +to the newly accepted event/address) and forward references (the newly accepted +event points to the orphan). Entries are removed only after a terminal policy +outcome. + +This durable tier is relay-input bounded: at most 1,024 events, 8 MiB total, +128 KiB per event, and 512 attempts in one triggered closure. Oldest entries +are evicted first. Evicted or oversized events are not marked as locally held, +so later broad synchronization may offer them again. Exact-ID hydration also +keeps dependency-pending cached IDs unresolved rather than misclassifying a +cached policy orphan as recovered. + See [Architecture: Rejected Events Index](architecture.md#rejected-events-index) and [`src/sync/rejected_index.rs`](../../src/sync/rejected_index.rs) for the design and implementation. diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 7331096..3e4d838 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -31,8 +31,6 @@ pub use metrics::SyncMetrics; // Re-export rejected index types pub use rejected_index::{EventType, RejectionReason}; -// Note: RejectedEventsIndex struct exists in rejected_index.rs but not yet used -// Current code still uses the simple HashSet type alias below // Re-export relay connection types pub use relay_connection::{ @@ -1772,12 +1770,14 @@ async fn run_rejected_index_cleanup( let (_, ann_cold_expired) = manager.rejected_events_index.cleanup_expired_for_type("announcement"); let (_, state_cold_expired) = manager.rejected_events_index.cleanup_expired_for_type("state"); let unrecoverable_expired = manager.rejected_events_index.cleanup_expired_unrecoverable(); + let related_expired = manager.rejected_events_index.cleanup_expired_related(); - if ann_cold_expired + state_cold_expired + unrecoverable_expired > 0 { + if ann_cold_expired + state_cold_expired + unrecoverable_expired + related_expired > 0 { tracing::info!( announcements = ann_cold_expired, states = state_cold_expired, unrecoverable = unrecoverable_expired, + related = related_expired, "Cleaned up expired entries from rejected events cold index" ); } @@ -1886,6 +1886,7 @@ async fn run_health_and_metrics_checker( ("queued_connection_attempts", manager.in_flight_connect_attempts.len()), ("purgatory_dependencies", manager.purgatory_dependency_attempts.len()), ("dependency_relays", manager.dependency_relay_deadlines.len()), + ("related_dependency_events", manager.rejected_events_index.related_len()), ("deferred_consolidations", manager.deferred_consolidations.relays.len()), ("descendant_rotations", manager.descendant_sync_rotations.len()), ] { @@ -2993,19 +2994,23 @@ impl SyncManager { if let Some(metrics) = metrics.as_ref() { metrics.record_hydration_events(&relay_url, "recovery", "delivered", 1); } - // Events already tracked as rejected count as recovered without - // re-processing: re-validation is owned by the rejected-index - // re-processing machinery, and unrecoverable IDs never revalidate. + // Permanent cached rejections account for the requested ID. + // Dependency-sensitive entries remain pending until their normal + // re-processing machinery observes the missing accepted event. if rejected_events_index.contains(&event.id) { + let dependency_pending = rejected_events_index.is_dependency_pending(&event.id); tracing::debug!( relay = %relay_url, event_id = %event.id, + dependency_pending, "Recovered missing event already tracked as rejected, skipping re-processing" ); if let Some(metrics) = metrics.as_ref() { metrics.record_hydration_events(&relay_url, "recovery", "rejected_cached", 1); } - recovered.insert(event.id); + if !dependency_pending { + recovered.insert(event.id); + } continue; } let result = Self::process_event_static( @@ -6364,6 +6369,30 @@ impl SyncManager { local_relay: &LocalRelay, rejected_events_index: &Arc, save_context: crate::nostr::persistence::SaveContext, + ) -> ProcessResult { + Self::process_event_static_inner( + event, + relay_url, + database, + write_policy, + local_relay, + rejected_events_index, + save_context, + true, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + async fn process_event_static_inner( + event: &Event, + relay_url: &str, + database: &SharedDatabase, + write_policy: &Nip34WritePolicy, + local_relay: &LocalRelay, + rejected_events_index: &Arc, + save_context: crate::nostr::persistence::SaveContext, + resolve_related_dependencies: bool, ) -> ProcessResult { use nostr_sdk::prelude::{WritePolicy, WritePolicyResult}; use std::net::{IpAddr, Ipv4Addr, SocketAddr}; @@ -6568,6 +6597,51 @@ impl SyncManager { } } + if resolve_related_dependencies { + let mut queue: VecDeque = rejected_events_index + .related_candidates_resolved_by(event) + .into_iter() + .collect(); + let mut attempted = HashSet::new(); + let mut saved = 0usize; + while let Some(candidate) = queue.pop_front() { + if attempted.len() >= rejected_index::RELATED_RETRY_CLOSURE_LIMIT + || !attempted.insert(candidate.id) + { + continue; + } + let outcome = Box::pin(Self::process_event_static_inner( + &candidate, + relay_url, + database, + write_policy, + local_relay, + rejected_events_index, + save_context, + false, + )) + .await; + if outcome.is_terminally_accounted() { + rejected_events_index.remove(&candidate.id); + } + if outcome == ProcessResult::Saved { + saved += 1; + queue.extend( + rejected_events_index.related_candidates_resolved_by(&candidate), + ); + } + } + if !attempted.is_empty() { + tracing::info!( + trigger_event = %event.id, + attempted = attempted.len(), + saved, + remaining = rejected_events_index.related_len(), + "Reprocessed dependency-sensitive related-event closure" + ); + } + } + ProcessResult::Saved } WritePolicyResult::Reject { @@ -6670,6 +6744,25 @@ impl SyncManager { "Synced event missing 'd' tag, tracked as unrecoverable by ID" ); } + } else if process_result == ProcessResult::Rejected(PolicyRejection::Restricted) + && message.contains( + "event must reference an accepted repository or accepted event", + ) + { + let (addressable_refs, event_refs) = + crate::nostr::policy::RelatedEventPolicy::extract_reference_tags(event); + let retained = rejected_events_index.add_related_from_relay( + event.clone(), + event_refs.into_iter().collect(), + addressable_refs.into_iter().collect(), + Some(relay_url.to_string()), + ); + tracing::debug!( + event_id = %event.id, + kind = %event.kind.as_u16(), + retained, + "Indexed dependency-sensitive related event for durable retry" + ); } process_result @@ -7978,6 +8071,82 @@ impl SyncManager { mod tests { use super::*; + #[tokio::test] + async fn accepted_dependency_reprocesses_a_synced_policy_orphan() { + let directory = tempfile::tempdir().expect("create test directory"); + let git_data_path = directory.path().join("git"); + let mut config = Config::for_testing(); + config.git_data_path = git_data_path.to_string_lossy().into_owned(); + config.relay_data_path = directory + .path() + .join("relay") + .to_string_lossy() + .into_owned(); + let purgatory = Arc::new(crate::purgatory::Purgatory::new(git_data_path)); + let runtime = crate::nostr::builder::create_relay( + &config, + purgatory, + crate::grasp06::receive::RepoInitLocks::default(), + ) + .await + .expect("create test relay runtime"); + let rejected = Arc::new(RejectedEventsIndex::new( + Duration::from_secs(120), + Duration::from_secs(604800), + )); + let keys = Keys::generate(); + let parent = EventBuilder::new(Kind::GitUserGraspList, "dependency") + .finalize(&keys) + .expect("build accepted dependency"); + let child = EventBuilder::new(Kind::TextNote, "dependent comment") + .tags([Tag::event(parent.id)]) + .finalize(&keys) + .expect("build dependent event"); + + let orphan_result = SyncManager::process_event_static( + &child, + "wss://source.example", + &runtime.stores.database, + &runtime.write_policy, + &runtime.relay, + &rejected, + crate::nostr::persistence::SaveContext::RelaySync, + ) + .await; + assert_eq!( + orphan_result, + ProcessResult::Rejected(PolicyRejection::Restricted) + ); + assert!(rejected.is_dependency_pending(&child.id)); + assert!(runtime + .stores + .database + .event_by_id(&child.id) + .await + .unwrap() + .is_none()); + + let parent_result = SyncManager::process_event_static( + &parent, + "wss://source.example", + &runtime.stores.database, + &runtime.write_policy, + &runtime.relay, + &rejected, + crate::nostr::persistence::SaveContext::RelaySync, + ) + .await; + assert_eq!(parent_result, ProcessResult::Saved); + assert!(runtime + .stores + .database + .event_by_id(&child.id) + .await + .unwrap() + .is_some()); + assert!(!rejected.contains(&child.id)); + } + #[test] fn partial_nip65_batch_retries_missing_authors_early() { let returned = Keys::generate().public_key(); diff --git a/src/sync/rejected_index.rs b/src/sync/rejected_index.rs index 1e21590..0b7da5b 100644 --- a/src/sync/rejected_index.rs +++ b/src/sync/rejected_index.rs @@ -94,6 +94,16 @@ use std::path::Path; use std::sync::{Arc, RwLock}; use std::time::{Duration, Instant, SystemTime}; +/// Hard bounds for dependency-sensitive related events retained across restarts. +/// +/// Full events are required for immediate re-processing, but relay-controlled +/// payloads must not turn dependency recovery into an unbounded memory or +/// checkpoint sink. Oversized events remain eligible for a later broad fetch. +const RELATED_MAX_ENTRIES: usize = 1_024; +const RELATED_MAX_SERIALIZED_BYTES: usize = 8 * 1024 * 1024; +const RELATED_MAX_EVENT_BYTES: usize = 128 * 1024; +pub(crate) const RELATED_RETRY_CLOSURE_LIMIT: usize = 512; + /// Type of event stored in the rejected events index #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] pub enum EventType { @@ -231,6 +241,182 @@ struct SerializableUnrecoverableIndex { entries: HashMap, } +/// A related event whose policy rejection can become valid when another event +/// or repository address is accepted. +#[derive(Debug, Clone)] +struct RelatedDependencyEntry { + event: Event, + event_refs: HashSet, + addressable_refs: HashSet, + relay_hints: HashSet, + rejected_at: Instant, + serialized_bytes: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SerializableRelatedDependencyEntry { + event: Event, + event_refs: HashSet, + addressable_refs: HashSet, + #[serde(default)] + relay_hints: HashSet, + rejected_at_offset_secs: u64, +} + +/// Bounded, durable full-event index for dependency-sensitive related events. +#[derive(Debug, Clone)] +struct RelatedDependencyIndex { + entries: Arc>>, + expiry_duration: Duration, +} + +impl RelatedDependencyIndex { + fn new(expiry_duration: Duration) -> Self { + Self { + entries: Arc::new(RwLock::new(HashMap::new())), + expiry_duration, + } + } + + fn event_coordinate(event: &Event) -> Option { + let kind = event.kind.as_u16(); + if (10_000..20_000).contains(&kind) { + return Some(format!("{}:{}", kind, event.pubkey.to_hex())); + } + if !(30_000..40_000).contains(&kind) { + return None; + } + let identifier = event + .tags + .iter() + .find(|tag| tag.kind() == "d") + .and_then(|tag| tag.content())?; + Some(format!("{}:{}:{}", kind, event.pubkey.to_hex(), identifier)) + } + + fn serialized_event_bytes(event: &Event) -> Option { + serde_json::to_vec(event).ok().map(|bytes| bytes.len()) + } + + fn add( + &self, + event: Event, + event_refs: HashSet, + addressable_refs: HashSet, + relay_url: Option, + ) -> bool { + let Some(serialized_bytes) = Self::serialized_event_bytes(&event) else { + return false; + }; + if serialized_bytes > RELATED_MAX_EVENT_BYTES { + return false; + } + + let mut entries = self.entries.write().unwrap(); + if let Some(existing) = entries.get_mut(&event.id) { + existing.relay_hints.extend(relay_url); + return true; + } + + entries.insert( + event.id, + RelatedDependencyEntry { + event, + event_refs, + addressable_refs, + relay_hints: relay_url.into_iter().collect(), + rejected_at: Instant::now(), + serialized_bytes, + }, + ); + Self::enforce_bounds(&mut entries); + true + } + + fn insert_restored(&self, entry: RelatedDependencyEntry) { + if entry.serialized_bytes > RELATED_MAX_EVENT_BYTES { + return; + } + let mut entries = self.entries.write().unwrap(); + entries.insert(entry.event.id, entry); + Self::enforce_bounds(&mut entries); + } + + fn enforce_bounds(entries: &mut HashMap) { + while entries.len() > RELATED_MAX_ENTRIES + || entries + .values() + .map(|entry| entry.serialized_bytes) + .sum::() + > RELATED_MAX_SERIALIZED_BYTES + { + let Some(oldest) = entries + .iter() + .min_by(|(left_id, left), (right_id, right)| { + left.rejected_at + .cmp(&right.rejected_at) + .then_with(|| left_id.cmp(right_id)) + }) + .map(|(event_id, _)| *event_id) + else { + break; + }; + entries.remove(&oldest); + } + } + + fn contains(&self, event_id: &EventId) -> bool { + self.entries + .read() + .unwrap() + .get(event_id) + .is_some_and(|entry| entry.rejected_at.elapsed() < self.expiry_duration) + } + + fn candidates_resolved_by(&self, accepted: &Event) -> Vec { + let (accepted_addressable_refs, accepted_event_refs) = + crate::nostr::policy::RelatedEventPolicy::extract_reference_tags(accepted); + let accepted_addressable_refs: HashSet<_> = accepted_addressable_refs.into_iter().collect(); + let accepted_event_refs: HashSet<_> = accepted_event_refs.into_iter().collect(); + let accepted_coordinate = Self::event_coordinate(accepted); + let now = Instant::now(); + let entries = self.entries.read().unwrap(); + let mut candidates: Vec<_> = entries + .values() + .filter(|entry| { + now.duration_since(entry.rejected_at) < self.expiry_duration + && (entry.event_refs.contains(&accepted.id) + || accepted_event_refs.contains(&entry.event.id) + || accepted_coordinate + .as_ref() + .is_some_and(|coordinate| entry.addressable_refs.contains(coordinate)) + || Self::event_coordinate(&entry.event).is_some_and(|coordinate| { + accepted_addressable_refs.contains(&coordinate) + })) + }) + .map(|entry| (entry.rejected_at, entry.event.id, entry.event.clone())) + .collect(); + candidates.sort_by(|left, right| left.0.cmp(&right.0).then_with(|| left.1.cmp(&right.1))); + candidates.into_iter().map(|(_, _, event)| event).collect() + } + + fn remove(&self, event_id: &EventId) { + self.entries.write().unwrap().remove(event_id); + } + + fn cleanup_expired(&self) -> usize { + let mut entries = self.entries.write().unwrap(); + let initial = entries.len(); + let now = Instant::now(); + entries.retain(|_, entry| now.duration_since(entry.rejected_at) < self.expiry_duration); + initial - entries.len() + } + + fn len(&self) -> usize { + self.entries.read().unwrap().len() + } +} + /// Complete rejected cache state for persistence /// /// Stores both hot cache and cold index with version and timestamp information. @@ -250,6 +436,9 @@ struct RejectedCacheState { /// Defaults to empty when restoring caches saved before this index existed. #[serde(default)] unrecoverable: SerializableUnrecoverableIndex, + /// Full related events awaiting an accepted reference dependency. + #[serde(default)] + related_dependencies: HashMap, } /// Hot cache: Stores full events for immediate re-processing @@ -636,6 +825,7 @@ pub struct RejectedEventsIndex { hot_cache: HotCache, cold_index: ColdIndex, unrecoverable: UnrecoverableIndex, + related_dependencies: RelatedDependencyIndex, metrics: Option, } @@ -646,6 +836,7 @@ impl std::fmt::Debug for RejectedEventsIndex { .field("hot_cache", &self.hot_cache) .field("cold_index", &self.cold_index) .field("unrecoverable", &self.unrecoverable) + .field("related_dependencies", &self.related_dependencies) .field("metrics", &self.metrics.is_some()) .finish() } @@ -663,6 +854,7 @@ impl RejectedEventsIndex { hot_cache: HotCache::new(hot_cache_duration), cold_index: ColdIndex::new(cold_index_duration), unrecoverable: UnrecoverableIndex::new(cold_index_duration), + related_dependencies: RelatedDependencyIndex::new(cold_index_duration), metrics: None, } } @@ -683,6 +875,7 @@ impl RejectedEventsIndex { hot_cache: HotCache::new(hot_cache_duration), cold_index: ColdIndex::new(cold_index_duration), unrecoverable: UnrecoverableIndex::new(cold_index_duration), + related_dependencies: RelatedDependencyIndex::new(cold_index_duration), metrics: Some(metrics), }; @@ -813,6 +1006,42 @@ impl RejectedEventsIndex { self.hot_cache.contains(event_id) || self.cold_index.contains(event_id) || self.unrecoverable.contains(event_id) + || self.related_dependencies.contains(event_id) + } + + /// Whether an indexed event is waiting on a dependency and must not be + /// treated as a terminally accounted hydration outcome. + pub fn is_dependency_pending(&self, event_id: &EventId) -> bool { + self.related_dependencies.contains(event_id) + || self + .cold_index + .entries + .read() + .unwrap() + .get(event_id) + .is_some_and(|entry| { + entry.reason == RejectionReason::MaintainerNotYetValid + && entry.rejected_at.elapsed() < self.cold_index.expiry_duration + }) + } + + /// Retain a policy-orphaned related event until either direction of its + /// reference relationship becomes accepted. + pub fn add_related_from_relay( + &self, + event: Event, + event_refs: HashSet, + addressable_refs: HashSet, + relay_url: Option, + ) -> bool { + self.related_dependencies + .add(event, event_refs, addressable_refs, relay_url) + } + + /// Return retained events whose backward or forward reference was made + /// valid by `accepted`. + pub fn related_candidates_resolved_by(&self, accepted: &Event) -> Vec { + self.related_dependencies.candidates_resolved_by(accepted) } /// Track a structurally unrecoverable event by ID alone. @@ -920,6 +1149,7 @@ impl RejectedEventsIndex { pub fn remove(&self, event_id: &EventId) { self.hot_cache.remove(event_id); self.cold_index.remove(event_id); + self.related_dependencies.remove(event_id); } /// Clean up expired entries from both tiers @@ -965,6 +1195,12 @@ impl RejectedEventsIndex { self.unrecoverable.cleanup_expired() } + /// Clean up related events whose dependency did not resolve within the + /// durable cold-index retention window. + pub fn cleanup_expired_related(&self) -> usize { + self.related_dependencies.cleanup_expired() + } + /// Get current number of entries in cold index pub fn cold_index_len(&self) -> usize { self.cold_index.len() @@ -975,6 +1211,11 @@ impl RejectedEventsIndex { self.unrecoverable.len() } + /// Get the number of retained dependency-sensitive related events. + pub fn related_len(&self) -> usize { + self.related_dependencies.len() + } + /// Get all rejected event IDs (from both hot cache and cold index) /// /// Used for excluding rejected events from negentropy sync. @@ -994,6 +1235,9 @@ impl RejectedEventsIndex { let unrecoverable_entries = self.unrecoverable.entries.read().unwrap(); ids.extend(unrecoverable_entries.keys().cloned()); + let related_entries = self.related_dependencies.entries.read().unwrap(); + ids.extend(related_entries.keys().cloned()); + ids } @@ -1018,6 +1262,7 @@ impl RejectedEventsIndex { let hot_entries = self.hot_cache.entries.read().unwrap(); let cold_entries = self.cold_index.entries.read().unwrap(); let unrecoverable_entries = self.unrecoverable.entries.read().unwrap(); + let related_entries = self.related_dependencies.entries.read().unwrap(); // Convert hot cache entries to serializable format let serializable_hot_entries: HashMap = hot_entries @@ -1074,6 +1319,25 @@ impl RejectedEventsIndex { }) .collect(); + let serializable_related_entries: HashMap = + related_entries + .iter() + .map(|(event_id, entry)| { + ( + *event_id, + SerializableRelatedDependencyEntry { + event: entry.event.clone(), + event_refs: entry.event_refs.clone(), + addressable_refs: entry.addressable_refs.clone(), + relay_hints: entry.relay_hints.clone(), + rejected_at_offset_secs: now + .duration_since(entry.rejected_at) + .as_secs(), + }, + ) + }) + .collect(); + // Create complete state let state = RejectedCacheState { version: 1, @@ -1090,6 +1354,7 @@ impl RejectedEventsIndex { expiry_duration_secs: self.unrecoverable.expiry_duration.as_secs(), entries: serializable_unrecoverable_entries, }, + related_dependencies: serializable_related_entries, }; // Replace the previous checkpoint only after the new snapshot is @@ -1202,6 +1467,28 @@ impl RejectedEventsIndex { drop(cold_entries); drop(unrecoverable_entries); + for (_, serializable_entry) in state.related_dependencies { + let total_offset = + Duration::from_secs(serializable_entry.rejected_at_offset_secs) + downtime; + if total_offset >= self.related_dependencies.expiry_duration { + continue; + } + let Some(serialized_bytes) = + RelatedDependencyIndex::serialized_event_bytes(&serializable_entry.event) + else { + continue; + }; + self.related_dependencies + .insert_restored(RelatedDependencyEntry { + event: serializable_entry.event, + event_refs: serializable_entry.event_refs, + addressable_refs: serializable_entry.addressable_refs, + relay_hints: serializable_entry.relay_hints, + rejected_at: now_instant - total_offset, + serialized_bytes, + }); + } + Ok(()) } } @@ -1209,7 +1496,9 @@ impl RejectedEventsIndex { #[cfg(test)] mod tests { use super::*; - use nostr_sdk::prelude::{EventBuilder, FinalizeUnsignedEvent, Keys, Kind, SignEvent}; + use nostr_sdk::prelude::{ + EventBuilder, FinalizeEvent, FinalizeUnsignedEvent, Keys, Kind, SignEvent, + }; async fn create_test_event() -> Event { let keys = Keys::generate(); @@ -2252,4 +2541,99 @@ mod tests { .get_maintainer_events(&event.pubkey, "test-repo", None); assert_eq!(events.len(), 0); } + + #[tokio::test] + async fn related_dependencies_resolve_in_both_reference_directions() { + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); + let keys = Keys::generate(); + let accepted_parent = create_test_event().await; + let backward_child = EventBuilder::new(Kind::TextNote, "backward child") + .tags([nostr_sdk::prelude::Tag::event(accepted_parent.id)]) + .finalize(&keys) + .expect("build backward child"); + assert!(index.add_related_from_relay( + backward_child.clone(), + HashSet::from([accepted_parent.id]), + HashSet::new(), + Some("wss://source.example".to_string()), + )); + + let backward = index.related_candidates_resolved_by(&accepted_parent); + assert_eq!( + backward.iter().map(|event| event.id).collect::>(), + vec![backward_child.id] + ); + + let forward_orphan = EventBuilder::new(Kind::TextNote, "forward orphan") + .finalize(&keys) + .expect("build forward orphan"); + assert!(index.add_related_from_relay( + forward_orphan.clone(), + HashSet::new(), + HashSet::new(), + None, + )); + let accepted_forward_ref = EventBuilder::new(Kind::TextNote, "accepted forward ref") + .tags([nostr_sdk::prelude::Tag::event(forward_orphan.id)]) + .finalize(&keys) + .expect("build accepted forward ref"); + + let forward = index.related_candidates_resolved_by(&accepted_forward_ref); + assert!(forward.iter().any(|event| event.id == forward_orphan.id)); + } + + #[tokio::test] + async fn related_dependencies_survive_checkpoint_restore() { + let directory = tempfile::tempdir().unwrap(); + let state_path = directory.path().join("rejected-events-cache.json"); + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); + let parent = create_test_event().await; + let keys = Keys::generate(); + let child = EventBuilder::new(Kind::TextNote, "durable child") + .tags([nostr_sdk::prelude::Tag::event(parent.id)]) + .finalize(&keys) + .expect("build durable child"); + assert!(index.add_related_from_relay( + child.clone(), + HashSet::from([parent.id]), + HashSet::new(), + Some("wss://source.example".to_string()), + )); + index.save_to_disk(&state_path).unwrap(); + + let restored = + RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); + restored.restore_from_disk(&state_path).unwrap(); + + assert!(restored.is_dependency_pending(&child.id)); + assert_eq!(restored.related_len(), 1); + assert_eq!( + restored + .related_candidates_resolved_by(&parent) + .into_iter() + .map(|event| event.id) + .collect::>(), + vec![child.id] + ); + } + + #[tokio::test] + async fn related_dependencies_evict_oldest_entry_at_the_count_bound() { + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); + let keys = Keys::generate(); + let oldest = EventBuilder::new(Kind::TextNote, "oldest") + .finalize(&keys) + .expect("build oldest event"); + assert!(index.add_related_from_relay(oldest.clone(), HashSet::new(), HashSet::new(), None,)); + + for sequence in 0..RELATED_MAX_ENTRIES { + let event = EventBuilder::new(Kind::TextNote, sequence.to_string()) + .finalize(&keys) + .expect("build bounded event"); + assert!(index.add_related_from_relay(event, HashSet::new(), HashSet::new(), None,)); + } + + assert_eq!(index.related_len(), RELATED_MAX_ENTRIES); + assert!(!index.contains(&oldest.id)); + } }