diff --git a/CHANGELOG.md b/CHANGELOG.md index f956646..605978e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,6 +30,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 requests. Regroup only explicit filter-count refusals, and ignore those for already-retired subscriptions so late responses cannot incorrectly pause replacement coverage for 24 hours. +- Reconcile NIP-77 against the actual local event inventory on the intended + relay, avoiding downloads of already-stored history. Keep reconciliation + read-only and incoming payloads subject to the normal write policy. - Reuse installed core live filters during historic retries, preventing duplicate subscriptions from exhausting relay budgets. - Retire partially admitted or incomplete historic batches on admission failure, diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 12dbaa9..372f753 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -692,6 +692,10 @@ The ngit-grasp relay implements **Proactive Sync of Nostr Events**, which synchr a large source catches up. Pending batches register request IDs before sends and retain an admission barrier until the worker finishes, so early EOSE cannot confirm partial coverage. Reset/disconnect cancels stale jobs +- **Local reconciliation inventory** comes from the policy-controlled database + and is passed explicitly to the intended relay's NIP-77 dry run. The SDK's + separate seen-ID tracker remains isolated: reconciliation neither persists + incoming payloads nor publishes local events - **Daily sync** with random 23-25h timer to detect state drift - **Filter consolidation** when incremental fragmentation exceeds the desired live-filter baseline by 70; rebuilds are deferred until in-flight batches diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 0318d84..3ed2eaa 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -327,6 +327,14 @@ Each layer creates one or more `PendingBatch` entries tracked in `PendingSyncInd **Why the double-check?** There's an async gap between receiving EOSE and the self-subscriber processing events to create Layer 2/3 filters. The 6-second wait (5s batch window + 1s buffer) ensures we don't prematurely mark sync complete while Layer 2/3 batches are being created. +NIP-77 compares the intended relay's inventory with matching IDs and timestamps +from the local policy-controlled database. These items are passed explicitly to +the relay's dry-run reconciliation API; the SDK's separate seen-ID tracker +contains no inventory. Reconciliation does not upload local events or save +remote payloads. Downloads still pass through the write policy. A local +inventory read failure aborts reconciliation rather than pretending the +database is empty. + **Batch Failure Tracking**: Semantic REQ+EOSE fallback is reserved for a material first-pass hydration incompatibility: at least 20 IDs were advertised and exact-ID REQ delivered no more than 10%. This preserves the fallback for relays that advertise inventory but cannot substantially serve it by ID, without disabling NIP-77 for small residuals caused by expiry, indexing lag, or concurrent deletion. Other incomplete batches retry the residual IDs once. If that retry makes no progress, the batch is marked as `failed = true` and NIP-77 remains enabled. This causes the relay to transition to `ConnectedHistoricSyncFailures` instead of `Connected`, signaling that live sync is active but historic sync is incomplete. The event IDs the relay failed to deliver are not dropped with the batch: they are registered for bounded background recovery (see "Missing-Event Recovery for Incomplete Batches" below), and a relay whose pending IDs are all eventually recovered — with no unrelated batch failures — is promoted back to `Connected`. **Metrics tracking**: The `ngit_sync_relay_connected` metric shows: diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 339690d..f7da78e 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -2903,20 +2903,47 @@ impl RelayConnection { self.ensure_query_admitted()?; // Use dry_run to only identify differences without downloading events self.ensure_current_session(&ledger_slot)?; - let sync_opts = SyncOptions::default().dry_run(); - let client = self.client.clone(); + // Keep the SDK's seen-ID tracker separate from policy-controlled storage. + // It has no reconciliation inventory; attaching the production database + // to the SDK would let incoming events bypass our write policy. + let items = match &self.database { + Some(database) => tokio::time::timeout( + NEGENTROPY_DIFF_TIMEOUT, + database.negentropy_items(filter.clone()), + ) + .await + .map_err(|_| "Timed out loading local negentropy inventory".to_string())? + .map_err(|error| format!("Failed to load local negentropy inventory: {error}"))?, + None => Vec::new(), + }; + let relay = self + .client + .relay(&self.url) + .await + .map_err(|error| error.to_string())? + .ok_or_else(|| format!("Relay is not registered: {}", self.url))?; let sync_task = async move { self.ensure_query_admitted()?; - match client.sync(filter).opts(sync_opts).await { - Ok(output) => { - if !output.failed.is_empty() { - Err(format!("Negentropy diff had failures: {:?}", output.failed)) - } else { - Ok(output.value) - } - } - Err(error) => Err(format!("Negentropy diff failed: {}", error)), - } + self.ensure_current_session(&ledger_slot)?; + let summary = relay + .sync(filter) + .items(items) + .opts(SyncOptions::default().dry_run()) + .await + .map_err(|error| error.to_string())?; + let url = relay.url().clone(); + let with_source = |ids: std::collections::HashSet| { + ids.into_iter() + .map(|id| (id, std::collections::HashSet::from([url.clone()]))) + .collect() + }; + Ok(nostr_sdk::client::SyncSummary { + local: summary.local, + remote: with_source(summary.remote), + sent: with_source(summary.sent), + received: with_source(summary.received), + send_failures: std::collections::HashMap::from([(url, summary.send_failures)]), + }) }; self.run_negentropy_diff_with_timeout(sync_task, NEGENTROPY_DIFF_TIMEOUT) @@ -3382,6 +3409,100 @@ mod tests { } } + #[tokio::test] + async fn negentropy_uses_local_inventory_without_bypassing_write_policy() { + use nostr_sdk::prelude::NostrDatabase; + use std::collections::HashSet; + + let remote_database = std::sync::Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let remote = + TestRelay::start(LocalRelayBuilder::default().database(remote_database.clone())).await; + let other = TestRelay::start(LocalRelayBuilder::default()).await; + let database: SharedDatabase = + std::sync::Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let keys = Keys::generate(); + let event = |content| { + EventBuilder::new(Kind::TextNote, content) + .finalize(&keys) + .unwrap() + }; + let common = event("already stored"); + let remote_only = event("must pass write policy"); + let local_only = event("must not be uploaded"); + let unrelated = event("on another relay"); + database.save_event(&common).await.unwrap(); + database.save_event(&local_only).await.unwrap(); + remote.add_event(common.clone()).await.unwrap(); + remote.add_event(remote_only.clone()).await.unwrap(); + other.add_event(unrelated.clone()).await.unwrap(); + let connection = RelayConnection::new_with_database( + remote.url().await.to_string(), + database.clone(), + None, + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.unwrap(); + let other_url = other.url().await; + connection + .client + .add_relay(other_url.clone()) + .await + .unwrap(); + connection + .client + .try_connect_relay(other_url, Duration::from_secs(3)) + .await + .unwrap(); + + let diff = tokio::time::timeout( + Duration::from_secs(5), + connection.negentropy_sync_diff(Filter::new().kind(Kind::TextNote)), + ) + .await + .unwrap() + .unwrap(); + assert_eq!( + diff.remote.keys().copied().collect::>(), + HashSet::from([remote_only.id]) + ); + assert_eq!(diff.local, HashSet::from([local_only.id])); + assert!(diff.sent.is_empty()); + assert!(diff.received.is_empty()); + assert!( + database + .event_by_id(&remote_only.id) + .await + .unwrap() + .is_none(), + "reconciliation must not save remote payloads outside write policy" + ); + assert!( + remote_database + .event_by_id(&local_only.id) + .await + .unwrap() + .is_none(), + "dry-run reconciliation must not publish local events" + ); + + // Ordinary SDK fetching must also leave the production store untouched. + let fetched = connection + .fetch_events(Filter::new().id(remote_only.id), Duration::from_secs(3)) + .await + .unwrap(); + assert_eq!(fetched, vec![remote_only.clone()]); + assert!(database + .event_by_id(&remote_only.id) + .await + .unwrap() + .is_none()); + + connection.disconnect().await; + remote.shutdown(); + other.shutdown(); + } + #[tokio::test] async fn fetch_events_targets_the_connections_exact_relay() { let configured = TestRelay::start(LocalRelayBuilder::default()).await;