1use sc_consensus::IncomingBlock;
22use sp_consensus::BlockOrigin;
23
24use crate::{
25 block_relay_protocol::{BlockDownloader, BlockResponseError},
26 service::network::NetworkServiceHandle,
27 strategy::{
28 chain_sync::validate_blocks, disconnected_peers::DisconnectedPeers, StrategyKey,
29 SyncingAction,
30 },
31 types::{BadPeer, SyncState, SyncStatus},
32 LOG_TARGET,
33};
34use codec::{Decode, Encode};
35use futures::{channel::oneshot, FutureExt};
36use log::{debug, error, trace, warn};
37use sc_network::{IfDisconnected, ProtocolName};
38use sc_network_common::sync::message::{
39 BlockAnnounce, BlockAttributes, BlockData, BlockRequest, Direction, FromBlock,
40};
41use sc_network_types::PeerId;
42use sp_blockchain::HeaderBackend;
43use sp_runtime::{
44 traits::{Block as BlockT, Header, NumberFor, Zero},
45 Justifications, SaturatedConversion,
46};
47use std::{any::Any, collections::HashMap, fmt, sync::Arc};
48
49const MIN_PEERS_TO_START_WARP_SYNC: usize = 3;
51
52pub struct EncodedProof(pub Vec<u8>);
54
55#[derive(Encode, Decode, Debug, Clone)]
57pub struct WarpProofRequest<B: BlockT> {
58 pub begin: B::Hash,
60}
61
62pub trait Verifier<Block: BlockT>: Send + Sync {
64 fn verify(
66 &mut self,
67 proof: &EncodedProof,
68 ) -> Result<VerificationResult<Block>, Box<dyn std::error::Error + Send + Sync>>;
69 fn next_proof_context(&self) -> Block::Hash;
71 fn status(&self) -> Option<String>;
73}
74
75pub enum VerificationResult<Block: BlockT> {
77 Partial(Vec<(Block::Header, Justifications)>),
79 Complete(Block::Header, Vec<(Block::Header, Justifications)>),
81}
82
83pub trait WarpSyncProvider<Block: BlockT>: Send + Sync {
85 fn generate(
88 &self,
89 start: Block::Hash,
90 ) -> Result<EncodedProof, Box<dyn std::error::Error + Send + Sync>>;
91 fn create_verifier(&self) -> Box<dyn Verifier<Block>>;
93}
94
95mod rep {
96 use sc_network::ReputationChange as Rep;
97
98 pub const UNEXPECTED_RESPONSE: Rep = Rep::new(-(1 << 29), "Unexpected response");
100
101 pub const BAD_WARP_PROOF: Rep = Rep::new(-(1 << 29), "Bad warp proof");
103
104 pub const NO_BLOCK: Rep = Rep::new(-(1 << 29), "No requested block data");
106
107 pub const NOT_REQUESTED: Rep = Rep::new(-(1 << 29), "Not requested block data");
109
110 pub const VERIFICATION_FAIL: Rep = Rep::new(-(1 << 29), "Block verification failed");
112
113 pub const BAD_MESSAGE: Rep = Rep::new(-(1 << 12), "Bad message");
115}
116
117#[derive(Clone, Eq, PartialEq, Debug)]
119pub enum WarpSyncPhase<Block: BlockT> {
120 AwaitingPeers { required_peers: usize },
122 DownloadingWarpProofs,
124 DownloadingTargetBlock,
126 DownloadingState,
128 ImportingState,
130 DownloadingBlocks(NumberFor<Block>),
132 Complete,
134}
135
136impl<Block: BlockT> fmt::Display for WarpSyncPhase<Block> {
137 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
138 match self {
139 Self::AwaitingPeers { required_peers } => {
140 write!(f, "Waiting for {required_peers} peers to be connected")
141 },
142 Self::DownloadingWarpProofs => write!(f, "Downloading finality proofs"),
143 Self::DownloadingTargetBlock => write!(f, "Downloading target block"),
144 Self::DownloadingState => write!(f, "Downloading state"),
145 Self::ImportingState => write!(f, "Importing state"),
146 Self::DownloadingBlocks(n) => write!(f, "Downloading block history (#{})", n),
147 Self::Complete => write!(f, "Warp sync is complete"),
148 }
149 }
150}
151
152#[derive(Clone, Eq, PartialEq, Debug)]
154pub struct WarpSyncProgress<Block: BlockT> {
155 pub phase: WarpSyncPhase<Block>,
157 pub total_bytes: u64,
159 pub status: Option<String>,
161}
162
163pub enum WarpSyncConfig<Block: BlockT> {
165 WithProvider(Arc<dyn WarpSyncProvider<Block>>),
167 WithTarget(<Block as BlockT>::Header),
171}
172
173enum Phase<B: BlockT> {
175 WaitingForPeers { warp_sync_provider: Arc<dyn WarpSyncProvider<B>> },
177 WarpProof { verifier: Box<dyn Verifier<B>> },
179 TargetBlock(B::Header),
181 Complete,
183}
184
185enum PeerState {
186 Available,
187 DownloadingProofs,
188 DownloadingTargetBlock,
189}
190
191impl PeerState {
192 fn is_available(&self) -> bool {
193 matches!(self, PeerState::Available)
194 }
195}
196
197struct Peer<B: BlockT> {
198 best_number: NumberFor<B>,
199 state: PeerState,
200}
201
202pub struct WarpSyncResult<B: BlockT> {
203 pub target_header: B::Header,
204 pub target_body: Option<Vec<B::Extrinsic>>,
205 pub target_justifications: Option<Justifications>,
206}
207
208pub struct WarpSync<B: BlockT> {
210 phase: Phase<B>,
211 total_proof_bytes: u64,
212 total_state_bytes: u64,
213 peers: HashMap<PeerId, Peer<B>>,
214 disconnected_peers: DisconnectedPeers,
215 protocol_name: Option<ProtocolName>,
216 block_downloader: Arc<dyn BlockDownloader<B>>,
217 actions: Vec<SyncingAction<B>>,
218 result: Option<WarpSyncResult<B>>,
219 min_peers_to_start_warp_sync: usize,
221}
222
223impl<B> WarpSync<B>
224where
225 B: BlockT,
226{
227 pub const STRATEGY_KEY: StrategyKey = StrategyKey::new("Warp");
229
230 pub fn new<Client>(
234 client: Arc<Client>,
235 warp_sync_config: WarpSyncConfig<B>,
236 protocol_name: Option<ProtocolName>,
237 block_downloader: Arc<dyn BlockDownloader<B>>,
238 min_peers_to_start_warp_sync: Option<usize>,
239 ) -> Self
240 where
241 Client: HeaderBackend<B> + 'static,
242 {
243 let min_peers_to_start_warp_sync =
244 min_peers_to_start_warp_sync.unwrap_or(MIN_PEERS_TO_START_WARP_SYNC);
245 if client.info().finalized_state.is_some() {
246 error!(
247 target: LOG_TARGET,
248 "Can't use warp sync mode with a partially synced database. Reverting to full sync mode."
249 );
250 return Self {
251 phase: Phase::Complete,
252 total_proof_bytes: 0,
253 total_state_bytes: 0,
254 peers: HashMap::new(),
255 disconnected_peers: DisconnectedPeers::new(),
256 protocol_name,
257 block_downloader,
258 actions: vec![SyncingAction::Finished],
259 result: None,
260 min_peers_to_start_warp_sync,
261 };
262 }
263
264 let phase = match warp_sync_config {
265 WarpSyncConfig::WithProvider(warp_sync_provider) => {
266 Phase::WaitingForPeers { warp_sync_provider }
267 },
268 WarpSyncConfig::WithTarget(target_header) => Phase::TargetBlock(target_header),
269 };
270
271 Self {
272 phase,
273 total_proof_bytes: 0,
274 total_state_bytes: 0,
275 peers: HashMap::new(),
276 disconnected_peers: DisconnectedPeers::new(),
277 protocol_name,
278 block_downloader,
279 actions: Vec::new(),
280 result: None,
281 min_peers_to_start_warp_sync,
282 }
283 }
284
285 pub fn add_peer(&mut self, peer_id: PeerId, _best_hash: B::Hash, best_number: NumberFor<B>) {
287 self.peers.insert(peer_id, Peer { best_number, state: PeerState::Available });
288
289 self.try_to_start_warp_sync();
290 }
291
292 pub fn remove_peer(&mut self, peer_id: &PeerId) {
294 if let Some(state) = self.peers.remove(peer_id) {
295 if !state.state.is_available() {
296 if let Some(bad_peer) =
297 self.disconnected_peers.on_disconnect_during_request(*peer_id)
298 {
299 self.actions.push(SyncingAction::DropPeer(bad_peer));
300 }
301 }
302 }
303 }
304
305 #[must_use]
309 pub fn on_validated_block_announce(
310 &mut self,
311 is_best: bool,
312 peer_id: PeerId,
313 announce: &BlockAnnounce<B::Header>,
314 ) -> Option<(B::Hash, NumberFor<B>)> {
315 is_best.then(|| {
316 let best_number = *announce.header.number();
317 let best_hash = announce.header.hash();
318 if let Some(ref mut peer) = self.peers.get_mut(&peer_id) {
319 peer.best_number = best_number;
320 }
321 (best_hash, best_number)
323 })
324 }
325
326 fn try_to_start_warp_sync(&mut self) {
328 let Phase::WaitingForPeers { warp_sync_provider } = &self.phase else { return };
329
330 if self.peers.len() < self.min_peers_to_start_warp_sync {
331 return;
332 }
333
334 let verifier = warp_sync_provider.create_verifier();
335 self.phase = Phase::WarpProof { verifier };
336 debug!(target: LOG_TARGET, "Started warp sync with {} peers.", self.peers.len());
337 }
338
339 pub fn on_generic_response(
340 &mut self,
341 peer_id: &PeerId,
342 protocol_name: ProtocolName,
343 response: Box<dyn Any + Send>,
344 ) {
345 if &protocol_name == self.block_downloader.protocol_name() {
346 let Ok(response) = response
347 .downcast::<(BlockRequest<B>, Result<Vec<BlockData<B>>, BlockResponseError>)>()
348 else {
349 warn!(target: LOG_TARGET, "Failed to downcast block response");
350 debug_assert!(false);
351 return;
352 };
353
354 let (request, response) = *response;
355 let blocks = match response {
356 Ok(blocks) => blocks,
357 Err(BlockResponseError::DecodeFailed(e)) => {
358 debug!(
359 target: LOG_TARGET,
360 "Failed to decode block response from peer {:?}: {:?}.",
361 peer_id,
362 e
363 );
364 self.actions.push(SyncingAction::DropPeer(BadPeer(*peer_id, rep::BAD_MESSAGE)));
365 return;
366 },
367 Err(BlockResponseError::ExtractionFailed(e)) => {
368 debug!(
369 target: LOG_TARGET,
370 "Failed to extract blocks from peer response {:?}: {:?}.",
371 peer_id,
372 e
373 );
374 self.actions.push(SyncingAction::DropPeer(BadPeer(*peer_id, rep::BAD_MESSAGE)));
375 return;
376 },
377 };
378
379 self.on_block_response(*peer_id, request, blocks);
380 } else {
381 let Ok(response) = response.downcast::<Vec<u8>>() else {
382 warn!(target: LOG_TARGET, "Failed to downcast warp sync response");
383 debug_assert!(false);
384 return;
385 };
386
387 self.on_warp_proof_response(peer_id, EncodedProof(*response));
388 }
389 }
390
391 pub fn on_warp_proof_response(&mut self, peer_id: &PeerId, response: EncodedProof) {
393 if let Some(peer) = self.peers.get_mut(peer_id) {
394 peer.state = PeerState::Available;
395 }
396
397 let Phase::WarpProof { verifier } = &mut self.phase else {
398 debug!(target: LOG_TARGET, "Unexpected warp proof response");
399 self.actions
400 .push(SyncingAction::DropPeer(BadPeer(*peer_id, rep::UNEXPECTED_RESPONSE)));
401 return;
402 };
403
404 let proof_to_incoming_block =
405 |(header, justifications): (B::Header, Justifications)| -> IncomingBlock<B> {
406 IncomingBlock {
407 hash: header.hash(),
408 header: Some(header),
409 body: None,
410 indexed_body: None,
411 justifications: Some(justifications),
412 origin: Some(*peer_id),
413 allow_missing_state: true,
416 skip_execution: true,
417 import_existing: false,
419 state: None,
420 }
421 };
422
423 match verifier.verify(&response) {
424 Err(e) => {
425 debug!(target: LOG_TARGET, "Bad warp proof response: {}", e);
426 self.actions
427 .push(SyncingAction::DropPeer(BadPeer(*peer_id, rep::BAD_WARP_PROOF)))
428 },
429 Ok(VerificationResult::Partial(proofs)) => {
430 debug!(target: LOG_TARGET, "Verified partial proof");
431 self.total_proof_bytes += response.0.len() as u64;
432 self.actions.push(SyncingAction::ImportBlocks {
433 origin: BlockOrigin::WarpSync,
434 blocks: proofs.into_iter().map(proof_to_incoming_block).collect(),
435 });
436 },
437 Ok(VerificationResult::Complete(header, proofs)) => {
438 debug!(
439 target: LOG_TARGET,
440 "Verified complete proof. Continuing with target block download: {} ({}).",
441 header.hash(),
442 header.number(),
443 );
444 self.total_proof_bytes += response.0.len() as u64;
445 self.phase = Phase::TargetBlock(header.clone());
446 let incoming_blocks: Vec<_> = proofs
447 .into_iter()
448 .map(proof_to_incoming_block)
449 .filter(|i| {
450 if header.number() != i.header.as_ref().unwrap().number() {
455 true
456 } else {
457 log::trace!(
458 target: LOG_TARGET,
459 "Filtered out target block: {} ({})",
460 header.hash(),
461 header.number()
462 );
463 false
464 }
465 })
466 .collect();
467 self.actions.push(SyncingAction::ImportBlocks {
468 origin: BlockOrigin::WarpSync,
469 blocks: incoming_blocks,
470 });
471 },
472 }
473 }
474
475 pub fn on_block_response(
477 &mut self,
478 peer_id: PeerId,
479 request: BlockRequest<B>,
480 blocks: Vec<BlockData<B>>,
481 ) {
482 if let Err(bad_peer) = self.on_block_response_inner(peer_id, request, blocks) {
483 self.actions.push(SyncingAction::DropPeer(bad_peer));
484 }
485 }
486
487 fn on_block_response_inner(
488 &mut self,
489 peer_id: PeerId,
490 request: BlockRequest<B>,
491 mut blocks: Vec<BlockData<B>>,
492 ) -> Result<(), BadPeer> {
493 if let Some(peer) = self.peers.get_mut(&peer_id) {
494 peer.state = PeerState::Available;
495 }
496
497 let Phase::TargetBlock(header) = &mut self.phase else {
498 debug!(target: LOG_TARGET, "Unexpected target block response from {peer_id}");
499 return Err(BadPeer(peer_id, rep::UNEXPECTED_RESPONSE));
500 };
501
502 if blocks.is_empty() {
503 debug!(
504 target: LOG_TARGET,
505 "Downloading target block failed: empty block response from {peer_id}",
506 );
507 return Err(BadPeer(peer_id, rep::NO_BLOCK));
508 }
509
510 if blocks.len() > 1 {
511 debug!(
512 target: LOG_TARGET,
513 "Too many blocks ({}) in warp target block response from {peer_id}",
514 blocks.len(),
515 );
516 return Err(BadPeer(peer_id, rep::NOT_REQUESTED));
517 }
518
519 validate_blocks::<B>(&blocks, &peer_id, Some(request))?;
520
521 let block = blocks.pop().expect("`blocks` len checked above; qed");
522
523 let Some(block_header) = &block.header else {
524 debug!(
525 target: LOG_TARGET,
526 "Downloading target block failed: missing header in response from {peer_id}.",
527 );
528 return Err(BadPeer(peer_id, rep::VERIFICATION_FAIL));
529 };
530
531 if block_header != header {
532 debug!(
533 target: LOG_TARGET,
534 "Downloading target block failed: different header in response from {peer_id}.",
535 );
536 return Err(BadPeer(peer_id, rep::VERIFICATION_FAIL));
537 }
538
539 if block.body.is_none() {
540 debug!(
541 target: LOG_TARGET,
542 "Downloading target block failed: missing body in response from {peer_id}.",
543 );
544 return Err(BadPeer(peer_id, rep::VERIFICATION_FAIL));
545 }
546
547 self.result = Some(WarpSyncResult {
548 target_header: header.clone(),
549 target_body: block.body,
550 target_justifications: block.justifications,
551 });
552 self.phase = Phase::Complete;
553 self.actions.push(SyncingAction::Finished);
554 Ok(())
555 }
556
557 fn schedule_next_peer(
559 &mut self,
560 new_state: PeerState,
561 min_best_number: Option<NumberFor<B>>,
562 ) -> Option<PeerId> {
563 let mut targets: Vec<_> = self.peers.values().map(|p| p.best_number).collect();
564 if targets.is_empty() {
565 return None;
566 }
567 targets.sort();
568 let median = targets[targets.len() / 2];
569 let threshold = std::cmp::max(median, min_best_number.unwrap_or(Zero::zero()));
570 for (peer_id, peer) in self.peers.iter_mut() {
573 if peer.state.is_available() &&
574 peer.best_number >= threshold &&
575 self.disconnected_peers.is_peer_available(peer_id)
576 {
577 peer.state = new_state;
578 return Some(*peer_id);
579 }
580 }
581 None
582 }
583
584 fn warp_proof_request(&mut self) -> Option<(PeerId, ProtocolName, WarpProofRequest<B>)> {
586 let Phase::WarpProof { verifier } = &self.phase else { return None };
587
588 let begin = verifier.next_proof_context();
590
591 if self
592 .peers
593 .values()
594 .any(|peer| matches!(peer.state, PeerState::DownloadingProofs))
595 {
596 return None;
598 }
599
600 let peer_id = self.schedule_next_peer(PeerState::DownloadingProofs, None)?;
601 trace!(target: LOG_TARGET, "New WarpProofRequest to {peer_id}, begin hash: {begin}.");
602
603 let request = WarpProofRequest { begin };
604
605 let Some(protocol_name) = self.protocol_name.clone() else {
606 warn!(
607 target: LOG_TARGET,
608 "Trying to send warp sync request when no protocol is configured {request:?}",
609 );
610 return None;
611 };
612
613 Some((peer_id, protocol_name, request))
614 }
615
616 fn target_block_request(&mut self) -> Option<(PeerId, BlockRequest<B>)> {
618 let Phase::TargetBlock(target_header) = &self.phase else { return None };
619
620 if self
621 .peers
622 .values()
623 .any(|peer| matches!(peer.state, PeerState::DownloadingTargetBlock))
624 {
625 return None;
627 }
628
629 let target_hash = target_header.hash();
631 let target_number = *target_header.number();
632
633 let peer_id =
634 self.schedule_next_peer(PeerState::DownloadingTargetBlock, Some(target_number))?;
635
636 trace!(
637 target: LOG_TARGET,
638 "New target block request to {peer_id}, target: {} ({}).",
639 target_hash,
640 target_number,
641 );
642
643 Some((
644 peer_id,
645 BlockRequest::<B> {
646 id: 0,
647 fields: BlockAttributes::HEADER |
648 BlockAttributes::BODY |
649 BlockAttributes::JUSTIFICATION,
650 from: FromBlock::Hash(target_hash),
651 direction: Direction::Ascending,
652 max: Some(1),
653 },
654 ))
655 }
656
657 pub fn progress(&self) -> WarpSyncProgress<B> {
659 match &self.phase {
660 Phase::WaitingForPeers { .. } => WarpSyncProgress {
661 phase: WarpSyncPhase::AwaitingPeers {
662 required_peers: self.min_peers_to_start_warp_sync,
663 },
664 total_bytes: self.total_proof_bytes,
665 status: None,
666 },
667 Phase::WarpProof { verifier } => WarpSyncProgress {
668 phase: WarpSyncPhase::DownloadingWarpProofs,
669 total_bytes: self.total_proof_bytes,
670 status: verifier.status(),
671 },
672 Phase::TargetBlock(_) => WarpSyncProgress {
673 phase: WarpSyncPhase::DownloadingTargetBlock,
674 total_bytes: self.total_proof_bytes,
675 status: None,
676 },
677 Phase::Complete => WarpSyncProgress {
678 phase: WarpSyncPhase::Complete,
679 total_bytes: self.total_proof_bytes + self.total_state_bytes,
680 status: None,
681 },
682 }
683 }
684
685 pub fn num_peers(&self) -> usize {
687 self.peers.len()
688 }
689
690 pub fn status(&self) -> SyncStatus<B> {
692 SyncStatus {
693 state: match &self.phase {
694 Phase::WaitingForPeers { .. } => SyncState::Downloading { target: Zero::zero() },
695 Phase::WarpProof { .. } => SyncState::Downloading { target: Zero::zero() },
696 Phase::TargetBlock(header) => SyncState::Downloading { target: *header.number() },
697 Phase::Complete => SyncState::Idle,
698 },
699 best_seen_block: match &self.phase {
700 Phase::WaitingForPeers { .. } => None,
701 Phase::WarpProof { .. } => None,
702 Phase::TargetBlock(header) => Some(*header.number()),
703 Phase::Complete => None,
704 },
705 num_peers: self.peers.len().saturated_into(),
706 queued_blocks: 0,
707 state_sync: None,
708 warp_sync: Some(self.progress()),
709 }
710 }
711
712 #[must_use]
714 pub fn actions(
715 &mut self,
716 network_service: &NetworkServiceHandle,
717 ) -> impl Iterator<Item = SyncingAction<B>> {
718 let warp_proof_request =
719 self.warp_proof_request().into_iter().map(|(peer_id, protocol_name, request)| {
720 trace!(
721 target: LOG_TARGET,
722 "Created `WarpProofRequest` to {}, request: {:?}.",
723 peer_id,
724 request,
725 );
726
727 let (tx, rx) = oneshot::channel();
728
729 network_service.start_request(
730 peer_id,
731 protocol_name,
732 request.encode(),
733 tx,
734 IfDisconnected::ImmediateError,
735 );
736
737 SyncingAction::StartRequest {
738 peer_id,
739 key: Self::STRATEGY_KEY,
740 request: async move {
741 Ok(rx.await?.and_then(|(response, protocol_name)| {
742 Ok((Box::new(response) as Box<dyn Any + Send>, protocol_name))
743 }))
744 }
745 .boxed(),
746 }
747 });
748 self.actions.extend(warp_proof_request);
749
750 let target_block_request =
751 self.target_block_request().into_iter().map(|(peer_id, request)| {
752 let downloader = self.block_downloader.clone();
753
754 SyncingAction::StartRequest {
755 peer_id,
756 key: Self::STRATEGY_KEY,
757 request: async move {
758 Ok(downloader.download_blocks(peer_id, request.clone()).await?.and_then(
759 |(response, protocol_name)| {
760 let decoded_response =
761 downloader.block_response_into_blocks(&request, response);
762 let result =
763 Box::new((request, decoded_response)) as Box<dyn Any + Send>;
764 Ok((result, protocol_name))
765 },
766 ))
767 }
768 .boxed(),
769 }
770 });
771 self.actions.extend(target_block_request);
772
773 std::mem::take(&mut self.actions).into_iter()
774 }
775
776 #[must_use]
778 pub fn take_result(&mut self) -> Option<WarpSyncResult<B>> {
779 self.result.take()
780 }
781}
782
783#[cfg(test)]
784mod test {
785 use super::*;
786 use crate::{mock::MockBlockDownloader, service::network::NetworkServiceProvider};
787 use sc_block_builder::BlockBuilderBuilder;
788 use sp_blockchain::{BlockStatus, Error as BlockchainError, HeaderBackend, Info};
789 use sp_core::H256;
790 use sp_runtime::{
791 traits::{Block as BlockT, Header as HeaderT, NumberFor},
792 ConsensusEngineId,
793 };
794 use std::{io::ErrorKind, sync::Arc};
795 use substrate_test_runtime_client::{
796 runtime::{Block, Hash},
797 BlockBuilderExt, DefaultTestClientBuilderExt, TestClientBuilder, TestClientBuilderExt,
798 };
799
800 pub const TEST_ENGINE_ID: ConsensusEngineId = *b"TEST";
801
802 mockall::mock! {
803 pub Client<B: BlockT> {}
804
805 impl<B: BlockT> HeaderBackend<B> for Client<B> {
806 fn header(&self, hash: B::Hash) -> Result<Option<B::Header>, BlockchainError>;
807 fn info(&self) -> Info<B>;
808 fn status(&self, hash: B::Hash) -> Result<BlockStatus, BlockchainError>;
809 fn number(
810 &self,
811 hash: B::Hash,
812 ) -> Result<Option<<<B as BlockT>::Header as HeaderT>::Number>, BlockchainError>;
813 fn hash(&self, number: NumberFor<B>) -> Result<Option<B::Hash>, BlockchainError>;
814 }
815 }
816
817 mockall::mock! {
818 pub WarpSyncProvider<B: BlockT> {}
819
820 impl<B: BlockT> super::WarpSyncProvider<B> for WarpSyncProvider<B> {
821 fn generate(
822 &self,
823 start: B::Hash,
824 ) -> Result<EncodedProof, Box<dyn std::error::Error + Send + Sync>>;
825 fn create_verifier(&self) -> Box<dyn super::Verifier<B>>;
826 }
827 }
828
829 mockall::mock! {
830 pub Verifier<B: BlockT> {}
831
832 impl<B: BlockT> super::Verifier<B> for Verifier<B> {
833 fn verify(
834 &mut self,
835 proof: &EncodedProof,
836 ) -> Result<VerificationResult<B>, Box<dyn std::error::Error + Send + Sync>>;
837 fn next_proof_context(&self) -> B::Hash;
838 fn status(&self) -> Option<String>;
839 }
840 }
841
842 fn mock_client_with_state() -> MockClient<Block> {
843 let mut client = MockClient::<Block>::new();
844 let genesis_hash = Hash::random();
845 client.expect_info().return_once(move || Info {
846 best_hash: genesis_hash,
847 best_number: 0,
848 genesis_hash,
849 finalized_hash: genesis_hash,
850 finalized_number: 0,
851 finalized_state: Some((genesis_hash, 0)),
853 number_leaves: 0,
854 block_gap: None,
855 });
856
857 client
858 }
859
860 fn mock_client_without_state() -> MockClient<Block> {
861 let mut client = MockClient::<Block>::new();
862 let genesis_hash = Hash::random();
863 client.expect_info().returning(move || Info {
864 best_hash: genesis_hash,
865 best_number: 0,
866 genesis_hash,
867 finalized_hash: genesis_hash,
868 finalized_number: 0,
869 finalized_state: None,
870 number_leaves: 0,
871 block_gap: None,
872 });
873
874 client
875 }
876
877 #[test]
878 fn warp_sync_with_provider_for_db_with_finalized_state_is_noop() {
879 let client = mock_client_with_state();
880 let provider = MockWarpSyncProvider::<Block>::new();
881 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
882 let mut warp_sync = WarpSync::new(
883 Arc::new(client),
884 config,
885 None,
886 Arc::new(MockBlockDownloader::new()),
887 None,
888 );
889
890 let network_provider = NetworkServiceProvider::new();
891 let network_handle = network_provider.handle();
892
893 let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
895 assert_eq!(actions.len(), 1);
896 assert!(matches!(actions[0], SyncingAction::Finished));
897
898 assert!(warp_sync.take_result().is_none());
900 }
901
902 #[test]
903 fn warp_sync_to_target_for_db_with_finalized_state_is_noop() {
904 let client = mock_client_with_state();
905 let config = WarpSyncConfig::WithTarget(<Block as BlockT>::Header::new(
906 1,
907 Default::default(),
908 Default::default(),
909 Default::default(),
910 Default::default(),
911 ));
912 let mut warp_sync = WarpSync::new(
913 Arc::new(client),
914 config,
915 None,
916 Arc::new(MockBlockDownloader::new()),
917 None,
918 );
919
920 let network_provider = NetworkServiceProvider::new();
921 let network_handle = network_provider.handle();
922
923 let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
925 assert_eq!(actions.len(), 1);
926 assert!(matches!(actions[0], SyncingAction::Finished));
927
928 assert!(warp_sync.take_result().is_none());
930 }
931
932 #[test]
933 fn warp_sync_with_provider_for_empty_db_doesnt_finish_instantly() {
934 let client = mock_client_without_state();
935 let provider = MockWarpSyncProvider::<Block>::new();
936 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
937 let mut warp_sync = WarpSync::new(
938 Arc::new(client),
939 config,
940 None,
941 Arc::new(MockBlockDownloader::new()),
942 None,
943 );
944
945 let network_provider = NetworkServiceProvider::new();
946 let network_handle = network_provider.handle();
947
948 assert_eq!(warp_sync.actions(&network_handle).count(), 0)
950 }
951
952 #[test]
953 fn warp_sync_to_target_for_empty_db_doesnt_finish_instantly() {
954 let client = mock_client_without_state();
955 let config = WarpSyncConfig::WithTarget(<Block as BlockT>::Header::new(
956 1,
957 Default::default(),
958 Default::default(),
959 Default::default(),
960 Default::default(),
961 ));
962 let mut warp_sync = WarpSync::new(
963 Arc::new(client),
964 config,
965 None,
966 Arc::new(MockBlockDownloader::new()),
967 None,
968 );
969
970 let network_provider = NetworkServiceProvider::new();
971 let network_handle = network_provider.handle();
972
973 assert_eq!(warp_sync.actions(&network_handle).count(), 0)
975 }
976
977 #[test]
978 fn warp_sync_is_started_only_when_there_is_enough_peers() {
979 let client = mock_client_without_state();
980 let mut provider = MockWarpSyncProvider::<Block>::new();
981 let mut verifier = MockVerifier::<Block>::new();
982 verifier.expect_next_proof_context().returning(|| Hash::random());
983 verifier
984 .expect_verify()
985 .returning(|_| unreachable!("verify should not be called in this test"));
986 provider.expect_create_verifier().return_once(move || Box::new(verifier));
987 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
988 let mut warp_sync = WarpSync::new(
989 Arc::new(client),
990 config,
991 None,
992 Arc::new(MockBlockDownloader::new()),
993 None,
994 );
995
996 for _ in 0..(MIN_PEERS_TO_START_WARP_SYNC - 1) {
998 warp_sync.add_peer(PeerId::random(), Hash::random(), 10);
999 assert!(matches!(warp_sync.phase, Phase::WaitingForPeers { .. }))
1000 }
1001
1002 warp_sync.add_peer(PeerId::random(), Hash::random(), 10);
1004 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }))
1005 }
1006
1007 #[test]
1008 fn no_peer_is_scheduled_if_no_peers_connected() {
1009 let client = mock_client_without_state();
1010 let provider = MockWarpSyncProvider::<Block>::new();
1011 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1012 let mut warp_sync = WarpSync::new(
1013 Arc::new(client),
1014 config,
1015 None,
1016 Arc::new(MockBlockDownloader::new()),
1017 None,
1018 );
1019
1020 assert!(warp_sync.schedule_next_peer(PeerState::DownloadingProofs, None).is_none());
1021 }
1022
1023 #[test]
1024 fn enough_peers_are_used_in_tests() {
1025 assert!(
1027 10 >= MIN_PEERS_TO_START_WARP_SYNC,
1028 "Tests must be updated to use that many initial peers.",
1029 );
1030 }
1031
1032 #[test]
1033 fn at_least_median_synced_peer_is_scheduled() {
1034 for _ in 0..100 {
1035 let client = mock_client_without_state();
1036 let mut provider = MockWarpSyncProvider::<Block>::new();
1037 let mut verifier = MockVerifier::<Block>::new();
1038 verifier.expect_next_proof_context().returning(|| Hash::random());
1039 verifier
1040 .expect_verify()
1041 .returning(|_| unreachable!("verify should not be called in this test"));
1042 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1043 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1044 let mut warp_sync = WarpSync::new(
1045 Arc::new(client),
1046 config,
1047 None,
1048 Arc::new(MockBlockDownloader::new()),
1049 None,
1050 );
1051
1052 for best_number in 1..11 {
1053 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1054 }
1055
1056 let peer_id = warp_sync.schedule_next_peer(PeerState::DownloadingProofs, None);
1057 assert!(warp_sync.peers.get(&peer_id.unwrap()).unwrap().best_number >= 6);
1058 }
1059 }
1060
1061 #[test]
1062 fn min_best_number_peer_is_scheduled() {
1063 for _ in 0..10 {
1064 let client = mock_client_without_state();
1065 let mut provider = MockWarpSyncProvider::<Block>::new();
1066 let mut verifier = MockVerifier::<Block>::new();
1067 verifier.expect_next_proof_context().returning(|| Hash::random());
1068 verifier
1069 .expect_verify()
1070 .returning(|_| unreachable!("verify should not be called in this test"));
1071 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1072 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1073 let mut warp_sync = WarpSync::new(
1074 Arc::new(client),
1075 config,
1076 None,
1077 Arc::new(MockBlockDownloader::new()),
1078 None,
1079 );
1080
1081 for best_number in 1..11 {
1082 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1083 }
1084
1085 let peer_id = warp_sync.schedule_next_peer(PeerState::DownloadingProofs, Some(10));
1086 assert!(warp_sync.peers.get(&peer_id.unwrap()).unwrap().best_number == 10);
1087 }
1088 }
1089
1090 #[test]
1091 fn backedoff_number_peer_is_not_scheduled() {
1092 let client = mock_client_without_state();
1093 let mut provider = MockWarpSyncProvider::<Block>::new();
1094 let mut verifier = MockVerifier::<Block>::new();
1095 verifier.expect_next_proof_context().returning(|| Hash::random());
1096 verifier
1097 .expect_verify()
1098 .returning(|_| unreachable!("verify should not be called in this test"));
1099 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1100 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1101 let mut warp_sync = WarpSync::new(
1102 Arc::new(client),
1103 config,
1104 None,
1105 Arc::new(MockBlockDownloader::new()),
1106 None,
1107 );
1108
1109 for best_number in 1..11 {
1110 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1111 }
1112
1113 let ninth_peer =
1114 *warp_sync.peers.iter().find(|(_, state)| state.best_number == 9).unwrap().0;
1115 let tenth_peer =
1116 *warp_sync.peers.iter().find(|(_, state)| state.best_number == 10).unwrap().0;
1117
1118 warp_sync.remove_peer(&tenth_peer);
1120 assert!(warp_sync.disconnected_peers.is_peer_available(&tenth_peer));
1121
1122 warp_sync.add_peer(tenth_peer, H256::random(), 10);
1123 let peer_id = warp_sync.schedule_next_peer(PeerState::DownloadingProofs, Some(10));
1124 assert_eq!(tenth_peer, peer_id.unwrap());
1125 warp_sync.remove_peer(&tenth_peer);
1126
1127 assert!(!warp_sync.disconnected_peers.is_peer_available(&tenth_peer));
1129
1130 warp_sync.add_peer(tenth_peer, H256::random(), 10);
1132 let peer_id: Option<PeerId> =
1133 warp_sync.schedule_next_peer(PeerState::DownloadingProofs, Some(10));
1134 assert!(peer_id.is_none());
1135
1136 let peer_id: Option<PeerId> =
1138 warp_sync.schedule_next_peer(PeerState::DownloadingProofs, Some(9));
1139 assert_eq!(ninth_peer, peer_id.unwrap());
1140 }
1141
1142 #[test]
1143 fn no_warp_proof_request_in_another_phase() {
1144 let client = mock_client_without_state();
1145 let mut provider = MockWarpSyncProvider::<Block>::new();
1146 let mut verifier = MockVerifier::<Block>::new();
1147 verifier.expect_next_proof_context().returning(|| Hash::random());
1148 verifier
1149 .expect_verify()
1150 .returning(|_| unreachable!("verify should not be called in this test"));
1151 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1152 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1153 let mut warp_sync = WarpSync::new(
1154 Arc::new(client),
1155 config,
1156 Some(ProtocolName::Static("")),
1157 Arc::new(MockBlockDownloader::new()),
1158 None,
1159 );
1160
1161 for best_number in 1..11 {
1163 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1164 }
1165
1166 warp_sync.phase = Phase::TargetBlock(<Block as BlockT>::Header::new(
1168 1,
1169 Default::default(),
1170 Default::default(),
1171 Default::default(),
1172 Default::default(),
1173 ));
1174
1175 assert!(warp_sync.warp_proof_request().is_none());
1177 }
1178
1179 #[test]
1180 fn warp_proof_request_starts_at_last_hash() {
1181 let client = mock_client_without_state();
1182 let mut provider = MockWarpSyncProvider::<Block>::new();
1183 let mut verifier = MockVerifier::<Block>::new();
1184 let known_last_hash = Hash::random();
1185 verifier.expect_next_proof_context().returning(move || known_last_hash);
1186 verifier
1187 .expect_verify()
1188 .returning(|_| unreachable!("verify should not be called in this test"));
1189 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1190 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1191 let mut warp_sync = WarpSync::new(
1192 Arc::new(client),
1193 config,
1194 Some(ProtocolName::Static("")),
1195 Arc::new(MockBlockDownloader::new()),
1196 None,
1197 );
1198
1199 for best_number in 1..11 {
1201 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1202 }
1203 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1204
1205 let (_peer_id, _protocol_name, request) = warp_sync.warp_proof_request().unwrap();
1206 assert_eq!(request.begin, known_last_hash);
1207 }
1208
1209 #[test]
1210 fn no_parallel_warp_proof_requests() {
1211 let client = mock_client_without_state();
1212 let mut provider = MockWarpSyncProvider::<Block>::new();
1213 let mut verifier = MockVerifier::<Block>::new();
1214 verifier.expect_next_proof_context().returning(|| Hash::random());
1215 verifier
1216 .expect_verify()
1217 .returning(|_| unreachable!("verify should not be called in this test"));
1218 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1219 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1220 let mut warp_sync = WarpSync::new(
1221 Arc::new(client),
1222 config,
1223 Some(ProtocolName::Static("")),
1224 Arc::new(MockBlockDownloader::new()),
1225 None,
1226 );
1227
1228 for best_number in 1..11 {
1230 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1231 }
1232 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1233
1234 assert!(warp_sync.warp_proof_request().is_some());
1236 assert!(warp_sync.warp_proof_request().is_none());
1238 }
1239
1240 #[test]
1241 fn bad_warp_proof_response_drops_peer() {
1242 let client = mock_client_without_state();
1243 let mut provider = MockWarpSyncProvider::<Block>::new();
1244 let mut verifier = MockVerifier::<Block>::new();
1245 verifier.expect_next_proof_context().returning(|| Hash::random());
1246 verifier.expect_verify().return_once(|_proof| {
1248 Err(Box::new(std::io::Error::new(ErrorKind::Other, "test-verification-failure")))
1249 });
1250 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1251 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1252 let mut warp_sync = WarpSync::new(
1253 Arc::new(client),
1254 config,
1255 Some(ProtocolName::Static("")),
1256 Arc::new(MockBlockDownloader::new()),
1257 None,
1258 );
1259
1260 for best_number in 1..11 {
1262 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1263 }
1264 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1265
1266 let network_provider = NetworkServiceProvider::new();
1267 let network_handle = network_provider.handle();
1268
1269 let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1271 assert_eq!(actions.len(), 1);
1272 let SyncingAction::StartRequest { peer_id: request_peer_id, .. } = actions[0] else {
1273 panic!("Invalid action");
1274 };
1275
1276 warp_sync.on_warp_proof_response(&request_peer_id, EncodedProof(Vec::new()));
1277
1278 let actions = std::mem::take(&mut warp_sync.actions);
1280 assert_eq!(actions.len(), 1);
1281 assert!(matches!(
1282 actions[0],
1283 SyncingAction::DropPeer(BadPeer(peer_id, _rep)) if peer_id == request_peer_id
1284 ));
1285 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1286 }
1287
1288 #[test]
1289 fn partial_warp_proof_doesnt_advance_phase() {
1290 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1291 let mut provider = MockWarpSyncProvider::<Block>::new();
1292 let target_block = BlockBuilderBuilder::new(&*client)
1293 .on_parent_block(client.chain_info().best_hash)
1294 .with_parent_block_number(client.chain_info().best_number)
1295 .build()
1296 .unwrap()
1297 .build()
1298 .unwrap()
1299 .block;
1300 let target_header = target_block.header().clone();
1301 let justifications = Justifications::new(vec![(TEST_ENGINE_ID, vec![1, 2, 3, 4, 5])]);
1302 let mut verifier = MockVerifier::<Block>::new();
1304 let context = client.info().genesis_hash;
1305 verifier.expect_next_proof_context().returning(move || context);
1306 let header_for_verify = target_header.clone();
1307 let just_for_verify = justifications.clone();
1308 verifier.expect_verify().return_once(move |_proof| {
1309 Ok(VerificationResult::Partial(vec![(
1310 header_for_verify.clone(),
1311 just_for_verify.clone(),
1312 )]))
1313 });
1314 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1315 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1316 let mut warp_sync = WarpSync::new(
1317 client,
1318 config,
1319 Some(ProtocolName::Static("")),
1320 Arc::new(MockBlockDownloader::new()),
1321 None,
1322 );
1323
1324 for best_number in 1..11 {
1326 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1327 }
1328 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1329
1330 let network_provider = NetworkServiceProvider::new();
1331 let network_handle = network_provider.handle();
1332
1333 let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1335 assert_eq!(actions.len(), 1);
1336 let SyncingAction::StartRequest { peer_id: request_peer_id, .. } = actions[0] else {
1337 panic!("Invalid action");
1338 };
1339
1340 warp_sync.on_warp_proof_response(&request_peer_id, EncodedProof(Vec::new()));
1341
1342 assert_eq!(warp_sync.actions.len(), 1);
1343 let SyncingAction::ImportBlocks { origin, mut blocks } = warp_sync.actions.pop().unwrap()
1344 else {
1345 panic!("Expected `ImportBlocks` action.");
1346 };
1347 assert_eq!(origin, BlockOrigin::WarpSync);
1348 assert_eq!(blocks.len(), 1);
1349 let import_block = blocks.pop().unwrap();
1350 assert_eq!(
1351 import_block,
1352 IncomingBlock {
1353 hash: target_header.hash(),
1354 header: Some(target_header),
1355 body: None,
1356 indexed_body: None,
1357 justifications: Some(justifications),
1358 origin: Some(request_peer_id),
1359 allow_missing_state: true,
1360 skip_execution: true,
1361 import_existing: false,
1362 state: None,
1363 }
1364 );
1365 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1366 }
1367
1368 #[test]
1369 fn complete_warp_proof_advances_phase() {
1370 use crate::{pending_responses::PendingResponses, service::network::ToServiceCommand};
1371 use futures::{executor::block_on, StreamExt};
1372
1373 sp_tracing::try_init_simple();
1375
1376 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1377 let mut provider = MockWarpSyncProvider::<Block>::new();
1378
1379 let warp_synced_header = <<Block as BlockT>::Header as HeaderT>::new(
1383 1,
1384 Default::default(),
1385 Default::default(),
1386 client.chain_info().best_hash,
1387 Default::default(),
1388 );
1389
1390 let target_header = <<Block as BlockT>::Header as HeaderT>::new(
1392 2,
1393 Default::default(),
1394 Default::default(),
1395 warp_synced_header.hash(),
1396 Default::default(),
1397 );
1398 let warp_justifications = Justifications::new(vec![(TEST_ENGINE_ID, vec![1, 2, 3, 4, 5])]);
1399 let target_justifications =
1400 Justifications::new(vec![(TEST_ENGINE_ID, vec![6, 7, 8, 9, 10])]);
1401
1402 let mut verifier = MockVerifier::<Block>::new();
1404 let context = client.info().genesis_hash;
1405 verifier.expect_next_proof_context().returning(move || context);
1406 let warp_synced_header_for_verify = warp_synced_header.clone();
1407 let warp_just_for_verify = warp_justifications.clone();
1408 let target_header_for_verify = target_header.clone();
1409 let target_just_for_verify = target_justifications.clone();
1410 verifier.expect_verify().return_once(move |_proof| {
1411 Ok(VerificationResult::Complete(
1412 target_header_for_verify.clone(),
1413 vec![
1414 (warp_synced_header_for_verify, warp_just_for_verify),
1415 (target_header_for_verify, target_just_for_verify),
1416 ],
1417 ))
1418 });
1419 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1420 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1421 let mut warp_sync = WarpSync::new(
1422 client,
1423 config,
1424 Some(ProtocolName::Static("")),
1425 Arc::new(MockBlockDownloader::new()),
1426 Some(1),
1427 );
1428
1429 let peer_id = PeerId::random();
1431 warp_sync.add_peer(peer_id, Hash::random(), 10);
1432 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1433
1434 let (tx, mut rx) = sc_utils::mpsc::tracing_unbounded("warp_test", 10);
1435 let network_handle = NetworkServiceHandle::new(tx);
1436 let mut pending = PendingResponses::new();
1437 let mut actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1438 assert_eq!(actions.len(), 1);
1439 let SyncingAction::StartRequest { peer_id: request_peer_id, key, request } =
1440 actions.pop().unwrap()
1441 else {
1442 panic!("Expected warp proof request.");
1443 };
1444 assert_eq!(request_peer_id, peer_id);
1445 pending.insert(request_peer_id, key, request);
1446 let ToServiceCommand::StartRequest(_, protocol, _, response_tx, _) =
1447 block_on(rx.next()).unwrap()
1448 else {
1449 panic!("Expected network request.");
1450 };
1451 response_tx.send(Ok((Vec::new(), protocol))).unwrap();
1452 let event = block_on(pending.next()).unwrap();
1453 assert_eq!(pending.len(), 0);
1455 let (response, _) = event.response.unwrap().unwrap();
1456 warp_sync.on_warp_proof_response(
1457 &event.peer_id,
1458 EncodedProof(*response.downcast::<Vec<u8>>().unwrap()),
1459 );
1460
1461 assert_eq!(warp_sync.actions.len(), 1);
1462 let SyncingAction::ImportBlocks { origin, mut blocks } = warp_sync.actions.pop().unwrap()
1463 else {
1464 panic!("Expected `ImportBlocks` action.");
1465 };
1466 assert_eq!(origin, BlockOrigin::WarpSync);
1467 assert_eq!(blocks.len(), 1);
1469 let import_block = blocks.pop().unwrap();
1470 assert_eq!(
1471 import_block,
1472 IncomingBlock {
1473 hash: warp_synced_header.hash(),
1474 header: Some(warp_synced_header),
1475 body: None,
1476 indexed_body: None,
1477 justifications: Some(warp_justifications),
1478 origin: Some(request_peer_id),
1479 allow_missing_state: true,
1480 skip_execution: true,
1481 import_existing: false,
1482 state: None,
1483 }
1484 );
1485 assert!(matches!(&warp_sync.phase, Phase::TargetBlock(header) if *header == target_header));
1486
1487 let mut actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1488 assert_eq!(actions.len(), 1);
1489 let SyncingAction::StartRequest { peer_id: target_peer, key, request } =
1490 actions.pop().unwrap()
1491 else {
1492 panic!("Expected target block request.");
1493 };
1494 assert_eq!(target_peer, request_peer_id);
1495 pending.insert(target_peer, key, request);
1497 assert_eq!(pending.len(), 1);
1498 assert_eq!(warp_sync.actions(&network_handle).count(), 0);
1499 }
1500
1501 #[test]
1502 fn no_target_block_requests_in_another_phase() {
1503 let client = mock_client_without_state();
1504 let mut provider = MockWarpSyncProvider::<Block>::new();
1505 let mut verifier = MockVerifier::<Block>::new();
1506 verifier.expect_next_proof_context().returning(|| Hash::random());
1507 verifier
1508 .expect_verify()
1509 .returning(|_| unreachable!("verify should not be called in this test"));
1510 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1511 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1512 let mut warp_sync = WarpSync::new(
1513 Arc::new(client),
1514 config,
1515 None,
1516 Arc::new(MockBlockDownloader::new()),
1517 None,
1518 );
1519
1520 for best_number in 1..11 {
1522 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1523 }
1524 assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1526
1527 assert!(warp_sync.target_block_request().is_none());
1529 }
1530
1531 #[test]
1532 fn target_block_request_is_correct() {
1533 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1534 let mut provider = MockWarpSyncProvider::<Block>::new();
1535 let mut verifier = MockVerifier::<Block>::new();
1536 let header_for_ctx = client.info().genesis_hash;
1537 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1538 let target_block = BlockBuilderBuilder::new(&*client)
1539 .on_parent_block(client.chain_info().best_hash)
1540 .with_parent_block_number(client.chain_info().best_number)
1541 .build()
1542 .unwrap()
1543 .build()
1544 .unwrap()
1545 .block;
1546 let target_header = target_block.header().clone();
1547 let header_for_verify = target_header.clone();
1549 verifier.expect_verify().return_once(move |_proof| {
1550 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1551 });
1552 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1553 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1554 let mut warp_sync =
1555 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1556
1557 for best_number in 1..11 {
1559 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1560 }
1561
1562 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1564
1565 let (_peer_id, request) = warp_sync.target_block_request().unwrap();
1566 assert_eq!(request.from, FromBlock::Hash(target_block.header().hash()));
1567 assert_eq!(
1568 request.fields,
1569 BlockAttributes::HEADER | BlockAttributes::BODY | BlockAttributes::JUSTIFICATION
1570 );
1571 assert_eq!(request.max, Some(1));
1572 }
1573
1574 #[test]
1575 fn externally_set_target_block_is_requested() {
1576 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1577 let target_block = BlockBuilderBuilder::new(&*client)
1578 .on_parent_block(client.chain_info().best_hash)
1579 .with_parent_block_number(client.chain_info().best_number)
1580 .build()
1581 .unwrap()
1582 .build()
1583 .unwrap()
1584 .block;
1585 let target_header = target_block.header().clone();
1586 let config = WarpSyncConfig::WithTarget(target_header);
1587 let mut warp_sync =
1588 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1589
1590 for best_number in 1..11 {
1592 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1593 }
1594
1595 assert!(matches!(warp_sync.phase, Phase::TargetBlock(_)));
1596
1597 let (_peer_id, request) = warp_sync.target_block_request().unwrap();
1598 assert_eq!(request.from, FromBlock::Hash(target_block.header().hash()));
1599 assert_eq!(
1600 request.fields,
1601 BlockAttributes::HEADER | BlockAttributes::BODY | BlockAttributes::JUSTIFICATION
1602 );
1603 assert_eq!(request.max, Some(1));
1604 }
1605
1606 #[test]
1607 fn no_parallel_target_block_requests() {
1608 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1609 let mut provider = MockWarpSyncProvider::<Block>::new();
1610 let mut verifier = MockVerifier::<Block>::new();
1611 let header_for_ctx = client.info().genesis_hash;
1612 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1613 let target_block = BlockBuilderBuilder::new(&*client)
1614 .on_parent_block(client.chain_info().best_hash)
1615 .with_parent_block_number(client.chain_info().best_number)
1616 .build()
1617 .unwrap()
1618 .build()
1619 .unwrap()
1620 .block;
1621 let target_header = target_block.header().clone();
1622 let header_for_verify = target_header.clone();
1624 verifier.expect_verify().return_once(move |_proof| {
1625 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1626 });
1627 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1628 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1629 let mut warp_sync =
1630 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1631
1632 for best_number in 1..11 {
1634 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1635 }
1636
1637 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1639
1640 assert!(warp_sync.target_block_request().is_some());
1642 assert!(warp_sync.target_block_request().is_none());
1644 }
1645
1646 #[test]
1647 fn target_block_response_with_no_blocks_drops_peer() {
1648 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1649 let mut provider = MockWarpSyncProvider::<Block>::new();
1650 let mut verifier = MockVerifier::<Block>::new();
1651 let header_for_ctx = client.info().genesis_hash;
1652 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1653 let target_block = BlockBuilderBuilder::new(&*client)
1654 .on_parent_block(client.chain_info().best_hash)
1655 .with_parent_block_number(client.chain_info().best_number)
1656 .build()
1657 .unwrap()
1658 .build()
1659 .unwrap()
1660 .block;
1661 let target_header = target_block.header().clone();
1662 let header_for_verify = target_header.clone();
1664 verifier.expect_verify().return_once(move |_proof| {
1665 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1666 });
1667 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1668 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1669 let mut warp_sync =
1670 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1671
1672 for best_number in 1..11 {
1674 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1675 }
1676
1677 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1679
1680 let (peer_id, request) = warp_sync.target_block_request().unwrap();
1681
1682 let response = Vec::new();
1684 assert!(matches!(
1686 warp_sync.on_block_response_inner(peer_id, request, response),
1687 Err(BadPeer(id, _rep)) if id == peer_id,
1688 ));
1689 }
1690
1691 #[test]
1692 fn target_block_response_with_extra_blocks_drops_peer() {
1693 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1694 let mut provider = MockWarpSyncProvider::<Block>::new();
1695 let mut verifier = MockVerifier::<Block>::new();
1696 let header_for_ctx = client.info().genesis_hash;
1697 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1698 let target_block = BlockBuilderBuilder::new(&*client)
1699 .on_parent_block(client.chain_info().best_hash)
1700 .with_parent_block_number(client.chain_info().best_number)
1701 .build()
1702 .unwrap()
1703 .build()
1704 .unwrap()
1705 .block;
1706
1707 let mut extra_block_builder = BlockBuilderBuilder::new(&*client)
1708 .on_parent_block(client.chain_info().best_hash)
1709 .with_parent_block_number(client.chain_info().best_number)
1710 .build()
1711 .unwrap();
1712 extra_block_builder
1713 .push_storage_change(vec![1, 2, 3], Some(vec![4, 5, 6]))
1714 .unwrap();
1715 let extra_block = extra_block_builder.build().unwrap().block;
1716
1717 let target_header = target_block.header().clone();
1718 let header_for_verify = target_header.clone();
1720 verifier.expect_verify().return_once(move |_proof| {
1721 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1722 });
1723 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1724 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1725 let mut warp_sync =
1726 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1727
1728 for best_number in 1..11 {
1730 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1731 }
1732
1733 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1735
1736 let (peer_id, request) = warp_sync.target_block_request().unwrap();
1737
1738 let response = vec![
1740 BlockData::<Block> {
1741 hash: target_block.header().hash(),
1742 header: Some(target_block.header().clone()),
1743 body: Some(target_block.extrinsics().iter().cloned().collect::<Vec<_>>()),
1744 indexed_body: None,
1745 receipt: None,
1746 message_queue: None,
1747 justification: None,
1748 justifications: None,
1749 },
1750 BlockData::<Block> {
1751 hash: extra_block.header().hash(),
1752 header: Some(extra_block.header().clone()),
1753 body: Some(extra_block.extrinsics().iter().cloned().collect::<Vec<_>>()),
1754 indexed_body: None,
1755 receipt: None,
1756 message_queue: None,
1757 justification: None,
1758 justifications: None,
1759 },
1760 ];
1761 assert!(matches!(
1763 warp_sync.on_block_response_inner(peer_id, request, response),
1764 Err(BadPeer(id, _rep)) if id == peer_id,
1765 ));
1766 }
1767
1768 #[test]
1769 fn target_block_response_with_wrong_block_drops_peer() {
1770 sp_tracing::try_init_simple();
1771
1772 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1773 let mut provider = MockWarpSyncProvider::<Block>::new();
1774 let mut verifier = MockVerifier::<Block>::new();
1775 let header_for_ctx = client.info().genesis_hash;
1776 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1777 let target_block = BlockBuilderBuilder::new(&*client)
1778 .on_parent_block(client.chain_info().best_hash)
1779 .with_parent_block_number(client.chain_info().best_number)
1780 .build()
1781 .unwrap()
1782 .build()
1783 .unwrap()
1784 .block;
1785
1786 let mut wrong_block_builder = BlockBuilderBuilder::new(&*client)
1787 .on_parent_block(client.chain_info().best_hash)
1788 .with_parent_block_number(client.chain_info().best_number)
1789 .build()
1790 .unwrap();
1791 wrong_block_builder
1792 .push_storage_change(vec![1, 2, 3], Some(vec![4, 5, 6]))
1793 .unwrap();
1794 let wrong_block = wrong_block_builder.build().unwrap().block;
1795
1796 let target_header = target_block.header().clone();
1797 let header_for_verify = target_header.clone();
1799 verifier.expect_verify().return_once(move |_proof| {
1800 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1801 });
1802 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1803 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1804 let mut warp_sync =
1805 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1806
1807 for best_number in 1..11 {
1809 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1810 }
1811
1812 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1814
1815 let (peer_id, request) = warp_sync.target_block_request().unwrap();
1816
1817 let response = vec![BlockData::<Block> {
1819 hash: wrong_block.header().hash(),
1820 header: Some(wrong_block.header().clone()),
1821 body: Some(wrong_block.extrinsics().iter().cloned().collect::<Vec<_>>()),
1822 indexed_body: None,
1823 receipt: None,
1824 message_queue: None,
1825 justification: None,
1826 justifications: None,
1827 }];
1828 assert!(matches!(
1830 warp_sync.on_block_response_inner(peer_id, request, response),
1831 Err(BadPeer(id, _rep)) if id == peer_id,
1832 ));
1833 }
1834
1835 #[test]
1836 fn correct_target_block_response_sets_strategy_result() {
1837 let client = Arc::new(TestClientBuilder::new().set_no_genesis().build());
1838 let mut provider = MockWarpSyncProvider::<Block>::new();
1839 let mut verifier = MockVerifier::<Block>::new();
1840 let header_for_ctx = client.info().genesis_hash;
1841 verifier.expect_next_proof_context().returning(move || header_for_ctx);
1842 let mut target_block_builder = BlockBuilderBuilder::new(&*client)
1843 .on_parent_block(client.chain_info().best_hash)
1844 .with_parent_block_number(client.chain_info().best_number)
1845 .build()
1846 .unwrap();
1847 target_block_builder
1848 .push_storage_change(vec![1, 2, 3], Some(vec![4, 5, 6]))
1849 .unwrap();
1850 let target_block = target_block_builder.build().unwrap().block;
1851 let target_header = target_block.header().clone();
1852 let header_for_verify = target_header.clone();
1854 verifier.expect_verify().return_once(move |_proof| {
1855 Ok(VerificationResult::Complete(header_for_verify, Default::default()))
1856 });
1857 provider.expect_create_verifier().return_once(move || Box::new(verifier));
1858 let config = WarpSyncConfig::WithProvider(Arc::new(provider));
1859 let mut warp_sync =
1860 WarpSync::new(client, config, None, Arc::new(MockBlockDownloader::new()), None);
1861
1862 for best_number in 1..11 {
1864 warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1865 }
1866
1867 warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1869
1870 let (peer_id, request) = warp_sync.target_block_request().unwrap();
1871
1872 let body = Some(target_block.extrinsics().iter().cloned().collect::<Vec<_>>());
1874 let justifications = Some(Justifications::from((*b"FRNK", Vec::new())));
1875 let response = vec![BlockData::<Block> {
1876 hash: target_block.header().hash(),
1877 header: Some(target_block.header().clone()),
1878 body: body.clone(),
1879 indexed_body: None,
1880 receipt: None,
1881 message_queue: None,
1882 justification: None,
1883 justifications: justifications.clone(),
1884 }];
1885
1886 assert!(warp_sync.on_block_response_inner(peer_id, request, response).is_ok());
1887
1888 let network_provider = NetworkServiceProvider::new();
1889 let network_handle = network_provider.handle();
1890
1891 let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1893 assert_eq!(actions.len(), 1);
1894 assert!(matches!(actions[0], SyncingAction::Finished));
1895
1896 let result = warp_sync.take_result().unwrap();
1898 assert_eq!(result.target_header, *target_block.header());
1899 assert_eq!(result.target_body, body);
1900 assert_eq!(result.target_justifications, justifications);
1901 }
1902}