fix(sync): separate delivered dependencies from missing history

An unrelated same-identifier state was delivered by relay.ngit.dev but retained as MaintainerNotYetValid on gitnostr.com. Hydration treated local authorization readiness as remote delivery failure, repeatedly fetched the immutable payload, and degraded the source relay's history status.

Account retained authorization dependencies separately from terminal policy outcomes during stream processing, EOSE reconciliation, and bounded recovery. Consult retained dependencies when the SDK suppresses duplicate event notifications. Keep hot/cold dependency state and relay hints for later membership-triggered reprocessing; do not broaden state authority or start Git fetching before authorization. Untracked retryable failures and persistence errors remain pending.

A wire-level manager regression covers first and cached delivery, in-flight and queued recovery, cold payload expiry, retained source hints, and absence of Git-fetch purgatory. Existing dependency revival and maintainer acceptance coverage remains intact. This assumes retained dependency metadata owns future authorization recovery; cache TTLs, protocol rules, and deployment configuration are unchanged.

Validation: 937 library tests and 114 sync integration tests passed (one ignored). After tightening the retryable-error boundary, all 421 sync unit tests and eight maintainer integration tests passed again. Strict all-target Clippy, formatting, and diff checks passed.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-09-25 08:07:23 +00:00
parent 75bea23318
commit 5d65cd5eaf
5 changed files with 343 additions and 73 deletions
+4
View File
@@ -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
+2 -1
View File
@@ -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
│
+16 -10
View File
@@ -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
+318 -60
View File
@@ -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<EventId> = 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<EventId> = match self
let mut satisfied: Vec<EventId> = 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]
+3 -2
View File
@@ -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