Skip to main content

reth_engine_tree/tree/
instrumented_state.rs

1//! Implements a state provider that tracks latency metrics.
2use alloy_primitives::{Address, StorageKey, StorageValue, B256};
3use metrics::{Gauge, Histogram};
4use reth_errors::ProviderResult;
5use reth_metrics::Metrics;
6use reth_primitives_traits::{Account, Bytecode, FastInstant as Instant};
7use reth_provider::EvmStateProvider;
8use std::{
9    sync::{
10        atomic::{AtomicU64, AtomicUsize, Ordering},
11        Arc,
12    },
13    time::Duration,
14};
15
16/// Nanoseconds per second
17const NANOS_PER_SEC: u32 = 1_000_000_000;
18
19/// An atomic version of [`Duration`], using an [`AtomicU64`] to store the total nanoseconds in the
20/// duration.
21#[derive(Debug, Default)]
22pub(crate) struct AtomicDuration {
23    /// The nanoseconds part of the duration
24    ///
25    /// We would have to accumulate 584 years of nanoseconds to overflow a u64, so this is
26    /// sufficiently large for our use case. We don't expect to be adding arbitrary durations to
27    /// this value.
28    nanos: AtomicU64,
29}
30
31impl AtomicDuration {
32    /// Returns the duration as a [`Duration`]
33    pub(crate) fn duration(&self) -> Duration {
34        let nanos = self.nanos.load(Ordering::Relaxed);
35        let seconds = nanos / NANOS_PER_SEC as u64;
36        let nanos = nanos % NANOS_PER_SEC as u64;
37        // `as u32` is ok because we did a mod by u32 const
38        Duration::new(seconds, nanos as u32)
39    }
40
41    /// Adds a [`Duration`] to the atomic duration.
42    pub(crate) fn add_duration(&self, duration: Duration) {
43        // this is `as_nanos` but without the `as u128` - we do not expect durations over 584 years
44        // as input here
45        let total_nanos =
46            duration.as_secs() * NANOS_PER_SEC as u64 + duration.subsec_nanos() as u64;
47        // add the nanoseconds part of the duration
48        self.nanos.fetch_add(total_nanos, Ordering::Relaxed);
49    }
50}
51
52/// A wrapper of a state provider and latency metrics.
53#[derive(Debug)]
54pub struct InstrumentedStateProvider<S> {
55    /// The state provider
56    state_provider: S,
57    /// Prometheus metrics for the instrumented state provider
58    metrics: StateProviderMetrics,
59    /// Shared fetch statistics, readable after the provider is consumed.
60    stats: Arc<StateProviderStats>,
61}
62
63impl<S> InstrumentedStateProvider<S>
64where
65    S: EvmStateProvider,
66{
67    /// Creates a new [`InstrumentedStateProvider`] from a state provider with the provided label
68    /// for metrics.
69    pub fn new(state_provider: S, source: &'static str) -> Self {
70        Self::with_stats(
71            state_provider,
72            StateProviderMetrics::with_source(source),
73            Arc::new(StateProviderStats::default()),
74        )
75    }
76
77    /// Creates a new [`InstrumentedStateProvider`] that writes into shared statistics.
78    pub(crate) const fn with_stats(
79        state_provider: S,
80        metrics: StateProviderMetrics,
81        stats: Arc<StateProviderStats>,
82    ) -> Self {
83        Self { state_provider, metrics, stats }
84    }
85
86    /// Returns a shared reference to the accumulated fetch statistics.
87    pub fn stats(&self) -> Arc<StateProviderStats> {
88        Arc::clone(&self.stats)
89    }
90}
91
92/// Metrics for the instrumented state provider
93#[derive(Metrics, Clone)]
94#[metrics(scope = "sync.state_provider")]
95pub(crate) struct StateProviderMetrics {
96    /// A histogram of the time it takes to get a storage value
97    storage_fetch_latency: Histogram,
98
99    /// A histogram of the time it takes to get a code value
100    code_fetch_latency: Histogram,
101
102    /// A histogram of the time it takes to get an account value
103    account_fetch_latency: Histogram,
104
105    /// A histogram of the total time we spend fetching storage over the lifetime of this state
106    /// provider
107    total_storage_fetch_latency: Histogram,
108
109    /// A gauge of the total time we spend fetching storage over the lifetime of this state
110    /// provider
111    total_storage_fetch_latency_gauge: Gauge,
112
113    /// A histogram of the total time we spend fetching code over the lifetime of this state
114    /// provider
115    total_code_fetch_latency: Histogram,
116
117    /// A gauge of the total time we spend fetching code over the lifetime of this state provider
118    total_code_fetch_latency_gauge: Gauge,
119
120    /// A histogram of the total time we spend fetching accounts over the lifetime of this state
121    /// provider
122    total_account_fetch_latency: Histogram,
123
124    /// A gauge of the total time we spend fetching accounts over the lifetime of this state
125    /// provider
126    total_account_fetch_latency_gauge: Gauge,
127}
128
129impl StateProviderMetrics {
130    /// Creates state-provider metrics with the given source label.
131    pub(crate) fn with_source(source: &'static str) -> Self {
132        Self::new_with_labels(&[("source", source)])
133    }
134
135    /// Records accumulated fetch latency totals.
136    pub(crate) fn record_totals(&self, stats: &StateProviderStats) {
137        let total_storage_fetch_latency = stats.total_storage_fetch_latency.duration();
138        self.total_storage_fetch_latency.record(total_storage_fetch_latency);
139        self.total_storage_fetch_latency_gauge.set(total_storage_fetch_latency.as_secs_f64());
140
141        let total_code_fetch_latency = stats.total_code_fetch_latency.duration();
142        self.total_code_fetch_latency.record(total_code_fetch_latency);
143        self.total_code_fetch_latency_gauge.set(total_code_fetch_latency.as_secs_f64());
144
145        let total_account_fetch_latency = stats.total_account_fetch_latency.duration();
146        self.total_account_fetch_latency.record(total_account_fetch_latency);
147        self.total_account_fetch_latency_gauge.set(total_account_fetch_latency.as_secs_f64());
148    }
149}
150
151impl<S: EvmStateProvider> EvmStateProvider for InstrumentedStateProvider<S> {
152    fn basic_account(&self, address: &Address) -> ProviderResult<Option<Account>> {
153        let start = Instant::now();
154        let res = self.state_provider.basic_account(address);
155        let elapsed = start.elapsed();
156        self.metrics.account_fetch_latency.record(elapsed);
157        self.stats.total_account_fetches.fetch_add(1, Ordering::Relaxed);
158        self.stats.total_account_fetch_latency.add_duration(elapsed);
159        res
160    }
161
162    fn storage(
163        &self,
164        account: Address,
165        storage_key: StorageKey,
166    ) -> ProviderResult<Option<StorageValue>> {
167        let start = Instant::now();
168        let res = self.state_provider.storage(account, storage_key);
169        let elapsed = start.elapsed();
170        self.metrics.storage_fetch_latency.record(elapsed);
171        self.stats.total_storage_fetches.fetch_add(1, Ordering::Relaxed);
172        self.stats.total_storage_fetch_latency.add_duration(elapsed);
173        res
174    }
175
176    fn bytecode_by_hash(&self, code_hash: &B256) -> ProviderResult<Option<Bytecode>> {
177        let start = Instant::now();
178        let res = self.state_provider.bytecode_by_hash(code_hash);
179        let elapsed = start.elapsed();
180        self.metrics.code_fetch_latency.record(elapsed);
181        self.stats.total_code_fetches.fetch_add(1, Ordering::Relaxed);
182        self.stats.total_code_fetch_latency.add_duration(elapsed);
183        self.stats.total_code_fetched_bytes.fetch_add(
184            res.as_ref()
185                .ok()
186                .and_then(|code| code.as_ref().map(|code| code.len()))
187                .unwrap_or_default(),
188            Ordering::Relaxed,
189        );
190        res
191    }
192
193    fn block_hash(&self, number: alloy_primitives::BlockNumber) -> ProviderResult<Option<B256>> {
194        self.state_provider.block_hash(number)
195    }
196}
197
198/// Accumulated fetch statistics from an [`InstrumentedStateProvider`].
199///
200/// Shared via `Arc` so statistics can be read after the provider is consumed.
201#[derive(Debug, Default)]
202pub struct StateProviderStats {
203    total_storage_fetches: AtomicUsize,
204    total_storage_fetch_latency: AtomicDuration,
205
206    total_code_fetches: AtomicUsize,
207    total_code_fetch_latency: AtomicDuration,
208    total_code_fetched_bytes: AtomicUsize,
209
210    total_account_fetches: AtomicUsize,
211    total_account_fetch_latency: AtomicDuration,
212}
213
214impl StateProviderStats {
215    /// Returns total number of storage fetches.
216    pub fn total_storage_fetches(&self) -> usize {
217        self.total_storage_fetches.load(Ordering::Relaxed)
218    }
219
220    /// Returns total time spent on storage fetches.
221    pub fn total_storage_fetch_latency(&self) -> Duration {
222        self.total_storage_fetch_latency.duration()
223    }
224
225    /// Returns total number of code fetches.
226    pub fn total_code_fetches(&self) -> usize {
227        self.total_code_fetches.load(Ordering::Relaxed)
228    }
229
230    /// Returns total time spent on code fetches.
231    pub fn total_code_fetch_latency(&self) -> Duration {
232        self.total_code_fetch_latency.duration()
233    }
234
235    /// Returns total amount of code fetched, in bytes.
236    pub fn total_code_fetched_bytes(&self) -> usize {
237        self.total_code_fetched_bytes.load(Ordering::Relaxed)
238    }
239
240    /// Returns total number of account fetches.
241    pub fn total_account_fetches(&self) -> usize {
242        self.total_account_fetches.load(Ordering::Relaxed)
243    }
244
245    /// Returns total time spent on account fetches.
246    pub fn total_account_fetch_latency(&self) -> Duration {
247        self.total_account_fetch_latency.duration()
248    }
249}