1use log::{debug, trace};
31use std::{
32 fmt,
33 time::{Duration, Instant},
34};
35
36use sp_consensus::{error::Error as ConsensusError, BlockOrigin};
37use sp_runtime::{
38 traits::{Block as BlockT, Header as _, NumberFor},
39 Justifications,
40};
41
42use crate::{
43 block_import::{
44 BlockCheckParams, BlockImport, BlockImportParams, ImportResult, ImportedAux, ImportedState,
45 JustificationImport, StateAction,
46 },
47 metrics::Metrics,
48};
49
50pub use basic_queue::BasicQueue;
51
52const LOG_TARGET: &str = "sync::import-queue";
53
54pub type DefaultImportQueue<Block> = BasicQueue<Block>;
58
59mod basic_queue;
60pub mod buffered_link;
61pub mod mock;
62
63pub type BoxBlockImport<B> = Box<dyn BlockImport<B, Error = ConsensusError> + Send + Sync>;
65
66pub type BoxJustificationImport<B> =
68 Box<dyn JustificationImport<B, Error = ConsensusError> + Send + Sync>;
69
70pub type RuntimeOrigin = sc_network_types::PeerId;
72
73#[derive(Debug, PartialEq, Eq, Clone)]
75pub struct IncomingBlock<B: BlockT> {
76 pub hash: <B as BlockT>::Hash,
78 pub header: Option<<B as BlockT>::Header>,
80 pub body: Option<Vec<<B as BlockT>::Extrinsic>>,
82 pub justifications: Option<Justifications>,
84 pub origin: Option<RuntimeOrigin>,
86 pub allow_missing_state: bool,
88 pub skip_execution: bool,
90 pub import_existing: bool,
92 pub state: Option<ImportedState<B>>,
94}
95
96#[async_trait::async_trait]
98pub trait Verifier<B: BlockT>: Send + Sync {
99 async fn verify(&self, block: BlockImportParams<B>) -> Result<BlockImportParams<B>, String>;
102}
103
104pub trait ImportQueueService<B: BlockT>: Send {
108 fn import_blocks(&mut self, origin: BlockOrigin, blocks: Vec<IncomingBlock<B>>);
111
112 fn import_justifications(
114 &mut self,
115 who: RuntimeOrigin,
116 hash: B::Hash,
117 number: NumberFor<B>,
118 justifications: Justifications,
119 );
120}
121
122#[async_trait::async_trait]
123pub trait ImportQueue<B: BlockT>: Send {
124 fn service(&self) -> Box<dyn ImportQueueService<B>>;
126
127 fn service_ref(&mut self) -> &mut dyn ImportQueueService<B>;
129
130 fn poll_actions(&mut self, cx: &mut futures::task::Context, link: &dyn Link<B>);
134
135 async fn run(self, link: &dyn Link<B>);
140}
141
142#[derive(Debug, PartialEq)]
144pub enum JustificationImportResult {
145 Success,
147
148 Failure,
150
151 OutdatedJustification,
153}
154
155pub trait Link<B: BlockT>: Send + Sync {
158 fn blocks_processed(
160 &self,
161 _imported: usize,
162 _count: usize,
163 _results: Vec<(BlockImportResult<B>, B::Hash)>,
164 ) {
165 }
166
167 fn justification_imported(
169 &self,
170 _who: RuntimeOrigin,
171 _hash: &B::Hash,
172 _number: NumberFor<B>,
173 _import_result: JustificationImportResult,
174 ) {
175 }
176
177 fn request_justification(&self, _hash: &B::Hash, _number: NumberFor<B>) {}
179}
180
181#[derive(Debug, PartialEq)]
183pub enum BlockImportStatus<BlockNumber: fmt::Debug + PartialEq> {
184 ImportedKnown(BlockNumber, Option<RuntimeOrigin>),
186 ImportedUnknown(BlockNumber, ImportedAux, Option<RuntimeOrigin>),
188}
189
190impl<BlockNumber: fmt::Debug + PartialEq> BlockImportStatus<BlockNumber> {
191 pub fn number(&self) -> &BlockNumber {
193 match self {
194 BlockImportStatus::ImportedKnown(n, _) |
195 BlockImportStatus::ImportedUnknown(n, _, _) => n,
196 }
197 }
198}
199
200#[derive(Debug, thiserror::Error)]
202pub enum BlockImportError {
203 #[error("block is missing a header (origin = {0:?})")]
205 IncompleteHeader(Option<RuntimeOrigin>),
206
207 #[error("block verification failed (origin = {0:?}): {1}")]
209 VerificationFailed(Option<RuntimeOrigin>, String),
210
211 #[error("bad block (origin = {0:?})")]
213 BadBlock(Option<RuntimeOrigin>),
214
215 #[error("block is missing parent state")]
217 MissingState,
218
219 #[error("block has an unknown parent")]
221 UnknownParent,
222
223 #[error("import has been cancelled")]
225 Cancelled,
226
227 #[error("consensus error: {0}")]
229 Other(ConsensusError),
230}
231
232type BlockImportResult<B> = Result<BlockImportStatus<NumberFor<B>>, BlockImportError>;
233
234pub async fn import_single_block<B: BlockT, V: Verifier<B>>(
236 import_handle: &mut impl BlockImport<B, Error = ConsensusError>,
237 block_origin: BlockOrigin,
238 block: IncomingBlock<B>,
239 verifier: &V,
240) -> BlockImportResult<B> {
241 match verify_single_block_metered(import_handle, block_origin, block, verifier, None).await? {
242 SingleBlockVerificationOutcome::Imported(import_status) => Ok(import_status),
243 SingleBlockVerificationOutcome::Verified(import_parameters) => {
244 import_single_block_metered(import_handle, import_parameters, None).await
245 },
246 }
247}
248
249fn import_handler<Block>(
250 number: NumberFor<Block>,
251 hash: Block::Hash,
252 parent_hash: Block::Hash,
253 block_origin: Option<RuntimeOrigin>,
254 import: Result<ImportResult, ConsensusError>,
255) -> Result<BlockImportStatus<NumberFor<Block>>, BlockImportError>
256where
257 Block: BlockT,
258{
259 match import {
260 Ok(ImportResult::AlreadyInChain) => {
261 trace!(target: LOG_TARGET, "Block already in chain {}: {:?}", number, hash);
262 Ok(BlockImportStatus::ImportedKnown(number, block_origin))
263 },
264 Ok(ImportResult::Imported(aux)) => {
265 Ok(BlockImportStatus::ImportedUnknown(number, aux, block_origin))
266 },
267 Ok(ImportResult::MissingState) => {
268 debug!(
269 target: LOG_TARGET,
270 "Parent state is missing for {}: {:?}, parent: {:?}", number, hash, parent_hash
271 );
272 Err(BlockImportError::MissingState)
273 },
274 Ok(ImportResult::UnknownParent) => {
275 debug!(
276 target: LOG_TARGET,
277 "Block with unknown parent {}: {:?}, parent: {:?}", number, hash, parent_hash
278 );
279 Err(BlockImportError::UnknownParent)
280 },
281 Ok(ImportResult::KnownBad) => {
282 debug!(target: LOG_TARGET, "Peer gave us a bad block {}: {:?}", number, hash);
283 Err(BlockImportError::BadBlock(block_origin))
284 },
285 Err(e) => {
286 debug!(target: LOG_TARGET, "Error importing block {}: {:?}: {}", number, hash, e);
287 Err(BlockImportError::Other(e))
288 },
289 }
290}
291
292pub(crate) enum SingleBlockVerificationOutcome<Block: BlockT> {
293 Imported(BlockImportStatus<NumberFor<Block>>),
295 Verified(SingleBlockImportParameters<Block>),
297}
298
299pub(crate) struct SingleBlockImportParameters<Block: BlockT> {
300 import_block: BlockImportParams<Block>,
301 hash: Block::Hash,
302 block_origin: Option<RuntimeOrigin>,
303 verification_time: Duration,
304}
305
306pub(crate) async fn verify_single_block_metered<B: BlockT, V: Verifier<B>>(
308 import_handle: &impl BlockImport<B, Error = ConsensusError>,
309 block_origin: BlockOrigin,
310 block: IncomingBlock<B>,
311 verifier: &V,
312 metrics: Option<&Metrics>,
313) -> Result<SingleBlockVerificationOutcome<B>, BlockImportError> {
314 let peer = block.origin;
315 let justifications = block.justifications;
316
317 let Some(header) = block.header else {
318 if let Some(ref peer) = peer {
319 debug!(target: LOG_TARGET, "Header {} was not provided by {peer} ", block.hash);
320 } else {
321 debug!(target: LOG_TARGET, "Header {} was not provided ", block.hash);
322 }
323 return Err(BlockImportError::IncompleteHeader(peer));
324 };
325
326 let number = *header.number();
327 let hash = block.hash;
328 let parent_hash = *header.parent_hash();
329
330 trace!(target: LOG_TARGET, "Block {number} ({hash}) has {:?} logs (origin: {:?})", header.digest().logs().len(), block_origin);
331
332 if matches!(block_origin, BlockOrigin::WarpSync) {
335 return Ok(SingleBlockVerificationOutcome::Verified(SingleBlockImportParameters {
336 import_block: BlockImportParams::new(block_origin, header),
337 hash: block.hash,
338 block_origin: peer,
339 verification_time: Duration::ZERO,
340 }));
341 }
342
343 match import_handler::<B>(
344 number,
345 hash,
346 parent_hash,
347 peer,
348 import_handle
349 .check_block(BlockCheckParams {
350 hash,
351 number,
352 parent_hash,
353 allow_missing_state: block.allow_missing_state,
354 import_existing: block.import_existing,
355 allow_missing_parent: block.state.is_some(),
356 })
357 .await,
358 )? {
359 BlockImportStatus::ImportedUnknown { .. } => (),
360 r => {
361 return Ok(SingleBlockVerificationOutcome::Imported(r));
363 },
364 }
365
366 let started = Instant::now();
367
368 let mut import_block = BlockImportParams::new(block_origin, header);
369 import_block.body = block.body;
370 import_block.justifications = justifications;
371 import_block.post_hash = Some(hash);
372 import_block.import_existing = block.import_existing;
373
374 if let Some(state) = block.state {
375 let changes = crate::block_import::StorageChanges::Import(state);
376 import_block.state_action = StateAction::ApplyChanges(changes);
377 } else if block.skip_execution {
378 import_block.state_action = StateAction::Skip;
379 } else if block.allow_missing_state {
380 import_block.state_action = StateAction::ExecuteIfPossible;
381 }
382
383 let import_block = verifier.verify(import_block).await.map_err(|msg| {
384 if let Some(ref peer) = peer {
385 trace!(
386 target: LOG_TARGET,
387 "Verifying {}({}) from {} failed: {}",
388 number,
389 hash,
390 peer,
391 msg
392 );
393 } else {
394 trace!(target: LOG_TARGET, "Verifying {}({}) failed: {}", number, hash, msg);
395 }
396 if let Some(metrics) = metrics {
397 metrics.report_verification(false, started.elapsed());
398 }
399 BlockImportError::VerificationFailed(peer, msg)
400 })?;
401
402 let verification_time = started.elapsed();
403 if let Some(metrics) = metrics {
404 metrics.report_verification(true, verification_time);
405 }
406
407 Ok(SingleBlockVerificationOutcome::Verified(SingleBlockImportParameters {
408 import_block,
409 hash,
410 block_origin: peer,
411 verification_time,
412 }))
413}
414
415pub(crate) async fn import_single_block_metered<Block: BlockT>(
416 import_handle: &mut impl BlockImport<Block, Error = ConsensusError>,
417 import_parameters: SingleBlockImportParameters<Block>,
418 metrics: Option<&Metrics>,
419) -> BlockImportResult<Block> {
420 let started = Instant::now();
421
422 let SingleBlockImportParameters { import_block, hash, block_origin, verification_time } =
423 import_parameters;
424
425 let number = *import_block.header.number();
426 let parent_hash = *import_block.header.parent_hash();
427
428 let imported = import_handle.import_block(import_block).await;
429 if let Some(metrics) = metrics {
430 metrics.report_verification_and_import(started.elapsed() + verification_time);
431 }
432
433 import_handler::<Block>(number, hash, parent_hash, block_origin, imported)
434}