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!(