1use crate::{
4 NodeResult,
5 eth::{
6 backend::validate::TransactionValidator, error::BlockchainError, fees::FeeHistoryService,
7 miner::Miner, pool::Pool,
8 },
9 filter::Filters,
10 mem::Backend,
11};
12use alloy_consensus::TxReceipt;
13use alloy_network::Network;
14use foundry_primitives::{FoundryReceiptEnvelope, FoundryTxEnvelope};
15use futures::{FutureExt, Stream, StreamExt};
16use std::{
17 collections::VecDeque,
18 pin::Pin,
19 sync::Arc,
20 task::{Context, Poll},
21};
22use tokio::{task::JoinHandle, time::Interval};
23
24pub struct NodeService<N: Network>
31where
32 N::ReceiptEnvelope: TxReceipt<Log = alloy_primitives::Log>,
33{
34 pool: Arc<Pool<N::TxEnvelope>>,
36 block_producer: BlockProducer<N>,
38 miner: Miner<N::TxEnvelope>,
40 fee_history: FeeHistoryService<N>,
42 filters: Filters<N>,
44 filter_eviction_interval: Interval,
46}
47
48impl<N: Network> NodeService<N>
49where
50 Backend<N>: TransactionValidator<N::TxEnvelope>,
51 N: Network<TxEnvelope = FoundryTxEnvelope, ReceiptEnvelope = FoundryReceiptEnvelope>,
52{
53 pub fn new(
54 pool: Arc<Pool<N::TxEnvelope>>,
55 backend: Arc<Backend<N>>,
56 miner: Miner<N::TxEnvelope>,
57 fee_history: FeeHistoryService<N>,
58 filters: Filters<N>,
59 ) -> Self {
60 let start = tokio::time::Instant::now() + filters.keep_alive();
61 let filter_eviction_interval = tokio::time::interval_at(start, filters.keep_alive());
62 Self {
63 block_producer: BlockProducer::new(backend, pool.clone(), miner.clone()),
64 pool,
65 miner,
66 fee_history,
67 filter_eviction_interval,
68 filters,
69 }
70 }
71}
72
73impl<N: Network> Future for NodeService<N>
74where
75 Backend<N>: TransactionValidator<N::TxEnvelope>,
76 N: Network<TxEnvelope = FoundryTxEnvelope, ReceiptEnvelope = FoundryReceiptEnvelope>,
77{
78 type Output = NodeResult<()>;
79
80 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
81 let pin = self.get_mut();
82
83 loop {
86 while let Poll::Ready(Some(result)) = pin.block_producer.poll_next_unpin(cx) {
88 match result {
89 BlockProduction::Mined(block_number) => {
90 trace!(target: "node", "mined block {block_number}");
91 }
92 BlockProduction::Failed(generation) => {
93 pin.miner.handle_failed_candidate(generation);
94 break;
95 }
96 BlockProduction::Skipped => {}
97 }
98 }
99
100 if pin.block_producer.is_idle()
103 && let Poll::Ready(work) = pin.miner.poll(&pin.pool, cx)
104 {
105 pin.block_producer.queued.push_back(work);
107 } else {
108 break;
110 }
111 }
112
113 let _ = pin.fee_history.poll_unpin(cx);
115
116 if pin.filter_eviction_interval.poll_tick(cx).is_ready() {
117 let filters = pin.filters.clone();
118
119 tokio::task::spawn(async move { filters.evict().await });
121 }
122
123 Poll::Pending
124 }
125}
126
127type MiningResult<N> = (Result<Option<u64>, BlockchainError>, Arc<Backend<N>>, u64);
128
129enum BlockProduction {
130 Mined(u64),
131 Failed(u64),
132 Skipped,
133}
134
135#[must_use = "streams do nothing unless polled"]
137struct BlockProducer<N: Network> {
138 idle_backend: Option<Arc<Backend<N>>>,
140 pool: Arc<Pool<N::TxEnvelope>>,
142 miner: Miner<N::TxEnvelope>,
144 block_mining: Option<JoinHandle<MiningResult<N>>>,
146 queued: VecDeque<crate::eth::miner::MiningWork<N::TxEnvelope>>,
148}
149
150impl<N: Network> BlockProducer<N>
151where
152 Backend<N>: TransactionValidator<N::TxEnvelope>,
153 N: Network<TxEnvelope = FoundryTxEnvelope, ReceiptEnvelope = FoundryReceiptEnvelope>,
154{
155 fn new(
156 backend: Arc<Backend<N>>,
157 pool: Arc<Pool<N::TxEnvelope>>,
158 miner: Miner<N::TxEnvelope>,
159 ) -> Self {
160 Self {
161 idle_backend: Some(backend),
162 pool,
163 miner,
164 block_mining: None,
165 queued: Default::default(),
166 }
167 }
168
169 fn is_idle(&self) -> bool {
170 self.idle_backend.is_some() && self.block_mining.is_none() && self.queued.is_empty()
171 }
172}
173
174impl<N: Network> Stream for BlockProducer<N>
175where
176 Backend<N>: TransactionValidator<N::TxEnvelope> + Send + Sync + 'static,
177 N: Network<TxEnvelope = FoundryTxEnvelope, ReceiptEnvelope = FoundryReceiptEnvelope> + 'static,
178{
179 type Item = BlockProduction;
180
181 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
182 let pin = self.get_mut();
183
184 if !pin.queued.is_empty() {
185 if let Some(backend) = pin.idle_backend.take() {
187 let work = pin.queued.pop_front().expect("not empty; qed");
188 let generation = work.generation;
189 let selected = work.transactions.len();
190 let pool = pin.pool.clone();
191 let miner = pin.miner.clone();
192
193 let handle = tokio::runtime::Handle::current();
196 let mining = tokio::task::spawn_blocking(move || {
197 handle.block_on(async move {
198 let mining_guard = backend.lock_mining_owned().await;
199 let transactions =
200 pool.ready_transactions().take(selected).collect::<Vec<_>>();
201 if selected > 0 && transactions.is_empty() {
202 return (Ok(None), backend, generation);
203 }
204 trace!(target: "miner", "creating new block");
205 let result = backend.mine_block_locked(transactions).await.map(|outcome| {
206 let block_number = outcome.block_number;
207 if pool.on_mined_block(outcome) {
208 miner.retry_ready_transactions();
209 }
210 trace!(target: "miner", "created new block: {block_number}");
211 Some(block_number)
212 });
213 drop(mining_guard);
214 (result, backend, generation)
215 })
216 });
217 pin.block_mining = Some(mining);
218 }
219 }
220
221 if let Some(mut mining) = pin.block_mining.take() {
222 if let Poll::Ready(res) = mining.poll_unpin(cx) {
223 return match res {
224 Ok((Ok(Some(block_number)), backend, _)) => {
225 pin.idle_backend = Some(backend);
226 Poll::Ready(Some(BlockProduction::Mined(block_number)))
227 }
228 Ok((Ok(None), backend, _)) => {
229 pin.idle_backend = Some(backend);
230 Poll::Ready(Some(BlockProduction::Skipped))
231 }
232 Ok((Err(error), backend, generation)) => {
233 pin.idle_backend = Some(backend);
234 pin.queued.clear();
235 warn!(target: "miner", %error, "failed to finalize block");
236 Poll::Ready(Some(BlockProduction::Failed(generation)))
237 }
238 Err(err) => {
239 panic!("miner task failed: {err}");
240 }
241 };
242 }
243 pin.block_mining = Some(mining)
244 }
245
246 Poll::Pending
247 }
248}