use std::time::Duration;
use alloy::{eips::BlockId, providers::Provider, rpc::types::Filter, sol_types::SolEventInterface};
use futures::{Stream, stream};
use crate::{
Chain,
abi::dex::Exchange::ExchangeEvents,
error::{DexError, ProviderError},
types,
};
pub type RawEvent = types::EventContext<ExchangeEvents>;
pub type RawBlockEvents = types::BlockEvents<RawEvent>;
pub fn raw<P, S, SFut>(
chain: &Chain,
provider: P,
from: types::StateInstant,
sleep: S,
) -> impl Stream<Item = Result<RawBlockEvents, DexError>>
where
P: Provider,
S: Fn(Duration) -> SFut + Copy,
SFut: Future<Output = ()>,
{
stream::unfold((provider, from.block_number()), move |(provider, mut block_num)| async move {
let filter = Filter::new()
.address(chain.exchange())
.from_block(block_num)
.to_block(block_num);
loop {
let result = futures::try_join!(
provider.get_block(BlockId::safe()).into_future(),
provider.get_block(BlockId::number(block_num)).into_future(),
provider.get_logs(&filter)
)
.map_err(ProviderError::from)
.and_then(|(safe_block, block, logs)| {
if safe_block.is_none_or(|sb| sb.header.number < block_num) {
return Err(ProviderError::InvalidRequest(
"block is not available yet".to_string(),
));
}
let block_header = block
.ok_or(ProviderError::InvalidRequest("block is not available yet".to_string()))?
.header;
let mut events = Vec::with_capacity(logs.len());
for log in &logs {
events.push(RawEvent::new(
log.transaction_hash.unwrap_or_default(),
log.transaction_index.unwrap_or_default(),
log.log_index.unwrap_or_default(),
ExchangeEvents::decode_log(&log.inner)
.map_err(ProviderError::from)?
.data,
));
}
events.sort_by_key(|e| e.log_index());
Ok(RawBlockEvents::new(
types::StateInstant::new(block_num, block_header.timestamp),
events,
))
});
if result.is_ok() {
block_num += 1;
return Some((result.map_err(DexError::Provider), (provider, block_num)));
}
if matches!(result, Err(ProviderError::InvalidRequest(_))) {
sleep(provider.client().poll_interval()).await;
continue;
}
return Some((result.map_err(DexError::Provider), (provider, block_num)));
}
})
}
#[cfg(test)]
mod tests {
use alloy::{
providers::ProviderBuilder, rpc::client::RpcClient, transports::layers::RetryBackoffLayer,
};
use futures::StreamExt;
use super::*;
use crate::Chain;
#[tokio::test]
async fn test_stream_recent_blocks() {
let client = RpcClient::builder()
.layer(RetryBackoffLayer::new(10, 100, 200))
.connect("https://testnet-rpc.monad.xyz")
.await
.unwrap();
client.set_poll_interval(Duration::from_millis(100));
let provider = ProviderBuilder::new().connect_client(client);
let testnet = Chain::testnet();
let mut block_num = provider.get_block_number().await.unwrap() + 1;
let stream =
raw(&testnet, provider, types::StateInstant::new(block_num, 0), tokio::time::sleep);
let block_results = stream.take(10).collect::<Vec<_>>().await;
for b in &block_results {
assert_eq!(b.as_ref().unwrap().instant().block_number(), block_num);
block_num += 1;
}
}
}