Skip to main content

perpl_sdk/stream/
raw.rs

1use std::time::Duration;
2
3use alloy::{eips::BlockId, providers::Provider, rpc::types::Filter, sol_types::SolEventInterface};
4use futures::{Stream, stream};
5
6use crate::{
7    Chain,
8    abi::dex::Exchange::ExchangeEvents,
9    error::{DexError, ProviderError},
10    types,
11};
12
13pub type RawEvent = types::EventContext<ExchangeEvents>;
14pub type RawBlockEvents = types::BlockEvents<RawEvent>;
15
16/// Returns stream of raw events emitted by the DEX smart contract,
17/// batched per block, starting from the specified block.
18///
19/// Polls logs via the given [`Provider`] to produce strictly continuous
20/// event sequence, with [`Provider`]-configured interval.
21///
22/// It is recommended to setup provider with
23/// [`alloy::transports::layers::FallbackLayer`]
24/// and/or [`alloy::transports::layers::RetryBackoffLayer`].
25///
26/// See [`crate::abi::dex::Exchange::ExchangeEvents`] for the list of possible
27/// events and corresponding details.
28///
29/// # Safety note
30///
31/// The returned stream is not cancellation-safe and should not be used within
32/// `select!`.
33pub fn raw<P, S, SFut>(
34    chain: &Chain,
35    provider: P,
36    from: types::StateInstant,
37    sleep: S,
38) -> impl Stream<Item = Result<RawBlockEvents, DexError>>
39where
40    P: Provider,
41    S: Fn(Duration) -> SFut + Copy,
42    SFut: Future<Output = ()>,
43{
44    stream::unfold((provider, from.block_number()), move |(provider, mut block_num)| async move {
45        let filter = Filter::new()
46            .address(chain.exchange())
47            .from_block(block_num)
48            .to_block(block_num);
49        loop {
50            // Anvil node, and maybe some RPC providers, produce empty response instead of
51            // error in case the block in the filter does not exist yet.
52            // Checking the block presence explicitly, and also checking requested block
53            // number against `safe` block tag as Monad RPC assumes `latest` ==
54            // `Proposed` since v0.13.0 and `Proposed` is not safe enough to
55            // preserve state consistency.
56            let result = futures::try_join!(
57                provider.get_block(BlockId::safe()).into_future(),
58                provider.get_block(BlockId::number(block_num)).into_future(),
59                provider.get_logs(&filter)
60            )
61            .map_err(ProviderError::from)
62            .and_then(|(safe_block, block, logs)| {
63                if safe_block.is_none_or(|sb| sb.header.number < block_num) {
64                    return Err(ProviderError::InvalidRequest(
65                        "block is not available yet".to_string(),
66                    ));
67                }
68                let block_header = block
69                    .ok_or(ProviderError::InvalidRequest("block is not available yet".to_string()))?
70                    .header;
71                let mut events = Vec::with_capacity(logs.len());
72                for log in &logs {
73                    events.push(RawEvent::new(
74                        log.transaction_hash.unwrap_or_default(),
75                        log.transaction_index.unwrap_or_default(),
76                        log.log_index.unwrap_or_default(),
77                        ExchangeEvents::decode_log(&log.inner)
78                            .map_err(ProviderError::from)?
79                            .data,
80                    ));
81                }
82                // Monad RPC does not guarantee logs are returned in block-internal order
83                // (eg block 68747089 from https://rpc-mainnet.monadinfra.com)
84                events.sort_by_key(|e| e.log_index());
85                Ok(RawBlockEvents::new(
86                    types::StateInstant::new(block_num, block_header.timestamp),
87                    events,
88                ))
89            });
90            if result.is_ok() {
91                block_num += 1;
92                return Some((result.map_err(DexError::Provider), (provider, block_num)));
93            }
94            if matches!(result, Err(ProviderError::InvalidRequest(_))) {
95                // Block is not available yet
96                sleep(provider.client().poll_interval()).await;
97                continue;
98            }
99            return Some((result.map_err(DexError::Provider), (provider, block_num)));
100        }
101    })
102}
103
104#[cfg(test)]
105mod tests {
106    use alloy::{
107        providers::ProviderBuilder, rpc::client::RpcClient, transports::layers::RetryBackoffLayer,
108    };
109    use futures::StreamExt;
110
111    use super::*;
112    use crate::Chain;
113
114    #[tokio::test]
115    async fn test_stream_recent_blocks() {
116        let client = RpcClient::builder()
117            .layer(RetryBackoffLayer::new(10, 100, 200))
118            .connect("https://testnet-rpc.monad.xyz")
119            .await
120            .unwrap();
121        client.set_poll_interval(Duration::from_millis(100));
122        let provider = ProviderBuilder::new().connect_client(client);
123
124        let testnet = Chain::testnet();
125        let mut block_num = provider.get_block_number().await.unwrap() + 1;
126        let stream =
127            raw(&testnet, provider, types::StateInstant::new(block_num, 0), tokio::time::sleep);
128        let block_results = stream.take(10).collect::<Vec<_>>().await;
129
130        for b in &block_results {
131            assert_eq!(b.as_ref().unwrap().instant().block_number(), block_num);
132            block_num += 1;
133        }
134    }
135}