loom_node_json_rpc/
node_mempool_actor.rs

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
use std::marker::PhantomData;

use alloy_network::{Ethereum, TransactionResponse};
use alloy_primitives::TxHash;
use alloy_provider::Provider;
use alloy_transport::Transport;
use futures::StreamExt;
use revm::DatabaseRef;
use tracing::error;

use loom_core_actors::{Actor, ActorResult, Broadcaster, Producer, WorkerResult};
use loom_core_actors_macros::*;
use loom_core_blockchain::Blockchain;
use loom_types_blockchain::MempoolTx;
use loom_types_events::{MessageMempoolDataUpdate, NodeMempoolDataUpdate};

/// Worker listens for new transactions in the node mempool and broadcasts [`MessageMempoolDataUpdate`].
pub async fn new_node_mempool_worker<P, T>(client: P, name: String, mempool_tx: Broadcaster<MessageMempoolDataUpdate>) -> WorkerResult
where
    T: Transport + Clone,
    P: Provider<T, Ethereum> + Send + Sync + 'static,
{
    let mempool_subscription = client.subscribe_full_pending_transactions().await?;
    let mut stream = mempool_subscription.into_stream();

    while let Some(tx) = stream.next().await {
        let tx_hash: TxHash = tx.tx_hash();
        let update_msg: MessageMempoolDataUpdate = MessageMempoolDataUpdate::new_with_source(
            NodeMempoolDataUpdate { tx_hash, mempool_tx: MempoolTx { tx: Some(tx), ..MempoolTx::default() } },
            name.clone(),
        );
        if let Err(e) = mempool_tx.send(update_msg).await {
            error!("mempool_tx.send error : {}", e);
            break;
        }
    }
    Ok(name)
}

#[derive(Producer)]
pub struct NodeMempoolActor<P, T> {
    name: &'static str,
    client: P,
    #[producer]
    mempool_tx: Option<Broadcaster<MessageMempoolDataUpdate>>,
    _t: PhantomData<T>,
}

impl<P, T> NodeMempoolActor<P, T>
where
    T: Transport + Clone,
    P: Provider<T, Ethereum> + Send + Sync + Clone + 'static,
{
    pub fn new(client: P) -> NodeMempoolActor<P, T> {
        NodeMempoolActor { client, name: "NodeMempoolActor", mempool_tx: None, _t: PhantomData }
    }

    pub fn with_name(self, name: String) -> Self {
        Self { name: Box::leak(name.into_boxed_str()), ..self }
    }

    fn get_name(&self) -> &'static str {
        self.name
    }

    pub fn on_bc<DB: DatabaseRef + Send + Sync + Clone + 'static>(self, bc: &Blockchain<DB>) -> Self {
        Self { mempool_tx: Some(bc.new_mempool_tx_channel()), ..self }
    }
}

impl<P, T> Actor for NodeMempoolActor<P, T>
where
    T: Transport + Clone,
    P: Provider<T, Ethereum> + Send + Sync + Clone + 'static,
{
    fn start(&self) -> ActorResult {
        let task =
            tokio::task::spawn(new_node_mempool_worker(self.client.clone(), self.name.to_string(), self.mempool_tx.clone().unwrap()));
        Ok(vec![task])
    }

    fn name(&self) -> &'static str {
        self.get_name()
    }
}