From db139ce3fe82816e5a60c149f6de8945d2d2b9b1 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 07 Jun 2026 18:47:59 +0000 Subject: [PATCH] [ingest] add more stats around signature verifs and message handling --- docs/api/firehose.md | 49 +++++++++++++++++++++++++++++++++++++++++++++++++ src/ingest/firehose_stats.rs | 321 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/ingest/relay.rs | 277 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------------------------- 3 file(s) changed, 606 insertion(s)(+), 41 deletion(s)(-) diff --git a/docs/api/firehose.md b/docs/api/firehose.md --- a/docs/api/firehose.md +++ b/docs/api/firehose.md @@ -101,7 +101,48 @@ "processed_messages": 999, "process_errors": 0, "commit_errors": 0, + "commit_messages": 998, + "sync_messages": 0, + "identity_messages": 1, + "account_messages": 1, + "repo_state_hits": 996, + "repo_state_misses": 2, + "repo_state_drops": 2, + "host_authority_checks": 1000, + "host_authority_authorized": 999, + "host_authority_stale": 1, + "host_authority_wrong": 0, + "host_authority_errors": 0, + "fetch_key_calls": 998, + "validate_commit_calls": 998, + "validate_commit_accepted": 998, + "validate_commit_stale": 0, + "validate_commit_sig_failures": 0, + "validate_commit_rejected": 0, + "validate_sync_calls": 0, + "validate_sync_accepted": 0, + "validate_sync_sig_failures": 0, + "validate_sync_rejected": 0, + "refresh_doc_calls": 1, + "new_account_calls": 2, + "resolve_doc_calls": 3, + "repo_status_probe_calls": 2, + "queue_emit_calls": 999, "process_message_micros": 100000, + "load_repo_state_micros": 10000, + "host_authority_micros": 5000, + "handle_commit_micros": 80000, + "handle_sync_micros": 0, + "handle_identity_micros": 1000, + "handle_account_micros": 1000, + "fetch_key_micros": 20000, + "validate_commit_micros": 50000, + "validate_sync_micros": 0, + "refresh_doc_micros": 1000, + "new_account_micros": 2000, + "resolve_doc_micros": 2000, + "repo_status_probe_micros": 1000, + "queue_emit_micros": 10000, "stage_counts_micros": 20000, "stage_and_commit_micros": 900000, "apply_counts_micros": 30000, @@ -109,6 +150,14 @@ "cursor_micros": 1000, "total_micros": 1100000, "max_process_message_micros": 1000, + "max_load_repo_state_micros": 200, + "max_host_authority_micros": 100, + "max_handle_commit_micros": 800, + "max_fetch_key_micros": 300, + "max_validate_commit_micros": 700, + "max_refresh_doc_micros": 1000, + "max_new_account_micros": 1000, + "max_queue_emit_micros": 200, "max_stage_and_commit_micros": 10000, "max_total_micros": 12000, "last_received_at": 1717239958, diff --git a/src/ingest/firehose_stats.rs b/src/ingest/firehose_stats.rs --- a/src/ingest/firehose_stats.rs +++ b/src/ingest/firehose_stats.rs @@ -296,6 +296,37 @@ } } +#[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, +} + #[derive(Default)] pub struct RelayShardStats { received_messages: AtomicU64, @@ -303,7 +334,48 @@ processed_messages: AtomicU64, process_errors: AtomicU64, commit_errors: AtomicU64, + commit_messages: AtomicU64, + sync_messages: AtomicU64, + identity_messages: AtomicU64, + account_messages: AtomicU64, + repo_state_hits: AtomicU64, + repo_state_misses: AtomicU64, + repo_state_drops: AtomicU64, + host_authority_checks: AtomicU64, + host_authority_authorized: AtomicU64, + host_authority_stale: AtomicU64, + host_authority_wrong: AtomicU64, + host_authority_errors: AtomicU64, + fetch_key_calls: AtomicU64, + validate_commit_calls: AtomicU64, + validate_commit_accepted: AtomicU64, + validate_commit_stale: AtomicU64, + validate_commit_sig_failures: AtomicU64, + validate_commit_rejected: AtomicU64, + validate_sync_calls: AtomicU64, + validate_sync_accepted: AtomicU64, + validate_sync_sig_failures: AtomicU64, + validate_sync_rejected: AtomicU64, + refresh_doc_calls: AtomicU64, + new_account_calls: AtomicU64, + resolve_doc_calls: AtomicU64, + repo_status_probe_calls: AtomicU64, + queue_emit_calls: AtomicU64, process_message_micros: AtomicU64, + load_repo_state_micros: AtomicU64, + host_authority_micros: AtomicU64, + handle_commit_micros: AtomicU64, + handle_sync_micros: AtomicU64, + handle_identity_micros: AtomicU64, + handle_account_micros: AtomicU64, + fetch_key_micros: AtomicU64, + validate_commit_micros: AtomicU64, + validate_sync_micros: AtomicU64, + refresh_doc_micros: AtomicU64, + new_account_micros: AtomicU64, + resolve_doc_micros: AtomicU64, + repo_status_probe_micros: AtomicU64, + queue_emit_micros: AtomicU64, stage_counts_micros: AtomicU64, stage_and_commit_micros: AtomicU64, apply_counts_micros: AtomicU64, @@ -311,6 +383,14 @@ cursor_micros: AtomicU64, total_micros: AtomicU64, max_process_message_micros: AtomicU64, + max_load_repo_state_micros: AtomicU64, + max_host_authority_micros: AtomicU64, + max_handle_commit_micros: AtomicU64, + max_fetch_key_micros: AtomicU64, + max_validate_commit_micros: AtomicU64, + max_refresh_doc_micros: AtomicU64, + max_new_account_micros: AtomicU64, + max_queue_emit_micros: AtomicU64, max_stage_and_commit_micros: AtomicU64, max_total_micros: AtomicU64, last_received_at: AtomicI64, @@ -330,6 +410,149 @@ pub fn record_info(&self) { self.info_messages.fetch_add(1, Ordering::Relaxed); + } + + pub fn record_message_kind(&self, kind: RelayMessageKind) { + match kind { + RelayMessageKind::Commit => self.commit_messages.fetch_add(1, Ordering::Relaxed), + RelayMessageKind::Sync => self.sync_messages.fetch_add(1, Ordering::Relaxed), + RelayMessageKind::Identity => self.identity_messages.fetch_add(1, Ordering::Relaxed), + RelayMessageKind::Account => self.account_messages.fetch_add(1, Ordering::Relaxed), + }; + } + + pub fn record_repo_state_load(&self, elapsed: Duration) { + add_duration_with_max( + &self.load_repo_state_micros, + &self.max_load_repo_state_micros, + elapsed, + ); + } + + pub fn record_repo_state_outcome(&self, outcome: RepoStateLoadOutcome) { + match outcome { + RepoStateLoadOutcome::Hit => self.repo_state_hits.fetch_add(1, Ordering::Relaxed), + RepoStateLoadOutcome::Miss => self.repo_state_misses.fetch_add(1, Ordering::Relaxed), + RepoStateLoadOutcome::Drop => self.repo_state_drops.fetch_add(1, Ordering::Relaxed), + }; + } + + pub fn record_host_authority(&self, elapsed: Duration, outcome: HostAuthorityStatsOutcome) { + self.host_authority_checks.fetch_add(1, Ordering::Relaxed); + match outcome { + HostAuthorityStatsOutcome::Authorized => self + .host_authority_authorized + .fetch_add(1, Ordering::Relaxed), + HostAuthorityStatsOutcome::WasStale => { + self.host_authority_stale.fetch_add(1, Ordering::Relaxed) + } + HostAuthorityStatsOutcome::WrongHost => { + self.host_authority_wrong.fetch_add(1, Ordering::Relaxed) + } + HostAuthorityStatsOutcome::Error => { + self.host_authority_errors.fetch_add(1, Ordering::Relaxed) + } + }; + add_duration_with_max( + &self.host_authority_micros, + &self.max_host_authority_micros, + elapsed, + ); + } + + pub fn record_handle_message(&self, kind: RelayMessageKind, elapsed: Duration) { + match kind { + RelayMessageKind::Commit => add_duration_with_max( + &self.handle_commit_micros, + &self.max_handle_commit_micros, + elapsed, + ), + RelayMessageKind::Sync => add_duration(&self.handle_sync_micros, elapsed), + RelayMessageKind::Identity => add_duration(&self.handle_identity_micros, elapsed), + RelayMessageKind::Account => add_duration(&self.handle_account_micros, elapsed), + } + } + + pub fn record_fetch_key(&self, elapsed: Duration) { + self.fetch_key_calls.fetch_add(1, Ordering::Relaxed); + add_duration_with_max(&self.fetch_key_micros, &self.max_fetch_key_micros, elapsed); + } + + pub fn record_validate_commit(&self, elapsed: Duration, outcome: ValidationStatsOutcome) { + self.validate_commit_calls.fetch_add(1, Ordering::Relaxed); + match outcome { + ValidationStatsOutcome::Accepted => self + .validate_commit_accepted + .fetch_add(1, Ordering::Relaxed), + ValidationStatsOutcome::Stale => { + self.validate_commit_stale.fetch_add(1, Ordering::Relaxed) + } + ValidationStatsOutcome::SigFailure => self + .validate_commit_sig_failures + .fetch_add(1, Ordering::Relaxed), + ValidationStatsOutcome::Rejected => self + .validate_commit_rejected + .fetch_add(1, Ordering::Relaxed), + }; + add_duration_with_max( + &self.validate_commit_micros, + &self.max_validate_commit_micros, + elapsed, + ); + } + + pub fn record_validate_sync(&self, elapsed: Duration, outcome: ValidationStatsOutcome) { + self.validate_sync_calls.fetch_add(1, Ordering::Relaxed); + match outcome { + ValidationStatsOutcome::Accepted => { + self.validate_sync_accepted.fetch_add(1, Ordering::Relaxed) + } + ValidationStatsOutcome::SigFailure => self + .validate_sync_sig_failures + .fetch_add(1, Ordering::Relaxed), + ValidationStatsOutcome::Stale | ValidationStatsOutcome::Rejected => { + self.validate_sync_rejected.fetch_add(1, Ordering::Relaxed) + } + }; + add_duration(&self.validate_sync_micros, elapsed); + } + + pub fn record_refresh_doc(&self, elapsed: Duration) { + self.refresh_doc_calls.fetch_add(1, Ordering::Relaxed); + add_duration_with_max( + &self.refresh_doc_micros, + &self.max_refresh_doc_micros, + elapsed, + ); + } + + pub fn record_new_account(&self, elapsed: Duration) { + self.new_account_calls.fetch_add(1, Ordering::Relaxed); + add_duration_with_max( + &self.new_account_micros, + &self.max_new_account_micros, + elapsed, + ); + } + + pub fn record_resolve_doc(&self, elapsed: Duration) { + self.resolve_doc_calls.fetch_add(1, Ordering::Relaxed); + add_duration(&self.resolve_doc_micros, elapsed); + } + + pub fn record_repo_status_probe(&self, elapsed: Duration) { + self.repo_status_probe_calls.fetch_add(1, Ordering::Relaxed); + add_duration(&self.repo_status_probe_micros, elapsed); + } + + #[cfg_attr(not(feature = "relay"), allow(dead_code))] + pub fn record_queue_emit(&self, elapsed: Duration) { + self.queue_emit_calls.fetch_add(1, Ordering::Relaxed); + add_duration_with_max( + &self.queue_emit_micros, + &self.max_queue_emit_micros, + elapsed, + ); } pub fn record_process_error(&self) { @@ -370,7 +593,48 @@ processed_messages: self.processed_messages.load(Ordering::Relaxed), process_errors: self.process_errors.load(Ordering::Relaxed), commit_errors: self.commit_errors.load(Ordering::Relaxed), + commit_messages: self.commit_messages.load(Ordering::Relaxed), + sync_messages: self.sync_messages.load(Ordering::Relaxed), + identity_messages: self.identity_messages.load(Ordering::Relaxed), + account_messages: self.account_messages.load(Ordering::Relaxed), + repo_state_hits: self.repo_state_hits.load(Ordering::Relaxed), + repo_state_misses: self.repo_state_misses.load(Ordering::Relaxed), + repo_state_drops: self.repo_state_drops.load(Ordering::Relaxed), + host_authority_checks: self.host_authority_checks.load(Ordering::Relaxed), + host_authority_authorized: self.host_authority_authorized.load(Ordering::Relaxed), + host_authority_stale: self.host_authority_stale.load(Ordering::Relaxed), + host_authority_wrong: self.host_authority_wrong.load(Ordering::Relaxed), + host_authority_errors: self.host_authority_errors.load(Ordering::Relaxed), + fetch_key_calls: self.fetch_key_calls.load(Ordering::Relaxed), + validate_commit_calls: self.validate_commit_calls.load(Ordering::Relaxed), + validate_commit_accepted: self.validate_commit_accepted.load(Ordering::Relaxed), + validate_commit_stale: self.validate_commit_stale.load(Ordering::Relaxed), + validate_commit_sig_failures: self.validate_commit_sig_failures.load(Ordering::Relaxed), + validate_commit_rejected: self.validate_commit_rejected.load(Ordering::Relaxed), + validate_sync_calls: self.validate_sync_calls.load(Ordering::Relaxed), + validate_sync_accepted: self.validate_sync_accepted.load(Ordering::Relaxed), + validate_sync_sig_failures: self.validate_sync_sig_failures.load(Ordering::Relaxed), + validate_sync_rejected: self.validate_sync_rejected.load(Ordering::Relaxed), + refresh_doc_calls: self.refresh_doc_calls.load(Ordering::Relaxed), + new_account_calls: self.new_account_calls.load(Ordering::Relaxed), + resolve_doc_calls: self.resolve_doc_calls.load(Ordering::Relaxed), + repo_status_probe_calls: self.repo_status_probe_calls.load(Ordering::Relaxed), + queue_emit_calls: self.queue_emit_calls.load(Ordering::Relaxed), process_message_micros: self.process_message_micros.load(Ordering::Relaxed), + load_repo_state_micros: self.load_repo_state_micros.load(Ordering::Relaxed), + host_authority_micros: self.host_authority_micros.load(Ordering::Relaxed), + handle_commit_micros: self.handle_commit_micros.load(Ordering::Relaxed), + handle_sync_micros: self.handle_sync_micros.load(Ordering::Relaxed), + handle_identity_micros: self.handle_identity_micros.load(Ordering::Relaxed), + handle_account_micros: self.handle_account_micros.load(Ordering::Relaxed), + fetch_key_micros: self.fetch_key_micros.load(Ordering::Relaxed), + validate_commit_micros: self.validate_commit_micros.load(Ordering::Relaxed), + validate_sync_micros: self.validate_sync_micros.load(Ordering::Relaxed), + refresh_doc_micros: self.refresh_doc_micros.load(Ordering::Relaxed), + new_account_micros: self.new_account_micros.load(Ordering::Relaxed), + resolve_doc_micros: self.resolve_doc_micros.load(Ordering::Relaxed), + repo_status_probe_micros: self.repo_status_probe_micros.load(Ordering::Relaxed), + queue_emit_micros: self.queue_emit_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), @@ -378,6 +642,14 @@ 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_load_repo_state_micros: self.max_load_repo_state_micros.load(Ordering::Relaxed), + max_host_authority_micros: self.max_host_authority_micros.load(Ordering::Relaxed), + max_handle_commit_micros: self.max_handle_commit_micros.load(Ordering::Relaxed), + max_fetch_key_micros: self.max_fetch_key_micros.load(Ordering::Relaxed), + max_validate_commit_micros: self.max_validate_commit_micros.load(Ordering::Relaxed), + max_refresh_doc_micros: self.max_refresh_doc_micros.load(Ordering::Relaxed), + max_new_account_micros: self.max_new_account_micros.load(Ordering::Relaxed), + max_queue_emit_micros: self.max_queue_emit_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), @@ -413,7 +685,48 @@ pub processed_messages: u64, pub process_errors: u64, pub commit_errors: u64, + pub commit_messages: u64, + pub sync_messages: u64, + pub identity_messages: u64, + pub account_messages: u64, + pub repo_state_hits: u64, + pub repo_state_misses: u64, + pub repo_state_drops: u64, + pub host_authority_checks: u64, + pub host_authority_authorized: u64, + pub host_authority_stale: u64, + pub host_authority_wrong: u64, + pub host_authority_errors: u64, + pub fetch_key_calls: u64, + pub validate_commit_calls: u64, + pub validate_commit_accepted: u64, + pub validate_commit_stale: u64, + pub validate_commit_sig_failures: u64, + pub validate_commit_rejected: u64, + pub validate_sync_calls: u64, + pub validate_sync_accepted: u64, + pub validate_sync_sig_failures: u64, + pub validate_sync_rejected: u64, + pub refresh_doc_calls: u64, + pub new_account_calls: u64, + pub resolve_doc_calls: u64, + pub repo_status_probe_calls: u64, + pub queue_emit_calls: u64, pub process_message_micros: u64, + pub load_repo_state_micros: u64, + pub host_authority_micros: u64, + pub handle_commit_micros: u64, + pub handle_sync_micros: u64, + pub handle_identity_micros: u64, + pub handle_account_micros: u64, + pub fetch_key_micros: u64, + pub validate_commit_micros: u64, + pub validate_sync_micros: u64, + pub refresh_doc_micros: u64, + pub new_account_micros: u64, + pub resolve_doc_micros: u64, + pub repo_status_probe_micros: u64, + pub queue_emit_micros: u64, pub stage_counts_micros: u64, pub stage_and_commit_micros: u64, pub apply_counts_micros: u64, @@ -421,6 +734,14 @@ pub cursor_micros: u64, pub total_micros: u64, pub max_process_message_micros: u64, + pub max_load_repo_state_micros: u64, + pub max_host_authority_micros: u64, + pub max_handle_commit_micros: u64, + pub max_fetch_key_micros: u64, + pub max_validate_commit_micros: u64, + pub max_refresh_doc_micros: u64, + pub max_new_account_micros: u64, + pub max_queue_emit_micros: u64, pub max_stage_and_commit_micros: u64, pub max_total_micros: u64, #[serde(skip_serializing_if = "Option::is_none")] diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -23,6 +23,10 @@ #[cfg(all(feature = "relay", feature = "jetstream"))] use crate::db::types::TrimmedDid; use crate::db::{self, CountDeltas, keys}; +#[cfg(feature = "firehose-diagnostics")] +use crate::ingest::firehose_stats::{ + HostAuthorityStatsOutcome, RelayMessageKind, RepoStateLoadOutcome, ValidationStatsOutcome, +}; use crate::ingest::stream::AccountStatus; #[cfg(feature = "relay")] use crate::ingest::stream::encode_frame; @@ -63,10 +67,46 @@ Some(repo_state) } +#[cfg(feature = "firehose-diagnostics")] +fn relay_message_kind(msg: &SubscribeReposMessage<'_>) -> Option { + match msg { + SubscribeReposMessage::Commit(_) => Some(RelayMessageKind::Commit), + SubscribeReposMessage::Sync(_) => Some(RelayMessageKind::Sync), + SubscribeReposMessage::Identity(_) => Some(RelayMessageKind::Identity), + SubscribeReposMessage::Account(_) => Some(RelayMessageKind::Account), + SubscribeReposMessage::Info(_) => None, + } +} + +#[cfg(feature = "firehose-diagnostics")] +fn commit_validation_outcome( + res: &std::result::Result, CommitValidationError>, +) -> ValidationStatsOutcome { + match res { + Ok(_) => ValidationStatsOutcome::Accepted, + Err(CommitValidationError::StaleRev) => ValidationStatsOutcome::Stale, + Err(CommitValidationError::SigFailure) => ValidationStatsOutcome::SigFailure, + Err(_) => ValidationStatsOutcome::Rejected, + } +} + +#[cfg(feature = "firehose-diagnostics")] +fn sync_validation_outcome( + res: &std::result::Result, +) -> ValidationStatsOutcome { + match res { + Ok(_) => ValidationStatsOutcome::Accepted, + Err(SyncValidationError::SigFailure) => ValidationStatsOutcome::SigFailure, + Err(_) => ValidationStatsOutcome::Rejected, + } +} + struct WorkerContext<'a> { verify_signatures: bool, state: &'a AppState, vctx: ValidationContext<'a>, + #[cfg(feature = "firehose-diagnostics")] + stats: Arc, batch: OwnedWriteBatch, count_deltas: CountDeltas, #[cfg(feature = "relay")] @@ -193,6 +233,8 @@ vctx: ValidationContext { opts: &validation_opts, }, + #[cfg(feature = "firehose-diagnostics")] + stats: shard_stats.clone(), batch: state.db.inner.batch(), count_deltas: CountDeltas::default(), #[cfg(feature = "relay")] @@ -346,7 +388,17 @@ } fn process_message(ctx: &mut WorkerContext, msg: WorkerMessage) -> Result<()> { - let Some(mut repo_state) = ctx.load_repo_state(&msg)? else { + #[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 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(()); }; let did = msg.msg.did().expect("already checked for did"); @@ -354,7 +406,20 @@ if let Some(host) = msg.firehose.host_str() && msg.is_pds { - let outcome = ctx.check_host_authority(did, &mut repo_state, host)?; + #[cfg(feature = "firehose-diagnostics")] + let authority_started = Instant::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 { + Ok(AuthorityOutcome::Authorized) => HostAuthorityStatsOutcome::Authorized, + Ok(AuthorityOutcome::WasStale) => HostAuthorityStatsOutcome::WasStale, + Ok(AuthorityOutcome::WrongHost { .. }) => HostAuthorityStatsOutcome::WrongHost, + Err(_) => HostAuthorityStatsOutcome::Error, + }, + ); + let outcome = outcome_result?; if let AuthorityOutcome::WrongHost { expected } = outcome { if !ctx.inc_error(host) { warn!(got = host, expected = %expected, "message rejected: wrong host"); @@ -367,19 +432,50 @@ match msg.msg { SubscribeReposMessage::Commit(commit) => { trace!("processing commit"); - Self::handle_commit(ctx, &mut repo_state, &msg.firehose, *commit) + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::now(); + let result = Self::handle_commit(ctx, &mut repo_state, &msg.firehose, *commit); + #[cfg(feature = "firehose-diagnostics")] + ctx.stats + .record_handle_message(RelayMessageKind::Commit, started.elapsed()); + result } SubscribeReposMessage::Sync(sync) => { debug!("processing sync"); - Self::handle_sync(ctx, &mut repo_state, &msg.firehose, *sync) + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::now(); + let result = Self::handle_sync(ctx, &mut repo_state, &msg.firehose, *sync); + #[cfg(feature = "firehose-diagnostics")] + ctx.stats + .record_handle_message(RelayMessageKind::Sync, started.elapsed()); + result } SubscribeReposMessage::Identity(identity) => { debug!("processing identity"); - Self::handle_identity(ctx, &mut repo_state, &msg.firehose, *identity, msg.is_pds) + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::now(); + let result = Self::handle_identity( + ctx, + &mut repo_state, + &msg.firehose, + *identity, + msg.is_pds, + ); + #[cfg(feature = "firehose-diagnostics")] + ctx.stats + .record_handle_message(RelayMessageKind::Identity, started.elapsed()); + result } SubscribeReposMessage::Account(account) => { debug!("processing account"); - Self::handle_account(ctx, &mut repo_state, &msg.firehose, *account, msg.is_pds) + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::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(RelayMessageKind::Account, started.elapsed()); + result } _ => Ok(()), } @@ -571,7 +667,11 @@ // 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 = Instant::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) => { repo_state.update_from_doc(doc); @@ -816,20 +916,32 @@ } fn refresh_doc(&mut self, did: &Did, repo_state: &mut RepoState) -> Result<()> { - let db = &self.state.db; - self.state.resolver.invalidate_sync(did); - let doc = Handle::current() - .block_on(self.state.resolver.resolve_doc(did)) - .map_err(|e| miette::miette!("{e}"))?; - repo_state.update_from_doc(doc); - repo_state.touch(); + #[cfg(feature = "firehose-diagnostics")] + let refresh_started = Instant::now(); + let result = (|| { + let db = &self.state.db; + self.state.resolver.invalidate_sync(did); + #[cfg(feature = "firehose-diagnostics")] + let resolve_started = Instant::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); + repo_state.touch(); - self.batch.insert( - &db.repos, - keys::repo_key(did), - db::ser_repo_state(repo_state)?, - ); - Ok(()) + self.batch.insert( + &db.repos, + keys::repo_key(did), + db::ser_repo_state(repo_state)?, + ); + Ok(()) + })(); + #[cfg(feature = "firehose-diagnostics")] + self.stats.record_refresh_doc(refresh_started.elapsed()); + result } fn validate_commit<'c>( @@ -839,7 +951,15 @@ ) -> Result>> { let did = &commit.repo; let key = self.fetch_key(did)?; - match self.vctx.validate_commit(commit, repo_state, key.as_ref()) { + #[cfg(feature = "firehose-diagnostics")] + let validate_started = Instant::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), + ); + match validation { Ok(v) => return Ok(Some(v)), Err(CommitValidationError::StaleRev) => { trace!("skipping replayed commit"); @@ -854,7 +974,15 @@ self.refresh_doc(did, repo_state)?; let key = self.fetch_key(did)?; - match self.vctx.validate_commit(commit, repo_state, key.as_ref()) { + #[cfg(feature = "firehose-diagnostics")] + let validate_started = Instant::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), + ); + match validation { Ok(v) => Ok(Some(v)), Err(e) => { debug!(err = %e, "commit rejected after key refresh"); @@ -870,7 +998,15 @@ ) -> Result> { let did = &sync.did; let key = self.fetch_key(did)?; - match self.vctx.validate_sync(sync, key.as_ref()) { + #[cfg(feature = "firehose-diagnostics")] + let validate_started = Instant::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), + ); + match validation { Ok(v) => return Ok(Some(v)), Err(SyncValidationError::SigFailure) => {} Err(e) => { @@ -881,7 +1017,15 @@ self.refresh_doc(did, repo_state)?; let key = self.fetch_key(did)?; - match self.vctx.validate_sync(sync, key.as_ref()) { + #[cfg(feature = "firehose-diagnostics")] + let validate_started = Instant::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), + ); + match validation { Ok(v) => Ok(Some(v)), Err(e) => { debug!(err = %e, "sync rejected after key refresh"); @@ -891,14 +1035,19 @@ } fn fetch_key(&self, did: &Did) -> Result>> { - if self.verify_signatures { - let key = Handle::current() + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::now(); + let result = if self.verify_signatures { + Handle::current() .block_on(self.state.resolver.resolve_signing_key(did)) - .map_err(|e| miette::miette!("{e}"))?; - Ok(Some(key)) + .map(Some) + .map_err(|e| miette::miette!("{e}")) } else { Ok(None) - } + }; + #[cfg(feature = "firehose-diagnostics")] + self.stats.record_fetch_key(started.elapsed()); + result } /// maps an inactive account status to the corresponding `RepoStatus`. @@ -973,6 +1122,9 @@ 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(RepoStateLoadOutcome::Drop); return Ok(None); } @@ -984,6 +1136,9 @@ .transpose()?; if let Some(repo_state) = repo_state_opt { + #[cfg(feature = "firehose-diagnostics")] + self.stats + .record_repo_state_outcome(RepoStateLoadOutcome::Hit); return Ok(Some(repo_state)); } @@ -993,7 +1148,12 @@ if filter.mode == crate::filter::FilterMode::Filter && !filter.signals.is_empty() { let commit = match &msg.msg { SubscribeReposMessage::Commit(c) => c, - _ => return Ok(None), + _ => { + #[cfg(feature = "firehose-diagnostics")] + self.stats + .record_repo_state_outcome(RepoStateLoadOutcome::Drop); + return Ok(None); + } }; let touches_signal = commit.ops.iter().any(|op| { op.path @@ -1011,24 +1171,37 @@ }); if !touches_signal { trace!(did = %did, "dropping commit, no signal-matching ops"); + #[cfg(feature = "firehose-diagnostics")] + self.stats + .record_repo_state_outcome(RepoStateLoadOutcome::Drop); return Ok(None); } } } debug!(did = %did, "discovered new account from firehose, queueing backfill"); + #[cfg(feature = "firehose-diagnostics")] + let new_account_started = Instant::now(); // resolve doc to initialize repo state self.state.resolver.invalidate_sync(did); + #[cfg(feature = "firehose-diagnostics")] + let resolve_started = Instant::now(); let doc = tokio::runtime::Handle::current() .block_on(self.state.resolver.resolve_doc(did)) - .into_diagnostic()?; + .into_diagnostic(); + #[cfg(feature = "firehose-diagnostics")] + self.stats.record_resolve_doc(resolve_started.elapsed()); + let doc = doc?; // if it's a PDS, verify it's the authoritative one if msg.is_pds { 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(RepoStateLoadOutcome::Drop); return Ok(None); } @@ -1036,14 +1209,22 @@ let count = self.state.db.get_count_sync(&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(RepoStateLoadOutcome::Drop); return Ok(None); } } } // try to get upstream status - let mut repo_state = tokio::runtime::Handle::current() - .block_on(self.check_repo_status(did, &doc.pds)) + #[cfg(feature = "firehose-diagnostics")] + let probe_started = Instant::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()); + let mut repo_state = repo_state .ok() .flatten() .unwrap_or_else(RepoState::backfilling); @@ -1073,20 +1254,34 @@ self.count_deltas.add("repos", 1); + #[cfg(feature = "firehose-diagnostics")] + { + self.stats + .record_repo_state_outcome(RepoStateLoadOutcome::Miss); + self.stats.record_new_account(new_account_started.elapsed()); + } + Ok(Some(repo_state)) } #[cfg(feature = "relay")] fn queue_emit(&mut self, make_frame: impl FnOnce(i64) -> Result) -> Result { - let db = &self.state.db; - let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); - let frame = make_frame(seq as i64)?; - self.batch - .insert(&db.relay_events, keys::relay_event_key(seq), frame.as_ref()); - self.pending_broadcasts - .push(RelayBroadcast::Ephemeral(seq, frame)); - self.pending_broadcasts.push(RelayBroadcast::Persisted(seq)); - Ok(seq) + #[cfg(feature = "firehose-diagnostics")] + let started = Instant::now(); + let result = (|| { + let db = &self.state.db; + let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); + let frame = make_frame(seq as i64)?; + self.batch + .insert(&db.relay_events, keys::relay_event_key(seq), frame.as_ref()); + self.pending_broadcasts + .push(RelayBroadcast::Ephemeral(seq, frame)); + self.pending_broadcasts.push(RelayBroadcast::Persisted(seq)); + Ok(seq) + })(); + #[cfg(feature = "firehose-diagnostics")] + self.stats.record_queue_emit(started.elapsed()); + result } } -- tangled.sh