mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
Merge #a26eee65: fix(sync): retry related policy orphans
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqs2ymhwv5e62ckys4n94ggqu07tl8f83t0hlk3wun2d9yvzx9w299szfklkk PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 CoverNote: Stacked on nostr:nevent1qqsqzt6cxx9aeaj7wq5tr3mdne5e32ztpcmgvhpksz2jdyk0m20fksspz3mhxue69uhhyetvv9ujumn8d96zuer9wctgxcg4 and implements the durable dependency-retry layer from nostr:nevent1qqsfdj4fegeq7zjpje7cht6cj9xcsw5a9h2lcx2ju0f4uj80gsvvyhcpz3mhxue69uhhyetvv9ujumn8d96zuer9wc70zgv8.\n\nRepository-related events can legitimately arrive before the accepted event or repository address that makes them admissible. This change retains those policy orphans in the existing crash-safe rejected-event checkpoint and retries them in both directions whenever a dependency is saved.\n\nThe retained queue is deliberately bounded: 1,024 events, 8 MiB of serialized event payloads, 128 KiB per event, seven-day expiry, and 512 attempts per triggered closure. Oldest-first eviction is deterministic. Only terminal processing outcomes remove an entry; still-dependent and persistence outcomes remain pending.\n\nLocal evidence:\n- cargo fmt --check: pass\n- cargo test --lib: 754 passed\n- cargo clippy --all-targets -- -D warnings: pass\n- restore, forward/backward dependency matching, count-bound eviction, and end-to-end child-before-parent recovery are covered\n\nArchive canary evidence for exact commit 08d38fde:\n- deployed as 2.1.2-08d38fde at 2026-08-13 23:42:37 UTC\n- active/running with zero restarts; production remained on the parent PR\n- 26,636 deliveries observed across relay.ngit.dev and gitnostr.com after five minutes: 26,587 duplicates, 32 purgatory, 16 saved, one tombstone, zero rejected outcomes, and zero persistence errors\n- the retained-state sampler exported related_dependency_events=0 and sampled pipeline queue depth remained zero\n- no candidate-local panic, fatal, or persistence failure; channel-closed lines at 23:42:36 belong to shutdown of the previous process\n\nProduction evidence:\n- promoted unchanged at 2026-08-13 23:51:56 UTC as 2.1.2-08d38fde\n- active/running with zero restarts; archive start time remained unchanged\n- first 6,450 live deliveries were all duplicates, with zero rejected/persistence outcomes and queue depth zero\n- no candidate-local panic, fatal, or persistence failure in the post-start journal\n\nReady for merge as the second layer of the stack. Recursive relay fetching and participant mailbox discovery remain deliberately separate follow-up layers.
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