diff --git a/src/sync/mod.rs b/src/sync/mod.rs index e69266f..894d015 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -430,6 +430,93 @@ pub enum ProcessResult { Rejected, } +/// A low-volume summary of the per-relay data lane. +/// +/// Sync bursts can contain tens of thousands of duplicates, so per-event logs +/// obscure whether time is spent waiting for the bounded channel or applying +/// policy. A periodic aggregate keeps production diagnosis cheap enough to +/// leave enabled while preserving both parts of that distinction. +#[derive(Debug)] +struct EventPipelineWindow { + started_at: std::time::Instant, + delivered: u64, + saved: u64, + duplicate: u64, + purgatory: u64, + rejected: u64, + queue_delay: std::time::Duration, + max_queue_delay: std::time::Duration, + processing_time: std::time::Duration, + max_processing_time: std::time::Duration, +} + +impl Default for EventPipelineWindow { + fn default() -> Self { + Self { + started_at: std::time::Instant::now(), + delivered: 0, + saved: 0, + duplicate: 0, + purgatory: 0, + rejected: 0, + queue_delay: std::time::Duration::ZERO, + max_queue_delay: std::time::Duration::ZERO, + processing_time: std::time::Duration::ZERO, + max_processing_time: std::time::Duration::ZERO, + } + } +} + +impl EventPipelineWindow { + const REPORT_INTERVAL: std::time::Duration = std::time::Duration::from_secs(30); + + fn record( + &mut self, + result: ProcessResult, + queue_delay: std::time::Duration, + processing_time: std::time::Duration, + ) { + self.delivered += 1; + match result { + ProcessResult::Saved => self.saved += 1, + ProcessResult::Duplicate => self.duplicate += 1, + ProcessResult::Purgatory => self.purgatory += 1, + ProcessResult::Rejected => self.rejected += 1, + } + self.queue_delay += queue_delay; + self.max_queue_delay = self.max_queue_delay.max(queue_delay); + self.processing_time += processing_time; + self.max_processing_time = self.max_processing_time.max(processing_time); + } + + fn report_if_due(&mut self, relay: &str, queue_depth: usize) { + let elapsed = self.started_at.elapsed(); + if elapsed < Self::REPORT_INTERVAL || self.delivered == 0 { + return; + } + + let delivered = self.delivered as f64; + tracing::info!( + relay, + window_seconds = elapsed.as_secs_f64(), + delivered = self.delivered, + saved = self.saved, + duplicate = self.duplicate, + purgatory = self.purgatory, + rejected = self.rejected, + 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, + average_processing_ms = self.processing_time.as_secs_f64() * 1000.0 / delivered, + max_processing_ms = self.max_processing_time.as_secs_f64() * 1000.0, + queue_depth, + queue_capacity = relay_connection::RELAY_EVENT_BUFFER_CAPACITY, + "Relay sync event-pipeline window" + ); + *self = Self::default(); + } +} + /// Statistics from re-processing events from hot cache #[derive(Debug, Clone, Default)] pub struct ReprocessingStats { @@ -3546,10 +3633,13 @@ impl SyncManager { tokio::spawn(async move { let mut disconnect_sent = false; + let mut pipeline_window = EventPipelineWindow::default(); while let Some(relay_event) = event_rx.recv().await { match relay_event { - RelayEvent::Event(event, subscription_id) => { + RelayEvent::Event(event, subscription_id, data_lane_arrival) => { + let queue_delay = data_lane_arrival.elapsed(); + let processing_started = std::time::Instant::now(); // Count raw deliveries before deduplication or write policy. Relays spend // their result allowance on every matching delivery, including events we // route to purgatory, reject, or have already stored; pagination must use @@ -3578,6 +3668,12 @@ impl SyncManager { relay = %relay_url_clone, "Skipping previously rejected announcement event" ); + pipeline_window.record( + ProcessResult::Rejected, + queue_delay, + processing_started.elapsed(), + ); + pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); continue; } @@ -3658,8 +3754,21 @@ impl SyncManager { } } } + pipeline_window.record(result, queue_delay, processing_started.elapsed()); + pipeline_window.report_if_due(&relay_url_clone, event_rx.len()); } - RelayEvent::EndOfStoredEvents(sub_id) => { + RelayEvent::EndOfStoredEvents(sub_id, data_lane_arrival) => { + let data_lane_delay = data_lane_arrival.elapsed(); + if data_lane_delay >= std::time::Duration::from_secs(5) { + tracing::info!( + relay = %relay_url_clone, + sub_id = %sub_id, + data_lane_delay_seconds = data_lane_delay.as_secs_f64(), + queue_depth = event_rx.len(), + queue_capacity = relay_connection::RELAY_EVENT_BUFFER_CAPACITY, + "EOSE reached sync manager after data-lane delay" + ); + } tracing::debug!( relay = %relay_url_clone, sub_id = %sub_id, @@ -7024,6 +7133,38 @@ impl SyncManager { mod tests { use super::*; + #[test] + fn event_pipeline_window_separates_queue_and_processing_costs() { + let mut window = EventPipelineWindow::default(); + window.record( + ProcessResult::Duplicate, + std::time::Duration::from_millis(12), + std::time::Duration::from_millis(3), + ); + window.record( + ProcessResult::Saved, + std::time::Duration::from_millis(8), + std::time::Duration::from_millis(7), + ); + + assert_eq!(window.delivered, 2); + assert_eq!(window.saved, 1); + assert_eq!(window.duplicate, 1); + assert_eq!(window.queue_delay, std::time::Duration::from_millis(20)); + assert_eq!( + window.max_queue_delay, + std::time::Duration::from_millis(12) + ); + assert_eq!( + window.processing_time, + std::time::Duration::from_millis(10) + ); + assert_eq!( + window.max_processing_time, + std::time::Duration::from_millis(7) + ); + } + #[test] fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { let keys = Keys::generate(); diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index eed345c..987ea86 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -490,10 +490,10 @@ where /// Events from a relay connection #[derive(Debug)] pub enum RelayEvent { - /// A new event was received (event, subscription_id) - Event(Box, SubscriptionId), + /// A new event was received (event, subscription_id, data-lane arrival). + Event(Box, SubscriptionId, std::time::Instant), /// End of stored events for a subscription - EndOfStoredEvents(SubscriptionId), + EndOfStoredEvents(SubscriptionId, std::time::Instant), /// NOTICE message from relay Notice(String), /// Connection was closed @@ -1655,6 +1655,7 @@ impl RelayConnection { event, subscription_id, } => { + let data_lane_arrival = std::time::Instant::now(); self.record_transient_req_event(&subscription_id); tracing::trace!( relay = %url, @@ -1663,7 +1664,11 @@ impl RelayConnection { "Received event" ); if event_sender - .send(RelayEvent::Event(Box::new(*event), subscription_id.clone())) + .send(RelayEvent::Event( + Box::new(*event), + subscription_id.clone(), + data_lane_arrival, + )) .await .is_err() { @@ -1673,11 +1678,15 @@ impl RelayConnection { } RelayNotification::Message { message } => match *message { RelayMessage::EndOfStoredEvents(sub_id) => { + let data_lane_arrival = std::time::Instant::now(); tracing::debug!(relay = %url, sub_id = ?sub_id, "Received EOSE"); // Convert Cow to owned SubscriptionId let owned_sub_id = sub_id.into_owned(); if event_sender - .send(RelayEvent::EndOfStoredEvents(owned_sub_id)) + .send(RelayEvent::EndOfStoredEvents( + owned_sub_id, + data_lane_arrival, + )) .await .is_err() {