mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
fix(sync): reconcile the intended relay against local inventory
The SDK client uses a seen-ID tracker with no reconciliation inventory, so NIP-77 requested already-stored history. Client-wide sync could also merge another registered relay's inventory into the connection's result. Pass matching IDs and timestamps from the policy-controlled database to the exact relay's dry-run API. Abort on local inventory read errors and retain admission, pacing, session checks and bounded reconciliation. The local inventory is a snapshot; events saved concurrently still need batch-level accounting. Keep SDK storage isolated so fetching cannot bypass write policy, and leave payload admission and recovery semantics unchanged. Validation: the wire regression failed before the fix and now passes, including exact-relay isolation, no uploads and no policy-bypassing writes. All 114 sync integration tests passed (one ignored); strict all-target Clippy passed. Assisted-by: GPT-6
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
+133
-12
@@ -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<EventId>| {
|
||||
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<_>>(),
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user