123456789use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::time::Duration;13use futures::StreamExt;141516use unique_runtime::RuntimeApi;171819use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};20use cumulus_client_consensus_common::ParachainConsensus;21use cumulus_client_network::build_block_announce_validator;22use cumulus_client_service::{23 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,24};25use cumulus_primitives_core::ParaId;262728use sc_client_api::ExecutorProvider;29use sc_executor::NativeElseWasmExecutor;30use sc_executor::NativeExecutionDispatch;31use sc_network::NetworkService;32use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};33use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};34use sp_consensus::SlotData;35use sp_keystore::SyncCryptoStorePtr;36use sp_runtime::traits::BlakeTwo256;37use substrate_prometheus_endpoint::Registry;38use sc_client_api::BlockchainEvents;394041use fc_rpc_core::types::FilterPool;42use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};434445type BlockNumber = u32;46type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;47pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;48type Hash = sp_core::H256;495051pub struct ParachainRuntimeExecutor;5253impl NativeExecutionDispatch for ParachainRuntimeExecutor {54 type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;5556 fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {57 unique_runtime::api::dispatch(method, data)58 }5960 fn native_version() -> sc_executor::NativeVersion {61 unique_runtime::native_version()62 }63}6465pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {66 let config_dir = config67 .base_path68 .as_ref()69 .map(|base_path| base_path.config_dir(config.chain_spec.id()))70 .unwrap_or_else(|| {71 BasePath::from_project("", "", "unique").config_dir(config.chain_spec.id())72 });73 let database_dir = config_dir.join("frontier").join("db");7475 Ok(Arc::new(fc_db::Backend::<Block>::new(76 &fc_db::DatabaseSettings {77 source: fc_db::DatabaseSettingsSrc::RocksDb {78 path: database_dir,79 cache_size: 0,80 },81 },82 )?))83}8485type ExecutorDispatch = ParachainRuntimeExecutor;8687type FullClient =88 sc_service::TFullClient<Block, RuntimeApi, NativeElseWasmExecutor<ExecutorDispatch>>;89type FullBackend = sc_service::TFullBackend<Block>;90type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;919293949596#[allow(clippy::type_complexity)]97pub fn new_partial<BIQ>(98 config: &Configuration,99 build_import_queue: BIQ,100) -> Result<101 PartialComponents<102 FullClient,103 FullBackend,104 FullSelectChain,105 sc_consensus::DefaultImportQueue<Block, FullClient>,106 sc_transaction_pool::FullPool<Block, FullClient>,107 (108 Option<Telemetry>,109 Option<FilterPool>,110 Arc<fc_db::Backend<Block>>,111 Option<TelemetryWorkerHandle>,112 ),113 >,114 sc_service::Error,115>116where117 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,118 ExecutorDispatch: NativeExecutionDispatch + 'static,119 BIQ: FnOnce(120 Arc<FullClient>,121 &Configuration,122 Option<TelemetryHandle>,123 &TaskManager,124 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,125{126 let _telemetry = config127 .telemetry_endpoints128 .clone()129 .filter(|x| !x.is_empty())130 .map(|endpoints| -> Result<_, sc_telemetry::Error> {131 let worker = TelemetryWorker::new(16)?;132 let telemetry = worker.handle().new_telemetry(endpoints);133 Ok((worker, telemetry))134 })135 .transpose()?;136137 let telemetry = config138 .telemetry_endpoints139 .clone()140 .filter(|x| !x.is_empty())141 .map(|endpoints| -> Result<_, sc_telemetry::Error> {142 let worker = TelemetryWorker::new(16)?;143 let telemetry = worker.handle().new_telemetry(endpoints);144 Ok((worker, telemetry))145 })146 .transpose()?;147148 let executor = NativeElseWasmExecutor::<ExecutorDispatch>::new(149 config.wasm_method,150 config.default_heap_pages,151 config.max_runtime_instances,152 );153154 let (client, backend, keystore_container, task_manager) =155 sc_service::new_full_parts::<Block, RuntimeApi, _>(156 config,157 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),158 executor,159 )?;160 let client = Arc::new(client);161162 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());163164 let telemetry = telemetry.map(|(worker, telemetry)| {165 task_manager166 .spawn_handle()167 .spawn("telemetry", None, worker.run());168 telemetry169 });170171 let select_chain = sc_consensus::LongestChain::new(backend.clone());172173 let transaction_pool = sc_transaction_pool::BasicPool::new_full(174 config.transaction_pool.clone(),175 config.role.is_authority().into(),176 config.prometheus_registry(),177 task_manager.spawn_essential_handle(),178 client.clone(),179 );180181 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));182183 let frontier_backend = open_frontier_backend(config)?;184185 let import_queue = build_import_queue(186 client.clone(),187 config,188 telemetry.as_ref().map(|telemetry| telemetry.handle()),189 &task_manager,190 )?;191192 let params = PartialComponents {193 backend,194 client,195 import_queue,196 keystore_container,197 task_manager,198 transaction_pool,199 select_chain,200 other: (201 telemetry,202 filter_pool,203 frontier_backend,204 telemetry_worker_handle,205 ),206 };207208 Ok(params)209}210211212213214#[sc_tracing::logging::prefix_logs_with("Parachain")]215async fn start_node_impl<BIQ, BIC>(216 parachain_config: Configuration,217 polkadot_config: Configuration,218 id: ParaId,219 build_import_queue: BIQ,220 build_consensus: BIC,221) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>222where223 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,224 ExecutorDispatch: NativeExecutionDispatch + 'static,225 BIQ: FnOnce(226 Arc<FullClient>,227 &Configuration,228 Option<TelemetryHandle>,229 &TaskManager,230 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,231 BIC: FnOnce(232 Arc<FullClient>,233 Option<&Registry>,234 Option<TelemetryHandle>,235 &TaskManager,236 &polkadot_service::NewFull<polkadot_service::Client>,237 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,238 Arc<NetworkService<Block, Hash>>,239 SyncCryptoStorePtr,240 bool,241 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,242{243 if matches!(parachain_config.role, Role::Light) {244 return Err("Light client not supported!".into());245 }246247 let parachain_config = prepare_node_config(parachain_config);248249 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;250 let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle) = params.other;251252 let relay_chain_full_node =253 cumulus_client_service::build_polkadot_full_node(polkadot_config, telemetry_worker_handle)254 .map_err(|e| match e {255 polkadot_service::Error::Sub(x) => x,256 s => format!("{}", s).into(),257 })?;258259 let client = params.client.clone();260 let backend = params.backend.clone();261 let block_announce_validator = build_block_announce_validator(262 relay_chain_full_node.client.clone(),263 id,264 Box::new(relay_chain_full_node.network.clone()),265 relay_chain_full_node.backend.clone(),266 );267268 let force_authoring = parachain_config.force_authoring;269 let validator = parachain_config.role.is_authority();270 let prometheus_registry = parachain_config.prometheus_registry().cloned();271 let transaction_pool = params.transaction_pool.clone();272 let mut task_manager = params.task_manager;273 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);274275 let (network, system_rpc_tx, start_network) =276 sc_service::build_network(sc_service::BuildNetworkParams {277 config: ¶chain_config,278 client: client.clone(),279 transaction_pool: transaction_pool.clone(),280 spawn_handle: task_manager.spawn_handle(),281 import_queue: import_queue.clone(),282 block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),283 warp_sync: None,284 })?;285286 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());287 let rpc_client = client.clone();288 let rpc_pool = transaction_pool.clone();289 let select_chain = params.select_chain.clone();290 let is_authority = parachain_config.role.clone().is_authority();291 let rpc_network = network.clone();292293 let rpc_frontier_backend = frontier_backend.clone();294 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {295 let full_deps = unique_rpc::FullDeps {296 backend: rpc_frontier_backend.clone(),297 deny_unsafe,298 client: rpc_client.clone(),299 pool: rpc_pool.clone(),300 graph: rpc_pool.pool().clone(),301 302 enable_dev_signer: false,303 filter_pool: filter_pool.clone(),304 network: rpc_network.clone(),305 select_chain: select_chain.clone(),306 is_authority,307 308 max_past_logs: 10000,309 };310311 Ok(unique_rpc::create_full::<_, _, _, _, RuntimeApi, _>(312 full_deps,313 subscription_executor.clone(),314 ))315 });316317 task_manager.spawn_essential_handle().spawn(318 "frontier-mapping-sync-worker",319 None,320 MappingSyncWorker::new(321 client.import_notification_stream(),322 Duration::new(6, 0),323 client.clone(),324 backend.clone(),325 frontier_backend.clone(),326 SyncStrategy::Normal,327 )328 .for_each(|()| futures::future::ready(())),329 );330331 sc_service::spawn_tasks(sc_service::SpawnTasksParams {332 rpc_extensions_builder,333 client: client.clone(),334 transaction_pool: transaction_pool.clone(),335 task_manager: &mut task_manager,336 config: parachain_config,337 keystore: params.keystore_container.sync_keystore(),338 backend: backend.clone(),339 network: network.clone(),340 system_rpc_tx,341 telemetry: telemetry.as_mut(),342 })?;343344 let announce_block = {345 let network = network.clone();346 Arc::new(move |hash, data| network.announce_block(hash, data))347 };348349 if validator {350 let parachain_consensus = build_consensus(351 client.clone(),352 prometheus_registry.as_ref(),353 telemetry.as_ref().map(|t| t.handle()),354 &task_manager,355 &relay_chain_full_node,356 transaction_pool,357 network,358 params.keystore_container.sync_keystore(),359 force_authoring,360 )?;361362 let spawner = task_manager.spawn_handle();363364 let params = StartCollatorParams {365 para_id: id,366 block_status: client.clone(),367 announce_block,368 client: client.clone(),369 task_manager: &mut task_manager,370 relay_chain_full_node,371 spawner,372 parachain_consensus,373 import_queue,374 };375376 start_collator(params).await?;377 } else {378 let params = StartFullNodeParams {379 client: client.clone(),380 announce_block,381 task_manager: &mut task_manager,382 para_id: id,383 relay_chain_full_node,384 };385386 start_full_node(params)?;387 }388389 start_network.start_network();390391 Ok((task_manager, client))392}393394395pub fn parachain_build_import_queue(396 client: Arc<FullClient>,397 config: &Configuration,398 telemetry: Option<TelemetryHandle>,399 task_manager: &TaskManager,400) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {401 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;402403 cumulus_client_consensus_aura::import_queue::<404 sp_consensus_aura::sr25519::AuthorityPair,405 _,406 _,407 _,408 _,409 _,410 _,411 >(cumulus_client_consensus_aura::ImportQueueParams {412 block_import: client.clone(),413 client: client.clone(),414 create_inherent_data_providers: move |_, _| async move {415 let time = sp_timestamp::InherentDataProvider::from_system_time();416417 let slot =418 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(419 *time,420 slot_duration.slot_duration(),421 );422423 Ok((time, slot))424 },425 registry: config.prometheus_registry(),426 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),427 spawner: &task_manager.spawn_essential_handle(),428 telemetry,429 })430 .map_err(Into::into)431}432433434pub async fn start_node(435 parachain_config: Configuration,436 polkadot_config: Configuration,437 id: ParaId,438) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {439 start_node_impl::<_, _>(440 parachain_config,441 polkadot_config,442 id,443 parachain_build_import_queue,444 |client,445 prometheus_registry,446 telemetry,447 task_manager,448 relay_chain_node,449 transaction_pool,450 sync_oracle,451 keystore,452 force_authoring| {453 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;454455 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(456 task_manager.spawn_handle(),457 client.clone(),458 transaction_pool,459 prometheus_registry,460 telemetry.clone(),461 );462463 let relay_chain_backend = relay_chain_node.backend.clone();464 let relay_chain_client = relay_chain_node.client.clone();465 Ok(build_aura_consensus::<466 sp_consensus_aura::sr25519::AuthorityPair,467 _,468 _,469 _,470 _,471 _,472 _,473 _,474 _,475 _,476 >(BuildAuraConsensusParams {477 proposer_factory,478 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {479 let parachain_inherent =480 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(481 relay_parent,482 &relay_chain_client,483 &*relay_chain_backend,484 &validation_data,485 id,486 );487 async move {488 let time = sp_timestamp::InherentDataProvider::from_system_time();489490 let slot =491 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(492 *time,493 slot_duration.slot_duration(),494 );495496 let parachain_inherent = parachain_inherent.ok_or_else(|| {497 Box::<dyn std::error::Error + Send + Sync>::from(498 "Failed to create parachain inherent",499 )500 })?;501 Ok((time, slot, parachain_inherent))502 }503 },504 block_import: client.clone(),505 relay_chain_client: relay_chain_node.client.clone(),506 relay_chain_backend: relay_chain_node.backend.clone(),507 para_client: client,508 backoff_authoring_blocks: Option::<()>::None,509 sync_oracle,510 keystore,511 force_authoring,512 slot_duration,513 514 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),515 telemetry,516 max_block_proposal_slot_portion: None,517 }))518 },519 )520 .await521}