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
16pub 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 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 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 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}