mod block_processor; mod config; use block_processor::listen_blocks; use config::StreamerConfig; use near_indexer::StreamerMessage; use crate::{ errors::ChainGatewayError, event_subscriber::consts::DEFAULT_NUMBER_OF_TRACKED_BLOCKS, primitives::FetchLatestFinalBlockInfo, }; use super::{ block_events::BlockUpdate, stats::{IndexerStats, indexer_logger}, subscriber::BlockEventSubscriptions, }; pub(crate) async fn start( block_event_subscriber: BlockEventSubscriptions, stream: tokio::sync::mpsc::Receiver, info_fetcher: impl FetchLatestFinalBlockInfo, ) -> Result, ChainGatewayError> { let StreamerConfig { buffer_size, block_events, } = block_event_subscriber.into(); let number_of_tracked_blocks = DEFAULT_NUMBER_OF_TRACKED_BLOCKS.max( buffer_size .try_into() .expect("usize is expected to fit into u64"), ); let (stats_tx, stats_rx) = tokio::sync::watch::channel(IndexerStats::new()); let (block_tx, block_rx) = tokio::sync::mpsc::channel(buffer_size); tokio::spawn(async move { if let Err(err) = listen_blocks( stream, block_events, stats_tx, block_tx, number_of_tracked_blocks, ) .await { tracing::error!(target: "chain gateway", "block event listener stopped: {err}"); } }); tokio::spawn(indexer_logger(stats_rx, info_fetcher)); Ok(block_rx) }