ghost_crab/
block_handler.rs1use crate::indexer::rpc_manager::Provider;
2use crate::indexer::templates::TemplateManager;
3use crate::latest_block_manager::LatestBlockManager;
4use alloy::providers::Provider as AlloyProvider;
5use alloy::rpc::types::eth::Block;
6use alloy::rpc::types::eth::BlockNumberOrTag;
7use alloy::transports::TransportError;
8use async_trait::async_trait;
9use ghost_crab_common::config::BlockHandlerConfig;
10use ghost_crab_common::config::ExecutionMode;
11use std::sync::Arc;
12use std::time::Duration;
13
14pub struct BlockContext {
15 pub provider: Provider,
16 pub templates: TemplateManager,
17 pub block_number: u64,
18}
19
20impl BlockContext {
21 pub async fn block(&self, hydrate: bool) -> Result<Option<Block>, TransportError> {
22 self.provider
23 .get_block_by_number(BlockNumberOrTag::Number(self.block_number), hydrate)
24 .await
25 }
26}
27
28pub type BlockHandlerInstance = Arc<Box<(dyn BlockHandler + Send + Sync)>>;
29
30#[async_trait]
31pub trait BlockHandler {
32 async fn handle(&self, params: BlockContext);
33 fn name(&self) -> String;
34}
35
36#[derive(Clone)]
37pub struct ProcessBlocksInput {
38 pub handler: BlockHandlerInstance,
39 pub templates: TemplateManager,
40 pub provider: Provider,
41 pub config: BlockHandlerConfig,
42}
43
44pub async fn process_blocks(
45 ProcessBlocksInput { handler, templates, provider, config }: ProcessBlocksInput,
46) -> Result<(), TransportError> {
47 let execution_mode = config.execution_mode.unwrap_or(ExecutionMode::Parallel);
48
49 let mut current_block = config.start_block;
50 let mut latest_block_manager =
51 LatestBlockManager::new(provider.clone(), Duration::from_secs(10));
52
53 loop {
54 let latest_block = latest_block_manager.get().await?;
55
56 if current_block >= latest_block {
57 tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
58 continue;
59 }
60
61 match execution_mode {
62 ExecutionMode::Parallel => {
63 let handler = handler.clone();
64 let provider = provider.clone();
65 let templates = templates.clone();
66
67 tokio::spawn(async move {
68 handler
69 .handle(BlockContext { provider, templates, block_number: current_block })
70 .await;
71 });
72 }
73 ExecutionMode::Serial => {
74 let templates = templates.clone();
75 let provider = provider.clone();
76 let templates = templates.clone();
77
78 handler
79 .handle(BlockContext { provider, templates, block_number: current_block })
80 .await;
81 }
82 }
83
84 current_block += config.step;
85 }
86}