arb_rpc/
arbdebug.rs

1//! `arbdebug_*` namespace — historical pricing + retryable queue
2//! introspection. Samples ArbOS state at each block in the requested
3//! range.
4
5use std::sync::Arc;
6
7use alloy_consensus::BlockHeader;
8use alloy_primitives::{Address, B256, StorageKey, U256};
9use alloy_rpc_types_eth::BlockNumberOrTag;
10use arb_storage::{
11    ARBOS_STATE_ADDRESS,
12    layout::{
13        L1_PRICING_SUBSPACE, L2_PRICING_SUBSPACE, RETRYABLES_SUBSPACE, ROOT_STORAGE_KEY,
14        derive_subspace_key, map_slot, subspace_slot,
15    },
16};
17use arbos::{
18    l1_pricing::{
19        AMORTIZED_COST_CAP_BIPS_OFFSET as L1_AMORTIZED_COST_CAP_BIPS_OFFSET,
20        EQUILIBRATION_UNITS_OFFSET as L1_EQUILIBRATION_UNITS_OFFSET,
21        FUNDS_DUE_FOR_REWARDS_OFFSET as L1_FUNDS_DUE_FOR_REWARDS_OFFSET,
22        INERTIA_OFFSET as L1_INERTIA_OFFSET,
23        L1_FEES_AVAILABLE_OFFSET as L1_L1_FEES_AVAILABLE_OFFSET,
24        LAST_SURPLUS_OFFSET as L1_LAST_SURPLUS_OFFSET,
25        LAST_UPDATE_TIME_OFFSET as L1_LAST_UPDATE_TIME_OFFSET,
26        PAY_REWARDS_TO_OFFSET as L1_PAY_REWARDS_TO_OFFSET,
27        PER_BATCH_GAS_COST_OFFSET as L1_PER_BATCH_GAS_COST_OFFSET,
28        PER_UNIT_REWARD_OFFSET as L1_PER_UNIT_REWARD_OFFSET,
29        PRICE_PER_UNIT_OFFSET as L1_PRICE_PER_UNIT_OFFSET,
30        UNITS_SINCE_OFFSET as L1_UNITS_SINCE_UPDATE_OFFSET,
31    },
32    l2_pricing::{
33        BACKLOG_TOLERANCE_OFFSET as L2_BACKLOG_TOLERANCE_OFFSET,
34        BASE_FEE_WEI_OFFSET as L2_BASE_FEE_OFFSET, GAS_BACKLOG_OFFSET as L2_GAS_BACKLOG_OFFSET,
35        MIN_BASE_FEE_WEI_OFFSET as L2_MIN_BASE_FEE_OFFSET,
36        PER_BLOCK_GAS_LIMIT_OFFSET as L2_PER_BLOCK_GAS_LIMIT_OFFSET,
37        PRICING_INERTIA_OFFSET as L2_PRICING_INERTIA_OFFSET,
38        SPEED_LIMIT_PER_SECOND_OFFSET as L2_SPEED_LIMIT_OFFSET,
39    },
40    retryables::{TIMEOUT_OFFSET, TIMEOUT_QUEUE_KEY},
41};
42use jsonrpsee::{
43    core::RpcResult,
44    proc_macros::rpc,
45    types::{ErrorObject, error::INTERNAL_ERROR_CODE},
46};
47use reth_provider::{BlockReaderIdExt, ReceiptProvider, StateProviderFactory};
48use serde::{Deserialize, Serialize};
49
50#[derive(Debug, Clone, Serialize, Deserialize)]
51#[serde(rename_all = "camelCase")]
52pub struct PricingModelHistory {
53    pub start: u64,
54    pub end: u64,
55    pub step: u64,
56    pub timestamp: Vec<u64>,
57    pub base_fee: Vec<U256>,
58    pub gas_backlog: Vec<u64>,
59    pub gas_used: Vec<u64>,
60    pub min_base_fee: U256,
61    pub speed_limit: u64,
62    pub per_block_gas_limit: u64,
63    pub per_tx_gas_limit: u64,
64    pub pricing_inertia: u64,
65    pub backlog_tolerance: u64,
66    pub l1_base_fee_estimate: Vec<U256>,
67    pub l1_last_surplus: Vec<U256>,
68    pub l1_funds_due: Vec<U256>,
69    pub l1_funds_due_for_rewards: Vec<U256>,
70    pub l1_units_since_update: Vec<u64>,
71    pub l1_last_update_time: Vec<u64>,
72    pub l1_equilibration_units: U256,
73    pub l1_per_batch_cost: i64,
74    pub l1_amortized_cost_cap_bips: u64,
75    pub l1_pricing_inertia: u64,
76    pub l1_per_unit_reward: u64,
77    pub l1_pay_reward_to: Address,
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize)]
81#[serde(rename_all = "camelCase")]
82pub struct TimeoutQueueHistory {
83    pub start: u64,
84    pub end: u64,
85    pub step: u64,
86    pub timestamp: Vec<u64>,
87    pub size: Vec<u64>,
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize)]
91#[serde(rename_all = "camelCase")]
92pub struct TimeoutQueue {
93    pub block_number: u64,
94    pub tickets: Vec<B256>,
95    pub timeouts: Vec<u64>,
96}
97
98#[rpc(server, namespace = "arbdebug")]
99pub trait ArbDebugApi {
100    #[method(name = "pricingModel")]
101    async fn pricing_model(&self, start: u64, end: u64) -> RpcResult<PricingModelHistory>;
102
103    #[method(name = "timeoutQueueHistory")]
104    async fn timeout_queue_history(&self, start: u64, end: u64) -> RpcResult<TimeoutQueueHistory>;
105
106    #[method(name = "timeoutQueue")]
107    async fn timeout_queue(&self, block_num: u64) -> RpcResult<TimeoutQueue>;
108}
109
110#[derive(Debug, Clone)]
111pub struct ArbDebugConfig {
112    /// Max samples per query. Zero disables arbdebug.
113    pub block_range_bound: u64,
114    /// Max tickets returned from `timeoutQueue`.
115    pub timeout_queue_bound: u64,
116}
117
118impl Default for ArbDebugConfig {
119    fn default() -> Self {
120        Self {
121            block_range_bound: 256,
122            timeout_queue_bound: 256,
123        }
124    }
125}
126
127pub struct ArbDebugHandler<Provider> {
128    provider: Provider,
129    config: Arc<ArbDebugConfig>,
130}
131
132impl<Provider: std::fmt::Debug> std::fmt::Debug for ArbDebugHandler<Provider> {
133    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
134        f.debug_struct("ArbDebugHandler")
135            .field("config", &self.config)
136            .finish_non_exhaustive()
137    }
138}
139
140impl<Provider: Clone> Clone for ArbDebugHandler<Provider> {
141    fn clone(&self) -> Self {
142        Self {
143            provider: self.provider.clone(),
144            config: self.config.clone(),
145        }
146    }
147}
148
149impl<Provider> ArbDebugHandler<Provider> {
150    pub fn new(provider: Provider, config: ArbDebugConfig) -> Self {
151        Self {
152            provider,
153            config: Arc::new(config),
154        }
155    }
156}
157
158fn internal_err(msg: impl std::fmt::Display) -> ErrorObject<'static> {
159    ErrorObject::owned(INTERNAL_ERROR_CODE, msg.to_string(), None::<()>)
160}
161
162/// Sample step-size: if the range exceeds `bound` blocks, step > 1 so the
163/// total number of samples stays ≤ bound.
164fn compute_step(start: u64, end: u64, bound: u64) -> (u64, u64, u64) {
165    let span = end.saturating_sub(start).saturating_add(1);
166    if span == 0 || bound == 0 {
167        return (start, 1, 0);
168    }
169    let step = if span > bound {
170        span.div_ceil(bound)
171    } else {
172        1
173    };
174    let samples = span.div_ceil(step).min(bound);
175    let first = end.saturating_sub(step.saturating_mul(samples.saturating_sub(1)));
176    (first, step, samples)
177}
178
179impl<Provider> ArbDebugHandler<Provider>
180where
181    Provider: StateProviderFactory + BlockReaderIdExt + ReceiptProvider + Clone + 'static,
182{
183    /// Total gas consumed in the given block, summed from the last
184    /// receipt's `cumulative_gas_used`.
185    fn block_gas_used(&self, block: u64) -> Result<u64, ErrorObject<'static>> {
186        use alloy_consensus::TxReceipt;
187        let receipts = self
188            .provider
189            .receipts_by_block(alloy_eips::BlockHashOrNumber::Number(block))
190            .map_err(internal_err)?
191            .unwrap_or_default();
192        Ok(receipts
193            .last()
194            .map(|r| r.cumulative_gas_used())
195            .unwrap_or(0))
196    }
197
198    fn check_enabled(&self) -> Result<(), ErrorObject<'static>> {
199        if self.config.block_range_bound == 0 {
200            return Err(internal_err("arbdebug disabled (block_range_bound = 0)"));
201        }
202        Ok(())
203    }
204
205    fn validate_range(&self, start: u64, end: u64) -> Result<(), ErrorObject<'static>> {
206        if start > end {
207            return Err(internal_err(format!(
208                "invalid range: start {start} > end {end}"
209            )));
210        }
211        Ok(())
212    }
213
214    fn header_timestamp(&self, block: u64) -> Result<u64, ErrorObject<'static>> {
215        let header = self
216            .provider
217            .sealed_header_by_number_or_tag(BlockNumberOrTag::Number(block))
218            .map_err(internal_err)?
219            .ok_or_else(|| internal_err(format!("block {block} not found")))?;
220        Ok(header.timestamp())
221    }
222
223    fn read_slot(&self, block: u64, slot: U256) -> Result<U256, ErrorObject<'static>> {
224        let state = self
225            .provider
226            .state_by_block_id(BlockNumberOrTag::Number(block).into())
227            .map_err(internal_err)?;
228        let k = StorageKey::from(B256::from(slot.to_be_bytes::<32>()));
229        Ok(state
230            .storage(ARBOS_STATE_ADDRESS, k)
231            .map_err(internal_err)?
232            .unwrap_or(U256::ZERO))
233    }
234
235    fn read_l1_field(&self, block: u64, offset: u64) -> Result<U256, ErrorObject<'static>> {
236        self.read_slot(block, subspace_slot(L1_PRICING_SUBSPACE, offset))
237    }
238
239    fn read_l2_field(&self, block: u64, offset: u64) -> Result<U256, ErrorObject<'static>> {
240        self.read_slot(block, subspace_slot(L2_PRICING_SUBSPACE, offset))
241    }
242
243    /// Storage key for the retryable timeout queue's body.
244    fn retryable_queue_storage_key() -> B256 {
245        let retryables = derive_subspace_key(ROOT_STORAGE_KEY, RETRYABLES_SUBSPACE);
246        derive_subspace_key(retryables.as_slice(), TIMEOUT_QUEUE_KEY)
247    }
248
249    /// Size of the timeout queue at the given block.
250    fn queue_size_at(&self, block: u64) -> Result<u64, ErrorObject<'static>> {
251        // Queue layout: the queue's own storage has offset-0 = next_put,
252        // offset-1 = next_get. `size = next_put - next_get`.
253        let qk = Self::retryable_queue_storage_key();
254        let put = self
255            .read_slot(block, map_slot(qk.as_slice(), 0))?
256            .try_into()
257            .unwrap_or(0u64);
258        let get = self
259            .read_slot(block, map_slot(qk.as_slice(), 1))?
260            .try_into()
261            .unwrap_or(0u64);
262        Ok(put.saturating_sub(get))
263    }
264
265    /// Enumerate `(ticket_id, timeout)` for every pending retryable
266    /// in the timeout queue at the given block, up to `max_entries`.
267    /// Reads are against the historic state for `block` via the
268    /// StateProvider — no ArbosState instantiation required.
269    fn queue_snapshot_at(
270        &self,
271        block: u64,
272        max_entries: usize,
273    ) -> Result<Vec<(B256, u64)>, ErrorObject<'static>> {
274        let qk = Self::retryable_queue_storage_key();
275        let put: u64 = self
276            .read_slot(block, map_slot(qk.as_slice(), 0))?
277            .try_into()
278            .unwrap_or(0);
279        let get: u64 = self
280            .read_slot(block, map_slot(qk.as_slice(), 1))?
281            .try_into()
282            .unwrap_or(0);
283        let retryables_key = derive_subspace_key(ROOT_STORAGE_KEY, RETRYABLES_SUBSPACE);
284        let mut out = Vec::new();
285        for idx in get..put {
286            if out.len() >= max_entries {
287                break;
288            }
289            let ticket_slot = map_slot(qk.as_slice(), idx);
290            let ticket_word = self.read_slot(block, ticket_slot)?;
291            let id = B256::from(ticket_word.to_be_bytes::<32>());
292            if id == B256::ZERO {
293                continue;
294            }
295            // Per-retryable storage is keyed by ticket_id as a subspace
296            // of the retryables subspace. Read the timeout field.
297            let ret_key = derive_subspace_key(retryables_key.as_slice(), id.as_slice());
298            let timeout: u64 = self
299                .read_slot(block, map_slot(ret_key.as_slice(), TIMEOUT_OFFSET))?
300                .try_into()
301                .unwrap_or(0);
302            if timeout == 0 {
303                continue;
304            }
305            out.push((id, timeout));
306        }
307        Ok(out)
308    }
309}
310
311#[async_trait::async_trait]
312impl<Provider> ArbDebugApiServer for ArbDebugHandler<Provider>
313where
314    Provider:
315        StateProviderFactory + BlockReaderIdExt + ReceiptProvider + Clone + Send + Sync + 'static,
316{
317    async fn pricing_model(&self, start: u64, end: u64) -> RpcResult<PricingModelHistory> {
318        self.check_enabled()?;
319        self.validate_range(start, end)?;
320        let (first, step, samples) = compute_step(start, end, self.config.block_range_bound);
321
322        let mut timestamp = Vec::with_capacity(samples as usize);
323        let mut base_fee = Vec::with_capacity(samples as usize);
324        let mut gas_backlog = Vec::with_capacity(samples as usize);
325        let mut gas_used = Vec::with_capacity(samples as usize);
326        let mut l1_base_fee_estimate = Vec::with_capacity(samples as usize);
327        let mut l1_last_surplus = Vec::with_capacity(samples as usize);
328        let mut l1_funds_due = Vec::with_capacity(samples as usize);
329        let mut l1_funds_due_for_rewards = Vec::with_capacity(samples as usize);
330        let mut l1_units_since_update = Vec::with_capacity(samples as usize);
331        let mut l1_last_update_time = Vec::with_capacity(samples as usize);
332
333        for i in 0..samples {
334            let b = first + step * i;
335            timestamp.push(self.header_timestamp(b)?);
336            base_fee.push(self.read_l2_field(b, L2_BASE_FEE_OFFSET)?);
337            gas_backlog.push(
338                self.read_l2_field(b, L2_GAS_BACKLOG_OFFSET)?
339                    .try_into()
340                    .unwrap_or(0u64),
341            );
342            gas_used.push(self.block_gas_used(b)?);
343            l1_base_fee_estimate.push(self.read_l1_field(b, L1_PRICE_PER_UNIT_OFFSET)?);
344            l1_last_surplus.push(self.read_l1_field(b, L1_LAST_SURPLUS_OFFSET)?);
345            l1_funds_due.push(self.read_l1_field(b, L1_L1_FEES_AVAILABLE_OFFSET)?);
346            l1_funds_due_for_rewards.push(self.read_l1_field(b, L1_FUNDS_DUE_FOR_REWARDS_OFFSET)?);
347            l1_units_since_update.push(
348                self.read_l1_field(b, L1_UNITS_SINCE_UPDATE_OFFSET)?
349                    .try_into()
350                    .unwrap_or(0u64),
351            );
352            l1_last_update_time.push(
353                self.read_l1_field(b, L1_LAST_UPDATE_TIME_OFFSET)?
354                    .try_into()
355                    .unwrap_or(0u64),
356            );
357        }
358
359        // Scalar fields — read once at `end`.
360        let min_base_fee = self.read_l2_field(end, L2_MIN_BASE_FEE_OFFSET)?;
361        let speed_limit = self
362            .read_l2_field(end, L2_SPEED_LIMIT_OFFSET)?
363            .try_into()
364            .unwrap_or(0u64);
365        let per_block_gas_limit = self
366            .read_l2_field(end, L2_PER_BLOCK_GAS_LIMIT_OFFSET)?
367            .try_into()
368            .unwrap_or(0u64);
369        let pricing_inertia = self
370            .read_l2_field(end, L2_PRICING_INERTIA_OFFSET)?
371            .try_into()
372            .unwrap_or(0u64);
373        let backlog_tolerance = self
374            .read_l2_field(end, L2_BACKLOG_TOLERANCE_OFFSET)?
375            .try_into()
376            .unwrap_or(0u64);
377        let l1_equilibration_units = self.read_l1_field(end, L1_EQUILIBRATION_UNITS_OFFSET)?;
378        let l1_per_batch_cost: i64 = self
379            .read_l1_field(end, L1_PER_BATCH_GAS_COST_OFFSET)?
380            .try_into()
381            .unwrap_or(0i64);
382        let l1_amortized_cost_cap_bips = self
383            .read_l1_field(end, L1_AMORTIZED_COST_CAP_BIPS_OFFSET)?
384            .try_into()
385            .unwrap_or(0u64);
386        let l1_pricing_inertia = self
387            .read_l1_field(end, L1_INERTIA_OFFSET)?
388            .try_into()
389            .unwrap_or(0u64);
390        let l1_per_unit_reward = self
391            .read_l1_field(end, L1_PER_UNIT_REWARD_OFFSET)?
392            .try_into()
393            .unwrap_or(0u64);
394        let l1_pay_reward_to = {
395            let word = self.read_l1_field(end, L1_PAY_REWARDS_TO_OFFSET)?;
396            Address::from_slice(&word.to_be_bytes::<32>()[12..])
397        };
398
399        Ok(PricingModelHistory {
400            start,
401            end,
402            step,
403            timestamp,
404            base_fee,
405            gas_backlog,
406            gas_used,
407            min_base_fee,
408            speed_limit,
409            per_block_gas_limit,
410            per_tx_gas_limit: 0,
411            pricing_inertia,
412            backlog_tolerance,
413            l1_base_fee_estimate,
414            l1_last_surplus,
415            l1_funds_due,
416            l1_funds_due_for_rewards,
417            l1_units_since_update,
418            l1_last_update_time,
419            l1_equilibration_units,
420            l1_per_batch_cost,
421            l1_amortized_cost_cap_bips,
422            l1_pricing_inertia,
423            l1_per_unit_reward,
424            l1_pay_reward_to,
425        })
426    }
427
428    async fn timeout_queue_history(&self, start: u64, end: u64) -> RpcResult<TimeoutQueueHistory> {
429        self.check_enabled()?;
430        self.validate_range(start, end)?;
431        let (first, step, samples) = compute_step(start, end, self.config.block_range_bound);
432
433        let mut timestamp = Vec::with_capacity(samples as usize);
434        let mut size = Vec::with_capacity(samples as usize);
435        for i in 0..samples {
436            let b = first + step * i;
437            timestamp.push(self.header_timestamp(b)?);
438            size.push(self.queue_size_at(b)?);
439        }
440        Ok(TimeoutQueueHistory {
441            start,
442            end,
443            step,
444            timestamp,
445            size,
446        })
447    }
448
449    async fn timeout_queue(&self, block_num: u64) -> RpcResult<TimeoutQueue> {
450        self.check_enabled()?;
451        let entries =
452            self.queue_snapshot_at(block_num, self.config.timeout_queue_bound as usize)?;
453        let (tickets, timeouts): (Vec<B256>, Vec<u64>) = entries.into_iter().unzip();
454        Ok(TimeoutQueue {
455            block_number: block_num,
456            tickets,
457            timeouts,
458        })
459    }
460}
461
462#[cfg(test)]
463mod tests {
464    use super::*;
465
466    #[test]
467    fn compute_step_span_fits_bound() {
468        let (first, step, samples) = compute_step(100, 109, 256);
469        assert_eq!(first, 100);
470        assert_eq!(step, 1);
471        assert_eq!(samples, 10);
472    }
473
474    #[test]
475    fn compute_step_span_exceeds_bound() {
476        let (first, step, samples) = compute_step(0, 9999, 100);
477        assert!(samples <= 100);
478        assert!(step >= 100);
479        // Should anchor last sample at `end`.
480        assert_eq!(first + step * (samples - 1), 9999);
481    }
482
483    #[test]
484    fn compute_step_single_block() {
485        let (first, step, samples) = compute_step(42, 42, 256);
486        assert_eq!(first, 42);
487        assert_eq!(step, 1);
488        assert_eq!(samples, 1);
489    }
490
491    #[test]
492    fn compute_step_zero_bound() {
493        let (_, step, samples) = compute_step(0, 10, 0);
494        assert_eq!(step, 1);
495        assert_eq!(samples, 0);
496    }
497}