reth_transaction_pool/validate/
task.rs1use 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
20type ValidationFuture = Pin<Box<dyn Future<Output = ()> + Send>>;
22
23type ValidationStream = ReceiverStream<ValidationFuture>;
25
26#[derive(Clone)]
32pub struct ValidationTask {
33 validation_jobs: Arc<Mutex<ValidationStream>>,
34}
35
36impl ValidationTask {
37 pub fn new() -> (ValidationJobSender, Self) {
41 Self::with_capacity(1)
42 }
43
44 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 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 pub async fn run(self) {
60 loop {
61 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#[derive(Debug)]
78pub struct ValidationJobSender {
79 tx: mpsc::Sender<Pin<Box<dyn Future<Output = ()> + Send>>>,
80 metrics: TxPoolValidatorMetrics,
81}
82
83impl ValidationJobSender {
84 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#[derive(Debug)]
98pub struct TransactionValidationTaskExecutor<V> {
99 pub validator: Arc<V>,
101 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
115impl TransactionValidationTaskExecutor<()> {
118 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 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 pub fn validator(&self) -> &V {
146 &self.validator
147 }
148}
149
150impl<Client, Tx, Evm> TransactionValidationTaskExecutor<EthTransactionValidator<Client, Tx, Evm>> {
151 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 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 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 pub fn spawn(validator: V, tasks: &Runtime, additional_tasks: usize) -> Self {
206 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
322struct 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 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 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 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 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 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}