From 3733846aaf0efeed947862408893812a3a5e2c9b Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 8 Aug 2026 15:03:55 +0000 Subject: [PATCH] fix(sync): defer permanently refused coverage Archive-mode production burn-in exposed a self-amplifying live-subscription loop: permanent CLOSED responses rebuilt complete coverage immediately, and each replacement generated more refusals. More than 11,000 CLOSED/rebuild cycles were observed in one minute while the relay connection itself remained healthy. Classify blocked, restricted, membership-required, and filter-incompatible responses into bounded categories. Preserve the WebSocket session, remove any rejected transient batch, pause new subscription work through the existing relay health tracker, and issue one coverage probe after the dead-relay 24-hour cadence. Export the bounded category as a counter and retain the raw reason only in logs. Correctness assumes these explicit policy messages are not transient rate limits; rate-limit-shaped messages retain their shorter adaptive cooldown. Repeated refusals do not extend the daily deadline, connection handshakes do not clear it, and expiry yields exactly one recovery probe. Adaptive filter-count learning and per-repository Git fetch single-flight are deliberately excluded. The former needs relay-implementation evidence; the latter is tracked separately. Validation: nix develop -c cargo test --lib passed all 669 tests; focused policy classification and daily-probe tests passed; nix develop -c cargo check --lib passed. Repository-wide cargo fmt --check still reports pre-existing formatting drift in unrelated integration fixtures. --- CHANGELOG.md | 4 + docs/explanation/grasp-02-proactive-sync.md | 11 +- docs/explanation/monitoring.md | 3 +- src/sync/health.rs | 100 ++++++++++++++++-- src/sync/metrics.rs | 22 +++- src/sync/mod.rs | 109 ++++++++++++++++++-- 6 files changed, 228 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bd15449..3340d4c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,6 +20,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Stop permanently refused live subscriptions from rebuilding themselves in a + tight loop. Blocked, restricted, membership-required, and filter-incompatible + responses retain the connection, retire the rejected coverage for 24 hours, + then receive one recovery probe; bounded metrics distinguish the policies. - Keep purgatory announcements, state events, and PR events alive while a concrete background Git sync for their repository is actively running. A large healthy clone can exceed the nominal 30-minute purgatory window; diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 94c80ec..f9a894a 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -89,7 +89,12 @@ pub struct RepoSyncNeeds { Ownership removal is deliberately passive. It does not close a healthy connection or replace a working live request. If a live request receives `CLOSED`, its replacement filters are derived from the current index and omit -obsolete items. When the connection itself ends naturally, confirmed state is +obsolete items. Definitive policy refusals (`blocked`, `restricted`, membership +required, or incompatible filters) are the exception: the connection remains +open, rejected coverage is not immediately recreated, and one recovery probe +runs after 24 hours. The bounded refusal category is exported as a metric while +the relay's full reason remains in logs. When the connection itself ends +naturally, confirmed state is reconciled before reconnect: shared relays reconnect with current items only, while an unreferenced relay is retired instead of reconnected. This preserves continuous live coverage while allowing stale ownership to drain at ordinary @@ -371,7 +376,9 @@ The sync system uses three background tasks that run continuously: 2. **Retry disconnected**: Calls `retry_disconnected_relays()` to attempt reconnection per health tracker backoff while either confirmed or desired `RepoSyncIndex` work remains -3. **Rate limit recovery**: Calls `check_rate_limit_recovery()` to clear expired rate limits +3. **Subscription recovery**: Calls `check_rate_limit_recovery()` to clear + expired short rate-limit cooldowns and schedule one probe after a 24-hour + policy-refusal pause 4. **Metrics update**: Updates Prometheus metrics with current health states **Why combined**: The 2-second interval provides good responsiveness for health changes while minimizing overhead. All operations are lightweight (index checks, no I/O except actual connection attempts). diff --git a/docs/explanation/monitoring.md b/docs/explanation/monitoring.md index 1a284b5..d9e0123 100644 --- a/docs/explanation/monitoring.md +++ b/docs/explanation/monitoring.md @@ -117,7 +117,8 @@ When GRASP-02 proactive sync is implemented, the following metrics will be added |--------|------|--------|-------------| | `ngit_sync_relay_connected` | Gauge | relay | Connection status (0=disconnected, 1=connecting, 2=syncing, 3=connected, 4=connected_historic_sync_failures) | | `ngit_sync_connection_attempts_total` | Counter | relay, result | Connection attempt outcomes | -| `ngit_sync_relay_status` | Gauge | relay | Health status (1=healthy, 2=disconnected, 3=degraded, 4=dead, 5=rate_limited) | +| `ngit_sync_relay_status` | Gauge | relay | Health status (1=healthy, 2=disconnected, 3=degraded, 4=dead, 5=rate_limited, 6=policy_limited) | +| `ngit_sync_policy_refusals_total` | Counter | relay, category | Subscription policy refusals using bounded categories; raw reasons remain in logs | | `ngit_sync_relay_failures` | Gauge | relay | Current consecutive failure count | | `ngit_sync_events_synced_total` | Counter | - | Events synced (newly saved events only) | | `ngit_sync_relays_tracked_total` | Gauge | - | Total relays discovered | diff --git a/src/sync/health.rs b/src/sync/health.rs index a91537f..067634f 100644 --- a/src/sync/health.rs +++ b/src/sync/health.rs @@ -28,6 +28,11 @@ const DEAD_THRESHOLD_HOURS: u64 = 24; /// How often dead relays are retried (once per 24 hours) const DEAD_RETRY_INTERVAL_HOURS: u64 = 24; +/// Policy-refusing relays are probed once per day without closing an otherwise +/// useful connection. This reuses the dead-relay cadence while keeping policy +/// refusal distinct from transport failure. +const POLICY_RETRY_INTERVAL_HOURS: u64 = 24; + /// Default maximum backoff duration in seconds (1 hour) const DEFAULT_MAX_BACKOFF_SECS: u64 = 3600; @@ -54,6 +59,8 @@ pub enum HealthState { Dead, /// Rate limited by relay, temporary cooldown active RateLimited, + /// Connected, but the relay currently refuses our subscription policy. + PolicyLimited, } impl std::fmt::Display for HealthState { @@ -64,6 +71,7 @@ impl std::fmt::Display for HealthState { HealthState::Degraded => write!(f, "degraded"), HealthState::Dead => write!(f, "dead"), HealthState::RateLimited => write!(f, "rate_limited"), + HealthState::PolicyLimited => write!(f, "policy_limited"), } } } @@ -75,6 +83,8 @@ pub struct RelayHealth { pub connected: bool, /// Has this relay sent us a rate-limiting NOTICE recently pub rate_limited: bool, + /// Deadline for retrying subscriptions refused by relay policy. + pub policy_retry_at: Option, /// Number of consecutive connection failures pub consecutive_failures: u32, /// Time of the first failure in the current failure streak @@ -105,14 +115,19 @@ impl RelayHealth { /// /// ## State Logic /// - /// 1. **RateLimited**: If the rate-limit cooldown hasn't expired - /// 2. **Dead**: 24+ hours of continuous failures - /// 3. **Degraded**: Active connection failures OR in stability period after recovery - /// 4. **Disconnected**: Not connected, but no recent failures or issues - /// 5. **Healthy**: Connected and stable (past stability period with no failures) + /// 1. **PolicyLimited**: If a relay-policy recovery probe is pending + /// 2. **RateLimited**: If the rate-limit cooldown hasn't expired + /// 3. **Dead**: 24+ hours of continuous failures + /// 4. **Degraded**: Active connection failures OR in stability period after recovery + /// 5. **Disconnected**: Not connected, but no recent failures or issues + /// 6. **Healthy**: Connected and stable (past stability period with no failures) pub fn state(&self) -> HealthState { let now = Instant::now(); + if self.policy_retry_at.is_some_and(|deadline| now < deadline) { + return HealthState::PolicyLimited; + } + // Check rate limiting first (highest priority) if self.rate_limited { if let Some(next_retry) = self.next_retry_at { @@ -180,6 +195,11 @@ impl RelayHealth { } } + pub fn is_policy_limited_now(&self) -> bool { + self.policy_retry_at + .is_some_and(|deadline| Instant::now() < deadline) + } + /// Get the consecutive failure count pub fn failure_count(&self) -> u32 { self.consecutive_failures @@ -465,6 +485,46 @@ impl RelayHealthTracker { ); } + /// Pause new subscription work after a definitive relay-policy refusal. + /// Existing connections remain available and one recovery probe is made + /// after the same 24-hour interval used for dead relays. + pub fn record_policy_refusal(&self, relay_url: &str) { + let now = Instant::now(); + let mut entry = self.health.entry(relay_url.to_string()).or_default(); + let health = entry.value_mut(); + if health + .policy_retry_at + .is_some_and(|deadline| now < deadline) + { + return; + } + health.policy_retry_at = + Some(now + Duration::from_secs(POLICY_RETRY_INTERVAL_HOURS * 3600)); + } + + pub fn is_subscription_paused(&self, relay_url: &str) -> bool { + self.health + .get(relay_url) + .is_some_and(|entry| entry.is_rate_limited_now() || entry.is_policy_limited_now()) + } + + /// Clear expired policy refusals and return relays due one recovery probe. + pub fn exit_expired_policy_refusals(&self) -> Vec { + let now = Instant::now(); + let mut recovered = Vec::new(); + for mut entry in self.health.iter_mut() { + let (url, health) = entry.pair_mut(); + if health + .policy_retry_at + .is_some_and(|deadline| now >= deadline) + { + health.policy_retry_at = None; + recovered.push(url.clone()); + } + } + recovered + } + /// Clear rate limiting state for a specific relay /// /// This clears the rate-limit episode without affecting connection status @@ -557,7 +617,9 @@ impl RelayHealthTracker { // Check state-based logic match health.state() { - HealthState::Healthy | HealthState::Disconnected => true, + HealthState::Healthy + | HealthState::Disconnected + | HealthState::PolicyLimited => true, HealthState::Degraded | HealthState::Dead | HealthState::RateLimited => { // Check if backoff/cooldown period has elapsed match health.next_retry_at { @@ -922,4 +984,30 @@ mod tests { assert_eq!(health.next_retry_at, deadline); assert!(health.connected); } + + #[test] + fn policy_refusal_keeps_connection_and_allows_one_daily_probe() { + let tracker = RelayHealthTracker::with_defaults(); + let relay = "wss://members-only.example"; + + tracker.record_success(relay); + tracker.record_policy_refusal(relay); + let first_deadline = tracker.get_health(relay).unwrap().policy_retry_at; + tracker.record_policy_refusal(relay); + + let health = tracker.get_health(relay).unwrap(); + assert!(health.connected); + assert_eq!(health.state(), HealthState::PolicyLimited); + assert_eq!(health.policy_retry_at, first_deadline); + assert!(tracker.is_subscription_paused(relay)); + + tracker.health.get_mut(relay).unwrap().policy_retry_at = + Some(Instant::now() - Duration::from_millis(1)); + assert_eq!( + tracker.exit_expired_policy_refusals(), + vec![relay.to_string()] + ); + assert!(!tracker.is_subscription_paused(relay)); + assert!(tracker.exit_expired_policy_refusals().is_empty()); + } } diff --git a/src/sync/metrics.rs b/src/sync/metrics.rs index 3df63cf..175b75c 100644 --- a/src/sync/metrics.rs +++ b/src/sync/metrics.rs @@ -28,6 +28,8 @@ pub struct SyncMetrics { relay_status: IntGaugeVec, /// Per-relay consecutive failure count relay_failures: IntGaugeVec, + /// Subscription policy refusals by relay and bounded category. + policy_refusals_total: IntCounterVec, // === Event metrics === /// Total events synced (newly saved events only) @@ -94,7 +96,7 @@ impl SyncMetrics { let relay_status = IntGaugeVec::new( Opts::new( "ngit_sync_relay_status", - "Relay health status (1=healthy, 2=disconnected, 3=degraded, 4=dead, 5=rate_limited)", + "Relay health status (1=healthy, 2=disconnected, 3=degraded, 4=dead, 5=rate_limited, 6=policy_limited)", ), &["relay"], )?; @@ -109,6 +111,15 @@ impl SyncMetrics { )?; registry.register(Box::new(relay_failures.clone()))?; + let policy_refusals_total = IntCounterVec::new( + Opts::new( + "ngit_sync_policy_refusals_total", + "Subscription refusals by relay and bounded policy category", + ), + &["relay", "category"], + )?; + registry.register(Box::new(policy_refusals_total.clone()))?; + // Event metrics let events_synced_total = IntCounter::with_opts(Opts::new( "ngit_sync_events_synced_total", @@ -223,6 +234,7 @@ impl SyncMetrics { connection_attempts_total, relay_status, relay_failures, + policy_refusals_total, events_synced_total, relays_tracked_total, relays_connected_total, @@ -297,6 +309,7 @@ impl SyncMetrics { /// - Degraded = 3 (connection problems or unstable after recovery) /// - Dead = 4 (24h+ of failures) /// - RateLimited = 5 (rate limit cooldown active) + /// - PolicyLimited = 6 (relay policy refusal; daily probe pending) /// /// # Arguments /// @@ -309,6 +322,7 @@ impl SyncMetrics { HealthState::Degraded => 3, HealthState::Dead => 4, HealthState::RateLimited => 5, + HealthState::PolicyLimited => 6, }; self.relay_status .with_label_values(&[relay]) @@ -357,6 +371,12 @@ impl SyncMetrics { .set(count as i64); } + pub fn record_policy_refusal(&self, relay: &str, category: &str) { + self.policy_refusals_total + .with_label_values(&[relay, category]) + .inc(); + } + /// Update dead relay count. pub fn update_dead_count(&self, count: i64) { self.relays_dead_total.set(count); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 1363189..14bf9fe 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -791,6 +791,46 @@ fn is_rate_limit_message(message: &str) -> bool { || message.contains("throttl") } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum PolicyRefusal { + Blocked, + Restricted, + MembershipRequired, + FilterIncompatible, +} + +impl PolicyRefusal { + fn label(self) -> &'static str { + match self { + Self::Blocked => "blocked", + Self::Restricted => "restricted", + Self::MembershipRequired => "membership_required", + Self::FilterIncompatible => "filter_incompatible", + } + } +} + +fn policy_refusal(message: &str) -> Option { + if is_rate_limit_message(message) { + return None; + } + let message = message.to_ascii_lowercase(); + if message.contains("invalid number of filters") + || message.contains("filter validation failed") + || message.contains("unsupported filter") + { + Some(PolicyRefusal::FilterIncompatible) + } else if message.contains("not a member") || message.contains("membership") { + Some(PolicyRefusal::MembershipRequired) + } else if message.contains("restricted") { + Some(PolicyRefusal::Restricted) + } else if message.contains("blocked") || message.contains("request rejected") { + Some(PolicyRefusal::Blocked) + } else { + None + } +} + fn subscription_state_byte_limit(message: &str) -> Option { let lower = message.to_ascii_lowercase(); let marker = "active subscriptions exceed max size "; @@ -1612,7 +1652,7 @@ impl SyncManager { // A NOTICE can arrive immediately before this page's EOSE. // Keep a sentinel in the batch and let a detached worker // resume the exact grouped page after the cooldown. - if self.health_tracker.is_rate_limited(relay_url) { + if self.health_tracker.is_subscription_paused(relay_url) { let deferred_sub_id = mark_deferred_pagination(batch, &sub_id); drop(pending); @@ -3028,7 +3068,7 @@ impl SyncManager { } // Step 2: Check if relay is rate-limited before creating new pending items - if self.health_tracker.is_rate_limited(&action.relay_url) { + if self.health_tracker.is_subscription_paused(&action.relay_url) { tracing::debug!( relay = %action.relay_url, full_repos = action.items.repos.len(), @@ -5310,6 +5350,31 @@ impl SyncManager { return; } + if let Some(category) = policy_refusal(reason) { + let removed_batch = { + let mut pending = self.pending_sync_index.write().await; + take_batch_containing_subscription(&mut pending, relay_url, &subscription_id) + }; + self.health_tracker.record_policy_refusal(relay_url); + if let Some(metrics) = &self.metrics { + metrics.record_policy_refusal(relay_url, category.label()); + metrics.record_health_state( + relay_url, + self.health_tracker.get_state(relay_url), + ); + } + tracing::warn!( + relay = %relay_url, + sub_id = %subscription_id, + category = category.label(), + reason, + retry_hours = 24, + pending_batch_removed = removed_batch.is_some(), + "Relay policy refused subscription; preserving connection and deferring coverage probe" + ); + return; + } + if let Some(generation) = live_generation { self.restore_live_coverage_after_closed(relay_url, generation) .await; @@ -5593,7 +5658,13 @@ impl SyncManager { use crate::sync::algorithms::{compute_actions, derive_relay_targets}; // Exit rate limiting for relays whose cooldown has expired - let relays_to_recover: Vec = self.health_tracker.exit_expired_rate_limits(); + let mut relays_to_recover: Vec = self.health_tracker.exit_expired_rate_limits(); + let policy_relays = self.health_tracker.exit_expired_policy_refusals(); + for relay in policy_relays { + if !relays_to_recover.contains(&relay) { + relays_to_recover.push(relay); + } + } if relays_to_recover.is_empty() { return; @@ -5605,13 +5676,7 @@ impl SyncManager { drop(repo_index); for relay_url in relays_to_recover { - tracing::info!( - relay = %relay_url, - "Rate limit cooldown expired, recovering" - ); - - // Clear rate limit state - self.health_tracker.clear_rate_limit(&relay_url); + tracing::info!(relay = %relay_url, "Subscription pause expired; probing coverage"); // A rate-limited CLOSED deliberately leaves live coverage down so // it cannot immediately recreate the rejected request burst. @@ -5644,7 +5709,7 @@ impl SyncManager { full_repo_count = action.items.repos.len(), state_only_repo_count = action.items.state_only_repos.len(), event_count = action.items.root_events.len(), - "Submitting recovered actions after rate limit" + "Submitting recovered actions after subscription pause" ); self.handle_new_sync_filters(action).await; } @@ -6742,6 +6807,28 @@ mod tests { assert!(!is_rate_limit_message("blocked: unsupported filter")); } + #[test] + fn policy_refusal_classifier_uses_bounded_categories() { + assert_eq!( + policy_refusal("restricted: this relay does not accept REQs"), + Some(PolicyRefusal::Restricted) + ); + assert_eq!( + policy_refusal("restricted: you are not a member of this relay"), + Some(PolicyRefusal::MembershipRequired) + ); + assert_eq!( + policy_refusal("blocked: Request rejected"), + Some(PolicyRefusal::Blocked) + ); + assert_eq!( + policy_refusal("ERROR: filter validation failed: invalid number of filters: 8"), + Some(PolicyRefusal::FilterIncompatible) + ); + assert_eq!(policy_refusal("rate-limited: too many queries"), None); + assert_eq!(policy_refusal("error: temporary backend failure"), None); + } + #[test] fn subscription_state_limit_parser_is_specific_and_extracts_bytes() { assert_eq!(