use std::sync::Arc;
use async_trait::async_trait;
use node::database::{self};
use node::{LongLivedService, Network};
use node_data::events::Event as ChainEvent;
use tokio::sync::broadcast;
use tokio::sync::mpsc::Receiver;
use tracing::error;
use crate::http::RuesEvent;
pub(crate) struct ChainEventStreamer {
pub node_receiver: Receiver<ChainEvent>,
pub rues_sender: broadcast::Sender<RuesEvent>,
}
#[async_trait]
impl<N: Network, DB: database::DB, VM: node::vm::VMExecution>
LongLivedService<N, DB, VM> for ChainEventStreamer
{
async fn execute(
&mut self,
_: Arc<tokio::sync::RwLock<N>>,
_: Arc<tokio::sync::RwLock<DB>>,
_: Arc<tokio::sync::RwLock<VM>>,
) -> anyhow::Result<usize> {
loop {
if let Some(msg) = self.node_receiver.recv().await {
if let Err(e) = self.rues_sender.send(msg.clone().into()) {
error!("Cannot send to rues {e:?}");
}
}
}
}
fn name(&self) -> &'static str {
"chain event streamer"
}
}