diff --git a/CHANGELOG.md b/CHANGELOG.md index e7d003b..0f9b6aa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Skip paused peers before scheduling purgatory dependency fetches, avoiding + redundant retries and warnings while preserving retry eligibility on recovery. + - Retire old relay subscriptions before reconnecting so SDK replay cannot accumulate untracked subscriptions or exhaust peers' retained-filter limits. - Remove rolled-back live subscriptions from the SDK registry so authentication diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index c6a1b22..9653db5 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -522,6 +522,11 @@ So concurrency is not a free scaling axis; it is the residual of the ledger: rust-nostr owns the internal `NEG-MSG` exchange, so the application cannot guarantee that each charged frame passes through its pacing gate. Historic recovery then uses the paced REQ path. +- Purgatory dependency polling excludes peers under rate-limit, policy, or + incomplete-request pauses before reserving event-ID retries. A paused-only + pass leaves those IDs immediately eligible when the pause ends; healthy + sources remain eligible for shared IDs. Connection admission still checks + pauses again before sending to cover changes after scheduling. - Failed rate-limit recovery probes double their pause from 65 seconds to at most one hour. Repeated notices during a pause do not extend it. Five unpaused minutes without another refusal reset the next delay; time spent diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 48956fb..398e5cf 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -6994,6 +6994,13 @@ impl SyncManager { &self, mut batches: HashMap, ) { + // Paused peers cannot serve these IDs yet. Do not reserve their retries + // or spawn fetches that the connection admission guard will reject. + batches.retain(|relay_url, batch| { + !batch.event_ids.is_empty() + && self.connections.contains_key(relay_url) + && !self.health_tracker.is_subscription_paused(relay_url) + }); if batches.is_empty() { return; } @@ -7011,9 +7018,7 @@ impl SyncManager { .event_ids .retain(|event_id| due_event_ids.contains(event_id)); } - batches.retain(|relay_url, batch| { - !batch.event_ids.is_empty() && self.connections.contains_key(relay_url) - }); + batches.retain(|_, batch| !batch.event_ids.is_empty()); let connections = self.connections.clone(); let database = self.database.clone(); @@ -9473,6 +9478,15 @@ mod tests { /// publication can straddle ticks and leave their retry deadlines offset. #[tokio::test] async fn unresolved_repositories_share_one_dependency_poll_per_relay() { + dependency_poll_regression(false).await; + } + + #[tokio::test] + async fn paused_dependency_peers_do_not_reserve_retries_or_block_healthy_peers() { + dependency_poll_regression(true).await; + } + + async fn dependency_poll_regression(exercise_pauses: bool) { use futures_util::{SinkExt, StreamExt}; use tokio_tungstenite::tungstenite::Message; @@ -9610,13 +9624,80 @@ mod tests { ); } - // No background manager loop runs: this pass sees both fresh entries. - tokio::time::timeout( - Duration::from_secs(5), - manager.sync_purgatory_announcements_to_index(), - ) - .await - .expect("maintenance must finish"); + if exercise_pauses { + let ids: HashSet<_> = expected_ids + .iter() + .map(|id| EventId::from_hex(id).unwrap()) + .collect(); + let batch = || DependencyRelayBatch { + event_ids: ids.clone(), + identifiers: HashSet::from(["dependency-poll".into()]), + }; + manager.health_tracker.record_rate_limit(&relay_url); + manager.spawn_batched_purgatory_dependency_refetch(HashMap::from([( + relay_url.clone(), + batch(), + )])); + assert!( + manager + .dependency_refetch_attempts + .lock() + .unwrap() + .is_empty(), + "a paused-only pass must not consume the next dependency retry" + ); + + // End the pause without waiting on a wall-clock cooldown. Recovery + // must be immediately eligible, even inside the ID retry interval. + manager.health_tracker.clear_rate_limit(&relay_url); + let mut batches = HashMap::from([(relay_url.clone(), batch())]); + let mut skipped_ids = Vec::new(); + for (n, url) in [ + "wss://policy.example", + "wss://timeout.example", + "wss://absent.example", + ] + .into_iter() + .enumerate() + { + let id = EventId::from_byte_array([n as u8 + 1; 32]); + skipped_ids.push(id); + let mut pending = batch(); // Shared IDs must still use the healthy peer. + pending.event_ids.insert(id); + batches.insert(url.into(), pending); + if n < 2 { + manager.connections.insert( + url.into(), + RelayConnection::new( + url.into(), + None, + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ), + ); + } + if n == 0 { + manager.health_tracker.record_policy_refusal(url); + } else if n == 1 { + manager.health_tracker.record_request_failure(url); + } + } + manager.spawn_batched_purgatory_dependency_refetch(batches); + let attempts = manager.dependency_refetch_attempts.lock().unwrap(); + assert!(ids.iter().all(|id| attempts.contains_key(id))); + assert!( + skipped_ids.iter().all(|id| !attempts.contains_key(id)), + "paused and absent peers must not reserve their unique dependency IDs" + ); + } else { + // No background manager loop runs: this pass sees both fresh entries. + tokio::time::timeout( + Duration::from_secs(5), + manager.sync_purgatory_announcements_to_index(), + ) + .await + .expect("maintenance must finish"); + } let query = tokio::time::timeout(Duration::from_secs(5), query_rx) .await .expect("dependency query deadline")