mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(sync): retry related policy orphans
Repository comments, reactions, zap requests, and other related events can arrive before the event or repository address that makes them admissible. Sync previously discarded these restricted policy orphans, making arrival order a permanent coverage gap. Persist full dependency-sensitive related events in a bounded seven-day index, match later accepted references in both directions, and reprocess an iterative bounded closure. Keep cached dependency failures pending in exact-ID hydration and expose retained count through the fixed-cardinality state metric. The index is capped at 1,024 events, 8 MiB total, 128 KiB per event, and 512 attempts per trigger. Recursive remote discovery and participant mailbox probing remain separate stacked changes. Validated with 754 library tests, focused forward/backward/persistence/eviction and policy reprocessing tests, formatting, and clippy across all targets with warnings denied.
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
+176
-7
@@ -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<RejectedEventsIndex>,
|
||||
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<RejectedEventsIndex>,
|
||||
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<Event> = 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();
|
||||
|
||||
+385
-1
@@ -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<EventId, SerializableUnrecoverableEntry>,
|
||||
}
|
||||
|
||||
/// 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<EventId>,
|
||||
addressable_refs: HashSet<String>,
|
||||
relay_hints: HashSet<String>,
|
||||
rejected_at: Instant,
|
||||
serialized_bytes: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
struct SerializableRelatedDependencyEntry {
|
||||
event: Event,
|
||||
event_refs: HashSet<EventId>,
|
||||
addressable_refs: HashSet<String>,
|
||||
#[serde(default)]
|
||||
relay_hints: HashSet<String>,
|
||||
rejected_at_offset_secs: u64,
|
||||
}
|
||||
|
||||
/// Bounded, durable full-event index for dependency-sensitive related events.
|
||||
#[derive(Debug, Clone)]
|
||||
struct RelatedDependencyIndex {
|
||||
entries: Arc<RwLock<HashMap<EventId, RelatedDependencyEntry>>>,
|
||||
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<String> {
|
||||
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<usize> {
|
||||
serde_json::to_vec(event).ok().map(|bytes| bytes.len())
|
||||
}
|
||||
|
||||
fn add(
|
||||
&self,
|
||||
event: Event,
|
||||
event_refs: HashSet<EventId>,
|
||||
addressable_refs: HashSet<String>,
|
||||
relay_url: Option<String>,
|
||||
) -> 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<EventId, RelatedDependencyEntry>) {
|
||||
while entries.len() > RELATED_MAX_ENTRIES
|
||||
|| entries
|
||||
.values()
|
||||
.map(|entry| entry.serialized_bytes)
|
||||
.sum::<usize>()
|
||||
> 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<Event> {
|
||||
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<EventId, SerializableRelatedDependencyEntry>,
|
||||
}
|
||||
|
||||
/// 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<super::metrics::SyncMetrics>,
|
||||
}
|
||||
|
||||
@@ -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<EventId>,
|
||||
addressable_refs: HashSet<String>,
|
||||
relay_url: Option<String>,
|
||||
) -> 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<Event> {
|
||||
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<EventId, SerializableHotCacheEntry> = hot_entries
|
||||
@@ -1074,6 +1319,25 @@ impl RejectedEventsIndex {
|
||||
})
|
||||
.collect();
|
||||
|
||||
let serializable_related_entries: HashMap<EventId, SerializableRelatedDependencyEntry> =
|
||||
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<_>>(),
|
||||
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<_>>(),
|
||||
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));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user