fix(sync): skip paused peers before reserving dependency retries

Purgatory polling repeatedly spawned exact-ID fetches against peers already
under subscription cooldown. Connection admission blocked the sends, but each
pass emitted avoidable warnings and reserved dependency IDs without querying.

Filter paused peers, absent connections and empty batches before reserving
retry IDs. Healthy sources remain eligible for shared IDs, and a paused-only
pass cannot delay recovery after the pause ends. Keep the final connection
admission guard for pauses that begin after scheduling; do not change cooldown
policy, in-flight requests or successful dependency processing.

Validation: the new regression fails before the fix and passes afterward.
It covers rate-limit recovery, policy/request pauses, absent connections and
healthy-peer wire delivery. All 10 dependency-related unit tests, workspace
all-target Clippy with warnings denied, formatting and diff checks pass.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-10-02 08:07:58 +00:00
parent 283a30af59
commit 9fd6d7c7f0
3 changed files with 99 additions and 10 deletions
+3
View File
@@ -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
@@ -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
+91 -10
View File
@@ -6994,6 +6994,13 @@ impl SyncManager {
&self,
mut batches: HashMap<String, DependencyRelayBatch>,
) {
// 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")