Skip to main content

ghost_crab/
block_handler.rs

1use 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}