From 88ef365397760332c5d50cd46150a95afff31f66 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 12:56:22 +0000 Subject: [PATCH] chore(sync): expose relay event-pipeline pressure The archive reconciliation against relay.ngit.dev hydrates 50,580 IDs and takes long enough that batch completion alone cannot distinguish slow peer delivery, bounded-channel backpressure, duplicate lookup cost, or write-policy work. Per-event logging would make the production signal less usable and add substantial overhead to the workload being measured. Carry the data-lane arrival instant with EVENT and EOSE messages. Emit one per-relay aggregate every 30 seconds with outcome counts, throughput, queue delay, processing time, and current queue depth; separately report EOSE messages delayed at least five seconds in the processor lane. A focused unit test verifies that queue and processing costs remain separate. These measurements deliberately observe rather than change concurrency, queue capacity, or processing order. Resource isolation for foreground relay traffic and adaptive client-side subscription ceilings remain follow-up design work. The timestamps begin when the processor-facing notification lane observes a message, so peer-side delivery time is outside their scope. Validation: nix develop -c cargo test --lib (688 passed). cargo fmt --check remains blocked by pre-existing formatting drift across the stacked branch under the current Rust 1.96 formatter; git diff --check passes. --- src/sync/mod.rs | 145 ++++++++++++++++++++++++++++++++++- src/sync/relay_connection.rs | 19 +++-- 2 files changed, 157 insertions(+), 7 deletions(-) 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() {