1use cumulus_client_cli::CollatorOptions;
23use cumulus_client_network::AssumeSybilResistance;
24use cumulus_client_pov_recovery::{PoVRecovery, RecoveryDelayRange, RecoveryHandle};
25use cumulus_client_proof_size_recording::load_proof_size_recording;
26use cumulus_primitives_core::ParaId;
27pub use cumulus_primitives_proof_size_hostfunction::storage_proof_size;
28use cumulus_relay_chain_inprocess_interface::build_inprocess_relay_chain;
29use cumulus_relay_chain_interface::{RelayChainInterface, RelayChainResult};
30use cumulus_relay_chain_minimal_node::build_minimal_relay_chain_node_with_rpc;
31use futures::{channel::mpsc, StreamExt};
32use polkadot_primitives::{CandidateEvent, CollatorPair, OccupiedCoreAssumption};
33use prometheus::{Histogram, HistogramOpts, Registry};
34use sc_client_api::{
35 AuxStore, Backend as BackendT, BlockBackend, BlockchainEvents, CallExecutor, ExecutorProvider,
36 Finalizer, ProofProvider, UsageProvider,
37};
38use sc_consensus::{
39 import_queue::{ImportQueue, ImportQueueService},
40 BlockImport,
41};
42use sc_network::{
43 config::SyncMode, request_responses::IncomingRequest, service::traits::NetworkService,
44 NetworkBackend,
45};
46use sc_network_sync::{strategy::chain_sync::GapSyncBodyPolicyProvider, SyncingService};
47use sc_network_transactions::TransactionsHandlerController;
48use sc_service::{
49 Configuration, SpawnEssentialTaskHandle, SpawnTaskHandle, TaskManager, WarpSyncConfig,
50};
51use sc_telemetry::{log, TelemetryWorkerHandle};
52use sc_tracing::block::TracingExecuteBlock;
53use sc_utils::mpsc::TracingUnboundedSender;
54use sp_api::{ApiExt, Core, ProofRecorder, ProvideRuntimeApi};
55use sp_blockchain::{HeaderBackend, HeaderMetadata};
56use sp_core::{traits::CallContext, Decode};
57use sp_runtime::{
58 traits::{Block as BlockT, HashingFor, Header},
59 SaturatedConversion, Saturating,
60};
61use sp_state_machine::OverlayedChanges;
62use sp_trie::proof_size_extension::{ProofSizeExt, ReplayProofSizeProvider};
63use std::{
64 cell::RefCell,
65 sync::Arc,
66 time::{Duration, Instant},
67};
68
69pub type ParachainHostFunctions = (
74 cumulus_primitives_proof_size_hostfunction::storage_proof_size::HostFunctions,
75 sp_io::SubstrateHostFunctions,
76 sp_crypto_ec_utils::HostFunctionsRfc163,
77);
78
79const RECOVERY_CHAN_SIZE: usize = 8;
83const LOG_TARGET_SYNC: &str = "sync::cumulus";
84
85pub enum DARecoveryProfile {
88 Collator,
90 FullNode,
93 Other(RecoveryDelayRange),
95}
96
97pub struct StartRelayChainTasksParams<'a, Block: BlockT, Client, RCInterface> {
99 pub client: Arc<Client>,
100 pub announce_block: Arc<dyn Fn(Block::Hash, Option<Vec<u8>>) + Send + Sync>,
101 pub para_id: ParaId,
102 pub relay_chain_interface: RCInterface,
103 pub task_manager: &'a mut TaskManager,
104 pub da_recovery_profile: DARecoveryProfile,
105 pub import_queue: Box<dyn ImportQueueService<Block>>,
106 pub relay_chain_slot_duration: Duration,
107 pub recovery_handle: Box<dyn RecoveryHandle>,
108 pub sync_service: Arc<SyncingService<Block>>,
109 pub prometheus_registry: Option<&'a Registry>,
110}
111
112pub fn start_relay_chain_tasks<Block, Client, Backend, RCInterface>(
122 StartRelayChainTasksParams {
123 client,
124 announce_block,
125 para_id,
126 task_manager,
127 da_recovery_profile,
128 relay_chain_interface,
129 import_queue,
130 relay_chain_slot_duration,
131 recovery_handle,
132 sync_service,
133 prometheus_registry,
134 }: StartRelayChainTasksParams<Block, Client, RCInterface>,
135) -> sc_service::error::Result<()>
136where
137 Block: BlockT,
138 Client: Finalizer<Block, Backend>
139 + UsageProvider<Block>
140 + HeaderBackend<Block>
141 + Send
142 + Sync
143 + BlockBackend<Block>
144 + BlockchainEvents<Block>
145 + 'static,
146 for<'a> &'a Client: BlockImport<Block>,
147 Backend: BackendT<Block> + 'static,
148 RCInterface: RelayChainInterface + Clone + 'static,
149{
150 let (recovery_chan_tx, recovery_chan_rx) = mpsc::channel(RECOVERY_CHAN_SIZE);
151
152 cumulus_client_consensus_common::spawn_parachain_consensus_tasks(
153 para_id,
154 client.clone(),
155 relay_chain_interface.clone(),
156 announce_block.clone(),
157 Some(recovery_chan_tx),
158 task_manager.spawn_essential_handle(),
159 );
160
161 let da_recovery_profile = match da_recovery_profile {
162 DARecoveryProfile::Collator => {
163 RecoveryDelayRange {
167 min: relay_chain_slot_duration / 2,
168 max: relay_chain_slot_duration,
169 }
170 },
171 DARecoveryProfile::FullNode => {
172 RecoveryDelayRange {
178 min: relay_chain_slot_duration * 25,
179 max: relay_chain_slot_duration * 50,
180 }
181 },
182 DARecoveryProfile::Other(profile) => profile,
183 };
184
185 let pov_recovery = PoVRecovery::new(
186 recovery_handle,
187 da_recovery_profile,
188 client.clone(),
189 import_queue,
190 relay_chain_interface.clone(),
191 para_id,
192 recovery_chan_rx,
193 sync_service.clone(),
194 );
195
196 task_manager
197 .spawn_essential_handle()
198 .spawn("cumulus-pov-recovery", None, pov_recovery.run());
199
200 let parachain_informant = parachain_informant::<Block, _>(
201 para_id,
202 relay_chain_interface.clone(),
203 client.clone(),
204 prometheus_registry.map(ParachainInformantMetrics::new).transpose()?,
205 );
206 task_manager
207 .spawn_handle()
208 .spawn("parachain-informant", None, parachain_informant);
209
210 Ok(())
211}
212
213pub fn prepare_node_config(mut parachain_config: Configuration) -> Configuration {
220 parachain_config.announce_block = false;
221 parachain_config.network.min_peers_to_start_warp_sync = Some(1);
224
225 parachain_config
226}
227
228pub async fn build_relay_chain_interface(
232 relay_chain_config: Configuration,
233 parachain_config: &Configuration,
234 telemetry_worker_handle: Option<TelemetryWorkerHandle>,
235 task_manager: &mut TaskManager,
236 collator_options: CollatorOptions,
237 hwbench: Option<sc_sysinfo::HwBench>,
238) -> RelayChainResult<(
239 Arc<dyn RelayChainInterface + 'static>,
240 Option<CollatorPair>,
241 Arc<dyn NetworkService>,
242 async_channel::Receiver<IncomingRequest>,
243)> {
244 match collator_options.relay_chain_mode {
245 cumulus_client_cli::RelayChainMode::Embedded => build_inprocess_relay_chain(
246 relay_chain_config,
247 parachain_config,
248 telemetry_worker_handle,
249 task_manager,
250 hwbench,
251 ),
252 cumulus_client_cli::RelayChainMode::ExternalRpc(rpc_target_urls) => {
253 build_minimal_relay_chain_node_with_rpc(
254 relay_chain_config,
255 parachain_config.prometheus_registry(),
256 task_manager,
257 rpc_target_urls,
258 )
259 .await
260 },
261 }
262}
263
264pub struct BuildNetworkParams<
266 'a,
267 Block: BlockT,
268 Client,
269 Network: NetworkBackend<Block, <Block as BlockT>::Hash>,
270 RCInterface,
271 IQ,
272> {
273 pub parachain_config: &'a Configuration,
274 pub net_config:
275 sc_network::config::FullNetworkConfiguration<Block, <Block as BlockT>::Hash, Network>,
276 pub client: Arc<Client>,
277 pub transaction_pool: Arc<sc_transaction_pool::TransactionPoolHandle<Block>>,
278 pub para_id: ParaId,
279 pub relay_chain_interface: RCInterface,
280 pub spawn_handle: SpawnTaskHandle,
281 pub spawn_essential_handle: SpawnEssentialTaskHandle,
282 pub import_queue: IQ,
283 pub metrics: sc_network::NotificationMetrics,
284 pub gap_sync_body_policy: Option<GapSyncBodyPolicyProvider>,
287}
288
289pub async fn build_network<'a, Block, Client, RCInterface, IQ, Network>(
291 BuildNetworkParams {
292 parachain_config,
293 net_config,
294 client,
295 transaction_pool,
296 para_id,
297 spawn_handle,
298 spawn_essential_handle,
299 relay_chain_interface,
300 import_queue,
301 metrics,
302 gap_sync_body_policy,
303 }: BuildNetworkParams<'a, Block, Client, Network, RCInterface, IQ>,
304) -> sc_service::error::Result<(
305 Arc<dyn NetworkService>,
306 TracingUnboundedSender<sc_rpc::system::Request<Block>>,
307 TransactionsHandlerController<Block::Hash>,
308 Arc<SyncingService<Block>>,
309 Option<sc_network_bitswap::BitswapHandle>,
310)>
311where
312 Block: BlockT,
313 Client: ProvideRuntimeApi<Block>
314 + HeaderMetadata<Block, Error = sp_blockchain::Error>
315 + sp_consensus::block_validation::Chain<Block>
316 + BlockBackend<Block>
317 + ProofProvider<Block>
318 + HeaderBackend<Block>
319 + BlockchainEvents<Block>
320 + 'static,
321 RCInterface: RelayChainInterface + Clone + 'static,
322 IQ: ImportQueue<Block> + 'static,
323 Network: NetworkBackend<Block, <Block as BlockT>::Hash>,
324{
325 let warp_sync_config = match parachain_config.network.sync_mode {
326 SyncMode::Warp => {
327 log::debug!(target: LOG_TARGET_SYNC, "waiting for announce block...");
328
329 let target_block =
330 wait_for_finalized_para_head::<Block, _>(para_id, relay_chain_interface.clone())
331 .await
332 .inspect_err(|e| {
333 log::error!(
334 target: LOG_TARGET_SYNC,
335 "Unable to determine parachain target block {:?}",
336 e
337 );
338 })?;
339 Some(WarpSyncConfig::WithTarget(target_block))
340 },
341 _ => None,
342 };
343
344 let block_announce_validator = Box::new(AssumeSybilResistance::allow_seconded_messages());
345
346 sc_service::build_network(sc_service::BuildNetworkParams {
347 config: parachain_config,
348 net_config,
349 client,
350 transaction_pool,
351 spawn_handle,
352 spawn_essential_handle,
353 import_queue,
354 block_announce_validator_builder: Some(Box::new(move |_| block_announce_validator)),
355 warp_sync_config,
356 block_relay: None,
357 metrics,
358 gap_sync_body_policy,
359 })
360}
361
362async fn wait_for_finalized_para_head<B, RCInterface>(
365 para_id: ParaId,
366 relay_chain_interface: RCInterface,
367) -> sc_service::error::Result<<B as BlockT>::Header>
368where
369 B: BlockT + 'static,
370 RCInterface: RelayChainInterface + Send + 'static,
371{
372 let mut imported_blocks = relay_chain_interface
373 .import_notification_stream()
374 .await
375 .map_err(|error| {
376 sc_service::Error::Other(format!(
377 "Relay chain import notification stream error when waiting for parachain head: \
378 {error}"
379 ))
380 })?
381 .fuse();
382 while imported_blocks.next().await.is_some() {
383 let is_syncing = relay_chain_interface
384 .is_major_syncing()
385 .await
386 .map_err(|e| format!("Unable to determine sync status: {e}"))?;
387
388 if !is_syncing {
389 let relay_chain_best_hash = relay_chain_interface
390 .finalized_block_hash()
391 .await
392 .map_err(|e| Box::new(e) as Box<_>)?;
393
394 let validation_data = relay_chain_interface
395 .persisted_validation_data(
396 relay_chain_best_hash,
397 para_id,
398 OccupiedCoreAssumption::TimedOut,
399 )
400 .await
401 .map_err(|e| format!("{e:?}"))?
402 .ok_or("Could not find parachain head in relay chain")?;
403
404 let finalized_header = B::Header::decode(&mut &validation_data.parent_head.0[..])
405 .map_err(|e| format!("Failed to decode parachain head: {e}"))?;
406
407 log::info!(
408 "๐ Received target parachain header #{} ({}) from the relay chain.",
409 finalized_header.number(),
410 finalized_header.hash()
411 );
412 return Ok(finalized_header);
413 }
414 }
415
416 Err("Stopping following imported blocks. Could not determine parachain target block".into())
417}
418
419async fn parachain_informant<Block: BlockT, Client>(
421 para_id: ParaId,
422 relay_chain_interface: impl RelayChainInterface + Clone,
423 client: Arc<Client>,
424 metrics: Option<ParachainInformantMetrics>,
425) where
426 Client: HeaderBackend<Block> + Send + Sync + 'static,
427{
428 let mut import_notifications = match relay_chain_interface.import_notification_stream().await {
429 Ok(import_notifications) => import_notifications,
430 Err(e) => {
431 log::error!("Failed to get import notification stream: {e:?}. Parachain informant will not run!");
432 return;
433 },
434 };
435 let mut last_backed_block_time: Option<Instant> = None;
436 while let Some(n) = import_notifications.next().await {
437 let candidate_events = match relay_chain_interface.candidate_events(n.hash()).await {
438 Ok(candidate_events) => candidate_events,
439 Err(e) => {
440 log::warn!("Failed to get candidate events for block {}: {e:?}", n.hash());
441 continue;
442 },
443 };
444 let mut backed_candidates = Vec::new();
445 let mut included_candidates = Vec::new();
446 let mut timed_out_candidates = Vec::new();
447 for event in candidate_events {
448 match event {
449 CandidateEvent::CandidateBacked(receipt, head, _, _) => {
450 if receipt.descriptor.para_id() != para_id {
451 continue;
452 }
453 let backed_block = match Block::Header::decode(&mut &head.0[..]) {
454 Ok(header) => header,
455 Err(e) => {
456 log::warn!(
457 "Failed to decode parachain header from backed block: {e:?}"
458 );
459 continue;
460 },
461 };
462 let backed_block_time = Instant::now();
463 if let Some(last_backed_block_time) = &last_backed_block_time {
464 let duration = backed_block_time.duration_since(*last_backed_block_time);
465 if let Some(metrics) = &metrics {
466 metrics.parachain_block_backed_duration.observe(duration.as_secs_f64());
467 }
468 }
469 last_backed_block_time = Some(backed_block_time);
470 backed_candidates.push(backed_block);
471 },
472 CandidateEvent::CandidateIncluded(receipt, head, _, _) => {
473 if receipt.descriptor.para_id() != para_id {
474 continue;
475 }
476 let included_block = match Block::Header::decode(&mut &head.0[..]) {
477 Ok(header) => header,
478 Err(e) => {
479 log::warn!(
480 "Failed to decode parachain header from included block: {e:?}"
481 );
482 continue;
483 },
484 };
485 let unincluded_segment_size =
486 client.info().best_number.saturating_sub(*included_block.number());
487 let unincluded_segment_size: u32 = unincluded_segment_size.saturated_into();
488 if let Some(metrics) = &metrics {
489 metrics.unincluded_segment_size.observe(unincluded_segment_size.into());
490 }
491 included_candidates.push(included_block);
492 },
493 CandidateEvent::CandidateTimedOut(receipt, head, _) => {
494 if receipt.descriptor.para_id() != para_id {
495 continue;
496 }
497 let timed_out_block = match Block::Header::decode(&mut &head.0[..]) {
498 Ok(header) => header,
499 Err(e) => {
500 log::warn!(
501 "Failed to decode parachain header from timed out block: {e:?}"
502 );
503 continue;
504 },
505 };
506 timed_out_candidates.push(timed_out_block);
507 },
508 }
509 }
510 let mut log_parts = Vec::new();
511 if !backed_candidates.is_empty() {
512 let backed_candidates = backed_candidates
513 .into_iter()
514 .map(|c| format!("#{} ({})", c.number(), c.hash()))
515 .collect::<Vec<_>>()
516 .join(", ");
517 log_parts.push(format!("backed: {}", backed_candidates));
518 };
519 if !included_candidates.is_empty() {
520 let included_candidates = included_candidates
521 .into_iter()
522 .map(|c| format!("#{} ({})", c.number(), c.hash()))
523 .collect::<Vec<_>>()
524 .join(", ");
525 log_parts.push(format!("included: {}", included_candidates));
526 };
527 if !timed_out_candidates.is_empty() {
528 let timed_out_candidates = timed_out_candidates
529 .into_iter()
530 .map(|c| format!("#{} ({})", c.number(), c.hash()))
531 .collect::<Vec<_>>()
532 .join(", ");
533 log_parts.push(format!("timed out: {}", timed_out_candidates));
534 };
535 if !log_parts.is_empty() {
536 log::info!(
537 "Update at relay chain block #{} ({}) - {}",
538 n.number(),
539 n.hash(),
540 log_parts.join(", ")
541 );
542 }
543 }
544}
545
546struct ParachainInformantMetrics {
547 parachain_block_backed_duration: Histogram,
549 unincluded_segment_size: Histogram,
551}
552
553impl ParachainInformantMetrics {
554 fn new(prometheus_registry: &Registry) -> prometheus::Result<Self> {
555 let parachain_block_authorship_duration = Histogram::with_opts(HistogramOpts::new(
556 "parachain_block_backed_duration",
557 "Time between parachain blocks getting backed by the relaychain",
558 ))?;
559 prometheus_registry.register(Box::new(parachain_block_authorship_duration.clone()))?;
560
561 let unincluded_segment_size = Histogram::with_opts(
562 HistogramOpts::new(
563 "parachain_unincluded_segment_size",
564 "Number of blocks between best block and last included block",
565 )
566 .buckets((0..=24).into_iter().map(|i| i as f64).collect()),
567 )?;
568 prometheus_registry.register(Box::new(unincluded_segment_size.clone()))?;
569
570 Ok(Self {
571 parachain_block_backed_duration: parachain_block_authorship_duration,
572 unincluded_segment_size,
573 })
574 }
575}
576
577pub struct ParachainTracingExecuteBlock<Client> {
582 client: Arc<Client>,
583}
584
585impl<Client> ParachainTracingExecuteBlock<Client> {
586 pub fn new(client: Arc<Client>) -> Self {
588 Self { client }
589 }
590}
591
592fn recorded_proof_size_ext<Block, Client>(
595 client: &Client,
596 hash: Block::Hash,
597 recorder: &ProofRecorder<Block>,
598) -> sp_blockchain::Result<ProofSizeExt>
599where
600 Block: BlockT,
601 Client: AuxStore,
602{
603 Ok(load_proof_size_recording(client, hash)?.map_or_else(
604 || ProofSizeExt::new(recorder.clone()),
605 |recordings| ProofSizeExt::new(ReplayProofSizeProvider::from(recordings)),
606 ))
607}
608
609impl<Block, Client> TracingExecuteBlock<Block> for ParachainTracingExecuteBlock<Client>
610where
611 Block: BlockT,
612 Client: ProvideRuntimeApi<Block>
613 + ExecutorProvider<Block>
614 + HeaderBackend<Block>
615 + AuxStore
616 + Send
617 + Sync,
618 Client::Api: Core<Block>,
619{
620 fn execute_block(&self, orig_hash: Block::Hash, block: Block) -> sp_blockchain::Result<()> {
621 let mut runtime_api = self.client.runtime_api();
622 let storage_proof_recorder = ProofRecorder::<Block>::default();
623
624 let proof_size_ext =
625 recorded_proof_size_ext::<Block, _>(&*self.client, orig_hash, &storage_proof_recorder)?;
626 runtime_api.register_extension(proof_size_ext);
627
628 runtime_api.record_proof_with_recorder(storage_proof_recorder);
629
630 runtime_api
631 .execute_block(*block.header().parent_hash(), block.into())
632 .map_err(Into::into)
633 }
634
635 fn call_recorded(
636 &self,
637 block: Block::Hash,
638 method: &str,
639 call_data: &[u8],
640 ) -> sp_blockchain::Result<Vec<u8>> {
641 let header = self
642 .client
643 .header(block)?
644 .ok_or_else(|| sp_blockchain::Error::UnknownBlock(format!("{block:?}")))?;
645 let at = *header.parent_hash();
646 let number = self
647 .client
648 .number(at)?
649 .ok_or_else(|| sp_blockchain::Error::UnknownBlock(format!("{at:?}")))?;
650 let storage_proof_recorder = ProofRecorder::<Block>::default();
651
652 let proof_size_ext =
653 recorded_proof_size_ext::<Block, _>(&*self.client, block, &storage_proof_recorder)?;
654 let mut extensions = self.client.execution_extensions().extensions(at, number);
655 extensions.register(proof_size_ext);
656
657 self.client.executor().contextual_call(
658 at,
659 method,
660 call_data,
661 &RefCell::new(OverlayedChanges::<HashingFor<Block>>::default()),
662 &Some(storage_proof_recorder),
663 CallContext::Offchain,
664 &RefCell::new(extensions),
665 )
666 }
667}