123456789use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::time::Duration;13use fc_rpc_core::types::FeeHistoryCache;14use futures::StreamExt;1516use unique_rpc::overrides_handle;1718use unique_runtime::RuntimeApi;192021use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};22use cumulus_client_consensus_common::ParachainConsensus;23use cumulus_client_network::build_block_announce_validator;24use cumulus_client_service::{25 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,26};27use cumulus_primitives_core::ParaId;282930use sc_client_api::ExecutorProvider;31use sc_executor::NativeElseWasmExecutor;32use sc_executor::NativeExecutionDispatch;33use sc_network::NetworkService;34use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};35use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};36use sp_consensus::SlotData;37use sp_keystore::SyncCryptoStorePtr;38use sp_runtime::traits::BlakeTwo256;39use substrate_prometheus_endpoint::Registry;40use sc_client_api::BlockchainEvents;414243use fc_rpc_core::types::FilterPool;44use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};454647type BlockNumber = u32;48type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;49pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;50type Hash = sp_core::H256;515253pub struct ParachainRuntimeExecutor;5455impl NativeExecutionDispatch for ParachainRuntimeExecutor {56 type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;5758 fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {59 unique_runtime::api::dispatch(method, data)60 }6162 fn native_version() -> sc_executor::NativeVersion {63 unique_runtime::native_version()64 }65}6667pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {68 let config_dir = config69 .base_path70 .as_ref()71 .map(|base_path| base_path.config_dir(config.chain_spec.id()))72 .unwrap_or_else(|| {73 BasePath::from_project("", "", "unique").config_dir(config.chain_spec.id())74 });75 let database_dir = config_dir.join("frontier").join("db");7677 Ok(Arc::new(fc_db::Backend::<Block>::new(78 &fc_db::DatabaseSettings {79 source: fc_db::DatabaseSettingsSrc::RocksDb {80 path: database_dir,81 cache_size: 0,82 },83 },84 )?))85}8687type ExecutorDispatch = ParachainRuntimeExecutor;8889type FullClient =90 sc_service::TFullClient<Block, RuntimeApi, NativeElseWasmExecutor<ExecutorDispatch>>;91type FullBackend = sc_service::TFullBackend<Block>;92type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;939495969798#[allow(clippy::type_complexity)]99pub fn new_partial<BIQ>(100 config: &Configuration,101 build_import_queue: BIQ,102) -> Result<103 PartialComponents<104 FullClient,105 FullBackend,106 FullSelectChain,107 sc_consensus::DefaultImportQueue<Block, FullClient>,108 sc_transaction_pool::FullPool<Block, FullClient>,109 (110 Option<Telemetry>,111 Option<FilterPool>,112 Arc<fc_db::Backend<Block>>,113 Option<TelemetryWorkerHandle>,114 FeeHistoryCache,115 ),116 >,117 sc_service::Error,118>119where120 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,121 ExecutorDispatch: NativeExecutionDispatch + 'static,122 BIQ: FnOnce(123 Arc<FullClient>,124 &Configuration,125 Option<TelemetryHandle>,126 &TaskManager,127 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,128{129 let _telemetry = config130 .telemetry_endpoints131 .clone()132 .filter(|x| !x.is_empty())133 .map(|endpoints| -> Result<_, sc_telemetry::Error> {134 let worker = TelemetryWorker::new(16)?;135 let telemetry = worker.handle().new_telemetry(endpoints);136 Ok((worker, telemetry))137 })138 .transpose()?;139140 let telemetry = config141 .telemetry_endpoints142 .clone()143 .filter(|x| !x.is_empty())144 .map(|endpoints| -> Result<_, sc_telemetry::Error> {145 let worker = TelemetryWorker::new(16)?;146 let telemetry = worker.handle().new_telemetry(endpoints);147 Ok((worker, telemetry))148 })149 .transpose()?;150151 let executor = NativeElseWasmExecutor::<ExecutorDispatch>::new(152 config.wasm_method,153 config.default_heap_pages,154 config.max_runtime_instances,155 );156157 let (client, backend, keystore_container, task_manager) =158 sc_service::new_full_parts::<Block, RuntimeApi, _>(159 config,160 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),161 executor,162 )?;163 let client = Arc::new(client);164165 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());166167 let telemetry = telemetry.map(|(worker, telemetry)| {168 task_manager169 .spawn_handle()170 .spawn("telemetry", None, worker.run());171 telemetry172 });173174 let select_chain = sc_consensus::LongestChain::new(backend.clone());175176 let transaction_pool = sc_transaction_pool::BasicPool::new_full(177 config.transaction_pool.clone(),178 config.role.is_authority().into(),179 config.prometheus_registry(),180 task_manager.spawn_essential_handle(),181 client.clone(),182 );183184 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));185186 let frontier_backend = open_frontier_backend(config)?;187188 let import_queue = build_import_queue(189 client.clone(),190 config,191 telemetry.as_ref().map(|telemetry| telemetry.handle()),192 &task_manager,193 )?;194 let fee_history_cache: FeeHistoryCache = Arc::new(Mutex::new(BTreeMap::new()));195196 let params = PartialComponents {197 backend,198 client,199 import_queue,200 keystore_container,201 task_manager,202 transaction_pool,203 select_chain,204 other: (205 telemetry,206 filter_pool,207 frontier_backend,208 telemetry_worker_handle,209 fee_history_cache,210 ),211 };212213 Ok(params)214}215216217218219#[sc_tracing::logging::prefix_logs_with("Parachain")]220async fn start_node_impl<BIQ, BIC>(221 parachain_config: Configuration,222 polkadot_config: Configuration,223 id: ParaId,224 build_import_queue: BIQ,225 build_consensus: BIC,226) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>227where228 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,229 ExecutorDispatch: NativeExecutionDispatch + 'static,230 BIQ: FnOnce(231 Arc<FullClient>,232 &Configuration,233 Option<TelemetryHandle>,234 &TaskManager,235 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,236 BIC: FnOnce(237 Arc<FullClient>,238 Option<&Registry>,239 Option<TelemetryHandle>,240 &TaskManager,241 &polkadot_service::NewFull<polkadot_service::Client>,242 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,243 Arc<NetworkService<Block, Hash>>,244 SyncCryptoStorePtr,245 bool,246 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,247{248 if matches!(parachain_config.role, Role::Light) {249 return Err("Light client not supported!".into());250 }251252 let parachain_config = prepare_node_config(parachain_config);253254 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;255 let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle, fee_history_cache) =256 params.other;257258 let relay_chain_full_node =259 cumulus_client_service::build_polkadot_full_node(polkadot_config, telemetry_worker_handle)260 .map_err(|e| match e {261 polkadot_service::Error::Sub(x) => x,262 s => format!("{}", s).into(),263 })?;264265 let client = params.client.clone();266 let backend = params.backend.clone();267 let block_announce_validator = build_block_announce_validator(268 relay_chain_full_node.client.clone(),269 id,270 Box::new(relay_chain_full_node.network.clone()),271 relay_chain_full_node.backend.clone(),272 );273274 let force_authoring = parachain_config.force_authoring;275 let validator = parachain_config.role.is_authority();276 let prometheus_registry = parachain_config.prometheus_registry().cloned();277 let transaction_pool = params.transaction_pool.clone();278 let mut task_manager = params.task_manager;279 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);280281 let (network, system_rpc_tx, start_network) =282 sc_service::build_network(sc_service::BuildNetworkParams {283 config: ¶chain_config,284 client: client.clone(),285 transaction_pool: transaction_pool.clone(),286 spawn_handle: task_manager.spawn_handle(),287 import_queue: import_queue.clone(),288 block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),289 warp_sync: None,290 })?;291292 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());293 let rpc_client = client.clone();294 let rpc_pool = transaction_pool.clone();295 let select_chain = params.select_chain.clone();296 let is_authority = parachain_config.role.clone().is_authority();297 let rpc_network = network.clone();298299 let rpc_frontier_backend = frontier_backend.clone();300301 let block_data_cache = Arc::new(fc_rpc::EthBlockDataCache::new(302 task_manager.spawn_handle(),303 overrides_handle(client.clone()),304 50,305 50,306 ));307308 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {309 let full_deps = unique_rpc::FullDeps {310 backend: rpc_frontier_backend.clone(),311 deny_unsafe,312 client: rpc_client.clone(),313 pool: rpc_pool.clone(),314 graph: rpc_pool.pool().clone(),315 316 enable_dev_signer: false,317 filter_pool: filter_pool.clone(),318 network: rpc_network.clone(),319 select_chain: select_chain.clone(),320 is_authority,321 322 max_past_logs: 10000,323 block_data_cache: block_data_cache.clone(),324 fee_history_cache: fee_history_cache.clone(),325 326 fee_history_limit: 2048,327 };328329 Ok(unique_rpc::create_full::<_, _, _, _, RuntimeApi, _>(330 full_deps,331 subscription_executor.clone(),332 ))333 });334335 task_manager.spawn_essential_handle().spawn(336 "frontier-mapping-sync-worker",337 None,338 MappingSyncWorker::new(339 client.import_notification_stream(),340 Duration::new(6, 0),341 client.clone(),342 backend.clone(),343 frontier_backend.clone(),344 SyncStrategy::Normal,345 )346 .for_each(|()| futures::future::ready(())),347 );348349 sc_service::spawn_tasks(sc_service::SpawnTasksParams {350 rpc_extensions_builder,351 client: client.clone(),352 transaction_pool: transaction_pool.clone(),353 task_manager: &mut task_manager,354 config: parachain_config,355 keystore: params.keystore_container.sync_keystore(),356 backend: backend.clone(),357 network: network.clone(),358 system_rpc_tx,359 telemetry: telemetry.as_mut(),360 })?;361362 let announce_block = {363 let network = network.clone();364 Arc::new(move |hash, data| network.announce_block(hash, data))365 };366367 if validator {368 let parachain_consensus = build_consensus(369 client.clone(),370 prometheus_registry.as_ref(),371 telemetry.as_ref().map(|t| t.handle()),372 &task_manager,373 &relay_chain_full_node,374 transaction_pool,375 network,376 params.keystore_container.sync_keystore(),377 force_authoring,378 )?;379380 let spawner = task_manager.spawn_handle();381382 let params = StartCollatorParams {383 para_id: id,384 block_status: client.clone(),385 announce_block,386 client: client.clone(),387 task_manager: &mut task_manager,388 relay_chain_full_node,389 spawner,390 parachain_consensus,391 import_queue,392 };393394 start_collator(params).await?;395 } else {396 let params = StartFullNodeParams {397 client: client.clone(),398 announce_block,399 task_manager: &mut task_manager,400 para_id: id,401 relay_chain_full_node,402 };403404 start_full_node(params)?;405 }406407 start_network.start_network();408409 Ok((task_manager, client))410}411412413pub fn parachain_build_import_queue(414 client: Arc<FullClient>,415 config: &Configuration,416 telemetry: Option<TelemetryHandle>,417 task_manager: &TaskManager,418) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {419 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;420421 cumulus_client_consensus_aura::import_queue::<422 sp_consensus_aura::sr25519::AuthorityPair,423 _,424 _,425 _,426 _,427 _,428 _,429 >(cumulus_client_consensus_aura::ImportQueueParams {430 block_import: client.clone(),431 client: client.clone(),432 create_inherent_data_providers: move |_, _| async move {433 let time = sp_timestamp::InherentDataProvider::from_system_time();434435 let slot =436 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(437 *time,438 slot_duration.slot_duration(),439 );440441 Ok((time, slot))442 },443 registry: config.prometheus_registry(),444 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),445 spawner: &task_manager.spawn_essential_handle(),446 telemetry,447 })448 .map_err(Into::into)449}450451452pub async fn start_node(453 parachain_config: Configuration,454 polkadot_config: Configuration,455 id: ParaId,456) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {457 start_node_impl::<_, _>(458 parachain_config,459 polkadot_config,460 id,461 parachain_build_import_queue,462 |client,463 prometheus_registry,464 telemetry,465 task_manager,466 relay_chain_node,467 transaction_pool,468 sync_oracle,469 keystore,470 force_authoring| {471 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;472473 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(474 task_manager.spawn_handle(),475 client.clone(),476 transaction_pool,477 prometheus_registry,478 telemetry.clone(),479 );480481 let relay_chain_backend = relay_chain_node.backend.clone();482 let relay_chain_client = relay_chain_node.client.clone();483 Ok(build_aura_consensus::<484 sp_consensus_aura::sr25519::AuthorityPair,485 _,486 _,487 _,488 _,489 _,490 _,491 _,492 _,493 _,494 >(BuildAuraConsensusParams {495 proposer_factory,496 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {497 let parachain_inherent =498 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(499 relay_parent,500 &relay_chain_client,501 &*relay_chain_backend,502 &validation_data,503 id,504 );505 async move {506 let time = sp_timestamp::InherentDataProvider::from_system_time();507508 let slot =509 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(510 *time,511 slot_duration.slot_duration(),512 );513514 let parachain_inherent = parachain_inherent.ok_or_else(|| {515 Box::<dyn std::error::Error + Send + Sync>::from(516 "Failed to create parachain inherent",517 )518 })?;519 Ok((time, slot, parachain_inherent))520 }521 },522 block_import: client.clone(),523 relay_chain_client: relay_chain_node.client.clone(),524 relay_chain_backend: relay_chain_node.backend.clone(),525 para_client: client,526 backoff_authoring_blocks: Option::<()>::None,527 sync_oracle,528 keystore,529 force_authoring,530 slot_duration,531 532 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),533 telemetry,534 max_block_proposal_slot_portion: None,535 }))536 },537 )538 .await539}