1use crate::{
22 common::{sliding_stat::DurationSlidingStats, STAT_SLIDING_WINDOW},
23 graph::ValidateTransactionPriority,
24 insert_and_log_throttled, LOG_TARGET, LOG_TARGET_STAT,
25};
26use async_trait::async_trait;
27use codec::Encode;
28use futures::future::{Future, FutureExt};
29use prometheus_endpoint::Registry as PrometheusRegistry;
30use sp_api::ApiExt;
31use sp_blockchain::TreeRoute;
32use sp_core::traits::SpawnEssentialNamed;
33use sp_runtime::{
34 generic::BlockId,
35 traits::{self, Block as BlockT},
36 transaction_validity::{TransactionSource, TransactionValidity},
37};
38use sp_transaction_pool::runtime_api::TaggedTransactionQueue;
39use std::{
40 marker::PhantomData,
41 pin::Pin,
42 sync::Arc,
43 time::{Duration, Instant},
44};
45use tokio::sync::{mpsc, oneshot, Mutex};
46
47use super::{
48 error::{self, Error},
49 metrics::{ApiMetrics, ApiMetricsExt},
50};
51use crate::graph;
52use tracing::{trace, warn, Level};
53
54pub struct FullChainApi<Client, Block> {
56 client: Arc<Client>,
57 _marker: PhantomData<Block>,
58 metrics: Option<Arc<ApiMetrics>>,
59 validation_pool_normal: mpsc::Sender<Pin<Box<dyn Future<Output = ()> + Send>>>,
60 validation_pool_maintained: mpsc::Sender<Pin<Box<dyn Future<Output = ()> + Send>>>,
61 validate_transaction_normal_stats: DurationSlidingStats,
62 validate_transaction_maintained_stats: DurationSlidingStats,
63}
64
65fn spawn_validation_pool_task(
67 name: &'static str,
68 receiver_normal: Arc<Mutex<mpsc::Receiver<Pin<Box<dyn Future<Output = ()> + Send>>>>>,
69 receiver_maintained: Arc<Mutex<mpsc::Receiver<Pin<Box<dyn Future<Output = ()> + Send>>>>>,
70 spawner: &impl SpawnEssentialNamed,
71 stats: DurationSlidingStats,
72 blocking_stats: DurationSlidingStats,
73) {
74 spawner.spawn_essential_blocking(
75 name,
76 Some("transaction-pool"),
77 async move {
78 loop {
79 let start = Instant::now();
80
81 let task = {
82 let receiver_maintained = receiver_maintained.clone();
83 let receiver_normal = receiver_normal.clone();
84 tokio::select! {
85 Some(task) = async {
86 receiver_maintained.lock().await.recv().await
87 } => { task }
88 Some(task) = async {
89 receiver_normal.lock().await.recv().await
90 } => { task }
91 else => {
92 return
93 }
94 }
95 };
96
97 let blocking_duration = {
98 let start = Instant::now();
99 task.await;
100 start.elapsed()
101 };
102
103 insert_and_log_throttled!(
104 Level::DEBUG,
105 target:LOG_TARGET_STAT,
106 prefix:format!("validate_transaction_inner_stats"),
107 stats,
108 start.elapsed().into()
109 );
110 insert_and_log_throttled!(
111 Level::DEBUG,
112 target:LOG_TARGET_STAT,
113 prefix:format!("validate_transaction_blocking_stats"),
114 blocking_stats,
115 blocking_duration.into()
116 );
117 trace!(target:LOG_TARGET, duration=?start.elapsed(), "spawn_validation_pool_task");
118 }
119 }
120 .boxed(),
121 );
122}
123
124impl<Client, Block> FullChainApi<Client, Block> {
125 pub fn new(
127 client: Arc<Client>,
128 prometheus: Option<&PrometheusRegistry>,
129 spawner: &impl SpawnEssentialNamed,
130 ) -> Self {
131 let stats = DurationSlidingStats::new(Duration::from_secs(STAT_SLIDING_WINDOW));
132 let blocking_stats = DurationSlidingStats::new(Duration::from_secs(STAT_SLIDING_WINDOW));
133
134 let metrics = prometheus.map(ApiMetrics::register).and_then(|r| match r {
135 Err(error) => {
136 warn!(
137 target: LOG_TARGET,
138 ?error,
139 "Failed to register transaction pool API Prometheus metrics"
140 );
141 None
142 },
143 Ok(api) => Some(Arc::new(api)),
144 });
145
146 let (sender, receiver) = mpsc::channel(1);
147 let (sender_maintained, receiver_maintained) = mpsc::channel(1);
148
149 let receiver = Arc::new(Mutex::new(receiver));
150 let receiver_maintained = Arc::new(Mutex::new(receiver_maintained));
151 spawn_validation_pool_task(
152 "transaction-pool-task-0",
153 receiver.clone(),
154 receiver_maintained.clone(),
155 spawner,
156 stats.clone(),
157 blocking_stats.clone(),
158 );
159 spawn_validation_pool_task(
160 "transaction-pool-task-1",
161 receiver,
162 receiver_maintained,
163 spawner,
164 stats.clone(),
165 blocking_stats.clone(),
166 );
167
168 FullChainApi {
169 client,
170 validation_pool_normal: sender,
171 validation_pool_maintained: sender_maintained,
172 _marker: Default::default(),
173 metrics,
174 validate_transaction_normal_stats: DurationSlidingStats::new(Duration::from_secs(
175 STAT_SLIDING_WINDOW,
176 )),
177 validate_transaction_maintained_stats: DurationSlidingStats::new(Duration::from_secs(
178 STAT_SLIDING_WINDOW,
179 )),
180 }
181 }
182}
183
184#[async_trait]
185impl<Client, Block> graph::ChainApi for FullChainApi<Client, Block>
186where
187 Block: BlockT,
188 Client: crate::ClientForTransactionPool<Block>,
189{
190 type Block = Block;
191 type Error = error::Error;
192
193 async fn block_body(
194 &self,
195 hash: Block::Hash,
196 ) -> Result<Option<Vec<<Self::Block as BlockT>::Extrinsic>>, Self::Error> {
197 self.client.block_body(hash).map_err(error::Error::from)
198 }
199
200 async fn validate_transaction(
201 &self,
202 at: <Self::Block as BlockT>::Hash,
203 source: TransactionSource,
204 uxt: graph::ExtrinsicFor<Self>,
205 validation_priority: ValidateTransactionPriority,
206 ) -> Result<TransactionValidity, Self::Error> {
207 let start = Instant::now();
208 let (tx, rx) = oneshot::channel();
209 let client = self.client.clone();
210 let (stats, validation_pool, prefix) =
211 if validation_priority == ValidateTransactionPriority::Maintained {
212 (
213 self.validate_transaction_maintained_stats.clone(),
214 self.validation_pool_maintained.clone(),
215 "validate_transaction_maintained_stats",
216 )
217 } else {
218 (
219 self.validate_transaction_normal_stats.clone(),
220 self.validation_pool_normal.clone(),
221 "validate_transaction_stats",
222 )
223 };
224 let metrics = self.metrics.clone();
225
226 metrics.report(|m| m.validations_scheduled.inc());
227
228 {
229 validation_pool
230 .send(
231 async move {
232 let res = validate_transaction_blocking(&*client, at, source, uxt);
233 let _ = tx.send(res);
234 metrics.report(|m| m.validations_finished.inc());
235 }
236 .boxed(),
237 )
238 .await
239 .map_err(|e| Error::RuntimeApi(format!("Validation pool down: {:?}", e)))?;
240 }
241
242 let validity = match rx.await {
243 Ok(r) => r,
244 Err(_) => Err(Error::RuntimeApi("Validation was canceled".into())),
245 };
246
247 insert_and_log_throttled!(
248 Level::DEBUG,
249 target:LOG_TARGET_STAT,
250 prefix:prefix,
251 stats,
252 start.elapsed().into()
253 );
254
255 validity
256 }
257
258 fn validate_transaction_blocking(
262 &self,
263 at: Block::Hash,
264 source: TransactionSource,
265 uxt: graph::ExtrinsicFor<Self>,
266 ) -> Result<TransactionValidity, Self::Error> {
267 validate_transaction_blocking(&*self.client, at, source, uxt)
268 }
269
270 fn block_id_to_number(
271 &self,
272 at: &BlockId<Self::Block>,
273 ) -> Result<Option<graph::NumberFor<Self>>, Self::Error> {
274 self.client.to_number(at).map_err(|e| Error::BlockIdConversion(e.to_string()))
275 }
276
277 fn block_id_to_hash(
278 &self,
279 at: &BlockId<Self::Block>,
280 ) -> Result<Option<graph::BlockHash<Self>>, Self::Error> {
281 self.client.to_hash(at).map_err(|e| Error::BlockIdConversion(e.to_string()))
282 }
283
284 fn hash_and_length(
285 &self,
286 ex: &graph::RawExtrinsicFor<Self>,
287 ) -> (graph::ExtrinsicHash<Self>, usize) {
288 ex.using_encoded(|x| (<traits::HashingFor<Block> as traits::Hash>::hash(x), x.len()))
289 }
290
291 fn block_header(
292 &self,
293 hash: <Self::Block as BlockT>::Hash,
294 ) -> Result<Option<<Self::Block as BlockT>::Header>, Self::Error> {
295 self.client.header(hash).map_err(Into::into)
296 }
297
298 fn tree_route(
299 &self,
300 from: <Self::Block as BlockT>::Hash,
301 to: <Self::Block as BlockT>::Hash,
302 ) -> Result<TreeRoute<Self::Block>, Self::Error> {
303 sp_blockchain::tree_route::<Block, Client>(&*self.client, from, to).map_err(Into::into)
304 }
305}
306
307fn validate_transaction_blocking<Client, Block>(
310 client: &Client,
311 at: Block::Hash,
312 source: TransactionSource,
313 uxt: graph::ExtrinsicFor<FullChainApi<Client, Block>>,
314) -> error::Result<TransactionValidity>
315where
316 Block: BlockT,
317 Client: crate::ClientForTransactionPool<Block>,
318{
319 let s = std::time::Instant::now();
320 let tx_hash = uxt.using_encoded(|x| <traits::HashingFor<Block> as traits::Hash>::hash(x));
321
322 let result = sp_tracing::within_span!(sp_tracing::Level::TRACE, "validate_transaction";
323 {
324 let runtime_api = client.runtime_api();
325 let api_version = sp_tracing::within_span! { sp_tracing::Level::TRACE, "check_version";
326 runtime_api
327 .api_version::<dyn TaggedTransactionQueue<Block>>(at)
328 .map_err(|e| Error::RuntimeApi(e.to_string()))?
329 .ok_or_else(|| Error::RuntimeApi(
330 format!("Could not find `TaggedTransactionQueue` api for block `{:?}`.", at)
331 ))
332 }?;
333
334 use sp_api::Core;
335
336 sp_tracing::within_span!(
337 sp_tracing::Level::TRACE, "runtime::validate_transaction";
338 {
339 if api_version >= 3 {
340 runtime_api.validate_transaction(at, source, (*uxt).clone(), at)
341 .map_err(|e| Error::RuntimeApi(e.to_string()))
342 } else {
343 let block_number = client.to_number(&BlockId::Hash(at))
344 .map_err(|e| Error::RuntimeApi(e.to_string()))?
345 .ok_or_else(||
346 Error::RuntimeApi(format!("Could not get number for block `{:?}`.", at))
347 )?;
348
349 runtime_api.initialize_block(at, &sp_runtime::traits::Header::new(
351 block_number + sp_runtime::traits::One::one(),
352 Default::default(),
353 Default::default(),
354 at,
355 Default::default()),
356 ).map_err(|e| Error::RuntimeApi(e.to_string()))?;
357
358 if api_version == 2 {
359 #[allow(deprecated)] runtime_api.validate_transaction_before_version_3(at, source, (*uxt).clone())
361 .map_err(|e| Error::RuntimeApi(e.to_string()))
362 } else {
363 #[allow(deprecated)] runtime_api.validate_transaction_before_version_2(at, (*uxt).clone())
365 .map_err(|e| Error::RuntimeApi(e.to_string()))
366 }
367 }
368 })
369 });
370 trace!(
371 target: LOG_TARGET,
372 ?tx_hash,
373 ?at,
374 duration = ?s.elapsed(),
375 "validate_transaction_blocking"
376 );
377 result
378}