From 245f3b1298ec9a8db098e9334ab823c0f4e40776 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 7 Jun 2026 19:15:46 +0300 Subject: [PATCH] [ingest] add stats for relay worker shards --- docs/api/firehose.md | 37 ++++++ src/api/firehose.rs | 16 ++- src/control/firehose.rs | 16 ++- src/control/mod.rs | 2 + src/ingest/firehose_stats.rs | 217 +++++++++++++++++++++++++++++++++++ src/ingest/relay.rs | 46 ++++++++ 6 files changed, 331 insertions(+), 3 deletions(-) diff --git a/docs/api/firehose.md b/docs/api/firehose.md index 85d3d06..d0160fe 100644 --- a/docs/api/firehose.md +++ b/docs/api/firehose.md @@ -86,6 +86,43 @@ all filters are exact-match and optional. multiple filters are combined with log | `failing` | filter by whether hydrant has recorded failures/backoff state for the source | | `throttled` | filter by whether the source is currently inside a retry backoff window | +## GET /firehose/diagnostics + +available only in builds compiled with the `firehose-diagnostics` feature. returns process-local diagnostics for firehose internals: + +```json +{ + "relay_worker": { + "shards": [ + { + "id": 0, + "received_messages": 1000, + "info_messages": 0, + "processed_messages": 999, + "process_errors": 0, + "commit_errors": 0, + "process_message_micros": 100000, + "stage_counts_micros": 20000, + "stage_and_commit_micros": 900000, + "apply_counts_micros": 30000, + "broadcast_micros": 20000, + "cursor_micros": 1000, + "total_micros": 1100000, + "max_process_message_micros": 1000, + "max_stage_and_commit_micros": 10000, + "max_total_micros": 12000, + "last_received_at": 1717239958, + "last_processed_at": 1717239958, + "last_seq": 1322, + "max_seq": 1322 + } + ] + } +} +``` + +poll this endpoint together with `/firehose/sources`. if source `send_wait_micros` is rising while relay-worker `stage_and_commit_micros` or `process_message_micros` dominates, the bottleneck is downstream of websocket reading rather than source connectivity. + ## GET /firehose/source fetch a single firehose source with the same runtime fields as `GET /firehose/sources`. diff --git a/src/api/firehose.rs b/src/api/firehose.rs index dd050bc..0b2645e 100644 --- a/src/api/firehose.rs +++ b/src/api/firehose.rs @@ -10,12 +10,17 @@ use url::Url; use crate::control::{FirehoseSourceInfo, Hydrant}; pub fn router() -> Router { - Router::new() + let router = Router::new() .route("/firehose/source", get(get_source)) .route("/firehose/sources", get(list_sources)) .route("/firehose/sources", post(add_source)) .route("/firehose/sources", delete(remove_source)) - .route("/firehose/cursors", delete(reset_cursor)) + .route("/firehose/cursors", delete(reset_cursor)); + + #[cfg(feature = "firehose-diagnostics")] + let router = router.route("/firehose/diagnostics", get(get_diagnostics)); + + router } #[derive(Debug, Deserialize)] @@ -103,6 +108,13 @@ pub async fn list_sources( Json(sources) } +#[cfg(feature = "firehose-diagnostics")] +pub async fn get_diagnostics( + State(hydrant): State, +) -> Json { + Json(hydrant.firehose.diagnostics()) +} + #[derive(Deserialize)] pub struct AddSourceRequest { pub url: Url, diff --git a/src/control/firehose.rs b/src/control/firehose.rs index b72eeea..6d4bfa5 100644 --- a/src/control/firehose.rs +++ b/src/control/firehose.rs @@ -11,7 +11,7 @@ use url::Url; use crate::config::FirehoseSource; use crate::db::{self, keys}; #[cfg(feature = "firehose-diagnostics")] -use crate::ingest::firehose_stats::FirehoseStatsSnapshot; +use crate::ingest::firehose_stats::{FirehoseStatsSnapshot, RelayWorkerStatsSnapshot}; use crate::ingest::{BufferTx, firehose::FirehoseIngestor}; use crate::state::AppState; @@ -74,6 +74,13 @@ pub struct FirehoseSourceInfo { pub pds: Option, } +/// feature-gated runtime diagnostics for firehose ingestion internals. +#[cfg(feature = "firehose-diagnostics")] +#[derive(Debug, Clone, serde::Serialize)] +pub struct FirehoseDiagnosticsInfo { + pub relay_worker: RelayWorkerStatsSnapshot, +} + /// runtime control over the firehose ingestor component. #[derive(Clone)] pub struct FirehoseHandle { @@ -252,6 +259,13 @@ impl FirehoseHandle { out } + #[cfg(feature = "firehose-diagnostics")] + pub fn diagnostics(&self) -> FirehoseDiagnosticsInfo { + FirehoseDiagnosticsInfo { + relay_worker: self.state.firehose_stats.relay_worker_snapshot(), + } + } + fn pds_info(&self, url: &Url, meta: &crate::pds_meta::PdsMeta) -> Option { let host = url.host_str()?; let seq = self diff --git a/src/control/mod.rs b/src/control/mod.rs index 7086a3f..369b008 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -24,6 +24,8 @@ mod jetstream; pub use jetstream::*; pub use filter::{FilterControl, FilterPatch, FilterSnapshot}; +#[cfg(feature = "firehose-diagnostics")] +pub use firehose::FirehoseDiagnosticsInfo; pub use firehose::{FirehoseHandle, FirehoseSourceInfo}; pub use pds::{PdsControl, PdsTierAssignment, PdsTierDefinition}; pub use repos::{ListedRecord, Record, RecordList, RepoHandle, RepoInfo, ReposControl}; diff --git a/src/ingest/firehose_stats.rs b/src/ingest/firehose_stats.rs index 0a224d5..49909c2 100644 --- a/src/ingest/firehose_stats.rs +++ b/src/ingest/firehose_stats.rs @@ -1,3 +1,4 @@ +use std::collections::BTreeMap; use std::sync::Arc; use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::time::Duration; @@ -9,6 +10,7 @@ use url::Url; #[derive(Default)] pub struct FirehoseStats { sources: scc::HashMap>, + relay_worker: RelayWorkerStats, } impl FirehoseStats { @@ -24,6 +26,14 @@ impl FirehoseStats { pub fn snapshot(&self, url: &Url) -> Option { self.sources.read_sync(url, |_, stats| stats.snapshot()) } + + pub fn relay_shard(&self, id: usize) -> Arc { + self.relay_worker.shard(id) + } + + pub fn relay_worker_snapshot(&self) -> RelayWorkerStatsSnapshot { + self.relay_worker.snapshot() + } } #[derive(Default)] @@ -261,6 +271,170 @@ pub struct FirehoseMessageStats { pub info: u64, } +#[derive(Default)] +pub struct RelayWorkerStats { + shards: Mutex>>, +} + +impl RelayWorkerStats { + fn shard(&self, id: usize) -> Arc { + let mut shards = self.shards.lock(); + shards + .entry(id) + .or_insert_with(|| Arc::new(RelayShardStats::default())) + .clone() + } + + fn snapshot(&self) -> RelayWorkerStatsSnapshot { + let shards = self + .shards + .lock() + .iter() + .map(|(&id, stats)| stats.snapshot(id)) + .collect(); + RelayWorkerStatsSnapshot { shards } + } +} + +#[derive(Default)] +pub struct RelayShardStats { + received_messages: AtomicU64, + info_messages: AtomicU64, + processed_messages: AtomicU64, + process_errors: AtomicU64, + commit_errors: AtomicU64, + process_message_micros: AtomicU64, + stage_counts_micros: AtomicU64, + stage_and_commit_micros: AtomicU64, + apply_counts_micros: AtomicU64, + broadcast_micros: AtomicU64, + cursor_micros: AtomicU64, + total_micros: AtomicU64, + max_process_message_micros: AtomicU64, + max_stage_and_commit_micros: AtomicU64, + max_total_micros: AtomicU64, + last_received_at: AtomicI64, + last_processed_at: AtomicI64, + last_error_at: AtomicI64, + last_seq: AtomicI64, + max_seq: AtomicI64, +} + +impl RelayShardStats { + pub fn record_received(&self, seq: i64) { + self.received_messages.fetch_add(1, Ordering::Relaxed); + self.last_received_at.store(now_ts(), Ordering::Relaxed); + self.last_seq.store(seq, Ordering::Relaxed); + self.max_seq.fetch_max(seq, Ordering::Relaxed); + } + + pub fn record_info(&self) { + self.info_messages.fetch_add(1, Ordering::Relaxed); + } + + pub fn record_process_error(&self) { + self.process_errors.fetch_add(1, Ordering::Relaxed); + self.last_error_at.store(now_ts(), Ordering::Relaxed); + } + + pub fn record_commit_error(&self) { + self.commit_errors.fetch_add(1, Ordering::Relaxed); + self.last_error_at.store(now_ts(), Ordering::Relaxed); + } + + pub fn record_processed(&self, timings: RelayShardTimings) { + self.processed_messages.fetch_add(1, Ordering::Relaxed); + self.last_processed_at.store(now_ts(), Ordering::Relaxed); + add_duration_with_max( + &self.process_message_micros, + &self.max_process_message_micros, + timings.process_message, + ); + add_duration(&self.stage_counts_micros, timings.stage_counts); + add_duration_with_max( + &self.stage_and_commit_micros, + &self.max_stage_and_commit_micros, + timings.stage_and_commit, + ); + add_duration(&self.apply_counts_micros, timings.apply_counts); + add_duration(&self.broadcast_micros, timings.broadcast); + add_duration(&self.cursor_micros, timings.cursor); + add_duration_with_max(&self.total_micros, &self.max_total_micros, timings.total); + } + + fn snapshot(&self, id: usize) -> RelayShardStatsSnapshot { + RelayShardStatsSnapshot { + id, + received_messages: self.received_messages.load(Ordering::Relaxed), + info_messages: self.info_messages.load(Ordering::Relaxed), + processed_messages: self.processed_messages.load(Ordering::Relaxed), + process_errors: self.process_errors.load(Ordering::Relaxed), + commit_errors: self.commit_errors.load(Ordering::Relaxed), + process_message_micros: self.process_message_micros.load(Ordering::Relaxed), + stage_counts_micros: self.stage_counts_micros.load(Ordering::Relaxed), + stage_and_commit_micros: self.stage_and_commit_micros.load(Ordering::Relaxed), + apply_counts_micros: self.apply_counts_micros.load(Ordering::Relaxed), + broadcast_micros: self.broadcast_micros.load(Ordering::Relaxed), + cursor_micros: self.cursor_micros.load(Ordering::Relaxed), + total_micros: self.total_micros.load(Ordering::Relaxed), + max_process_message_micros: self.max_process_message_micros.load(Ordering::Relaxed), + max_stage_and_commit_micros: self.max_stage_and_commit_micros.load(Ordering::Relaxed), + max_total_micros: self.max_total_micros.load(Ordering::Relaxed), + last_received_at: nonzero_i64(&self.last_received_at), + last_processed_at: nonzero_i64(&self.last_processed_at), + last_error_at: nonzero_i64(&self.last_error_at), + last_seq: nonzero_i64(&self.last_seq), + max_seq: nonzero_i64(&self.max_seq), + } + } +} + +#[derive(Default)] +pub struct RelayShardTimings { + pub process_message: Duration, + pub stage_counts: Duration, + pub stage_and_commit: Duration, + pub apply_counts: Duration, + pub broadcast: Duration, + pub cursor: Duration, + pub total: Duration, +} + +#[derive(Debug, Clone, Serialize)] +pub struct RelayWorkerStatsSnapshot { + pub shards: Vec, +} + +#[derive(Debug, Clone, Serialize)] +pub struct RelayShardStatsSnapshot { + pub id: usize, + pub received_messages: u64, + pub info_messages: u64, + pub processed_messages: u64, + pub process_errors: u64, + pub commit_errors: u64, + pub process_message_micros: u64, + pub stage_counts_micros: u64, + pub stage_and_commit_micros: u64, + pub apply_counts_micros: u64, + pub broadcast_micros: u64, + pub cursor_micros: u64, + pub total_micros: u64, + pub max_process_message_micros: u64, + pub max_stage_and_commit_micros: u64, + pub max_total_micros: u64, + #[serde(skip_serializing_if = "Option::is_none")] + pub last_received_at: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub last_processed_at: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub last_error_at: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub last_seq: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub max_seq: Option, +} + fn now_ts() -> i64 { chrono::Utc::now().timestamp() } @@ -269,6 +443,17 @@ fn duration_micros(duration: Duration) -> u64 { duration.as_micros().try_into().unwrap_or(u64::MAX) } +fn add_duration(total: &AtomicU64, duration: Duration) { + let micros = duration_micros(duration); + total.fetch_add(micros, Ordering::Relaxed); +} + +fn add_duration_with_max(total: &AtomicU64, max: &AtomicU64, duration: Duration) { + let micros = duration_micros(duration); + total.fetch_add(micros, Ordering::Relaxed); + max.fetch_max(micros, Ordering::Relaxed); +} + fn nonzero_i64(atomic: &AtomicI64) -> Option { let value = atomic.load(Ordering::Relaxed); (value != 0).then_some(value) @@ -300,4 +485,36 @@ mod tests { assert_eq!(snapshot.max_seq, Some(101)); assert_eq!(snapshot.message_kinds.commit, 1); } + + #[test] + fn snapshot_records_relay_worker_progress() { + let stats = RelayShardStats::default(); + + stats.record_received(200); + stats.record_process_error(); + stats.record_processed(RelayShardTimings { + process_message: Duration::from_micros(10), + stage_counts: Duration::from_micros(20), + stage_and_commit: Duration::from_micros(30), + apply_counts: Duration::from_micros(40), + broadcast: Duration::from_micros(50), + cursor: Duration::from_micros(60), + total: Duration::from_micros(210), + }); + + let snapshot = stats.snapshot(3); + assert_eq!(snapshot.id, 3); + assert_eq!(snapshot.received_messages, 1); + assert_eq!(snapshot.processed_messages, 1); + assert_eq!(snapshot.process_errors, 1); + assert_eq!(snapshot.last_seq, Some(200)); + assert_eq!(snapshot.max_seq, Some(200)); + assert_eq!(snapshot.process_message_micros, 10); + assert_eq!(snapshot.stage_counts_micros, 20); + assert_eq!(snapshot.stage_and_commit_micros, 30); + assert_eq!(snapshot.apply_counts_micros, 40); + assert_eq!(snapshot.broadcast_micros, 50); + assert_eq!(snapshot.cursor_micros, 60); + assert_eq!(snapshot.total_micros, 210); + } } diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index da3c556..7673e4e 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -2,6 +2,8 @@ use std::collections::HashMap; use std::sync::Arc; #[cfg(feature = "relay")] use std::sync::atomic::Ordering; +#[cfg(feature = "firehose-diagnostics")] +use std::time::Instant; use fjall::OwnedWriteBatch; @@ -182,6 +184,8 @@ impl RelayWorker { let span = info_span!("worker_shard", shard = id); let _entered = span.clone().entered(); debug!("relay shard started"); + #[cfg(feature = "firehose-diagnostics")] + let shard_stats = state.firehose_stats.relay_shard(id); let mut ctx = WorkerContext { verify_signatures, @@ -204,8 +208,12 @@ impl RelayWorker { }; while let Some(msg) = rx.blocking_recv() { + #[cfg(feature = "firehose-diagnostics")] + let message_started = Instant::now(); let IngestMessage::Firehose { url, is_pds, msg } = msg; if let SubscribeReposMessage::Info(inf) = msg { + #[cfg(feature = "firehose-diagnostics")] + shard_stats.record_info(); match inf.name { InfoName::OutdatedCursor => {} InfoName::Other(name) => { @@ -230,23 +238,37 @@ impl RelayWorker { SubscribeReposMessage::Sync(s) => (s.did.clone(), s.seq), _ => continue, }; + #[cfg(feature = "firehose-diagnostics")] + shard_stats.record_received(seq); let firehose = msg.firehose.clone(); let _span = info_span!("relay", did = %did, firehose = %firehose, seq = %seq).entered(); + #[cfg(feature = "firehose-diagnostics")] + let process_started = Instant::now(); if let Err(e) = Self::process_message(&mut ctx, msg) { + #[cfg(feature = "firehose-diagnostics")] + shard_stats.record_process_error(); error!(did = %did, err = %e, "relay shard: error processing message"); } + #[cfg(feature = "firehose-diagnostics")] + let process_message = process_started.elapsed(); let mut batch = std::mem::replace(&mut ctx.batch, ctx.state.db.inner.batch()); + #[cfg(feature = "firehose-diagnostics")] + let stage_counts_started = Instant::now(); let reservation = ctx .state .db .stage_count_deltas(&mut batch, &ctx.count_deltas); + #[cfg(feature = "firehose-diagnostics")] + let stage_counts = stage_counts_started.elapsed(); #[cfg(all(feature = "relay", feature = "jetstream"))] let mut jetstream_broadcasts = Vec::new(); + #[cfg(feature = "firehose-diagnostics")] + let stage_and_commit_started = Instant::now(); #[cfg(all(feature = "relay", feature = "jetstream"))] let res = { let _lock = ctx.state.db.jetstream_lock.lock(); @@ -266,15 +288,25 @@ impl RelayWorker { #[cfg(not(all(feature = "relay", feature = "jetstream")))] let res = batch.commit(); + #[cfg(feature = "firehose-diagnostics")] + let stage_and_commit = stage_and_commit_started.elapsed(); if let Err(e) = res { + #[cfg(feature = "firehose-diagnostics")] + shard_stats.record_commit_error(); error!(shard = id, err = %e, "relay shard: failed to commit batch"); drop(reservation); continue; } + #[cfg(feature = "firehose-diagnostics")] + let apply_counts_started = Instant::now(); ctx.state.db.apply_count_deltas(&ctx.count_deltas); drop(reservation); + #[cfg(feature = "firehose-diagnostics")] + let apply_counts = apply_counts_started.elapsed(); + #[cfg(feature = "firehose-diagnostics")] + let broadcast_started = Instant::now(); #[cfg(feature = "relay")] for broadcast in ctx.pending_broadcasts.drain(..) { let _ = state.db.relay_broadcast_tx.send(broadcast); @@ -287,15 +319,29 @@ impl RelayWorker { for msg in ctx.pending_hook_messages.drain(..) { let _ = ctx.hook.blocking_send(msg); } + #[cfg(feature = "firehose-diagnostics")] + let broadcast = broadcast_started.elapsed(); // advance cursor for this firehose only if we are the terminal consumer (relay mode) // in events mode, FirehoseWorker will advance the cursor after processing + #[cfg(feature = "firehose-diagnostics")] + let cursor_started = Instant::now(); #[cfg(feature = "relay")] { ctx.state .firehose_cursors .peek_with(&firehose, |_, c| c.store(seq, Ordering::SeqCst)); } + #[cfg(feature = "firehose-diagnostics")] + shard_stats.record_processed(crate::ingest::firehose_stats::RelayShardTimings { + process_message, + stage_counts, + stage_and_commit, + apply_counts, + broadcast, + cursor: cursor_started.elapsed(), + total: message_started.elapsed(), + }); } } -- 2.51.2