From 6d73f0b6f781b66febb148eec39452d2550a0120 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 11 Jul 2026 18:01:36 +0300 Subject: [PATCH] [ingest] make firehose diagnostics a zero-cost facade firehose_stats now always compiles: shared kind enums and timings live in kinds.rs; the real atomic-counter impls are gated behind firehose-diagnostics; a noop mirror of ZST types with inline empty recorders takes their place otherwise. StatsInstant (Instant when enabled, ZST returning Duration::ZERO when not) keeps timing call sites unconditional at zero cost. ~100 inline cfg sites across the ingest hot path (firehose.rs, relay worker/context/handlers) removed; the read side (snapshot endpoints) stays gated at its composition roots. --- src/ingest/firehose.rs | 32 ++-------- src/ingest/firehose_stats.rs | 22 ++++++- src/ingest/firehose_stats/kinds.rs | 67 ++++++++++++++++++++ src/ingest/firehose_stats/noop.rs | 99 ++++++++++++++++++++++++++++++ src/ingest/firehose_stats/relay.rs | 47 ++------------ src/ingest/mod.rs | 1 - src/ingest/relay.rs | 6 -- src/ingest/relay/context.rs | 64 +++++-------------- src/ingest/relay/handlers.rs | 4 +- src/ingest/relay/worker.rs | 62 +++++-------------- src/state.rs | 3 - 11 files changed, 225 insertions(+), 182 deletions(-) create mode 100644 src/ingest/firehose_stats/kinds.rs create mode 100644 src/ingest/firehose_stats/noop.rs diff --git a/src/ingest/firehose.rs b/src/ingest/firehose.rs index 530a502..a36aac5 100644 --- a/src/ingest/firehose.rs +++ b/src/ingest/firehose.rs @@ -119,7 +119,6 @@ fn classify_websocket_error(err: &tokio_websockets::Error) -> FirehoseFailure { } } -#[cfg(feature = "firehose-diagnostics")] fn message_stats(msg: &SubscribeReposMessage<'_>) -> (&'static str, Option) { match msg { SubscribeReposMessage::Commit(commit) => ("commit", Some(commit.seq)), @@ -149,7 +148,6 @@ pub struct FirehoseIngestor { _verify_signatures: bool, throttle: ThrottleHandle, max_failures: usize, - #[cfg(feature = "firehose-diagnostics")] stats: Arc, } @@ -165,7 +163,6 @@ impl FirehoseIngestor { max_failures: usize, ) -> Self { let throttle = state.throttler.get_handle(&relay_host).await; - #[cfg(feature = "firehose-diagnostics")] let stats = state.firehose_stats.handle(&relay_host).await; Self { state, @@ -177,7 +174,6 @@ impl FirehoseIngestor { _verify_signatures: verify_signatures, throttle, max_failures, - #[cfg(feature = "firehose-diagnostics")] stats, } } @@ -208,7 +204,6 @@ impl FirehoseIngestor { Some(c) => info!(cursor = %c, "resuming from cursor"), None => info!("no cursor found, live tailing"), } - #[cfg(feature = "firehose-diagnostics")] self.stats.record_connect_attempt(start_cursor); let host_status = self.is_pds.then(|| { @@ -229,7 +224,6 @@ impl FirehoseIngestor { Ok(s) => s, Err(e) => { let failure = classify_firehose_error(&e); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_connect_error(failure.kind); let secs = match self.on_failure(&failure).await { Some(secs) => secs, @@ -265,7 +259,6 @@ impl FirehoseIngestor { elapsed_ms = connect_started.elapsed().as_millis(), "firehose connected" ); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_connected(connect_started.elapsed()); let mut marked_active = false; let active_sleep_secs = if cfg!(debug_assertions) { 1 } else { 60 }; @@ -279,15 +272,11 @@ impl FirehoseIngestor { Ok(b) => b, Err(e) => break Err(e), }; - #[cfg(feature = "firehose-diagnostics")] self.stats.record_frame(bytes.len()); match decode_frame(&bytes) { Ok(msg) => { - #[cfg(feature = "firehose-diagnostics")] - { - let (kind, seq) = message_stats(&msg); - self.stats.record_decoded(kind, seq); - } + let (kind, seq) = message_stats(&msg); + self.stats.record_decoded(kind, seq); if self.is_pds { let tier = { let meta = self.state.pds_meta.load(); @@ -299,13 +288,11 @@ impl FirehoseIngestor { self.state.tier_policy.resolve(host, override_name) }; let accounts = self.state.db.get_count(&count_key).await; - #[cfg(feature = "firehose-diagnostics")] - let throttle_started = Instant::now(); + let throttle_started = crate::ingest::firehose_stats::StatsInstant::now(); let mut disabled = false; loop { tokio::select! { _ = self.throttle.wait_for_allow(accounts, &tier) => { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_throttle_wait(throttle_started.elapsed()); break; } @@ -389,13 +376,11 @@ impl FirehoseIngestor { // reconnect quickly. a malicious source could abuse this but they could equally // just drop the TCP connection, which would hit the normal backoff path anyway. Err(FirehoseError::StreamClosed { code: 1001, reason }) => { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_stream_error("stream_closed"); debug!(reason = %reason, "host gone away"); tokio::time::sleep(Duration::from_secs(1)).await; } Err(FirehoseError::FutureCursor) => { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_stream_error("future_cursor"); if self.is_pds && let Err(e) = self.set_host_status(HostStatus::Idle) @@ -414,7 +399,6 @@ impl FirehoseIngestor { .map_or(Cow::Borrowed(""), Cow::Borrowed); let failure = FirehoseFailure::new("relay_error", format!("{error}: {message}")); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_stream_error(failure.kind); error!(err = %error, "relay sent error: {message}"); let secs = match self.on_failure(&failure).await { @@ -444,7 +428,6 @@ impl FirehoseIngestor { } Err(e) => { let failure = classify_firehose_error(&e); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_stream_error(failure.kind); let secs = match self.on_failure(&failure).await { Some(secs) => secs, @@ -542,26 +525,22 @@ impl FirehoseIngestor { _ => return, }; - #[cfg(feature = "firehose-diagnostics")] - let should_process_started = Instant::now(); + let should_process_started = crate::ingest::firehose_stats::StatsInstant::now(); let process = self .should_process(did) .await .inspect_err(|e| error!(did = %did, err = %e, "failed to check if we should process")) .unwrap_or(false); - #[cfg(feature = "firehose-diagnostics")] self.stats .record_should_process(should_process_started.elapsed()); if !process { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_skipped(); trace!(did = %did, "skipping: not in filter"); return; } trace!(did = %did, "forwarding message to ingest buffer"); - #[cfg(feature = "firehose-diagnostics")] - let send_started = Instant::now(); + let send_started = crate::ingest::firehose_stats::StatsInstant::now(); let res = self .buffer_tx .send(IngestMessage::Firehose { @@ -570,7 +549,6 @@ impl FirehoseIngestor { msg: msg.into_static(), }) .await; - #[cfg(feature = "firehose-diagnostics")] { let elapsed = send_started.elapsed(); if res.is_ok() { diff --git a/src/ingest/firehose_stats.rs b/src/ingest/firehose_stats.rs index 497bb7d..68dea85 100644 --- a/src/ingest/firehose_stats.rs +++ b/src/ingest/firehose_stats.rs @@ -1,40 +1,56 @@ +mod kinds; +#[cfg(not(feature = "firehose-diagnostics"))] +mod noop; +#[cfg(feature = "firehose-diagnostics")] mod relay; +#[cfg(feature = "firehose-diagnostics")] mod source; +pub use kinds::*; +#[cfg(not(feature = "firehose-diagnostics"))] +pub use noop::*; +#[cfg(feature = "firehose-diagnostics")] pub use relay::{ - HostAuthorityStatsOutcome, RelayMessageKind, RelayShardStats, RelayShardTimings, - RelayWorkerStats, RelayWorkerStatsSnapshot, RepoStateLoadOutcome, ValidationStatsOutcome, + RelayShardStats, RelayWorkerStats, RelayWorkerStatsSnapshot, }; +#[cfg(feature = "firehose-diagnostics")] pub use source::{FirehoseSourceStats, FirehoseStats, FirehoseStatsSnapshot}; +#[cfg(feature = "firehose-diagnostics")] use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; +#[cfg(feature = "firehose-diagnostics")] use std::time::Duration; +#[cfg(feature = "firehose-diagnostics")] fn now_ts() -> i64 { chrono::Utc::now().timestamp() } +#[cfg(feature = "firehose-diagnostics")] fn duration_micros(duration: Duration) -> u64 { duration.as_micros().try_into().unwrap_or(u64::MAX) } +#[cfg(feature = "firehose-diagnostics")] fn add_duration(total: &AtomicU64, duration: Duration) { let micros = duration_micros(duration); total.fetch_add(micros, Ordering::Relaxed); } +#[cfg(feature = "firehose-diagnostics")] 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); } +#[cfg(feature = "firehose-diagnostics")] fn nonzero_i64(atomic: &AtomicI64) -> Option { let value = atomic.load(Ordering::Relaxed); (value != 0).then_some(value) } -#[cfg(test)] +#[cfg(all(test, feature = "firehose-diagnostics"))] mod tests { use super::*; diff --git a/src/ingest/firehose_stats/kinds.rs b/src/ingest/firehose_stats/kinds.rs new file mode 100644 index 0000000..a67e5b6 --- /dev/null +++ b/src/ingest/firehose_stats/kinds.rs @@ -0,0 +1,67 @@ +use std::time::Duration; + +#[derive(Clone, Copy, Debug)] +pub enum RelayMessageKind { + Commit, + Sync, + Identity, + Account, +} + +#[derive(Clone, Copy, Debug)] +pub enum RepoStateLoadOutcome { + Hit, + Miss, + Drop, +} + +#[derive(Clone, Copy, Debug)] +pub enum HostAuthorityStatsOutcome { + Authorized, + WasStale, + WrongHost, + Error, +} + +#[derive(Clone, Copy, Debug)] +pub enum ValidationStatsOutcome { + Accepted, + Stale, + SigFailure, + Rejected, +} + +#[cfg_attr(not(feature = "firehose-diagnostics"), allow(dead_code))] +#[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, +} + +/// instant used for diagnostics timing. compiles to a zero-sized no-op when +/// the `firehose-diagnostics` feature is disabled, so timing call sites stay +/// unconditional without paying for `Instant::now`. +#[cfg(feature = "firehose-diagnostics")] +pub(crate) use std::time::Instant as StatsInstant; + +#[cfg(not(feature = "firehose-diagnostics"))] +#[derive(Clone, Copy)] +pub(crate) struct StatsInstant; + +#[cfg(not(feature = "firehose-diagnostics"))] +impl StatsInstant { + #[inline(always)] + pub(crate) fn now() -> Self { + Self + } + + #[inline(always)] + pub(crate) fn elapsed(&self) -> Duration { + Duration::ZERO + } +} diff --git a/src/ingest/firehose_stats/noop.rs b/src/ingest/firehose_stats/noop.rs new file mode 100644 index 0000000..4835389 --- /dev/null +++ b/src/ingest/firehose_stats/noop.rs @@ -0,0 +1,99 @@ +//! zero-cost mirror of the diagnostics stats surface, used when the +//! `firehose-diagnostics` feature is disabled. every recorder is an inline +//! no-op so ingest call sites stay unconditional. + +use std::sync::Arc; +use std::time::Duration; +use url::Url; + +use super::{ + HostAuthorityStatsOutcome, RelayMessageKind, RelayShardTimings, RepoStateLoadOutcome, + ValidationStatsOutcome, +}; + +#[derive(Default)] +pub struct FirehoseStats; + +impl FirehoseStats { + #[inline(always)] + pub async fn handle(&self, _url: &Url) -> Arc { + Arc::new(FirehoseSourceStats) + } + + #[inline(always)] + pub fn relay_shard(&self, _id: usize) -> Arc { + Arc::new(RelayShardStats) + } +} + +#[derive(Default)] +pub struct FirehoseSourceStats; + +impl FirehoseSourceStats { + #[inline(always)] + pub fn record_connect_attempt(&self, _cursor: Option) {} + #[inline(always)] + pub fn record_connected(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_connect_error(&self, _kind: &'static str) {} + #[inline(always)] + pub fn record_stream_error(&self, _kind: &'static str) {} + #[inline(always)] + pub fn record_frame(&self, _len: usize) {} + #[inline(always)] + pub fn record_decoded(&self, _kind: &'static str, _seq: Option) {} + #[inline(always)] + pub fn record_throttle_wait(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_should_process(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_skipped(&self) {} + #[inline(always)] + pub fn record_forwarded(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_forward_error(&self, _elapsed: Duration) {} +} + +#[derive(Default)] +pub struct RelayShardStats; + +// mirrors the full diagnostics surface; per-combo unused recorders are expected +#[allow(dead_code)] +impl RelayShardStats { + #[inline(always)] + pub fn record_received(&self, _seq: i64) {} + #[inline(always)] + pub fn record_info(&self) {} + #[inline(always)] + pub fn record_message_kind(&self, _kind: RelayMessageKind) {} + #[inline(always)] + pub fn record_repo_state_load(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_repo_state_outcome(&self, _outcome: RepoStateLoadOutcome) {} + #[inline(always)] + pub fn record_host_authority(&self, _elapsed: Duration, _outcome: HostAuthorityStatsOutcome) {} + #[inline(always)] + pub fn record_handle_message(&self, _kind: RelayMessageKind, _elapsed: Duration) {} + #[inline(always)] + pub fn record_fetch_key(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_validate_commit(&self, _elapsed: Duration, _outcome: ValidationStatsOutcome) {} + #[inline(always)] + pub fn record_validate_sync(&self, _elapsed: Duration, _outcome: ValidationStatsOutcome) {} + #[inline(always)] + pub fn record_refresh_doc(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_new_account(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_resolve_doc(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_repo_status_probe(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_queue_emit(&self, _elapsed: Duration) {} + #[inline(always)] + pub fn record_process_error(&self) {} + #[inline(always)] + pub fn record_commit_error(&self) {} + #[inline(always)] + pub fn record_processed(&self, _timings: RelayShardTimings) {} +} diff --git a/src/ingest/firehose_stats/relay.rs b/src/ingest/firehose_stats/relay.rs index 008f6b4..c2d8e46 100644 --- a/src/ingest/firehose_stats/relay.rs +++ b/src/ingest/firehose_stats/relay.rs @@ -5,38 +5,10 @@ use std::sync::Arc; use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::time::Duration; -use super::{add_duration, add_duration_with_max, nonzero_i64, now_ts}; - -#[derive(Clone, Copy, Debug)] -pub enum RelayMessageKind { - Commit, - Sync, - Identity, - Account, -} - -#[derive(Clone, Copy, Debug)] -pub enum RepoStateLoadOutcome { - Hit, - Miss, - Drop, -} - -#[derive(Clone, Copy, Debug)] -pub enum HostAuthorityStatsOutcome { - Authorized, - WasStale, - WrongHost, - Error, -} - -#[derive(Clone, Copy, Debug)] -pub enum ValidationStatsOutcome { - Accepted, - Stale, - SigFailure, - Rejected, -} +use super::{ + HostAuthorityStatsOutcome, RelayMessageKind, RelayShardTimings, RepoStateLoadOutcome, + ValidationStatsOutcome, add_duration, add_duration_with_max, nonzero_i64, now_ts, +}; #[derive(Default)] pub struct RelayWorkerStats { @@ -398,17 +370,6 @@ impl RelayShardStats { } } -#[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, diff --git a/src/ingest/mod.rs b/src/ingest/mod.rs index 1d9ec7a..76f4933 100644 --- a/src/ingest/mod.rs +++ b/src/ingest/mod.rs @@ -1,5 +1,4 @@ pub mod firehose; -#[cfg(feature = "firehose-diagnostics")] pub mod firehose_stats; #[cfg(feature = "indexer")] pub mod indexer; diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index 55e3928..08804bb 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -5,11 +5,8 @@ use jacquard_api::com_atproto::sync::get_repo_status::{ }; use smol_str::SmolStr; -#[cfg(feature = "firehose-diagnostics")] use crate::ingest::firehose_stats::{RelayMessageKind, ValidationStatsOutcome}; -#[cfg(feature = "firehose-diagnostics")] use crate::ingest::stream::SubscribeReposMessage; -#[cfg(feature = "firehose-diagnostics")] use crate::ingest::validation::{ CommitValidationError, SyncValidationError, ValidatedCommit, ValidatedSync, }; @@ -50,7 +47,6 @@ pub(crate) fn map_repo_status_probe( Some(repo_state) } -#[cfg(feature = "firehose-diagnostics")] pub(crate) fn relay_message_kind(msg: &SubscribeReposMessage<'_>) -> Option { match msg { SubscribeReposMessage::Commit(_) => Some(RelayMessageKind::Commit), @@ -61,7 +57,6 @@ pub(crate) fn relay_message_kind(msg: &SubscribeReposMessage<'_>) -> Option, CommitValidationError>, ) -> ValidationStatsOutcome { @@ -73,7 +68,6 @@ pub(crate) fn commit_validation_outcome( } } -#[cfg(feature = "firehose-diagnostics")] pub(crate) fn sync_validation_outcome( res: &std::result::Result, ) -> ValidationStatsOutcome { diff --git a/src/ingest/relay/context.rs b/src/ingest/relay/context.rs index 4dcae3c..aaad979 100644 --- a/src/ingest/relay/context.rs +++ b/src/ingest/relay/context.rs @@ -1,5 +1,4 @@ use std::collections::HashMap; -#[cfg(feature = "firehose-diagnostics")] use std::sync::Arc; use std::time::Instant; @@ -35,14 +34,13 @@ use crate::types::RelayBroadcast; #[cfg(feature = "relay")] use std::sync::atomic::Ordering; -#[cfg(feature = "firehose-diagnostics")] use super::{commit_validation_outcome, sync_validation_outcome}; +use crate::ingest::firehose_stats::StatsInstant; pub(crate) struct WorkerContext<'a> { pub(crate) verify_signatures: bool, pub(crate) state: &'a AppState, pub(crate) vctx: ValidationContext<'a>, - #[cfg(feature = "firehose-diagnostics")] pub(crate) stats: Arc, pub(crate) batch: OwnedWriteBatch, pub(crate) count_deltas: CountDeltas, @@ -146,17 +144,14 @@ impl WorkerContext<'_> { } pub(crate) fn refresh_doc(&mut self, did: &Did, repo_state: &mut RepoState) -> Result<()> { - #[cfg(feature = "firehose-diagnostics")] - let refresh_started = Instant::now(); + let refresh_started = StatsInstant::now(); let result = (|| { let db = &self.state.db; self.state.resolver.invalidate_sync(did); - #[cfg(feature = "firehose-diagnostics")] - let resolve_started = Instant::now(); + let resolve_started = StatsInstant::now(); let doc = Handle::current() .block_on(self.state.resolver.resolve_doc(did)) .map_err(|e| miette::miette!("{e}")); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_resolve_doc(resolve_started.elapsed()); let doc = doc?; repo_state.update_from_doc(doc); @@ -169,7 +164,6 @@ impl WorkerContext<'_> { ); Ok(()) })(); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_refresh_doc(refresh_started.elapsed()); result } @@ -181,10 +175,8 @@ impl WorkerContext<'_> { ) -> Result>> { let did = &commit.repo; let key = self.fetch_key(did)?; - #[cfg(feature = "firehose-diagnostics")] - let validate_started = Instant::now(); + let validate_started = StatsInstant::now(); let validation = self.vctx.validate_commit(commit, repo_state, key.as_ref()); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_validate_commit( validate_started.elapsed(), commit_validation_outcome(&validation), @@ -204,10 +196,8 @@ impl WorkerContext<'_> { self.refresh_doc(did, repo_state)?; let key = self.fetch_key(did)?; - #[cfg(feature = "firehose-diagnostics")] - let validate_started = Instant::now(); + let validate_started = StatsInstant::now(); let validation = self.vctx.validate_commit(commit, repo_state, key.as_ref()); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_validate_commit( validate_started.elapsed(), commit_validation_outcome(&validation), @@ -228,10 +218,8 @@ impl WorkerContext<'_> { ) -> Result> { let did = &sync.did; let key = self.fetch_key(did)?; - #[cfg(feature = "firehose-diagnostics")] - let validate_started = Instant::now(); + let validate_started = StatsInstant::now(); let validation = self.vctx.validate_sync(sync, key.as_ref()); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_validate_sync( validate_started.elapsed(), sync_validation_outcome(&validation), @@ -247,10 +235,8 @@ impl WorkerContext<'_> { self.refresh_doc(did, repo_state)?; let key = self.fetch_key(did)?; - #[cfg(feature = "firehose-diagnostics")] - let validate_started = Instant::now(); + let validate_started = StatsInstant::now(); let validation = self.vctx.validate_sync(sync, key.as_ref()); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_validate_sync( validate_started.elapsed(), sync_validation_outcome(&validation), @@ -265,8 +251,7 @@ impl WorkerContext<'_> { } fn fetch_key(&self, did: &Did) -> Result>> { - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = if self.verify_signatures { Handle::current() .block_on(self.state.resolver.resolve_signing_key(did)) @@ -275,7 +260,6 @@ impl WorkerContext<'_> { } else { Ok(None) }; - #[cfg(feature = "firehose-diagnostics")] self.stats.record_fetch_key(started.elapsed()); result } @@ -356,7 +340,6 @@ impl WorkerContext<'_> { if metadata.is_some_and(|m| !m.tracked) { trace!(did = %did, "ignoring message, repo is explicitly untracked"); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, ); @@ -371,7 +354,6 @@ impl WorkerContext<'_> { .transpose()?; if let Some(repo_state) = repo_state_opt { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Hit, ); @@ -385,7 +367,6 @@ impl WorkerContext<'_> { let commit = match &msg.msg { SubscribeReposMessage::Commit(c) => c, _ => { - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, ); @@ -408,7 +389,6 @@ impl WorkerContext<'_> { }); if !touches_signal { trace!(did = %did, "dropping commit, no signal-matching ops"); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, ); @@ -418,17 +398,14 @@ impl WorkerContext<'_> { } debug!(did = %did, "discovered new account from firehose, queueing backfill"); - #[cfg(feature = "firehose-diagnostics")] - let new_account_started = Instant::now(); + let new_account_started = StatsInstant::now(); // resolve doc to initialize repo state self.state.resolver.invalidate_sync(did); - #[cfg(feature = "firehose-diagnostics")] - let resolve_started = Instant::now(); + let resolve_started = StatsInstant::now(); let doc = tokio::runtime::Handle::current() .block_on(self.state.resolver.resolve_doc(did)) .into_diagnostic(); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_resolve_doc(resolve_started.elapsed()); let doc = doc?; @@ -437,7 +414,6 @@ impl WorkerContext<'_> { let pds_host = doc.pds.host_str().map(|h| h.to_string()); if pds_host.as_deref() != msg.firehose.host_str() { warn!(did = %did, got = ?pds_host, expected = ?msg.firehose.host_str(), "message rejected: wrong host for new account"); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, ); @@ -451,7 +427,6 @@ impl WorkerContext<'_> { .get_count_sync(&keys::pds_account_count_key(host)); if self.state.is_over_account_limit(host, count) { warn!(did = %did, host, count, "account limit reached for host, dropping new account"); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_state_outcome( crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, ); @@ -463,11 +438,9 @@ impl WorkerContext<'_> { #[cfg(any(not(feature = "relay"), feature = "indexer"))] let mut repo_state = { // try to get upstream status - #[cfg(feature = "firehose-diagnostics")] - let probe_started = Instant::now(); + let probe_started = StatsInstant::now(); let repo_state = tokio::runtime::Handle::current().block_on(self.check_repo_status(did, &doc.pds)); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_repo_status_probe(probe_started.elapsed()); repo_state .ok() @@ -503,13 +476,10 @@ impl WorkerContext<'_> { self.count_deltas.add_repos(1); - #[cfg(feature = "firehose-diagnostics")] - { - self.stats.record_repo_state_outcome( - crate::ingest::firehose_stats::RepoStateLoadOutcome::Miss, - ); - self.stats.record_new_account(new_account_started.elapsed()); - } + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Miss, + ); + self.stats.record_new_account(new_account_started.elapsed()); Ok(Some(repo_state)) } @@ -519,8 +489,7 @@ impl WorkerContext<'_> { &mut self, make_frame: impl FnOnce(i64) -> Result, ) -> Result { - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = (|| { let db = &self.state.db; let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); @@ -532,7 +501,6 @@ impl WorkerContext<'_> { self.pending_broadcasts.push(RelayBroadcast::Persisted(seq)); Ok(seq) })(); - #[cfg(feature = "firehose-diagnostics")] self.stats.record_queue_emit(started.elapsed()); result } diff --git a/src/ingest/relay/handlers.rs b/src/ingest/relay/handlers.rs index b4303f2..7e9a085 100644 --- a/src/ingest/relay/handlers.rs +++ b/src/ingest/relay/handlers.rs @@ -211,10 +211,8 @@ impl RelayWorker { // or if there is no handle specified if is_pds || identity.handle.is_none() { ctx.state.resolver.invalidate_sync(&identity.did); - #[cfg(feature = "firehose-diagnostics")] - let resolve_started = std::time::Instant::now(); + let resolve_started = crate::ingest::firehose_stats::StatsInstant::now(); let doc = Handle::current().block_on(ctx.state.resolver.resolve_doc(&identity.did)); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_resolve_doc(resolve_started.elapsed()); match doc { Ok(doc) => { diff --git a/src/ingest/relay/worker.rs b/src/ingest/relay/worker.rs index a566c9d..6d60602 100644 --- a/src/ingest/relay/worker.rs +++ b/src/ingest/relay/worker.rs @@ -1,8 +1,7 @@ use std::sync::Arc; #[cfg(feature = "relay")] use std::sync::atomic::Ordering; -#[cfg(feature = "firehose-diagnostics")] -use std::time::Instant; +use crate::ingest::firehose_stats::StatsInstant; use miette::{IntoDiagnostic, Result}; use tokio::runtime::Handle; @@ -17,7 +16,6 @@ use crate::state::AppState; use super::{AuthorityOutcome, WorkerContext}; -#[cfg(feature = "firehose-diagnostics")] use super::relay_message_kind; pub struct WorkerMessage { @@ -120,7 +118,6 @@ 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 { @@ -129,7 +126,6 @@ impl RelayWorker { vctx: crate::ingest::validation::ValidationContext { opts: &validation_opts, }, - #[cfg(feature = "firehose-diagnostics")] stats: shard_stats.clone(), batch: state.db.inner.batch(), count_deltas: CountDeltas::default(), @@ -147,11 +143,9 @@ impl RelayWorker { }; while let Some(msg) = rx.blocking_recv() { - #[cfg(feature = "firehose-diagnostics")] - let message_started = Instant::now(); + let message_started = StatsInstant::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 => {} @@ -179,37 +173,30 @@ 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(); + let process_started = StatsInstant::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 stage_counts_started = StatsInstant::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(); + let stage_and_commit_started = StatsInstant::now(); #[cfg(all(feature = "relay", feature = "jetstream"))] let res = { let _lock = ctx.state.db.jetstream_lock.lock(); @@ -229,25 +216,20 @@ 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(); + let apply_counts_started = StatsInstant::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(); + let broadcast_started = StatsInstant::now(); #[cfg(feature = "relay")] for broadcast in ctx.pending_broadcasts.drain(..) { let _ = state.db.relay_broadcast_tx.send(broadcast); @@ -260,20 +242,17 @@ 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(); + let cursor_started = StatsInstant::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, @@ -287,15 +266,12 @@ impl RelayWorker { } fn process_message(ctx: &mut WorkerContext, msg: WorkerMessage) -> Result<()> { - #[cfg(feature = "firehose-diagnostics")] if let Some(kind) = relay_message_kind(&msg.msg) { ctx.stats.record_message_kind(kind); } - #[cfg(feature = "firehose-diagnostics")] - let load_started = Instant::now(); + let load_started = StatsInstant::now(); let repo_state_result = ctx.load_repo_state(&msg); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_repo_state_load(load_started.elapsed()); let Some(mut repo_state) = repo_state_result? else { return Ok(()); @@ -305,10 +281,8 @@ impl RelayWorker { if let Some(host) = msg.firehose.host_str() && msg.is_pds { - #[cfg(feature = "firehose-diagnostics")] - let authority_started = Instant::now(); + let authority_started = StatsInstant::now(); let outcome_result = ctx.check_host_authority(did, &mut repo_state, host); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_host_authority( authority_started.elapsed(), match &outcome_result { @@ -337,10 +311,8 @@ impl RelayWorker { match msg.msg { SubscribeReposMessage::Commit(commit) => { trace!("processing commit"); - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = Self::handle_commit(ctx, &mut repo_state, &msg.firehose, *commit); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_handle_message( crate::ingest::firehose_stats::RelayMessageKind::Commit, started.elapsed(), @@ -349,10 +321,8 @@ impl RelayWorker { } SubscribeReposMessage::Sync(sync) => { debug!("processing sync"); - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = Self::handle_sync(ctx, &mut repo_state, &msg.firehose, *sync); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_handle_message( crate::ingest::firehose_stats::RelayMessageKind::Sync, started.elapsed(), @@ -361,8 +331,7 @@ impl RelayWorker { } SubscribeReposMessage::Identity(identity) => { debug!("processing identity"); - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = Self::handle_identity( ctx, &mut repo_state, @@ -370,7 +339,6 @@ impl RelayWorker { *identity, msg.is_pds, ); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_handle_message( crate::ingest::firehose_stats::RelayMessageKind::Identity, started.elapsed(), @@ -379,11 +347,9 @@ impl RelayWorker { } SubscribeReposMessage::Account(account) => { debug!("processing account"); - #[cfg(feature = "firehose-diagnostics")] - let started = Instant::now(); + let started = StatsInstant::now(); let result = Self::handle_account(ctx, &mut repo_state, &msg.firehose, *account, msg.is_pds); - #[cfg(feature = "firehose-diagnostics")] ctx.stats.record_handle_message( crate::ingest::firehose_stats::RelayMessageKind::Account, started.elapsed(), diff --git a/src/state.rs b/src/state.rs index 760e014..0c40afb 100644 --- a/src/state.rs +++ b/src/state.rs @@ -10,7 +10,6 @@ use tokio::sync::Notify; use tokio::sync::{Semaphore, watch}; use url::Url; -#[cfg(feature = "firehose-diagnostics")] use crate::ingest::firehose_stats::FirehoseStats; #[cfg(feature = "relay")] use crate::pds_daily_limit::PdsDailyLimit; @@ -33,7 +32,6 @@ pub struct AppState { pub(crate) pds_daily_limit: PdsDailyLimit, pub(crate) tier_policy: TierPolicy, pub firehose_cursors: scc::HashIndex, - #[cfg(feature = "firehose-diagnostics")] pub(crate) firehose_stats: FirehoseStats, pub firehose_enabled: watch::Sender, #[cfg(feature = "indexer")] @@ -114,7 +112,6 @@ impl AppState { pds_meta, tier_policy: config.tier_policy.clone(), firehose_cursors: relay_cursors, - #[cfg(feature = "firehose-diagnostics")] firehose_stats: FirehoseStats::default(), #[cfg(feature = "indexer")] backfill_notify: Notify::new(), -- 2.51.2