Skip to main content

reth_transaction_pool/validate/
task.rs

1//! A validation service for transactions.
2
3use crate::{
4    blobstore::BlobStore,
5    metrics::TxPoolValidatorMetrics,
6    validate::{EthTransactionValidatorBuilder, TransactionValidatorError},
7    EthTransactionValidator, PoolTransaction, TransactionOrigin, TransactionValidationOutcome,
8    TransactionValidator,
9};
10use futures_util::{lock::Mutex, StreamExt};
11use reth_chainspec::{ChainSpecProvider, EthereumHardforks};
12use reth_evm::ConfigureEvm;
13use reth_primitives_traits::{HeaderTy, SealedBlock};
14use reth_storage_api::BlockReaderIdExt;
15use reth_tasks::Runtime;
16use std::{future::Future, pin::Pin, sync::Arc};
17use tokio::sync::{mpsc, oneshot};
18use tokio_stream::wrappers::ReceiverStream;
19
20/// Represents a future outputting unit type and is sendable.
21type ValidationFuture = Pin<Box<dyn Future<Output = ()> + Send>>;
22
23/// Represents a stream of validation futures.
24type ValidationStream = ReceiverStream<ValidationFuture>;
25
26/// A service that performs validation jobs.
27///
28/// This listens for incoming validation jobs and executes them.
29///
30/// This should be spawned as a task: [`ValidationTask::run`]
31#[derive(Clone)]
32pub struct ValidationTask {
33    validation_jobs: Arc<Mutex<ValidationStream>>,
34}
35
36impl ValidationTask {
37    /// Creates a new cloneable task pair.
38    ///
39    /// The sender sends new (transaction) validation tasks to an available validation task.
40    pub fn new() -> (ValidationJobSender, Self) {
41        Self::with_capacity(1)
42    }
43
44    /// Creates a new cloneable task pair with the given channel capacity.
45    pub fn with_capacity(capacity: usize) -> (ValidationJobSender, Self) {
46        let (tx, rx) = mpsc::channel(capacity);
47        let metrics = TxPoolValidatorMetrics::default();
48        (ValidationJobSender { tx, metrics }, Self::with_receiver(rx))
49    }
50
51    /// Creates a new task with the given receiver.
52    pub fn with_receiver(jobs: mpsc::Receiver<Pin<Box<dyn Future<Output = ()> + Send>>>) -> Self {
53        Self { validation_jobs: Arc::new(Mutex::new(ReceiverStream::new(jobs))) }
54    }
55
56    /// Executes all new validation jobs that come in.
57    ///
58    /// This will run as long as the channel is alive and is expected to be spawned as a task.
59    pub async fn run(self) {
60        loop {
61            // Release the shared receiver before running the job, so other workers
62            // can dequeue validations while this worker is busy.
63            let task = { self.validation_jobs.lock().await.next().await };
64            let Some(task) = task else { break };
65            task.await;
66        }
67    }
68}
69
70impl std::fmt::Debug for ValidationTask {
71    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
72        f.debug_struct("ValidationTask").field("validation_jobs", &"...").finish()
73    }
74}
75
76/// A sender new type for sending validation jobs to [`ValidationTask`].
77#[derive(Debug)]
78pub struct ValidationJobSender {
79    tx: mpsc::Sender<Pin<Box<dyn Future<Output = ()> + Send>>>,
80    metrics: TxPoolValidatorMetrics,
81}
82
83impl ValidationJobSender {
84    /// Sends the given job to the validation task.
85    pub async fn send(
86        &self,
87        job: Pin<Box<dyn Future<Output = ()> + Send>>,
88    ) -> Result<(), TransactionValidatorError> {
89        self.metrics.inflight_validation_jobs.increment(1);
90        let _guard = DecrementPendingOnDrop(&self.metrics.inflight_validation_jobs);
91        self.tx.send(job).await.map_err(|_| TransactionValidatorError::ValidationServiceUnreachable)
92    }
93}
94
95/// A [`TransactionValidator`] implementation that validates ethereum transaction.
96/// This validator is non-blocking, all validation work is done in a separate task.
97#[derive(Debug)]
98pub struct TransactionValidationTaskExecutor<V> {
99    /// The validator that will validate transactions on a separate task.
100    pub validator: Arc<V>,
101    /// The bounded multi-producer queue supplies backpressure without serializing
102    /// producers behind an additional lock.
103    pub to_validation_task: Arc<ValidationJobSender>,
104}
105
106impl<V> Clone for TransactionValidationTaskExecutor<V> {
107    fn clone(&self) -> Self {
108        Self {
109            validator: self.validator.clone(),
110            to_validation_task: self.to_validation_task.clone(),
111        }
112    }
113}
114
115// === impl TransactionValidationTaskExecutor ===
116
117impl TransactionValidationTaskExecutor<()> {
118    /// Convenience method to create a [`EthTransactionValidatorBuilder`]
119    pub fn eth_builder<Client, Evm>(
120        client: Client,
121        evm_config: Evm,
122    ) -> EthTransactionValidatorBuilder<Client, Evm>
123    where
124        Client: ChainSpecProvider<ChainSpec: EthereumHardforks>
125            + BlockReaderIdExt<Header = HeaderTy<Evm::Primitives>>,
126        Evm: ConfigureEvm,
127    {
128        EthTransactionValidatorBuilder::new(client, evm_config)
129    }
130}
131
132impl<V> TransactionValidationTaskExecutor<V> {
133    /// Maps the given validator to a new type.
134    pub fn map<F, T>(self, mut f: F) -> TransactionValidationTaskExecutor<T>
135    where
136        F: FnMut(V) -> T,
137    {
138        TransactionValidationTaskExecutor {
139            validator: Arc::new(f(Arc::into_inner(self.validator).unwrap())),
140            to_validation_task: self.to_validation_task,
141        }
142    }
143
144    /// Returns the validator.
145    pub fn validator(&self) -> &V {
146        &self.validator
147    }
148}
149
150impl<Client, Tx, Evm> TransactionValidationTaskExecutor<EthTransactionValidator<Client, Tx, Evm>> {
151    /// Creates a new instance for the given client
152    ///
153    /// This will spawn a single validation tasks that performs the actual validation.
154    /// See [`TransactionValidationTaskExecutor::eth_with_additional_tasks`]
155    pub fn eth<S: BlobStore>(client: Client, evm_config: Evm, blob_store: S, tasks: Runtime) -> Self
156    where
157        Client: ChainSpecProvider<ChainSpec: EthereumHardforks>
158            + BlockReaderIdExt<Header = HeaderTy<Evm::Primitives>>,
159        Evm: ConfigureEvm,
160    {
161        Self::eth_with_additional_tasks(client, evm_config, blob_store, tasks, 0)
162    }
163
164    /// Creates a new instance for the given client
165    ///
166    /// By default this will enable support for:
167    ///   - shanghai
168    ///   - eip1559
169    ///   - eip2930
170    ///
171    /// This will always spawn a validation task that performs the actual validation. It will spawn
172    /// `num_additional_tasks` additional tasks.
173    pub fn eth_with_additional_tasks<S: BlobStore>(
174        client: Client,
175        evm_config: Evm,
176        blob_store: S,
177        tasks: Runtime,
178        num_additional_tasks: usize,
179    ) -> Self
180    where
181        Client: ChainSpecProvider<ChainSpec: EthereumHardforks>
182            + BlockReaderIdExt<Header = HeaderTy<Evm::Primitives>>,
183        Evm: ConfigureEvm,
184    {
185        EthTransactionValidatorBuilder::new(client, evm_config)
186            .with_additional_tasks(num_additional_tasks)
187            .build_with_tasks(tasks, blob_store)
188    }
189}
190
191impl<V> TransactionValidationTaskExecutor<V> {
192    /// Creates a new executor instance with the given validator for transaction validation.
193    ///
194    /// Initializes the executor with the provided validator and sets up communication for
195    /// validation tasks.
196    pub fn new(validator: V) -> (Self, ValidationTask) {
197        let (tx, task) = ValidationTask::new();
198        (Self { validator: Arc::new(validator), to_validation_task: Arc::new(tx) }, task)
199    }
200
201    /// Creates a new executor and spawns the validation tasks on the given runtime.
202    ///
203    /// This spawns `additional_tasks` extra blocking tasks plus one critical blocking task
204    /// for the validation service.
205    pub fn spawn(validator: V, tasks: &Runtime, additional_tasks: usize) -> Self {
206        // Buffer two jobs per worker so workers can dequeue without waiting for
207        // a producer to refill the channel after every handoff.
208        let (tx, task) =
209            ValidationTask::with_capacity(additional_tasks.saturating_add(1).saturating_mul(2));
210
211        for _ in 0..additional_tasks {
212            let task = task.clone();
213            tasks.spawn_blocking_task(async move {
214                task.run().await;
215            });
216        }
217
218        tasks.spawn_critical_blocking_task("transaction-validation-service", async move {
219            task.run().await;
220        });
221
222        Self { validator: Arc::new(validator), to_validation_task: Arc::new(tx) }
223    }
224}
225
226impl<V> TransactionValidator for TransactionValidationTaskExecutor<V>
227where
228    V: TransactionValidator + 'static,
229{
230    type Transaction = <V as TransactionValidator>::Transaction;
231    type Block = V::Block;
232
233    async fn validate_transaction(
234        &self,
235        origin: TransactionOrigin,
236        transaction: Self::Transaction,
237    ) -> TransactionValidationOutcome<Self::Transaction> {
238        let hash = *transaction.hash();
239        let (tx, rx) = oneshot::channel();
240        {
241            let res = {
242                let validator = self.validator.clone();
243                let fut = Box::pin(async move {
244                    let res = validator.validate_transaction(origin, transaction).await;
245                    let _ = tx.send(res);
246                });
247                self.to_validation_task.send(fut).await
248            };
249            if res.is_err() {
250                return TransactionValidationOutcome::Error(
251                    hash,
252                    Box::new(TransactionValidatorError::ValidationServiceUnreachable),
253                );
254            }
255        }
256
257        match rx.await {
258            Ok(res) => res,
259            Err(_) => TransactionValidationOutcome::Error(
260                hash,
261                Box::new(TransactionValidatorError::ValidationServiceUnreachable),
262            ),
263        }
264    }
265
266    async fn validate_transactions(
267        &self,
268        transactions: impl IntoIterator<Item = (TransactionOrigin, Self::Transaction), IntoIter: Send>
269            + Send,
270    ) -> Vec<TransactionValidationOutcome<Self::Transaction>> {
271        let transactions: Vec<_> = transactions.into_iter().collect();
272        let hashes: Vec<_> = transactions.iter().map(|(_, tx)| *tx.hash()).collect();
273        let (tx, rx) = oneshot::channel();
274        {
275            let res = {
276                let validator = self.validator.clone();
277                let fut = Box::pin(async move {
278                    let res = validator.validate_transactions(transactions).await;
279                    let _ = tx.send(res);
280                });
281                self.to_validation_task.send(fut).await
282            };
283            if res.is_err() {
284                return validation_service_error_outcomes(hashes)
285            }
286        }
287        match rx.await {
288            Ok(res) => res,
289            Err(_) => validation_service_error_outcomes(hashes),
290        }
291    }
292
293    async fn validate_transactions_with_origin(
294        &self,
295        origin: TransactionOrigin,
296        transactions: impl IntoIterator<Item = Self::Transaction, IntoIter: Send> + Send,
297    ) -> Vec<TransactionValidationOutcome<Self::Transaction>> {
298        let transactions: Vec<_> = transactions.into_iter().collect();
299        let hashes: Vec<_> = transactions.iter().map(|tx| *tx.hash()).collect();
300        let (tx, rx) = oneshot::channel();
301        let validator = self.validator.clone();
302        let fut = Box::pin(async move {
303            let res = validator.validate_transactions_with_origin(origin, transactions).await;
304            let _ = tx.send(res);
305        });
306
307        if self.to_validation_task.send(fut).await.is_err() {
308            return validation_service_error_outcomes(hashes)
309        }
310
311        match rx.await {
312            Ok(res) => res,
313            Err(_) => validation_service_error_outcomes(hashes),
314        }
315    }
316
317    fn on_new_head_block(&self, new_tip_block: &SealedBlock<Self::Block>) {
318        self.validator.on_new_head_block(new_tip_block)
319    }
320}
321
322/// Decrements the pending-send count even if the send future is cancelled.
323struct DecrementPendingOnDrop<'a>(&'a metrics::Gauge);
324
325impl Drop for DecrementPendingOnDrop<'_> {
326    fn drop(&mut self) {
327        self.0.decrement(1);
328    }
329}
330
331#[inline]
332fn validation_service_error_outcomes<T: PoolTransaction>(
333    hashes: Vec<alloy_primitives::TxHash>,
334) -> Vec<TransactionValidationOutcome<T>> {
335    hashes
336        .into_iter()
337        .map(|hash| {
338            TransactionValidationOutcome::Error(
339                hash,
340                Box::new(TransactionValidatorError::ValidationServiceUnreachable),
341            )
342        })
343        .collect()
344}
345
346#[cfg(test)]
347mod tests {
348    use super::*;
349    use crate::{
350        test_utils::MockTransaction,
351        validate::{TransactionValidationOutcome, ValidTransaction},
352        TransactionOrigin,
353    };
354    use alloy_primitives::{Address, U256};
355    use metrics::atomics::AtomicU64;
356    use std::sync::atomic::Ordering;
357
358    #[tokio::test]
359    async fn validation_send_metrics_track_cancellation_and_completion() {
360        for close_channel in [false, true] {
361            let (mut sender, task) = ValidationTask::new();
362            let gauge = Arc::new(AtomicU64::new(0));
363            sender.metrics.inflight_validation_jobs = metrics::Gauge::from_arc(gauge.clone());
364            let pending_sends = || f64::from_bits(gauge.load(Ordering::Relaxed));
365
366            sender.send(Box::pin(async {})).await.unwrap();
367            assert_eq!(pending_sends(), 0.0);
368
369            // Keep the channel full so every producer waits for capacity.
370            let mut producers =
371                (0..8).map(|_| Box::pin(sender.send(Box::pin(async {})))).collect::<Vec<_>>();
372            for producer in &mut producers {
373                assert!(futures_util::poll!(producer.as_mut()).is_pending());
374            }
375            assert_eq!(pending_sends(), 8.0);
376
377            let survivor = producers.remove(0);
378            drop(producers);
379            assert_eq!(pending_sends(), 1.0);
380
381            if close_channel {
382                drop(task);
383            } else {
384                // Free a slot for the remaining producer.
385                drop(task.validation_jobs.lock().await.next().await.unwrap());
386            }
387            assert_eq!(survivor.await.is_err(), close_channel);
388            assert_eq!(pending_sends(), 0.0);
389        }
390    }
391
392    #[tokio::test]
393    async fn cloned_workers_validate_while_another_job_is_blocked() {
394        let (sender, task) = ValidationTask::new();
395        let first_worker = tokio::spawn(task.clone().run());
396        let second_worker = tokio::spawn(task.run());
397        let (started_tx, started_rx) = oneshot::channel();
398        let (release_tx, release_rx) = oneshot::channel();
399        sender
400            .send(Box::pin(async move {
401                started_tx.send(()).unwrap();
402                release_rx.await.unwrap();
403            }))
404            .await
405            .unwrap();
406        started_rx.await.unwrap();
407
408        // The first validation remains blocked. A second configured worker must
409        // independently dequeue and finish this job without releasing the first.
410        let (completed_tx, completed_rx) = oneshot::channel();
411        sender
412            .send(Box::pin(async move {
413                completed_tx.send(()).unwrap();
414            }))
415            .await
416            .unwrap();
417        tokio::time::timeout(std::time::Duration::from_secs(1), completed_rx)
418            .await
419            .expect("second worker cannot dequeue while first validation holds receiver lock")
420            .unwrap();
421
422        release_tx.send(()).unwrap();
423        drop(sender);
424        // Closing the bounded channel still shuts every worker down cleanly.
425        tokio::time::timeout(std::time::Duration::from_secs(1), async {
426            first_worker.await.unwrap();
427            second_worker.await.unwrap();
428        })
429        .await
430        .unwrap();
431    }
432
433    #[derive(Debug)]
434    struct NoopValidator;
435
436    impl TransactionValidator for NoopValidator {
437        type Transaction = MockTransaction;
438        type Block = reth_ethereum_primitives::Block;
439
440        async fn validate_transaction(
441            &self,
442            _origin: TransactionOrigin,
443            transaction: Self::Transaction,
444        ) -> TransactionValidationOutcome<Self::Transaction> {
445            TransactionValidationOutcome::Valid {
446                balance: U256::ZERO,
447                state_nonce: 0,
448                bytecode_hash: None,
449                transaction: ValidTransaction::Valid(transaction),
450                propagate: false,
451                authorities: Some(Vec::<Address>::new()),
452            }
453        }
454    }
455
456    #[tokio::test]
457    async fn executor_new_spawns_and_validates_single() {
458        let validator = NoopValidator;
459        let (executor, task) = TransactionValidationTaskExecutor::new(validator);
460        tokio::spawn(task.run());
461        let tx = MockTransaction::legacy();
462        let out = executor.validate_transaction(TransactionOrigin::External, tx).await;
463        assert!(out.is_valid());
464    }
465
466    #[tokio::test]
467    async fn executor_new_spawns_and_validates_batch() {
468        let validator = NoopValidator;
469        let (executor, task) = TransactionValidationTaskExecutor::new(validator);
470        tokio::spawn(task.run());
471        let txs = vec![
472            (TransactionOrigin::External, MockTransaction::legacy()),
473            (TransactionOrigin::Local, MockTransaction::legacy()),
474        ];
475        let out = executor.validate_transactions(txs).await;
476        assert_eq!(out.len(), 2);
477        assert!(out.iter().all(|o| o.is_valid()));
478    }
479
480    #[tokio::test]
481    async fn cloned_executors_share_bounded_queue() {
482        let (executor, task) = TransactionValidationTaskExecutor::new(NoopValidator);
483        let mut submissions = tokio::task::JoinSet::new();
484
485        let first_executor = executor.clone();
486        let mut first_submission = Box::pin(async move {
487            first_executor
488                .validate_transaction(TransactionOrigin::Local, MockTransaction::legacy())
489                .await
490        });
491        // Poll once to fill the bounded queue before starting workers.
492        assert!(futures_util::poll!(first_submission.as_mut()).is_pending());
493        assert_eq!(executor.to_validation_task.tx.capacity(), 0);
494
495        submissions.spawn(async move {
496            let outcome = first_submission.await;
497            assert!(outcome.is_valid());
498        });
499
500        for index in 1..64 {
501            let executor = executor.clone();
502            submissions.spawn(async move {
503                let transaction = MockTransaction::legacy();
504                let outcomes = match index % 3 {
505                    0 => vec![
506                        executor.validate_transaction(TransactionOrigin::Local, transaction).await,
507                    ],
508                    1 => {
509                        executor
510                            .validate_transactions(vec![(TransactionOrigin::Local, transaction)])
511                            .await
512                    }
513                    _ => {
514                        executor
515                            .validate_transactions_with_origin(
516                                TransactionOrigin::Local,
517                                vec![transaction],
518                            )
519                            .await
520                    }
521                };
522                assert_eq!(outcomes.len(), 1);
523                assert!(outcomes[0].is_valid());
524            });
525        }
526
527        let first_worker = tokio::spawn(task.clone().run());
528        let second_worker = tokio::spawn(task.run());
529        tokio::time::timeout(std::time::Duration::from_secs(5), async {
530            while let Some(result) = submissions.join_next().await {
531                result.unwrap();
532            }
533            drop(executor);
534            first_worker.await.unwrap();
535            second_worker.await.unwrap();
536        })
537        .await
538        .expect("concurrent producers and workers must drain and shut down");
539    }
540
541    #[tokio::test]
542    async fn configured_workers_have_bounded_handoff_capacity() {
543        let runtime = Runtime::test();
544        for (additional_tasks, expected_capacity) in [(0, 2), (1, 4), (8, 18)] {
545            let executor =
546                TransactionValidationTaskExecutor::spawn(NoopValidator, &runtime, additional_tasks);
547            assert_eq!(executor.to_validation_task.tx.max_capacity(), expected_capacity);
548            let outcome = tokio::time::timeout(
549                std::time::Duration::from_secs(5),
550                executor.validate_transaction(TransactionOrigin::Local, MockTransaction::legacy()),
551            )
552            .await
553            .expect("configured validation workers must process jobs");
554            assert!(outcome.is_valid());
555        }
556    }
557
558    #[derive(Debug)]
559    struct SameOriginBatchValidator;
560
561    impl TransactionValidator for SameOriginBatchValidator {
562        type Transaction = MockTransaction;
563        type Block = reth_ethereum_primitives::Block;
564
565        async fn validate_transaction(
566            &self,
567            _origin: TransactionOrigin,
568            _transaction: Self::Transaction,
569        ) -> TransactionValidationOutcome<Self::Transaction> {
570            panic!("same-origin batches must use the batch validator")
571        }
572
573        async fn validate_transactions_with_origin(
574            &self,
575            origin: TransactionOrigin,
576            transactions: impl IntoIterator<Item = Self::Transaction, IntoIter: Send> + Send,
577        ) -> Vec<TransactionValidationOutcome<Self::Transaction>> {
578            transactions
579                .into_iter()
580                .map(|transaction| TransactionValidationOutcome::Valid {
581                    balance: U256::ZERO,
582                    state_nonce: 0,
583                    bytecode_hash: None,
584                    transaction: ValidTransaction::Valid(transaction),
585                    propagate: origin.is_local(),
586                    authorities: None,
587                })
588                .collect()
589        }
590    }
591
592    #[tokio::test]
593    async fn executor_forwards_same_origin_batches() {
594        let (executor, task) = TransactionValidationTaskExecutor::new(SameOriginBatchValidator);
595        tokio::spawn(task.run());
596
597        let transactions = vec![MockTransaction::legacy(), MockTransaction::eip1559()];
598        let expected_hashes = transactions.iter().map(|tx| *tx.hash()).collect::<Vec<_>>();
599        let outcomes = executor
600            .validate_transactions_with_origin(TransactionOrigin::Local, transactions)
601            .await;
602
603        assert_eq!(outcomes.len(), expected_hashes.len());
604        assert!(outcomes.into_iter().zip(expected_hashes).all(|(outcome, expected_hash)| {
605            matches!(
606                outcome,
607                TransactionValidationOutcome::Valid { transaction, propagate: true, .. }
608                    if transaction.hash() == &expected_hash
609            )
610        }));
611    }
612}