1use 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 pub block_range_bound: u64,
114 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
162fn 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 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 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 fn queue_size_at(&self, block: u64) -> Result<u64, ErrorObject<'static>> {
251 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 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 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 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 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}