diff --git a/CHANGELOG.md b/CHANGELOG.md index b2e3858..f956646 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,6 +36,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 timeout, or closure before EOSE. Release their subscriptions, bound SDK retention with an auto-close deadline, and leave coverage unconfirmed for recovery instead of retaining stalled work or accepting partial history. + Count delivered events retained for authorization/dependency recovery as + received, avoiding repeated payload fetches and false remote-history failures. + Preserve their cached IDs and relay hints for later maintainer changes without + fetching Git data before state authorization. - Pace descendant fallback cycles with a one-minute refresh delay, retaining overlap and immediate baselines for changed frontiers. Report missing semantic fallback metadata as scheduled recovery rather than a subscription-creation diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 34884d3..12dbaa9 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -830,7 +830,8 @@ Exact-ID Recovery Negentropy Sync │ - └──▶ Exclude Cold Index IDs from "missing events" calculation + ├──▶ Exclude Cold Index IDs from "missing events" calculation + └──▶ Count delivered dependency-pending events without removing their cache entries Related Event Rejected as an Orphan │ diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index e796764..0318d84 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -928,13 +928,15 @@ maintenance timer: (30s base doubling up to 15min; sub-second in `NGIT_TEST`), one in-flight attempt per relay, at most 300 IDs per fetch. One persistently incomplete relay cannot starve other relays or later batches. -- **Outcome-aware**: a successful attempt clears only IDs with a durable - terminal explanation: saved, already stored, in purgatory, validly - tombstoned, blocked, or permanently invalid. Restricted, policy-error, - unknown, and persistence-error outcomes remain pending because a dependency - or transient fault may clear. Duplicate incomplete responses merge into the - existing pending set. IDs that arrive by other means (live sync, user - submission) are cleared on the next tick without consuming attempt budget. +- **Outcome-aware**: transport completion includes events saved, already stored, + in purgatory, validly tombstoned, blocked, permanently invalid, or retained in + the rejected-events dependency index. A delivered state awaiting maintainer + authorization is not a remote delivery failure. Its dependency entry remains + available for later reprocessing; transport completion does not authorize it. + Untracked retryable policy failures and persistence errors remain pending. + Duplicate incomplete responses merge into the existing pending set. IDs + stored or retained by policy tracking through another path are cleared on the + next tick without consuming attempt budget. - **Explicit expiry**: after 12 consecutive zero-progress attempts the relay's pending IDs are dropped with a warning, and the relay stays in `ConnectedHistoricSyncFailures` until the daily sync re-discovers the gap. @@ -1003,9 +1005,13 @@ 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. +so later broad synchronization may offer them again. Exact-ID hydration counts +retained dependency-pending events as delivered, including cached rejections, +without clearing their dependency entries or treating them as locally accepted. +Unrelated same-identifier state events therefore do not trigger repeated payload +fetches or degrade the source relay's history status. Git fetching still requires +normal state authorization; a future confirmed maintainer relationship can +activate an older state through hot-cache reprocessing or cold-ID recovery. See [Architecture: Rejected Events Index](architecture.md#rejected-events-index) and [`src/sync/rejected_index.rs`](../../src/sync/rejected_index.rs) for the diff --git a/src/sync/mod.rs b/src/sync/mod.rs index af5b970..9f3387e 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -640,11 +640,25 @@ impl ProcessResult { matches!(self, Self::Rejected(_)) } + /// Whether transport recovery is finished, independently of local authorization. + /// Retained dependencies can be revived by announcement changes; fetching the + /// same immutable payload again cannot resolve their missing authority. + fn is_hydration_accounted( + self, + event_id: &EventId, + rejected_events_index: &RejectedEventsIndex, + ) -> bool { + self.is_terminally_accounted() + || (self == Self::Rejected(PolicyRejection::Restricted) + && rejected_events_index.is_dependency_pending(event_id)) + } + /// Whether this result gives a durable explanation for the requested ID. /// + /// This controls dependency-cache removal, not transport completion. /// Restricted, server-error, and unknown policy failures can become valid - /// after dependencies arrive or a transient fault clears, so exact-ID - /// recovery keeps them pending. Invalid/blocked results are permanent for + /// after dependencies arrive or a transient fault clears, so dependency + /// reprocessing keeps them pending. Invalid/blocked results are permanent for /// the event, while purgatory and tombstones are explicit terminal states. fn is_terminally_accounted(self) -> bool { matches!( @@ -3086,8 +3100,17 @@ impl SyncManager { // For negentropy batches, check if all requested events were received if batch.sync_method == SyncMethod::Negentropy { if let (Some(requested), Some(received)) = - (&batch.requested_event_ids, &batch.received_event_ids) + (&batch.requested_event_ids, &mut batch.received_event_ids) { + // The SDK can suppress repeated Event notifications for payloads + // it already holds. Retained policy decisions still prove these + // IDs were delivered, even without another stream notification. + received.extend( + requested + .iter() + .copied() + .filter(|id| self.rejected_events_index.is_dependency_pending(id)), + ); let missing: Vec = requested.difference(received).cloned().collect(); if !missing.is_empty() { @@ -3400,7 +3423,7 @@ impl SyncManager { /// reconciliation but failed to deliver on exact-ID fetches. /// /// Runs on the sync maintenance timer. For each relay with pending IDs: - /// 1. IDs already present locally (live sync, user submission) are + /// 1. IDs already stored or retained by policy/dependency tracking are /// cleared promptly without consuming attempt budget. /// 2. When the relay's backoff deadline has passed and it has a live /// connection, one bounded exact-ID fetch is spawned. Network I/O runs @@ -3424,7 +3447,7 @@ impl SyncManager { }; // 1. Clear IDs satisfied by other means. - let satisfied: Vec = match self + let mut satisfied: Vec = match self .database .query(Filter::new().ids(pending.iter().copied())) .await @@ -3432,6 +3455,12 @@ impl SyncManager { Ok(events) => events.into_iter().map(|event| event.id).collect(), Err(_) => Vec::new(), }; + satisfied.extend( + pending + .iter() + .copied() + .filter(|id| self.rejected_events_index.is_dependency_pending(id)), + ); if !satisfied.is_empty() { let outcome = self .missing_event_recovery @@ -3441,7 +3470,7 @@ impl SyncManager { tracing::info!( relay = %relay_url, satisfied = satisfied.len(), - "Pending missing events satisfied by local arrivals" + "Pending missing events accounted for by local storage or retained policy decisions" ); if let Some(outcome) = outcome { Self::apply_recovery_outcome( @@ -3510,8 +3539,8 @@ impl SyncManager { /// /// Runs outside the sync actor lock. Every event the relay returns is /// passed through the normal write policy. An ID counts as recovered only - /// once it has a durable terminal outcome; transient persistence and - /// dependency-sensitive policy failures remain pending. + /// once it has a terminal outcome or is retained for dependency recovery. + /// Untracked policy failures and persistence errors remain pending. #[allow(clippy::too_many_arguments)] async fn run_missing_event_recovery_attempt( relay_url: String, @@ -3568,9 +3597,8 @@ impl SyncManager { if let Some(metrics) = metrics.as_ref() { metrics.record_hydration_events(&relay_url, "recovery", "delivered", 1); } - // Permanent cached rejections account for the requested ID. - // Dependency-sensitive entries remain pending until their normal - // re-processing machinery observes the missing accepted event. + // A retained policy decision proves delivery. Keep dependency-sensitive + // entries for authorization-triggered reprocessing, not payload retries. if rejected_events_index.contains(&event.id) { let dependency_pending = rejected_events_index.is_dependency_pending(&event.id); tracing::debug!( @@ -3582,9 +3610,7 @@ impl SyncManager { if let Some(metrics) = metrics.as_ref() { metrics.record_hydration_events(&relay_url, "recovery", "rejected_cached", 1); } - if !dependency_pending { - recovered.insert(event.id); - } + recovered.insert(event.id); continue; } let result = Self::process_event_static( @@ -3617,7 +3643,7 @@ impl SyncManager { "Recovered missing event was not persisted" ); } - if result.is_terminally_accounted() { + if result.is_hydration_accounted(&event.id, &rejected_events_index) { recovered.insert(event.id); } } @@ -4543,50 +4569,25 @@ impl SyncManager { } } - // Skip events we've already rejected (announcements only) - if (event.kind == Kind::GitRepoAnnouncement + // Cached decisions still account for delivery to this batch. + // Keep dependency entries available for later authorization changes. + let result = if (event.kind == Kind::GitRepoAnnouncement || event.kind == Kind::RepoState) && rejected_events_index.contains(&event.id) { - tracing::trace!( - event_id = %event.id, - kind = %event.kind.as_u16(), - relay = %relay_url_clone, - "Skipping previously rejected announcement event" - ); - pipeline_window.record( - ProcessResult::Rejected(PolicyRejection::PreviouslyRejected), - queue_delay, - processing_started.elapsed(), - ); - if let Some(ref metrics) = metrics_clone { - metrics.record_hydration_events( - &relay_url_clone, - "stream", - "delivered", - 1, - ); - metrics.record_hydration_events( - &relay_url_clone, - "stream", - "rejected_cached", - 1, - ); - } - pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); - continue; - } - - let result = Self::process_event_static( - &event, - &relay_url_clone, - &database, - &write_policy, - &local_relay, - &rejected_events_index, - crate::nostr::persistence::SaveContext::RelaySync, - ) - .await; + ProcessResult::Rejected(PolicyRejection::PreviouslyRejected) + } else { + Self::process_event_static( + &event, + &relay_url_clone, + &database, + &write_policy, + &local_relay, + &rejected_events_index, + crate::nostr::persistence::SaveContext::RelaySync, + ) + .await + }; if let Some(ref metrics) = metrics_clone { metrics.record_hydration_events( &relay_url_clone, @@ -4648,10 +4649,9 @@ impl SyncManager { } } - // Track received event IDs for negentropy batches. Unlike REQ+EOSE - // pagination above, negentropy completion is concerned with events that - // were actually saved or already present locally. - if result.is_terminally_accounted() { + // Account for delivered events saved locally or retained under policy. + // Dependency readiness is separate from the source's transport health. + if result.is_hydration_accounted(&event.id, &rejected_events_index) { let mut pending = pending_sync_index.write().await; if let Some(batches) = pending.get_mut(&relay_url_clone) { for batch in batches.iter_mut() { @@ -9580,6 +9580,258 @@ mod tests { manager.shutdown().await; } + #[tokio::test] + async fn delivered_unauthorized_state_completes_hydration_without_losing_dependency() { + use futures_util::{SinkExt, StreamExt}; + use tokio_tungstenite::tungstenite::Message; + + let directory = tempfile::tempdir().unwrap(); + let mut config = Config::for_testing(); + let git_path = directory.path().join("git"); + config.git_data_path = git_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_path.clone())); + let runtime = crate::nostr::builder::create_relay( + &config, + purgatory.clone(), + crate::grasp06::receive::RepoInitLocks::default(), + None, + ) + .await + .unwrap(); + let owner = Keys::generate(); + let stranger = Keys::generate(); + // Seed a hosted coordinate with the same identifier but no relationship + // to the other author. This is the production test3 collision. + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags([ + Tag::identifier("test3"), + Tag::custom("clone", ["https://service.example/test3.git"]), + ]) + .finalize(&owner) + .unwrap(); + runtime + .stores + .database + .save_event(&announcement) + .await + .unwrap(); + let event = EventBuilder::new(Kind::RepoState, "") + .tags([ + Tag::identifier("test3"), + Tag::custom( + "refs/heads/main", + ["1111111111111111111111111111111111111111"], + ), + Tag::custom("HEAD", ["ref: refs/heads/main"]), + ]) + .finalize(&stranger) + .unwrap(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let served_event = event.clone(); + let mut server = tokio::task::JoinSet::new(); + server.spawn(async move { + let (socket, _) = listener.accept().await.unwrap(); + let mut socket = tokio_tungstenite::accept_async(socket).await.unwrap(); + while let Some(Ok(frame)) = socket.next().await { + if !frame.is_text() { + continue; + } + let message: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + if message[0] == "REQ" { + socket + .send(Message::Text( + serde_json::json!(["EVENT", message[1], served_event]) + .to_string() + .into(), + )) + .await + .unwrap(); + socket + .send(Message::Text( + serde_json::json!(["EOSE", message[1]]).to_string().into(), + )) + .await + .unwrap(); + } + } + }); + let mut manager = SyncManager::new( + None, + "service.example".into(), + runtime.stores.database.clone(), + runtime.write_policy, + runtime.relay, + &config, + git_path.clone(), + None, + None, + None, + ); + // Expire full payloads immediately: cold metadata must still prevent + // transport retries and retain the ID/hint for later membership changes. + manager.rejected_events_index = Arc::new(RejectedEventsIndex::new( + Duration::ZERO, + Duration::from_secs(3600), + )); + let connection = RelayConnection::new( + url.clone(), + None, + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.unwrap(); + manager.connections.insert(url.clone(), connection.clone()); + manager.nip65_discovery_only_relays.insert(url.clone()); + let (disconnect_tx, _disconnect_rx) = tokio::sync::mpsc::channel(8); + let (eose_tx, mut eose_rx) = lifecycle_notification_channel(); + let (closed_tx, _closed_rx) = lifecycle_notification_channel(); + manager.disconnect_tx = Some(disconnect_tx); + manager.eose_tx = Some(eose_tx); + manager.subscription_closed_tx = Some(closed_tx); + manager.handle_connect_or_reconnect(&url).await; + for number in 0..2 { + let sub = SubscriptionId::new(format!("state-{number}")); + manager.pending_sync_index.write().await.insert( + url.clone(), + vec![PendingBatch { + batch_id: number, + purpose: PendingBatchPurpose::Core, + items: PendingItems::default(), + outstanding_subs: HashSet::from([sub.clone()]), + sync_method: SyncMethod::Negentropy, + pagination_state: HashMap::new(), + requested_event_ids: Some(HashSet::from([event.id])), + received_event_ids: Some(HashSet::new()), + initial_hydration_counts: None, + retry_count: 0, + failed: false, + }], + ); + connection + .subscribe_filter_with_id( + Filter::new().id(event.id), + TransientRequestClass::NegentropyRetry, + sub.clone(), + ) + .await + .unwrap(); + let eose = tokio::time::timeout(Duration::from_secs(5), eose_rx.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!(eose.sub_id, sub); + if number == 0 { + assert_eq!( + manager.pending_sync_index.read().await[&url][0].received_event_ids, + Some(HashSet::from([event.id])), + "a delivered dependency must satisfy initial hydration" + ); + } + manager.handle_eose(&url, eose.sub_id).await; + assert!(!manager.pending_sync_index.read().await.contains_key(&url)); + assert!(manager.missing_event_recovery.lock().unwrap().pending_ids(&url).is_none(), + "cached delivery must not create missing-event recovery even if the SDK suppresses its notification"); + } + assert!(manager + .rejected_events_index + .is_dependency_pending(&event.id)); + assert!(runtime + .stores + .database + .event_by_id(&event.id) + .await + .unwrap() + .is_none()); + assert!( + purgatory.find_state("test3").is_empty(), + "unauthorized state must not enter Git-fetch purgatory" + ); + assert!(!git_path + .join(stranger.public_key().to_bech32().unwrap()) + .exists()); + + // An already-running exact-ID recovery must also finish when its + // response is a retained dependency, without deleting that dependency. + let now = Instant::now(); + let attempt = { + let mut recovery = manager.missing_event_recovery.lock().unwrap(); + recovery.register(&url, 8, [event.id], false, now); + recovery + .begin_attempt(&url, now + Duration::from_secs(60)) + .unwrap() + }; + tokio::time::timeout( + Duration::from_secs(5), + SyncManager::run_missing_event_recovery_attempt( + url.clone(), + connection.clone(), + attempt, + manager.database.clone(), + manager.write_policy.clone(), + manager.local_relay.clone(), + manager.rejected_events_index.clone(), + manager.missing_event_recovery.clone(), + manager.relay_sync_index.clone(), + None, + ), + ) + .await + .unwrap(); + assert!(manager + .missing_event_recovery + .lock() + .unwrap() + .pending_ids(&url) + .is_none()); + assert!(manager + .rejected_events_index + .is_dependency_pending(&event.id)); + + manager.missing_event_recovery.lock().unwrap().register( + &url, + 9, + [event.id], + false, + Instant::now(), + ); + manager.tick_missing_event_recovery().await; + assert!(manager + .missing_event_recovery + .lock() + .unwrap() + .pending_ids(&url) + .is_none()); + assert!( + manager + .rejected_events_index + .is_dependency_pending(&event.id), + "transport completion must not remove the authorization dependency" + ); + let (ids, hot) = manager.rejected_events_index.dependency_candidates( + &stranger.public_key(), + "test3", + Some(rejected_index::EventType::State), + ); + assert_eq!(ids, vec![event.id]); + assert!(hot.is_empty()); + assert!(manager + .rejected_events_index + .dependency_relay_hints( + &stranger.public_key(), + "test3", + Some(rejected_index::EventType::State) + ) + .contains(&url)); + manager.shutdown().await; + } + #[tokio::test] async fn accepted_dependency_reprocesses_a_synced_policy_orphan() { let directory = tempfile::tempdir().expect("create test directory"); @@ -9628,6 +9880,8 @@ mod tests { ProcessResult::Rejected(PolicyRejection::Restricted) ); assert!(rejected.is_dependency_pending(&child.id)); + assert!(orphan_result.is_hydration_accounted(&child.id, &rejected)); + assert!(!orphan_result.is_terminally_accounted()); assert!(runtime .stores .database @@ -9849,6 +10103,10 @@ mod tests { assert_eq!(restricted.hydration_outcome(), "rejected_restricted"); assert!(!restricted.is_terminally_accounted()); assert!(!persistence_error.is_terminally_accounted()); + let rejected = RejectedEventsIndex::new(Duration::ZERO, Duration::from_secs(3600)); + let id = EventId::from_byte_array([42; 32]); + assert!(!restricted.is_hydration_accounted(&id, &rejected)); + assert!(!persistence_error.is_hydration_accounted(&id, &rejected)); } #[test] diff --git a/src/sync/rejected_index.rs b/src/sync/rejected_index.rs index 829a7c3..a26d395 100644 --- a/src/sync/rejected_index.rs +++ b/src/sync/rejected_index.rs @@ -1009,8 +1009,9 @@ impl RejectedEventsIndex { || 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. + /// Whether an indexed event is waiting on a dependency and must remain + /// available for reprocessing. Its payload has already been delivered, so + /// transport recovery can finish without treating the event as accepted. pub fn is_dependency_pending(&self, event_id: &EventId) -> bool { self.related_dependencies.contains(event_id) || self