referrerpolicy=no-referrer-when-downgrade

sc_network_sync/strategy/
warp.rs

1// This file is part of Substrate.
2
3// Copyright (C) Parity Technologies (UK) Ltd.
4// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
5
6// This program is free software: you can redistribute it and/or modify
7// it under the terms of the GNU General Public License as published by
8// the Free Software Foundation, either version 3 of the License, or
9// (at your option) any later version.
10
11// This program is distributed in the hope that it will be useful,
12// but WITHOUT ANY WARRANTY; without even the implied warranty of
13// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
14// GNU General Public License for more details.
15
16// You should have received a copy of the GNU General Public License
17// along with this program. If not, see <https://www.gnu.org/licenses/>.
18
19//! Warp syncing strategy. Bootstraps chain by downloading warp proofs and state.
20
21use 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
49/// Number of peers that need to be connected before warp sync is started.
50const MIN_PEERS_TO_START_WARP_SYNC: usize = 3;
51
52/// Scale-encoded warp sync proof response.
53pub struct EncodedProof(pub Vec<u8>);
54
55/// Warp sync request
56#[derive(Encode, Decode, Debug, Clone)]
57pub struct WarpProofRequest<B: BlockT> {
58	/// Start collecting proofs from this block.
59	pub begin: B::Hash,
60}
61
62/// Verifier for warp sync proofs. Each verifier operates in a specific context.
63pub trait Verifier<Block: BlockT>: Send + Sync {
64	/// Verify a warp sync proof.
65	fn verify(
66		&mut self,
67		proof: &EncodedProof,
68	) -> Result<VerificationResult<Block>, Box<dyn std::error::Error + Send + Sync>>;
69	/// Hash to be used as the starting point for the next proof request.
70	fn next_proof_context(&self) -> Block::Hash;
71	/// Get status text for progress reporting
72	fn status(&self) -> Option<String>;
73}
74
75/// Proof verification result.
76pub enum VerificationResult<Block: BlockT> {
77	/// Proof is valid, but the target was not reached.
78	Partial(Vec<(Block::Header, Justifications)>),
79	/// Target finality is proved.
80	Complete(Block::Header, Vec<(Block::Header, Justifications)>),
81}
82
83/// Warp sync backend. Handles retrieving and verifying warp sync proofs.
84pub trait WarpSyncProvider<Block: BlockT>: Send + Sync {
85	/// Generate proof starting at given block hash. The proof is accumulated until maximum proof
86	/// size is reached.
87	fn generate(
88		&self,
89		start: Block::Hash,
90	) -> Result<EncodedProof, Box<dyn std::error::Error + Send + Sync>>;
91	/// Create a verifier for warp sync proofs.
92	fn create_verifier(&self) -> Box<dyn Verifier<Block>>;
93}
94
95mod rep {
96	use sc_network::ReputationChange as Rep;
97
98	/// Unexpected response received form a peer
99	pub const UNEXPECTED_RESPONSE: Rep = Rep::new(-(1 << 29), "Unexpected response");
100
101	/// Peer provided invalid warp proof data
102	pub const BAD_WARP_PROOF: Rep = Rep::new(-(1 << 29), "Bad warp proof");
103
104	/// Peer did not provide us with advertised block data.
105	pub const NO_BLOCK: Rep = Rep::new(-(1 << 29), "No requested block data");
106
107	/// Reputation change for peers which send us non-requested block data.
108	pub const NOT_REQUESTED: Rep = Rep::new(-(1 << 29), "Not requested block data");
109
110	/// Reputation change for peers which send us a block which we fail to verify.
111	pub const VERIFICATION_FAIL: Rep = Rep::new(-(1 << 29), "Block verification failed");
112
113	/// We received a message that failed to decode.
114	pub const BAD_MESSAGE: Rep = Rep::new(-(1 << 12), "Bad message");
115}
116
117/// Reported warp sync phase.
118#[derive(Clone, Eq, PartialEq, Debug)]
119pub enum WarpSyncPhase<Block: BlockT> {
120	/// Waiting for peers to connect.
121	AwaitingPeers { required_peers: usize },
122	/// Downloading and verifying warp proofs.
123	DownloadingWarpProofs,
124	/// Downloading target block.
125	DownloadingTargetBlock,
126	/// Downloading state data.
127	DownloadingState,
128	/// Importing state.
129	ImportingState,
130	/// Downloading block history.
131	DownloadingBlocks(NumberFor<Block>),
132	/// Warp sync is complete.
133	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/// Reported warp sync progress.
153#[derive(Clone, Eq, PartialEq, Debug)]
154pub struct WarpSyncProgress<Block: BlockT> {
155	/// Estimated download percentage.
156	pub phase: WarpSyncPhase<Block>,
157	/// Total bytes downloaded so far.
158	pub total_bytes: u64,
159	/// Optional status text from the verifier.
160	pub status: Option<String>,
161}
162
163/// Warp sync configuration as accepted by [`WarpSync`].
164pub enum WarpSyncConfig<Block: BlockT> {
165	/// Standard warp sync for the chain.
166	WithProvider(Arc<dyn WarpSyncProvider<Block>>),
167	/// Skip downloading proofs and use provided header of the state that should be downloaded.
168	///
169	/// It is expected that the header provider ensures that the header is trusted.
170	WithTarget(<Block as BlockT>::Header),
171}
172
173/// Warp sync phase used by warp sync state machine.
174enum Phase<B: BlockT> {
175	/// Waiting for enough peers to connect.
176	WaitingForPeers { warp_sync_provider: Arc<dyn WarpSyncProvider<B>> },
177	/// Downloading warp proofs.
178	WarpProof { verifier: Box<dyn Verifier<B>> },
179	/// Downloading target block.
180	TargetBlock(B::Header),
181	/// Warp sync is complete.
182	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
208/// Warp sync state machine. Accumulates warp proofs and state.
209pub 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	/// Number of peers that need to be connected before warp sync is started.
220	min_peers_to_start_warp_sync: usize,
221}
222
223impl<B> WarpSync<B>
224where
225	B: BlockT,
226{
227	/// Strategy key used by warp sync.
228	pub const STRATEGY_KEY: StrategyKey = StrategyKey::new("Warp");
229
230	/// Create a new instance. When passing a warp sync provider we will be checking for proof and
231	/// authorities. Alternatively we can pass a target block when we want to skip downloading
232	/// proofs, in this case we will continue polling until the target block is known.
233	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	/// Notify that a new peer has connected.
286	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	/// Notify that a peer has disconnected.
293	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	/// Submit a validated block announcement.
306	///
307	/// Returns new best hash & best number of the peer if they are updated.
308	#[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			// Let `SyncingEngine` know that we should update the peer info.
322			(best_hash, best_number)
323		})
324	}
325
326	/// Start warp sync as soon as we have enough peers.
327	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	/// Process warp proof response.
392	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					// We are still in warp sync, so we don't have the state. This means
414					// we also can't execute the block.
415					allow_missing_state: true,
416					skip_execution: true,
417					// Shouldn't already exist in the database.
418					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						// We need target block with state and warp sync does not provide this.
451						// That's why we filter out target block here, otherwise oncoming state sync
452						// (which comes after warp sync) will abort because target block is already
453						// imported.
454						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	/// Process (target) block response.
476	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	/// Reserve a peer for a request assigning `new_state`.
558	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		// Find a random peer that is synced as much as peer majority and is above
571		// `min_best_number`.
572		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	/// Produce warp proof request.
585	fn warp_proof_request(&mut self) -> Option<(PeerId, ProtocolName, WarpProofRequest<B>)> {
586		let Phase::WarpProof { verifier } = &self.phase else { return None };
587
588		// Copy verifier context early to cut the borrowing tie.
589		let begin = verifier.next_proof_context();
590
591		if self
592			.peers
593			.values()
594			.any(|peer| matches!(peer.state, PeerState::DownloadingProofs))
595		{
596			// Only one warp proof request at a time is possible.
597			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	/// Produce target block request.
617	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			// Only one target block request at a time is possible.
626			return None;
627		}
628
629		// Cut the borrowing tie.
630		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	/// Returns warp sync estimated progress (stage, bytes received).
658	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	/// Get the number of peers known to warp sync.
686	pub fn num_peers(&self) -> usize {
687		self.peers.len()
688	}
689
690	/// Returns the current sync status.
691	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	/// Get actions that should be performed by the owner on [`WarpSync`]'s behalf
713	#[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	/// Take the result of finished warp sync, returning `None` if the sync was unsuccessful.
777	#[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			// We need some finalized state to render warp sync impossible.
852			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		// Warp sync instantly finishes
894		let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
895		assert_eq!(actions.len(), 1);
896		assert!(matches!(actions[0], SyncingAction::Finished));
897
898		// ... with no result.
899		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		// Warp sync instantly finishes
924		let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
925		assert_eq!(actions.len(), 1);
926		assert!(matches!(actions[0], SyncingAction::Finished));
927
928		// ... with no result.
929		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		// No actions are emitted.
949		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		// No actions are emitted.
974		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		// Warp sync is not started when there is not enough peers.
997		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		// Now we have enough peers and warp sync is started.
1003		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		// Tests below use 10 peers. Fail early if it's less than a threshold for warp sync.
1026		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		// Disconnecting a peer without an inflight request has no effect on persistent states.
1119		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		// Peer is backed off.
1128		assert!(!warp_sync.disconnected_peers.is_peer_available(&tenth_peer));
1129
1130		// No peer available for 10'th best block because of the backoff.
1131		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		// Other requests can still happen.
1137		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		// Make sure we have enough peers to make a request.
1162		for best_number in 1..11 {
1163			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1164		}
1165
1166		// Manually set to another phase.
1167		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		// No request is made.
1176		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		// Make sure we have enough peers to make a request.
1200		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		// Make sure we have enough peers to make requests.
1229		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		// First request is made.
1235		assert!(warp_sync.warp_proof_request().is_some());
1236		// Second request is not made.
1237		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		// Warp proof verification fails.
1247		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		// Make sure we have enough peers to make a request.
1261		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		// Consume `SendWarpProofRequest` action.
1270		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		// We only interested in already generated actions, not new requests.
1279		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		// Warp proof is partial.
1303		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		// Make sure we have enough peers to make a request.
1325		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		// Consume `SendWarpProofRequest` action.
1334		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		// Initialize logging
1374		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		// Building both blocks manually. Can't use BlockBuilderBuilder to build one block on top of
1380		// another, if no genesis is set (required for warp sync to make sure finalized_state is not
1381		// set)
1382		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		// Build target_block on top of warp_synced_block
1391		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		// Warp proof is complete and contains both blocks.
1403		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		// One peer must serve both the proof and the subsequent target-block request.
1430		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		// Receiving the proof already frees this peer's request slot.
1454		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		// Only warp_synced_block should be in the import list; target_block is filtered out
1468		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		// No explicit cancellation or implicit removal is needed between these requests.
1496		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		// Make sure we have enough peers to make a request.
1521		for best_number in 1..11 {
1522			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1523		}
1524		// We are not in `Phase::TargetBlock`
1525		assert!(matches!(warp_sync.phase, Phase::WarpProof { .. }));
1526
1527		// No request is made.
1528		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		// Warp proof is complete.
1548		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		// Make sure we have enough peers to make a request.
1558		for best_number in 1..11 {
1559			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1560		}
1561
1562		// Manually set `TargetBlock` phase.
1563		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		// Make sure we have enough peers to make a request.
1591		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		// Warp proof is complete.
1623		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		// Make sure we have enough peers to make a request.
1633		for best_number in 1..11 {
1634			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1635		}
1636
1637		// Manually set `TargetBlock` phase.
1638		warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1639
1640		// First target block request is made.
1641		assert!(warp_sync.target_block_request().is_some());
1642		// No parallel request is made.
1643		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		// Warp proof is complete.
1663		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		// Make sure we have enough peers to make a request.
1673		for best_number in 1..11 {
1674			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1675		}
1676
1677		// Manually set `TargetBlock` phase.
1678		warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1679
1680		let (peer_id, request) = warp_sync.target_block_request().unwrap();
1681
1682		// Empty block response received.
1683		let response = Vec::new();
1684		// Peer is dropped.
1685		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		// Warp proof is complete.
1719		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		// Make sure we have enough peers to make a request.
1729		for best_number in 1..11 {
1730			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1731		}
1732
1733		// Manually set `TargetBlock` phase.
1734		warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1735
1736		let (peer_id, request) = warp_sync.target_block_request().unwrap();
1737
1738		// Block response with extra blocks received.
1739		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		// Peer is dropped.
1762		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		// Warp proof is complete.
1798		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		// Make sure we have enough peers to make a request.
1808		for best_number in 1..11 {
1809			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1810		}
1811
1812		// Manually set `TargetBlock` phase.
1813		warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1814
1815		let (peer_id, request) = warp_sync.target_block_request().unwrap();
1816
1817		// Wrong block received.
1818		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		// Peer is dropped.
1829		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		// Warp proof is complete.
1853		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		// Make sure we have enough peers to make a request.
1863		for best_number in 1..11 {
1864			warp_sync.add_peer(PeerId::random(), Hash::random(), best_number);
1865		}
1866
1867		// Manually set `TargetBlock` phase.
1868		warp_sync.phase = Phase::TargetBlock(target_block.header().clone());
1869
1870		let (peer_id, request) = warp_sync.target_block_request().unwrap();
1871
1872		// Correct block received.
1873		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		// Strategy finishes.
1892		let actions = warp_sync.actions(&network_handle).collect::<Vec<_>>();
1893		assert_eq!(actions.len(), 1);
1894		assert!(matches!(actions[0], SyncingAction::Finished));
1895
1896		// With correct result.
1897		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}