diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 695346c..ca536f1 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -898,9 +898,12 @@ maintenance timer: (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. -- **Progress-aware**: a successful attempt clears only the IDs actually - recovered and resets the backoff; duplicate incomplete responses merge into - the existing pending set. IDs that arrive by other means (live sync, user +- **Outcome-aware**: a successful attempt clears only IDs with a durable + terminal explanation: saved, already stored, in purgatory, validly + tombstoned, blocked, or permanently invalid. Restricted, policy-error, + unknown, and persistence-error outcomes remain pending because a dependency + or transient fault may clear. Duplicate incomplete responses merge into the + existing pending set. IDs that arrive by other means (live sync, user submission) 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 @@ -1182,6 +1185,10 @@ The [`SyncMetrics`](src/sync/metrics.rs:18) module provides comprehensive monito ### Event Metrics - `ngit_sync_events_synced_total`: Total events synced (newly saved events only, not duplicates or rejected) +- `ngit_sync_hydration_events_total{relay,phase,outcome}`: Remote hydration + work by fixed phase and outcome. Outcomes distinguish requests, deliveries, source + non-delivery, saved/duplicate/purgatory/tombstone states, bounded policy + rejection classes, and persistence failures. ### Summary Metrics diff --git a/docs/explanation/monitoring.md b/docs/explanation/monitoring.md index d9e0123..7d6d939 100644 --- a/docs/explanation/monitoring.md +++ b/docs/explanation/monitoring.md @@ -121,6 +121,7 @@ When GRASP-02 proactive sync is implemented, the following metrics will be added | `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_hydration_events_total` | Counter | relay, phase, outcome | Remote hydration requests, deliveries, and bounded persistence/admission outcomes; phase is `stream` or `recovery` | | `ngit_sync_relays_tracked_total` | Gauge | - | Total relays discovered | | `ngit_sync_relays_connected_total` | Gauge | - | Currently connected relay count | | `ngit_sync_relays_dead_total` | Gauge | - | Relays marked as dead | @@ -181,6 +182,14 @@ sum(rate(ngit_sync_connection_attempts_total{result="success"}[1h])) # Event sync rate (newly saved events) rate(ngit_sync_events_synced_total[5m]) +# Exact-ID responses that a source did not deliver +sum by (relay) (rate(ngit_sync_hydration_events_total{phase="recovery",outcome="not_delivered"}[15m])) + +# Delivered events that were not made servable +sum by (relay, outcome) ( + rate(ngit_sync_hydration_events_total{outcome=~"purgatory|rejected_.*|persistence_error"}[15m]) +) + # Relays with high failure counts (potential issues) topk(10, ngit_sync_relay_failures) diff --git a/src/sync/metrics.rs b/src/sync/metrics.rs index fe063b3..74df374 100644 --- a/src/sync/metrics.rs +++ b/src/sync/metrics.rs @@ -34,6 +34,8 @@ pub struct SyncMetrics { // === Event metrics === /// Total events synced (newly saved events only) events_synced_total: IntCounter, + /// Historic/live hydration deliveries and terminal processing outcomes. + hydration_events_total: IntCounterVec, // === Summary metrics === /// Total relays discovered and tracked @@ -127,6 +129,15 @@ impl SyncMetrics { ))?; registry.register(Box::new(events_synced_total.clone()))?; + let hydration_events_total = IntCounterVec::new( + Opts::new( + "ngit_sync_hydration_events_total", + "Hydration events by relay, phase, and bounded delivery or persistence outcome", + ), + &["relay", "phase", "outcome"], + )?; + registry.register(Box::new(hydration_events_total.clone()))?; + // Summary metrics let relays_tracked_total = IntGauge::with_opts(Opts::new( "ngit_sync_relays_tracked_total", @@ -236,6 +247,7 @@ impl SyncMetrics { relay_failures, policy_refusals_total, events_synced_total, + hydration_events_total, relays_tracked_total, relays_connected_total, relays_dead_total, @@ -402,6 +414,18 @@ impl SyncMetrics { self.events_synced_total.inc(); } + /// Record hydration work using a fixed outcome vocabulary supplied by the + /// sync manager. `count` allows exact-ID attempts to account for a batch + /// without performing one Prometheus update per absent response. + pub fn record_hydration_events(&self, relay: &str, phase: &str, outcome: &str, count: usize) { + if count == 0 { + return; + } + self.hydration_events_total + .with_label_values(&[relay, phase, outcome]) + .inc_by(count as u64); + } + // === Summary Recording Methods === /// Set the total tracked relay count. @@ -643,6 +667,25 @@ mod tests { metrics.record_synced_event(); metrics.record_synced_event(); metrics.record_synced_event(); + + metrics.record_hydration_events("wss://relay.example", "recovery", "requested", 3); + metrics.record_hydration_events("wss://relay.example", "recovery", "saved", 2); + metrics.record_hydration_events("wss://relay.example", "recovery", "not_delivered", 1); + + let hydration = registry + .gather() + .into_iter() + .find(|family| family.name() == "ngit_sync_hydration_events_total") + .expect("hydration outcome metric must be exported"); + assert_eq!(hydration.get_metric().len(), 3); + assert_eq!( + hydration + .get_metric() + .iter() + .map(|metric| metric.get_counter().value() as u64) + .sum::(), + 6 + ); } #[test] diff --git a/src/sync/mod.rs b/src/sync/mod.rs index a796af4..7331096 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -432,8 +432,64 @@ pub enum ProcessResult { Duplicate, /// Event added to Purgatory Purgatory, - /// Event rejected by write policy - Rejected, + /// Event is absent by a valid persisted deletion or vanish request. + Tombstoned, + /// Event rejected by write policy. + Rejected(PolicyRejection), + /// The database could not be read or an accepted event could not be saved. + PersistenceError, +} + +/// Bounded admission-rejection classes used by hydration accounting. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PolicyRejection { + Blocked, + Invalid, + Restricted, + Error, + Other, + PreviouslyRejected, +} + +impl ProcessResult { + fn hydration_outcome(self) -> &'static str { + match self { + Self::Saved => "saved", + Self::Duplicate => "duplicate", + Self::Purgatory => "purgatory", + Self::Tombstoned => "tombstoned", + Self::Rejected(PolicyRejection::Blocked) => "rejected_blocked", + Self::Rejected(PolicyRejection::Invalid) => "rejected_invalid", + Self::Rejected(PolicyRejection::Restricted) => "rejected_restricted", + Self::Rejected(PolicyRejection::Error) => "rejected_error", + Self::Rejected(PolicyRejection::Other) => "rejected_other", + Self::Rejected(PolicyRejection::PreviouslyRejected) => "rejected_cached", + Self::PersistenceError => "persistence_error", + } + } + + fn is_rejected(self) -> bool { + matches!(self, Self::Rejected(_)) + } + + /// Whether this result gives a durable explanation for the requested ID. + /// + /// Restricted, server-error, and unknown policy failures can become valid + /// after dependencies arrive or a transient fault clears, so exact-ID + /// recovery keeps them pending. Invalid/blocked results are permanent for + /// the event, while purgatory and tombstones are explicit terminal states. + fn is_terminally_accounted(self) -> bool { + matches!( + self, + Self::Saved + | Self::Duplicate + | Self::Purgatory + | Self::Tombstoned + | Self::Rejected(PolicyRejection::Blocked) + | Self::Rejected(PolicyRejection::Invalid) + | Self::Rejected(PolicyRejection::PreviouslyRejected) + ) + } } /// A low-volume summary of the per-relay data lane. @@ -449,7 +505,9 @@ struct EventPipelineWindow { saved: u64, duplicate: u64, purgatory: u64, + tombstoned: u64, rejected: u64, + persistence_error: u64, queue_delay: std::time::Duration, max_queue_delay: std::time::Duration, processing_time: std::time::Duration, @@ -464,7 +522,9 @@ impl Default for EventPipelineWindow { saved: 0, duplicate: 0, purgatory: 0, + tombstoned: 0, rejected: 0, + persistence_error: 0, queue_delay: std::time::Duration::ZERO, max_queue_delay: std::time::Duration::ZERO, processing_time: std::time::Duration::ZERO, @@ -487,7 +547,9 @@ impl EventPipelineWindow { ProcessResult::Saved => self.saved += 1, ProcessResult::Duplicate => self.duplicate += 1, ProcessResult::Purgatory => self.purgatory += 1, - ProcessResult::Rejected => self.rejected += 1, + ProcessResult::Tombstoned => self.tombstoned += 1, + ProcessResult::Rejected(_) => self.rejected += 1, + ProcessResult::PersistenceError => self.persistence_error += 1, } self.queue_delay += queue_delay; self.max_queue_delay = self.max_queue_delay.max(queue_delay); @@ -509,7 +571,9 @@ impl EventPipelineWindow { saved = self.saved, duplicate = self.duplicate, purgatory = self.purgatory, + tombstoned = self.tombstoned, rejected = self.rejected, + persistence_error = self.persistence_error, events_per_second = delivered / elapsed.as_secs_f64(), average_queue_delay_ms = self.queue_delay.as_secs_f64() * 1000.0 / delivered, max_queue_delay_ms = self.max_queue_delay.as_secs_f64() * 1000.0, @@ -2870,9 +2934,9 @@ impl SyncManager { /// One bounded exact-ID recovery fetch against a single relay. /// /// Runs outside the sync actor lock. Every event the relay returns is - /// passed through the normal write policy; an ID counts as recovered once - /// the relay has delivered the event, regardless of the policy verdict - /// (rejected events have their own re-processing machinery). + /// passed through the normal write policy. An ID counts as recovered only + /// once it has a durable terminal outcome; transient persistence and + /// dependency-sensitive policy failures remain pending. #[allow(clippy::too_many_arguments)] async fn run_missing_event_recovery_attempt( relay_url: String, @@ -2893,6 +2957,9 @@ impl SyncManager { requested = requested.len(), "Retrying events missing from an incomplete historic sync batch" ); + if let Some(metrics) = metrics.as_ref() { + metrics.record_hydration_events(&relay_url, "recovery", "requested", requested.len()); + } let events = match connection .fetch_events( @@ -2914,11 +2981,18 @@ impl SyncManager { } }; + let mut delivered = HashSet::new(); let mut recovered = HashSet::new(); for event in events { if !requested.contains(&event.id) { continue; } + if !delivered.insert(event.id) { + continue; + } + if let Some(metrics) = metrics.as_ref() { + metrics.record_hydration_events(&relay_url, "recovery", "delivered", 1); + } // Events already tracked as rejected count as recovered without // re-processing: re-validation is owned by the rejected-index // re-processing machinery, and unrecoverable IDs never revalidate. @@ -2928,6 +3002,9 @@ impl SyncManager { event_id = %event.id, "Recovered missing event already tracked as rejected, skipping re-processing" ); + if let Some(metrics) = metrics.as_ref() { + metrics.record_hydration_events(&relay_url, "recovery", "rejected_cached", 1); + } recovered.insert(event.id); continue; } @@ -2941,14 +3018,38 @@ impl SyncManager { crate::nostr::persistence::SaveContext::RelaySync, ) .await; - if result == ProcessResult::Rejected { + if let Some(metrics) = metrics.as_ref() { + metrics.record_hydration_events( + &relay_url, + "recovery", + result.hydration_outcome(), + 1, + ); + if result == ProcessResult::Saved { + metrics.record_synced_event(); + } + } + if result.is_rejected() || result == ProcessResult::PersistenceError { tracing::debug!( relay = %relay_url, event_id = %event.id, - "Recovered missing event was rejected by the write policy" + outcome = result.hydration_outcome(), + terminal = result.is_terminally_accounted(), + "Recovered missing event was not persisted" ); } - recovered.insert(event.id); + if result.is_terminally_accounted() { + recovered.insert(event.id); + } + } + + if let Some(metrics) = metrics.as_ref() { + metrics.record_hydration_events( + &relay_url, + "recovery", + "not_delivered", + requested.len().saturating_sub(delivered.len()), + ); } let outcome = missing_event_recovery.lock().unwrap().complete_attempt( @@ -3832,10 +3933,24 @@ impl SyncManager { "Skipping previously rejected announcement event" ); pipeline_window.record( - ProcessResult::Rejected, + ProcessResult::Rejected(PolicyRejection::PreviouslyRejected), queue_delay, processing_started.elapsed(), ); + if let Some(ref metrics) = metrics_clone { + metrics.record_hydration_events( + &relay_url_clone, + "stream", + "delivered", + 1, + ); + metrics.record_hydration_events( + &relay_url_clone, + "stream", + "rejected_cached", + 1, + ); + } pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); continue; } @@ -3850,6 +3965,20 @@ impl SyncManager { crate::nostr::persistence::SaveContext::RelaySync, ) .await; + if let Some(ref metrics) = metrics_clone { + metrics.record_hydration_events( + &relay_url_clone, + "stream", + "delivered", + 1, + ); + metrics.record_hydration_events( + &relay_url_clone, + "stream", + result.hydration_outcome(), + 1, + ); + } // Only record metric when event is actually saved if result == ProcessResult::Saved { if let Some(ref metrics) = metrics_clone { @@ -3900,7 +4029,7 @@ impl SyncManager { // Track received event IDs for negentropy batches. Unlike REQ+EOSE // pagination above, negentropy completion is concerned with events that // were actually saved or already present locally. - if result == ProcessResult::Saved || result == ProcessResult::Duplicate { + if result.is_terminally_accounted() { let mut pending = pending_sync_index.write().await; if let Some(batches) = pending.get_mut(&relay_url_clone) { for batch in batches.iter_mut() { @@ -5880,7 +6009,7 @@ impl SyncManager { ) .await; - if result != ProcessResult::Rejected { + if result.is_terminally_accounted() { rejected_events_index.remove(&event.id); dependency_refetch_attempts .lock() @@ -6161,7 +6290,15 @@ impl SyncManager { write_policy.purgatory().enqueue_sync_immediate(&identifier); } } - ProcessResult::Rejected => { + ProcessResult::Tombstoned => { + stats.duplicate += 1; + tracing::debug!( + event_id = %event.id, + "{} remains absent under deletion semantics", + context + ); + } + ProcessResult::Rejected(_) | ProcessResult::PersistenceError => { stats.rejected += 1; tracing::warn!( event_id = %event.id, @@ -6173,7 +6310,7 @@ impl SyncManager { } } - if reprocess_result != ProcessResult::Rejected { + if reprocess_result.is_terminally_accounted() { rejected_events_index.remove(&event.id); } } @@ -6190,6 +6327,35 @@ impl SyncManager { /// - Broadcast to WebSocket subscribers via notify_event (enables recursive relay discovery) /// /// Returns `ProcessResult` to indicate whether the event was saved, duplicate, or rejected. + fn classify_policy_result( + prefix: &nostr_sdk::prelude::MachineReadablePrefix, + message: &str, + status: bool, + ) -> ProcessResult { + use nostr_sdk::prelude::MachineReadablePrefix; + + if status { + return if prefix.as_str() == "purgatory" || message.contains("purgatory") { + ProcessResult::Purgatory + } else { + ProcessResult::Duplicate + }; + } + + if message == "this event is deleted" || message == "this pubkey has requested to vanish" { + return ProcessResult::Tombstoned; + } + + let rejection = match prefix { + MachineReadablePrefix::Blocked => PolicyRejection::Blocked, + MachineReadablePrefix::Invalid => PolicyRejection::Invalid, + MachineReadablePrefix::Restricted => PolicyRejection::Restricted, + MachineReadablePrefix::Error => PolicyRejection::Error, + _ => PolicyRejection::Other, + }; + ProcessResult::Rejected(rejection) + } + async fn process_event_static( event: &Event, relay_url: &str, @@ -6209,7 +6375,7 @@ impl SyncManager { } Err(e) => { tracing::warn!(event_id = %event.id, error = %e, "Database error checking event"); - return ProcessResult::Rejected; + return ProcessResult::PersistenceError; } Ok(None) => {} // Continue processing } @@ -6229,7 +6395,7 @@ impl SyncManager { error = %e, "Failed to save synced event" ); - return ProcessResult::Rejected; + return ProcessResult::PersistenceError; } // Broadcast to WebSocket subscribers (enables recursive relay discovery) @@ -6405,23 +6571,31 @@ impl SyncManager { ProcessResult::Saved } WritePolicyResult::Reject { - message, status, .. + prefix, + message, + status, } => { - if status { + let process_result = Self::classify_policy_result(&prefix, &message, status); + if matches!( + process_result, + ProcessResult::Purgatory | ProcessResult::Duplicate + ) { tracing::debug!( event_id = %event.id, kind = %event.kind.as_u16(), + outcome = process_result.hydration_outcome(), reason = %message, - "Event added to purgatory" + "Event accepted without main-database persistence" ); // Note: git data sync for state events is triggered by the policy // layer when adding to purgatory (via start_state_sync) - ProcessResult::Purgatory + process_result } else { tracing::debug!( event_id = %event.id, relay = %relay_url, kind = %event.kind.as_u16(), + outcome = process_result.hydration_outcome(), reason = %message, "Event rejected by write policy" ); @@ -6498,7 +6672,7 @@ impl SyncManager { } } - ProcessResult::Rejected + process_result } } } @@ -7899,6 +8073,50 @@ mod tests { ); } + #[test] + fn hydration_outcomes_separate_terminal_and_retryable_failures() { + use nostr::message::relay::SingleWord; + use nostr_sdk::prelude::MachineReadablePrefix; + + let purgatory = MachineReadablePrefix::Custom( + SingleWord::from_static("purgatory").expect("valid custom prefix"), + ); + assert_eq!( + SyncManager::classify_policy_result( + &purgatory, + "won't be served until git data arrives", + true, + ), + ProcessResult::Purgatory + ); + assert_eq!( + SyncManager::classify_policy_result( + &MachineReadablePrefix::Invalid, + "this event is deleted", + false, + ), + ProcessResult::Tombstoned + ); + + let invalid = SyncManager::classify_policy_result( + &MachineReadablePrefix::Invalid, + "malformed event", + false, + ); + let restricted = SyncManager::classify_policy_result( + &MachineReadablePrefix::Restricted, + "dependency not accepted yet", + false, + ); + let persistence_error = ProcessResult::PersistenceError; + + assert_eq!(invalid.hydration_outcome(), "rejected_invalid"); + assert!(invalid.is_terminally_accounted()); + assert_eq!(restricted.hydration_outcome(), "rejected_restricted"); + assert!(!restricted.is_terminally_accounted()); + assert!(!persistence_error.is_terminally_accounted()); + } + #[test] fn lifecycle_inbox_accepts_production_sized_terminal_burst() { let (tx, mut rx) = lifecycle_notification_channel::(); @@ -8460,7 +8678,7 @@ mod tests { .custom_created_at(Timestamp::from_secs(10)) .finalize(&keys) .expect("build rejected event"), - ProcessResult::Rejected, + ProcessResult::Rejected(PolicyRejection::Restricted), ), ]; @@ -8469,7 +8687,7 @@ mod tests { pagination.record_event(&event); assert!(matches!( policy_result, - ProcessResult::Purgatory | ProcessResult::Rejected + ProcessResult::Purgatory | ProcessResult::Rejected(_) )); }