Worker listens for new transactions in the node mempool and broadcasts [`MessageMempoolDataUpdate`].
(client: P, name: String, mempool_tx: Broadcaster<MessageMempoolDataUpdate>)
| 13 | |
| 14 | /// Worker listens for new transactions in the node mempool and broadcasts [`MessageMempoolDataUpdate`]. |
| 15 | pub async fn new_node_mempool_worker<P>(client: P, name: String, mempool_tx: Broadcaster<MessageMempoolDataUpdate>) -> WorkerResult |
| 16 | where |
| 17 | P: Provider<Ethereum> + Send + Sync + 'static, |
| 18 | { |
| 19 | let mempool_subscription = client.subscribe_full_pending_transactions().await?; |
| 20 | let mut stream = mempool_subscription.into_stream(); |
| 21 | |
| 22 | while let Some(tx) = stream.next().await { |
| 23 | let tx_hash: TxHash = tx.tx_hash(); |
| 24 | let update_msg: MessageMempoolDataUpdate = MessageMempoolDataUpdate::new_with_source( |
| 25 | NodeMempoolDataUpdate { tx_hash, mempool_tx: MempoolTx { tx: Some(tx), ..MempoolTx::default() } }, |
| 26 | name.clone(), |
| 27 | ); |
| 28 | if let Err(e) = mempool_tx.send(update_msg) { |
| 29 | error!("mempool_tx.send error : {}", e); |
| 30 | break; |
| 31 | } |
| 32 | } |
| 33 | Ok(name) |
| 34 | } |
| 35 | |
| 36 | #[derive(Producer)] |
| 37 | pub struct NodeMempoolActor<P> { |