diff --git a/CHANGELOG.md b/CHANGELOG.md index 7c26b00..2eaa1aa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -43,7 +43,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 recovery before declaring history incomplete, avoiding repeated payload fetches and false remote-history failures. Preserve dependency IDs and relay hints for later maintainer changes without fetching Git data before state - authorization. Include bounded ID samples in background recovery logs. + authorization. Bound recovery per event ID so progress on new events cannot + renew older failures' retry budgets, rotate pending IDs fairly, and include + bounded ID samples in background recovery logs. - 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 cb51b88..70d0265 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -700,6 +700,9 @@ The ngit-grasp relay implements **Proactive Sync of Nostr Events**, which synchr even when the SDK suppresses repeat event notifications. Local lookups use bounded ID chunks without holding the pending-batch lock; absent IDs and failed storage lookups remain eligible for recovery +- **Missing-ID budgets** expire each ID after twelve unsuccessful fetches. + Per-relay requests rotate pending IDs fairly and retain backoff for unresolved + IDs even when unrelated IDs recover; expiry cannot falsely restore health - **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 51b1c5f..fc5fa00 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -942,8 +942,10 @@ maintenance timer: connection, outside the sync actor lock. - **Bounded and backed off**: attempts are per relay with exponential backoff (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. + attempt per relay, at most 300 IDs per fetch. Pending IDs rotate in FIFO + order so an unavailable chunk cannot starve other IDs, and new arrivals + cannot continually jump ahead of older work. The delay follows the highest + failed-fetch count among remaining IDs; unrelated recoveries cannot reset it. Each attempt logs at most five requested IDs alongside its count, allowing residual events to be inspected without dumping the whole inventory. - **Outcome-aware**: transport completion includes events saved, already stored, @@ -955,8 +957,10 @@ maintenance timer: 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 +- **Explicit expiry**: each ID expires after 12 unsuccessful fetches of that + ID, regardless of progress on other IDs. Newer and unattempted IDs retain + their own budgets. Expiry is logged and prevents health restoration even + if the remaining IDs later recover; the relay stays in `ConnectedHistoricSyncFailures` until the daily sync re-discovers the gap. Attempts against a disconnected relay are deferred, not counted, so an unavailable relay neither expires its pending work nor loops tightly. diff --git a/src/sync/missing_events.rs b/src/sync/missing_events.rs index 0f1befa..6da0631 100644 --- a/src/sync/missing_events.rs +++ b/src/sync/missing_events.rs @@ -16,10 +16,10 @@ //! //! - Attempts are scheduled per relay with exponential backoff //! (30s base doubling up to 15min in production; sub-second in `NGIT_TEST`). -//! - Any progress (at least one ID recovered) resets the backoff; the pending -//! set shrinks monotonically, so this cannot loop forever. -//! - After [`MAX_ATTEMPTS_WITHOUT_PROGRESS`] consecutive attempts with zero -//! progress the relay's pending IDs are expired: they are dropped with an +//! - Each ID has its own failed-fetch budget. Recovering newer IDs cannot +//! reset the delay or budget of older unresolved IDs. +//! - After [`MAX_ATTEMPTS_PER_ID`] unsuccessful fetches an ID expires, +//! without dropping newer or unattempted IDs. Expiry is logged with an //! explicit warning and the relay remains in //! `ConnectedHistoricSyncFailures` until the next daily sync re-discovers //! the gap. Nothing is persisted across restarts - a restart re-runs @@ -31,7 +31,7 @@ //! observed for that relay (the entry is then marked as unable to restore //! health, so the degraded status stays visible). -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::time::{Duration, Instant}; use nostr_sdk::prelude::EventId; @@ -40,8 +40,8 @@ use nostr_sdk::prelude::EventId; /// used by historic sync. Larger pending sets recover across several attempts. pub const MAX_RECOVERY_IDS_PER_ATTEMPT: usize = 300; -/// Consecutive zero-progress attempts before a relay's pending IDs expire. -pub const MAX_ATTEMPTS_WITHOUT_PROGRESS: u32 = 12; +/// Unsuccessful fetches before an individual pending ID expires. +pub const MAX_ATTEMPTS_PER_ID: u32 = 12; /// Cap on remembered source batch IDs (for logging only). const MAX_TRACKED_SOURCE_BATCHES: usize = 16; @@ -76,9 +76,11 @@ fn backoff_for(attempts_without_progress: u32) -> Duration { #[derive(Debug)] struct RelayRecoveryState { /// Event IDs the relay reported but has not yet delivered. - missing: HashSet, - /// Consecutive attempts that recovered nothing. Reset on progress. - attempts_without_progress: u32, + missing: HashMap, + /// Round-robin admission prevents either old or newly arriving IDs starving. + queue: VecDeque, + /// IDs in the current fetch; only these consume an attempt on failure. + attempted_ids: Vec, /// Total attempts issued, for logging. total_attempts: u32, /// Earliest time the next attempt may run. @@ -133,7 +135,14 @@ pub enum AttemptOutcome { remaining: usize, next_attempt_in: Duration, }, - /// The zero-progress attempt budget is exhausted; pending IDs dropped. + /// Some IDs exhausted their budget; younger pending IDs keep recovering. + SomeExpired { + expired: usize, + remaining: usize, + attempts: u32, + next_attempt_in: Duration, + }, + /// Every remaining ID exhausted its own fetch budget. Expired { remaining: usize, attempts: u32 }, } @@ -162,8 +171,9 @@ impl MissingEventRecoveryIndex { self.relays .entry(relay_url.to_string()) .or_insert_with(|| RelayRecoveryState { - missing: HashSet::new(), - attempts_without_progress: 0, + missing: HashMap::new(), + queue: VecDeque::new(), + attempted_ids: Vec::new(), total_attempts: 0, next_attempt_at: now + base_backoff(), in_flight: false, @@ -182,7 +192,12 @@ impl MissingEventRecoveryIndex { } let before = entry.missing.len(); - entry.missing.extend(ids); + for id in ids { + if let std::collections::hash_map::Entry::Vacant(slot) = entry.missing.entry(id) { + slot.insert(0); + entry.queue.push_back(id); + } + } RegisterOutcome { newly_added: entry.missing.len() - before, pending_total: entry.missing.len(), @@ -217,7 +232,7 @@ impl MissingEventRecoveryIndex { if entry.in_flight { return None; } - Some(entry.missing.iter().copied().collect()) + Some(entry.missing.keys().copied().collect()) } /// Remove IDs that were satisfied outside recovery (live sync, user @@ -235,9 +250,7 @@ impl MissingEventRecoveryIndex { } let recovered = before - entry.missing.len(); entry.total_recovered += recovered; - if recovered > 0 { - entry.attempts_without_progress = 0; - } + entry.queue.retain(|id| entry.missing.contains_key(id)); if entry.missing.is_empty() { let entry = self .relays @@ -260,13 +273,15 @@ impl MissingEventRecoveryIndex { } entry.in_flight = true; entry.total_attempts += 1; + entry.attempted_ids = entry + .queue + .iter() + .take(MAX_RECOVERY_IDS_PER_ATTEMPT) + .copied() + .collect(); + entry.queue.rotate_left(entry.attempted_ids.len()); Some(RecoveryAttempt { - ids: entry - .missing - .iter() - .take(MAX_RECOVERY_IDS_PER_ATTEMPT) - .copied() - .collect(), + ids: entry.attempted_ids.clone(), attempt_number: entry.total_attempts, }) } @@ -297,7 +312,7 @@ impl MissingEventRecoveryIndex { entry.in_flight = false; let before = entry.missing.len(); - entry.missing.retain(|id| !recovered.contains(id)); + entry.missing.retain(|id, _| !recovered.contains(id)); let recovered_count = before - entry.missing.len(); entry.total_recovered += recovered_count; @@ -313,30 +328,49 @@ impl MissingEventRecoveryIndex { }); } + for id in entry.attempted_ids.drain(..) { + if let Some(failures) = entry.missing.get_mut(&id) { + *failures += 1; + } + } + let before_expiry = entry.missing.len(); + entry + .missing + .retain(|_, failures| *failures < MAX_ATTEMPTS_PER_ID); + let expired = before_expiry - entry.missing.len(); + entry.queue.retain(|id| entry.missing.contains_key(id)); + if expired > 0 { + // Recovering the younger IDs cannot repair the abandoned gap. + entry.can_restore_health = false; + } + if entry.missing.is_empty() { + let entry = self + .relays + .remove(relay_url) + .expect("entry present just above"); + return Some(AttemptOutcome::Expired { + remaining: expired, + attempts: entry.total_attempts, + }); + } + + let delay = backoff_for(entry.missing.values().copied().max().unwrap_or(0)); + entry.next_attempt_at = now + delay; + if expired > 0 { + return Some(AttemptOutcome::SomeExpired { + expired, + remaining: entry.missing.len(), + attempts: entry.total_attempts, + next_attempt_in: delay, + }); + } if recovered_count > 0 { - entry.attempts_without_progress = 0; - let delay = backoff_for(0); - entry.next_attempt_at = now + delay; return Some(AttemptOutcome::PartiallyRecovered { recovered: recovered_count, remaining: entry.missing.len(), next_attempt_in: delay, }); } - - entry.attempts_without_progress += 1; - if entry.attempts_without_progress >= MAX_ATTEMPTS_WITHOUT_PROGRESS { - let entry = self - .relays - .remove(relay_url) - .expect("entry present just above"); - return Some(AttemptOutcome::Expired { - remaining: entry.missing.len(), - attempts: entry.total_attempts, - }); - } - let delay = backoff_for(entry.attempts_without_progress); - entry.next_attempt_at = now + delay; Some(AttemptOutcome::RetryScheduled { attempt: entry.total_attempts, remaining: entry.missing.len(), @@ -420,7 +454,7 @@ mod tests { } #[test] - fn successful_attempt_clears_only_recovered_ids_and_resets_backoff() { + fn successful_attempt_preserves_unresolved_ids_backoff() { let mut index = MissingEventRecoveryIndex::default(); let now = registered(&mut index, &[id(1), id(2), id(3)]); @@ -441,15 +475,96 @@ mod tests { AttemptOutcome::PartiallyRecovered { recovered: 1, remaining: 2, - next_attempt_in: backoff_for(0), + next_attempt_in: backoff_for(3), }, - "progress must clear only the recovered ID and reset the backoff" + "progress must not reset the unresolved IDs backoff" ); let mut remaining = index.pending_ids("ws://relay").unwrap(); remaining.sort(); assert_eq!(remaining, vec![id(1), id(3)]); } + #[test] + fn fresh_progress_cannot_renew_an_old_ids_budget() { + let mut index = MissingEventRecoveryIndex::default(); + let now = registered(&mut index, &[id(1)]); + for attempt in 1..=MAX_ATTEMPTS_PER_ID { + index.register("ws://relay", 2, [id(2), id(3)], false, now); + index.clear_satisfied("ws://relay", &[id(2)]); + index.begin_attempt("ws://relay", due(now)).unwrap(); + let outcome = index + .complete_attempt("ws://relay", &HashSet::from([id(3)]), now) + .unwrap(); + if attempt == MAX_ATTEMPTS_PER_ID { + assert_eq!( + outcome, + AttemptOutcome::Expired { + remaining: 1, + attempts: attempt + } + ); + } else { + assert_eq!( + outcome, + AttemptOutcome::PartiallyRecovered { + recovered: 1, + remaining: 1, + next_attempt_in: backoff_for(attempt), + } + ); + } + } + } + + #[test] + fn expiry_preserves_younger_ids_without_restoring_health() { + let mut index = MissingEventRecoveryIndex::default(); + let now = registered(&mut index, &[id(1)]); + for _ in 1..MAX_ATTEMPTS_PER_ID { + index.begin_attempt("ws://relay", due(now)).unwrap(); + index.complete_attempt("ws://relay", &HashSet::new(), now); + } + index.begin_attempt("ws://relay", due(now)).unwrap(); + // Registration during a fetch must not spend the new ID's budget. + index.register("ws://relay", 2, [id(2)], false, now); + assert_eq!( + index.complete_attempt("ws://relay", &HashSet::new(), now), + Some(AttemptOutcome::SomeExpired { + expired: 1, + remaining: 1, + attempts: MAX_ATTEMPTS_PER_ID, + next_attempt_in: backoff_for(0), + }) + ); + assert_eq!(index.pending_ids("ws://relay"), Some(vec![id(2)])); + index.begin_attempt("ws://relay", due(now)).unwrap(); + assert!(matches!( + index.complete_attempt("ws://relay", &HashSet::from([id(2)]), now), + Some(AttemptOutcome::FullyRecovered { + can_restore_health: false, + .. + }) + )); + } + + #[test] + fn large_inventory_does_not_starve_unattempted_ids() { + let mut index = MissingEventRecoveryIndex::default(); + let ids: Vec<_> = (0..MAX_RECOVERY_IDS_PER_ATTEMPT + 1) + .map(|n| { + let mut bytes = [0; 32]; + bytes[..8].copy_from_slice(&(n as u64).to_be_bytes()); + EventId::from_slice(&bytes).unwrap() + }) + .collect(); + let now = registered(&mut index, &ids); + let first = index.begin_attempt("ws://relay", due(now)).unwrap(); + let untouched = *ids.iter().find(|id| !first.ids.contains(id)).unwrap(); + index.complete_attempt("ws://relay", &HashSet::new(), now); + let second = index.begin_attempt("ws://relay", due(now)).unwrap(); + assert_eq!(second.ids[0], untouched); + } + #[test] fn zero_progress_attempts_back_off_exponentially_with_a_cap() { assert!(backoff_for(1) > backoff_for(0)); @@ -462,7 +577,7 @@ mod tests { let mut index = MissingEventRecoveryIndex::default(); let now = registered(&mut index, &[id(1), id(2)]); - for attempt in 1..MAX_ATTEMPTS_WITHOUT_PROGRESS { + for attempt in 1..MAX_ATTEMPTS_PER_ID { index.begin_attempt("ws://relay", due(now)).unwrap(); let outcome = index .complete_attempt("ws://relay", &HashSet::new(), now) @@ -481,7 +596,7 @@ mod tests { outcome, AttemptOutcome::Expired { remaining: 2, - attempts: MAX_ATTEMPTS_WITHOUT_PROGRESS, + attempts: MAX_ATTEMPTS_PER_ID, } ); assert!( diff --git a/src/sync/mod.rs b/src/sync/mod.rs index af107e8..48956fb 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -3786,6 +3786,18 @@ impl SyncManager { "No missing events recovered - retry scheduled with backoff" ); } + AttemptOutcome::SomeExpired { + expired, + remaining, + attempts, + next_attempt_in, + } => { + tracing::warn!( + relay = %relay_url, expired, remaining, attempts, + next_attempt_in_secs = next_attempt_in.as_secs_f64(), + "Missing-event IDs exhausted their fetch budget; continuing recovery of younger IDs" + ); + } AttemptOutcome::Expired { remaining, attempts, @@ -3794,7 +3806,7 @@ impl SyncManager { relay = %relay_url, remaining, attempts, - "Missing-event recovery expired by policy after repeated zero-progress attempts - relay remains ConnectedHistoricSyncFailures until daily sync" + "Missing-event recovery expired after per-ID fetch budgets were exhausted - relay remains ConnectedHistoricSyncFailures until daily sync" ); } }