1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
use crate::handler::{Context, HandlerConfig};
use crate::latest_block_manager::LatestBlockManager;
use alloy::primitives::Address;
use alloy::providers::Provider;
use alloy::rpc::types::eth::Filter;
use serde::Deserialize;

#[derive(Clone, Copy, Deserialize, Debug)]
#[serde(rename_all = "lowercase")]
pub enum ExecutionMode {
    Parallel,
    Serial,
}

pub async fn process_logs(
    HandlerConfig {
        start_block,
        step,
        address,
        handler,
        provider,
        templates,
        execution_mode,
    }: HandlerConfig,
) {
    let mut current_block = start_block;
    let event_signature = handler.get_event_signature();
    let address = address.parse::<Address>().unwrap();

    let mut block_manager = LatestBlockManager::new(1000, provider.clone());

    loop {
        let mut end_block = current_block + step;
        let latest_block = block_manager.get().await;

        if end_block > latest_block {
            end_block = latest_block;
        }

        if current_block >= end_block {
            tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
            continue;
        }

        let source = handler.get_source();

        println!(
            "[{}] Processing logs from {} to {}",
            source, current_block, end_block
        );

        let filter = Filter::new()
            .address(address.clone())
            .event(&event_signature)
            .from_block(current_block)
            .to_block(end_block);

        let logs = provider.get_logs(&filter).await.unwrap();

        match execution_mode {
            ExecutionMode::Parallel => {
                for log in logs {
                    let handler = handler.clone();
                    let provider = provider.clone();
                    let templates = templates.clone();

                    tokio::spawn(async move {
                        handler
                            .handle(Context {
                                log,
                                provider,
                                templates,
                                contract_address: address
                            })
                            .await;
                    });
                }
            }
            ExecutionMode::Serial => {
                for log in logs {
                    let templates = templates.clone();
                    let provider = provider.clone();
                    let templates = templates.clone();

                    handler
                        .handle(Context {
                            log,
                            provider,
                            templates,
                            contract_address: address
                        })
                        .await;
                }
            }
        }

        current_block = end_block;
    }
}