1use crate::tree::{error::InsertBlockFatalError, TreeOutcome};
2use alloy_rpc_types_engine::{PayloadStatus, PayloadStatusEnum};
3use reth_engine_primitives::{ForkchoiceStatus, OnForkChoiceUpdated};
4use reth_errors::ProviderError;
5use reth_evm::metrics::ExecutorMetrics;
6use reth_execution_types::BlockExecutionOutput;
7use reth_metrics::{
8 metrics::{Counter, Gauge, Histogram},
9 thread::{ThreadResourceUsage, ThreadResourceUsageDelta},
10 Metrics,
11};
12use reth_primitives_traits::{constants::gas_units::MEGAGAS, FastInstant as Instant};
13use reth_trie::updates::TrieUpdates;
14use std::time::Duration;
15
16const GAS_BUCKET_THRESHOLDS: [u64; 5] =
19 [5 * MEGAGAS, 10 * MEGAGAS, 20 * MEGAGAS, 30 * MEGAGAS, 40 * MEGAGAS];
20
21const NUM_GAS_BUCKETS: usize = GAS_BUCKET_THRESHOLDS.len() + 1;
23
24#[derive(Debug, Default)]
26pub struct EngineApiMetrics {
27 pub engine: EngineMetrics,
29 pub executor: ExecutorMetrics,
31 pub block_validation: BlockValidationMetrics,
33 pub tree: TreeMetrics,
35 #[allow(dead_code)]
37 pub(crate) bal: BalMetrics,
38 pub(crate) execution_gas_buckets: ExecutionGasBucketMetrics,
40 pub(crate) block_validation_gas_buckets: BlockValidationGasBucketMetrics,
42}
43
44impl EngineApiMetrics {
45 pub fn record_block_execution<R>(
50 &self,
51 output: &BlockExecutionOutput<R>,
52 execution_duration: Duration,
53 ) {
54 let execution_secs = execution_duration.as_secs_f64();
55 let gas_used = output.result.gas_used;
56
57 self.executor.gas_processed_total.increment(gas_used);
59 self.executor.gas_per_second.set(gas_used as f64 / execution_secs);
60 self.executor.gas_used_histogram.record(gas_used as f64);
61 self.executor.execution_histogram.record(execution_secs);
62 self.executor.execution_duration.set(execution_secs);
63
64 let accounts = output.state.state.len();
66 let storage_slots =
67 output.state.state.values().map(|account| account.storage.len()).sum::<usize>();
68 let bytecodes = output.state.contracts.len();
69
70 self.executor.accounts_updated_histogram.record(accounts as f64);
71 self.executor.storage_slots_updated_histogram.record(storage_slots as f64);
72 self.executor.bytecodes_updated_histogram.record(bytecodes as f64);
73 }
74
75 pub const fn executor_metrics(&self) -> &ExecutorMetrics {
77 &self.executor
78 }
79
80 pub fn record_pre_execution(&self, elapsed: Duration) {
82 self.executor.pre_execution_histogram.record(elapsed);
83 }
84
85 pub fn record_post_execution(&self, elapsed: Duration) {
87 self.executor.post_execution_histogram.record(elapsed);
88 }
89
90 pub fn record_block_execution_gas_bucket(&self, gas_used: u64, elapsed: Duration) {
92 let idx = GasBucketMetrics::bucket_index(gas_used);
93 self.execution_gas_buckets.buckets[idx]
94 .execution_gas_bucket_histogram
95 .record(elapsed.as_secs_f64());
96 }
97
98 pub fn record_state_root_gas_bucket(&self, gas_used: u64, elapsed_secs: f64) {
100 let idx = GasBucketMetrics::bucket_index(gas_used);
101 self.block_validation_gas_buckets.buckets[idx]
102 .state_root_gas_bucket_histogram
103 .record(elapsed_secs);
104 }
105
106 pub fn record_transaction_wait(&self, elapsed: Duration) {
108 self.executor.transaction_wait_histogram.record(elapsed);
109 }
110
111 pub fn record_transaction_execution(&self, elapsed: Duration) {
113 self.executor.transaction_execution_histogram.record(elapsed);
114 }
115}
116
117#[derive(Metrics)]
119#[metrics(scope = "blockchain_tree")]
120pub struct TreeMetrics {
121 pub canonical_chain_height: Gauge,
123 #[metric(skip)]
125 pub reorgs: ReorgMetrics,
126 pub latest_reorg_depth: Gauge,
128 pub safe_block_height: Gauge,
130 pub finalized_block_height: Gauge,
132}
133
134#[derive(Debug)]
136pub struct ReorgMetrics {
137 pub head: Counter,
139 pub safe: Counter,
141 pub finalized: Counter,
143}
144
145impl Default for ReorgMetrics {
146 fn default() -> Self {
147 Self {
148 head: metrics::counter!("blockchain_tree_reorgs", "commitment" => "head"),
149 safe: metrics::counter!("blockchain_tree_reorgs", "commitment" => "safe"),
150 finalized: metrics::counter!("blockchain_tree_reorgs", "commitment" => "finalized"),
151 }
152 }
153}
154
155#[derive(Metrics)]
157#[metrics(scope = "consensus.engine.beacon")]
158pub struct EngineMetrics {
159 #[metric(skip)]
161 pub(crate) forkchoice_updated: ForkchoiceUpdatedMetrics,
162 #[metric(skip)]
164 pub(crate) new_payload: NewPayloadStatusMetrics,
165 pub(crate) executed_blocks: Gauge,
167 pub(crate) inserted_already_executed_blocks: Counter,
169 pub(crate) pipeline_runs: Counter,
171 pub(crate) executed_new_block_cache_miss: Counter,
173 pub(crate) persistence_duration: Histogram,
175 pub(crate) backpressure_active: Gauge,
177 pub(crate) backpressure_stall_duration: Histogram,
179 pub(crate) failed_new_payload_response_deliveries: Counter,
185 pub(crate) failed_forkchoice_updated_response_deliveries: Counter,
187 pub(crate) block_insert_total_duration: Histogram,
189}
190
191#[derive(Metrics)]
193#[metrics(scope = "consensus.engine.beacon")]
194pub(crate) struct ForkchoiceUpdatedMetrics {
195 #[metric(skip)]
197 pub(crate) latest_finish_at: Option<Instant>,
198 #[metric(skip)]
200 pub(crate) latest_start_at: Option<Instant>,
201 pub(crate) forkchoice_updated_messages: Counter,
203 pub(crate) forkchoice_with_attributes_updated_messages: Counter,
205 pub(crate) forkchoice_updated_valid: Counter,
208 pub(crate) forkchoice_updated_invalid: Counter,
211 pub(crate) forkchoice_updated_syncing: Counter,
214 pub(crate) forkchoice_updated_error: Counter,
217 pub(crate) forkchoice_updated_latency: Histogram,
219 pub(crate) forkchoice_updated_last: Gauge,
221 pub(crate) new_payload_forkchoice_updated_time_diff: Histogram,
223 pub(crate) time_between_forkchoice_updated: Histogram,
226 pub(crate) forkchoice_updated_interval: Histogram,
229}
230
231impl ForkchoiceUpdatedMetrics {
232 pub(crate) fn update_response_metrics(
234 &mut self,
235 start: Instant,
236 latest_new_payload_at: &mut Option<Instant>,
237 has_attrs: bool,
238 result: &Result<TreeOutcome<OnForkChoiceUpdated>, ProviderError>,
239 ) {
240 let finish = Instant::now();
241 let elapsed = finish - start;
242
243 if let Some(prev_finish) = self.latest_finish_at {
244 self.time_between_forkchoice_updated.record(start - prev_finish);
245 }
246 if let Some(prev_start) = self.latest_start_at {
247 self.forkchoice_updated_interval.record(start - prev_start);
248 }
249 self.latest_finish_at = Some(finish);
250 self.latest_start_at = Some(start);
251
252 match result {
253 Ok(outcome) => match outcome.outcome.forkchoice_status() {
254 ForkchoiceStatus::Valid => self.forkchoice_updated_valid.increment(1),
255 ForkchoiceStatus::Invalid => self.forkchoice_updated_invalid.increment(1),
256 ForkchoiceStatus::Syncing => self.forkchoice_updated_syncing.increment(1),
257 },
258 Err(_) => self.forkchoice_updated_error.increment(1),
259 }
260 self.forkchoice_updated_messages.increment(1);
261 if has_attrs {
262 self.forkchoice_with_attributes_updated_messages.increment(1);
263 }
264 self.forkchoice_updated_latency.record(elapsed);
265 self.forkchoice_updated_last.set(elapsed);
266 if let Some(latest_new_payload_at) = latest_new_payload_at.take() {
267 self.new_payload_forkchoice_updated_time_diff.record(start - latest_new_payload_at);
268 }
269 }
270}
271
272#[derive(Clone, Metrics)]
274#[metrics(scope = "consensus.engine.beacon")]
275pub(crate) struct NewPayloadGasBucketMetrics {
276 pub(crate) new_payload_gas_bucket_latency: Histogram,
278 pub(crate) new_payload_gas_bucket_gas_per_second: Histogram,
280}
281
282#[derive(Debug)]
284pub(crate) struct GasBucketMetrics {
285 buckets: [NewPayloadGasBucketMetrics; NUM_GAS_BUCKETS],
286}
287
288impl Default for GasBucketMetrics {
289 fn default() -> Self {
290 Self {
291 buckets: std::array::from_fn(|i| {
292 let label = Self::bucket_label(i);
293 NewPayloadGasBucketMetrics::new_with_labels(&[("gas_bucket", label)])
294 }),
295 }
296 }
297}
298
299impl GasBucketMetrics {
300 fn record(&self, gas_used: u64, elapsed: Duration) {
301 let idx = Self::bucket_index(gas_used);
302 self.buckets[idx].new_payload_gas_bucket_latency.record(elapsed);
303 self.buckets[idx]
304 .new_payload_gas_bucket_gas_per_second
305 .record(gas_used as f64 / elapsed.as_secs_f64());
306 }
307
308 pub(crate) fn bucket_index(gas_used: u64) -> usize {
310 GAS_BUCKET_THRESHOLDS
311 .iter()
312 .position(|&threshold| gas_used < threshold)
313 .unwrap_or(GAS_BUCKET_THRESHOLDS.len())
314 }
315
316 pub(crate) fn bucket_label(index: usize) -> String {
318 if index == 0 {
319 let hi = GAS_BUCKET_THRESHOLDS[0] / MEGAGAS;
320 format!("<{hi}M")
321 } else if index < GAS_BUCKET_THRESHOLDS.len() {
322 let lo = GAS_BUCKET_THRESHOLDS[index - 1] / MEGAGAS;
323 let hi = GAS_BUCKET_THRESHOLDS[index] / MEGAGAS;
324 format!("{lo}-{hi}M")
325 } else {
326 let lo = GAS_BUCKET_THRESHOLDS[GAS_BUCKET_THRESHOLDS.len() - 1] / MEGAGAS;
327 format!(">{lo}M")
328 }
329 }
330}
331
332#[derive(Clone, Metrics)]
334#[metrics(scope = "sync.execution")]
335pub(crate) struct ExecutionGasBucketSeries {
336 pub(crate) execution_gas_bucket_histogram: Histogram,
338}
339
340#[derive(Debug)]
342pub(crate) struct ExecutionGasBucketMetrics {
343 buckets: [ExecutionGasBucketSeries; NUM_GAS_BUCKETS],
344}
345
346impl Default for ExecutionGasBucketMetrics {
347 fn default() -> Self {
348 Self {
349 buckets: std::array::from_fn(|i| {
350 let label = GasBucketMetrics::bucket_label(i);
351 ExecutionGasBucketSeries::new_with_labels(&[("gas_bucket", label)])
352 }),
353 }
354 }
355}
356
357#[derive(Clone, Metrics)]
359#[metrics(scope = "sync.block_validation")]
360pub(crate) struct BlockValidationGasBucketSeries {
361 pub(crate) state_root_gas_bucket_histogram: Histogram,
363}
364
365#[derive(Debug)]
367pub(crate) struct BlockValidationGasBucketMetrics {
368 buckets: [BlockValidationGasBucketSeries; NUM_GAS_BUCKETS],
369}
370
371impl Default for BlockValidationGasBucketMetrics {
372 fn default() -> Self {
373 Self {
374 buckets: std::array::from_fn(|i| {
375 let label = GasBucketMetrics::bucket_label(i);
376 BlockValidationGasBucketSeries::new_with_labels(&[("gas_bucket", label)])
377 }),
378 }
379 }
380}
381
382#[derive(Metrics)]
384#[metrics(scope = "consensus.engine.beacon")]
385pub(crate) struct NewPayloadStatusMetrics {
386 #[metric(skip)]
388 pub(crate) latest_finish_at: Option<Instant>,
389 #[metric(skip)]
391 pub(crate) latest_start_at: Option<Instant>,
392 #[metric(skip)]
394 pub(crate) gas_bucket: GasBucketMetrics,
395 #[metric(skip)]
397 thread_resource_usage: NewPayloadThreadResourceMetrics,
398 pub(crate) new_payload_messages: Counter,
400 pub(crate) new_payload_valid: Counter,
403 pub(crate) new_payload_invalid: Counter,
406 pub(crate) new_payload_syncing: Counter,
409 pub(crate) new_payload_accepted: Counter,
412 pub(crate) new_payload_error: Counter,
415 pub(crate) new_payload_total_gas: Histogram,
417 pub(crate) new_payload_total_gas_last: Gauge,
419 pub(crate) new_payload_gas_per_second: Histogram,
421 pub(crate) new_payload_gas_per_second_last: Gauge,
423 pub(crate) new_payload_latency: Histogram,
425 pub(crate) new_payload_last: Gauge,
427 pub(crate) time_between_new_payloads: Histogram,
429 pub(crate) new_payload_interval: Histogram,
431 pub(crate) forkchoice_updated_new_payload_time_diff: Histogram,
433}
434
435impl NewPayloadStatusMetrics {
436 pub(crate) fn measure_thread_resource_usage(&self) -> NewPayloadThreadResourceGuard {
438 self.thread_resource_usage.measure()
439 }
440
441 pub(crate) fn update_response_metrics(
443 &mut self,
444 start: Instant,
445 latest_forkchoice_updated_at: &mut Option<Instant>,
446 result: &Result<TreeOutcome<PayloadStatus>, InsertBlockFatalError>,
447 gas_used: u64,
448 ) {
449 let finish = Instant::now();
450 let elapsed = finish - start;
451
452 if let Some(prev_finish) = self.latest_finish_at {
453 self.time_between_new_payloads.record(start - prev_finish);
454 }
455 if let Some(prev_start) = self.latest_start_at {
456 self.new_payload_interval.record(start - prev_start);
457 }
458 self.latest_finish_at = Some(finish);
459 self.latest_start_at = Some(start);
460 match result {
461 Ok(outcome) => match outcome.outcome.status {
462 PayloadStatusEnum::Valid => {
463 self.new_payload_valid.increment(1);
464 if !outcome.already_seen {
465 self.new_payload_total_gas.record(gas_used as f64);
466 self.new_payload_total_gas_last.set(gas_used as f64);
467 let gas_per_second = gas_used as f64 / elapsed.as_secs_f64();
468 self.new_payload_gas_per_second.record(gas_per_second);
469 self.new_payload_gas_per_second_last.set(gas_per_second);
470
471 self.new_payload_latency.record(elapsed);
472 self.new_payload_last.set(elapsed);
473 self.gas_bucket.record(gas_used, elapsed);
474 }
475 }
476 PayloadStatusEnum::Syncing => self.new_payload_syncing.increment(1),
477 PayloadStatusEnum::Accepted => self.new_payload_accepted.increment(1),
478 PayloadStatusEnum::Invalid { .. } => self.new_payload_invalid.increment(1),
479 },
480 Err(_) => self.new_payload_error.increment(1),
481 }
482 self.new_payload_messages.increment(1);
483 if let Some(latest_forkchoice_updated_at) = latest_forkchoice_updated_at.take() {
484 self.forkchoice_updated_new_payload_time_diff
485 .record(start - latest_forkchoice_updated_at);
486 }
487 }
488}
489
490#[derive(Clone, Metrics)]
492#[metrics(scope = "consensus.engine.beacon")]
493struct NewPayloadThreadResourceMetrics {
494 new_payload_thread_user_cpu_seconds: Histogram,
496 new_payload_thread_system_cpu_seconds: Histogram,
498 new_payload_thread_minor_page_faults: Histogram,
500 new_payload_thread_major_page_faults: Histogram,
502 new_payload_thread_voluntary_context_switches: Histogram,
504 new_payload_thread_involuntary_context_switches: Histogram,
506 new_payload_thread_block_input_operations: Histogram,
508 new_payload_thread_block_output_operations: Histogram,
510}
511
512impl NewPayloadThreadResourceMetrics {
513 fn measure(&self) -> NewPayloadThreadResourceGuard {
514 let metrics = self.clone();
515 let start = ThreadResourceUsage::now();
516 NewPayloadThreadResourceGuard { start, metrics }
517 }
518
519 fn record(&self, usage: &ThreadResourceUsageDelta) {
520 self.new_payload_thread_user_cpu_seconds.record(usage.user_cpu_time);
521 self.new_payload_thread_system_cpu_seconds.record(usage.system_cpu_time);
522 self.new_payload_thread_minor_page_faults.record(usage.minor_page_faults as f64);
523 self.new_payload_thread_major_page_faults.record(usage.major_page_faults as f64);
524 self.new_payload_thread_voluntary_context_switches
525 .record(usage.voluntary_context_switches as f64);
526 self.new_payload_thread_involuntary_context_switches
527 .record(usage.involuntary_context_switches as f64);
528 self.new_payload_thread_block_input_operations.record(usage.block_input_operations as f64);
529 self.new_payload_thread_block_output_operations
530 .record(usage.block_output_operations as f64);
531 }
532}
533
534pub(crate) struct NewPayloadThreadResourceGuard {
536 start: ThreadResourceUsage,
537 metrics: NewPayloadThreadResourceMetrics,
538}
539
540impl Drop for NewPayloadThreadResourceGuard {
541 fn drop(&mut self) {
542 if let Some(usage) = self.start.elapsed() {
543 self.metrics.record(&usage);
544 }
545 }
546}
547
548#[allow(dead_code)]
552#[derive(Metrics, Clone)]
553#[metrics(scope = "execution.block_access_list")]
554pub(crate) struct BalMetrics {
555 pub(crate) size_bytes: Gauge,
557 pub(crate) valid_total: Counter,
559 pub(crate) invalid_total: Counter,
561 pub(crate) validation_time_seconds: Histogram,
563 pub(crate) account_changes: Gauge,
565 pub(crate) storage_changes: Gauge,
567 pub(crate) balance_changes: Gauge,
569 pub(crate) nonce_changes: Gauge,
571 pub(crate) code_changes: Gauge,
573}
574
575#[derive(Metrics, Clone)]
577#[metrics(scope = "sync.block_validation")]
578pub struct BlockValidationMetrics {
579 pub state_root_storage_tries_updated_total: Counter,
581 pub state_root_task_fallback_success_total: Counter,
583 pub state_root_task_timeout_total: Counter,
585 pub state_root_duration: Gauge,
587 pub state_root_histogram: Histogram,
589 pub deferred_trie_compute_duration: Histogram,
591 pub payload_validation_duration: Gauge,
593 pub payload_validation_histogram: Histogram,
595 pub spawn_payload_processor: Histogram,
597 pub post_execution_validation_duration: Histogram,
599 pub total_duration: Histogram,
601 pub hashed_post_state_size: Histogram,
603 pub trie_updates_sorted_size: Histogram,
605}
606
607impl BlockValidationMetrics {
608 pub fn record_state_root(&self, trie_output: &TrieUpdates, elapsed_as_secs: f64) {
610 self.state_root_storage_tries_updated_total
611 .increment(trie_output.storage_tries_ref().len() as u64);
612 self.state_root_duration.set(elapsed_as_secs);
613 self.state_root_histogram.record(elapsed_as_secs);
614 }
615
616 pub fn record_payload_validation(&self, elapsed_as_secs: f64) {
619 self.payload_validation_duration.set(elapsed_as_secs);
620 self.payload_validation_histogram.record(elapsed_as_secs);
621 }
622}
623
624#[derive(Metrics)]
626#[metrics(scope = "blockchain_tree.block_buffer")]
627pub(crate) struct BlockBufferMetrics {
628 pub blocks: Gauge,
630}
631
632#[cfg(test)]
633mod tests {
634 use super::*;
635 use alloy_eips::eip7685::Requests;
636 use metrics_util::debugging::{DebuggingRecorder, Snapshotter};
637 use reth_ethereum_primitives::Receipt;
638 use reth_execution_types::BlockExecutionResult;
639 use reth_revm::db::BundleState;
640
641 fn setup_test_recorder() -> Snapshotter {
642 let recorder = DebuggingRecorder::new();
643 let snapshotter = recorder.snapshotter();
644 recorder.install().unwrap();
645 snapshotter
646 }
647
648 #[test]
649 fn test_record_block_execution_metrics() {
650 let snapshotter = setup_test_recorder();
651 let metrics = EngineApiMetrics::default();
652
653 metrics.executor.gas_processed_total.increment(0);
655 metrics.executor.gas_per_second.set(0.0);
656 metrics.executor.gas_used_histogram.record(0.0);
657
658 let output = BlockExecutionOutput::<Receipt> {
659 state: BundleState::default(),
660 result: BlockExecutionResult {
661 receipts: vec![],
662 requests: Requests::default(),
663 gas_used: 21000,
664 blob_gas_used: 0,
665 },
666 };
667
668 metrics.record_block_execution(&output, Duration::from_millis(100));
669 metrics.engine.new_payload.thread_resource_usage.record(&ThreadResourceUsageDelta {
670 user_cpu_time: Duration::from_millis(1),
671 system_cpu_time: Duration::from_millis(2),
672 minor_page_faults: 3,
673 major_page_faults: 4,
674 voluntary_context_switches: 5,
675 involuntary_context_switches: 6,
676 block_input_operations: 7,
677 block_output_operations: 8,
678 });
679
680 let snapshot = snapshotter.snapshot().into_vec();
681
682 let mut found_execution_metrics = false;
684 let mut found_thread_resource_metrics = false;
685 for (key, _unit, _desc, _value) in snapshot {
686 let metric_name = key.key().name();
687 if metric_name.starts_with("sync.execution") {
688 found_execution_metrics = true;
689 }
690 if metric_name == "consensus.engine.beacon.new_payload_thread_major_page_faults" {
691 found_thread_resource_metrics = true;
692 }
693 }
694
695 assert!(found_execution_metrics, "Expected to find sync.execution metrics");
696 assert!(found_thread_resource_metrics, "Expected to find thread resource metrics");
697 }
698}