1234567891011121314151617181920use std::sync::Arc;21use std::sync::Mutex;22use std::collections::BTreeMap;23use std::time::Duration;24use fc_rpc_core::types::FeeHistoryCache;25use futures::StreamExt;2627use unique_rpc::overrides_handle;2829#[cfg(feature = "unique-runtime")]30use unique_runtime as runtime;3132#[cfg(feature = "quartz-runtime")]33use quartz_runtime as runtime;3435#[cfg(feature = "opal-runtime")]36use opal_runtime as runtime;3738use runtime::RuntimeApi;394041use cumulus_client_consensus_aura::{AuraConsensus, BuildAuraConsensusParams, SlotProportion};42use cumulus_client_consensus_common::ParachainConsensus;43use cumulus_client_service::{44 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,45};46use cumulus_client_network::BlockAnnounceValidator;47use cumulus_primitives_core::ParaId;48use cumulus_relay_chain_interface::RelayChainInterface;49use cumulus_relay_chain_local::build_relay_chain_interface;505152use sc_client_api::ExecutorProvider;53use sc_executor::NativeElseWasmExecutor;54use sc_executor::NativeExecutionDispatch;55use sc_network::NetworkService;56use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};57use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};58use sp_consensus::SlotData;59use sp_keystore::SyncCryptoStorePtr;60use sp_runtime::traits::BlakeTwo256;61use substrate_prometheus_endpoint::Registry;62use sc_client_api::BlockchainEvents;636465use fc_rpc_core::types::FilterPool;66use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};676869type BlockNumber = u32;70type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;71pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;72type Hash = sp_core::H256;737475pub struct ParachainRuntimeExecutor;7677impl NativeExecutionDispatch for ParachainRuntimeExecutor {78 type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;7980 fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {81 runtime::api::dispatch(method, data)82 }8384 fn native_version() -> sc_executor::NativeVersion {85 runtime::native_version()86 }87}8889pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {90 let config_dir = config91 .base_path92 .as_ref()93 .map(|base_path| base_path.config_dir(config.chain_spec.id()))94 .unwrap_or_else(|| {95 BasePath::from_project("", "", "unique").config_dir(config.chain_spec.id())96 });97 let database_dir = config_dir.join("frontier").join("db");9899 Ok(Arc::new(fc_db::Backend::<Block>::new(100 &fc_db::DatabaseSettings {101 source: fc_db::DatabaseSettingsSrc::RocksDb {102 path: database_dir,103 cache_size: 0,104 },105 },106 )?))107}108109type ExecutorDispatch = ParachainRuntimeExecutor;110111type FullClient =112 sc_service::TFullClient<Block, RuntimeApi, NativeElseWasmExecutor<ExecutorDispatch>>;113type FullBackend = sc_service::TFullBackend<Block>;114type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;115116117118119120#[allow(clippy::type_complexity)]121pub fn new_partial<BIQ>(122 config: &Configuration,123 build_import_queue: BIQ,124) -> Result<125 PartialComponents<126 FullClient,127 FullBackend,128 FullSelectChain,129 sc_consensus::DefaultImportQueue<Block, FullClient>,130 sc_transaction_pool::FullPool<Block, FullClient>,131 (132 Option<Telemetry>,133 Option<FilterPool>,134 Arc<fc_db::Backend<Block>>,135 Option<TelemetryWorkerHandle>,136 FeeHistoryCache,137 ),138 >,139 sc_service::Error,140>141where142 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,143 ExecutorDispatch: NativeExecutionDispatch + 'static,144 BIQ: FnOnce(145 Arc<FullClient>,146 &Configuration,147 Option<TelemetryHandle>,148 &TaskManager,149 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,150{151 let _telemetry = config152 .telemetry_endpoints153 .clone()154 .filter(|x| !x.is_empty())155 .map(|endpoints| -> Result<_, sc_telemetry::Error> {156 let worker = TelemetryWorker::new(16)?;157 let telemetry = worker.handle().new_telemetry(endpoints);158 Ok((worker, telemetry))159 })160 .transpose()?;161162 let telemetry = config163 .telemetry_endpoints164 .clone()165 .filter(|x| !x.is_empty())166 .map(|endpoints| -> Result<_, sc_telemetry::Error> {167 let worker = TelemetryWorker::new(16)?;168 let telemetry = worker.handle().new_telemetry(endpoints);169 Ok((worker, telemetry))170 })171 .transpose()?;172173 let executor = NativeElseWasmExecutor::<ExecutorDispatch>::new(174 config.wasm_method,175 config.default_heap_pages,176 config.max_runtime_instances,177 config.runtime_cache_size,178 );179180 let (client, backend, keystore_container, task_manager) =181 sc_service::new_full_parts::<Block, RuntimeApi, _>(182 config,183 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),184 executor,185 )?;186 let client = Arc::new(client);187188 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());189190 let telemetry = telemetry.map(|(worker, telemetry)| {191 task_manager192 .spawn_handle()193 .spawn("telemetry", None, worker.run());194 telemetry195 });196197 let select_chain = sc_consensus::LongestChain::new(backend.clone());198199 let transaction_pool = sc_transaction_pool::BasicPool::new_full(200 config.transaction_pool.clone(),201 config.role.is_authority().into(),202 config.prometheus_registry(),203 task_manager.spawn_essential_handle(),204 client.clone(),205 );206207 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));208209 let frontier_backend = open_frontier_backend(config)?;210211 let import_queue = build_import_queue(212 client.clone(),213 config,214 telemetry.as_ref().map(|telemetry| telemetry.handle()),215 &task_manager,216 )?;217 let fee_history_cache: FeeHistoryCache = Arc::new(Mutex::new(BTreeMap::new()));218219 let params = PartialComponents {220 backend,221 client,222 import_queue,223 keystore_container,224 task_manager,225 transaction_pool,226 select_chain,227 other: (228 telemetry,229 filter_pool,230 frontier_backend,231 telemetry_worker_handle,232 fee_history_cache,233 ),234 };235236 Ok(params)237}238239240241242#[sc_tracing::logging::prefix_logs_with("Parachain")]243async fn start_node_impl<BIQ, BIC>(244 parachain_config: Configuration,245 polkadot_config: Configuration,246 id: ParaId,247 build_import_queue: BIQ,248 build_consensus: BIC,249) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>250where251 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,252 ExecutorDispatch: NativeExecutionDispatch + 'static,253 BIQ: FnOnce(254 Arc<FullClient>,255 &Configuration,256 Option<TelemetryHandle>,257 &TaskManager,258 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,259 BIC: FnOnce(260 Arc<FullClient>,261 Option<&Registry>,262 Option<TelemetryHandle>,263 &TaskManager,264 Arc<dyn RelayChainInterface>,265 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,266 Arc<NetworkService<Block, Hash>>,267 SyncCryptoStorePtr,268 bool,269 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,270{271 if matches!(parachain_config.role, Role::Light) {272 return Err("Light client not supported!".into());273 }274275 let parachain_config = prepare_node_config(parachain_config);276277 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;278 let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle, fee_history_cache) =279 params.other;280281 let client = params.client.clone();282 let backend = params.backend.clone();283 let mut task_manager = params.task_manager;284285 let (relay_chain_interface, collator_key) =286 build_relay_chain_interface(polkadot_config, telemetry_worker_handle, &mut task_manager)287 .map_err(|e| match e {288 polkadot_service::Error::Sub(x) => x,289 s => format!("{}", s).into(),290 })?;291292 let block_announce_validator = BlockAnnounceValidator::new(relay_chain_interface.clone(), id);293294 let force_authoring = parachain_config.force_authoring;295 let validator = parachain_config.role.is_authority();296 let prometheus_registry = parachain_config.prometheus_registry().cloned();297 let transaction_pool = params.transaction_pool.clone();298 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);299300 let (network, system_rpc_tx, start_network) =301 sc_service::build_network(sc_service::BuildNetworkParams {302 config: ¶chain_config,303 client: client.clone(),304 transaction_pool: transaction_pool.clone(),305 spawn_handle: task_manager.spawn_handle(),306 import_queue: import_queue.clone(),307 block_announce_validator_builder: Some(Box::new(|_| {308 Box::new(block_announce_validator)309 })),310 warp_sync: None,311 })?;312313 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());314 let rpc_client = client.clone();315 let rpc_pool = transaction_pool.clone();316 let select_chain = params.select_chain.clone();317 let rpc_network = network.clone();318319 let rpc_frontier_backend = frontier_backend.clone();320321 let block_data_cache = Arc::new(fc_rpc::EthBlockDataCache::new(322 task_manager.spawn_handle(),323 overrides_handle(client.clone()),324 50,325 50,326 ));327328 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {329 let full_deps = unique_rpc::FullDeps {330 backend: rpc_frontier_backend.clone(),331 deny_unsafe,332 client: rpc_client.clone(),333 pool: rpc_pool.clone(),334 graph: rpc_pool.pool().clone(),335 336 enable_dev_signer: false,337 filter_pool: filter_pool.clone(),338 network: rpc_network.clone(),339 select_chain: select_chain.clone(),340 is_authority: validator,341 342 max_past_logs: 10000,343 block_data_cache: block_data_cache.clone(),344 fee_history_cache: fee_history_cache.clone(),345 346 fee_history_limit: 2048,347 };348349 Ok(unique_rpc::create_full::<_, _, _, _, RuntimeApi, _>(350 full_deps,351 subscription_executor.clone(),352 ))353 });354355 task_manager.spawn_essential_handle().spawn(356 "frontier-mapping-sync-worker",357 None,358 MappingSyncWorker::new(359 client.import_notification_stream(),360 Duration::new(6, 0),361 client.clone(),362 backend.clone(),363 frontier_backend.clone(),364 SyncStrategy::Normal,365 )366 .for_each(|()| futures::future::ready(())),367 );368369 sc_service::spawn_tasks(sc_service::SpawnTasksParams {370 rpc_extensions_builder,371 client: client.clone(),372 transaction_pool: transaction_pool.clone(),373 task_manager: &mut task_manager,374 config: parachain_config,375 keystore: params.keystore_container.sync_keystore(),376 backend: backend.clone(),377 network: network.clone(),378 system_rpc_tx,379 telemetry: telemetry.as_mut(),380 })?;381382 let announce_block = {383 let network = network.clone();384 Arc::new(move |hash, data| network.announce_block(hash, data))385 };386387 let relay_chain_slot_duration = Duration::from_secs(6);388389 if validator {390 let parachain_consensus = build_consensus(391 client.clone(),392 prometheus_registry.as_ref(),393 telemetry.as_ref().map(|t| t.handle()),394 &task_manager,395 relay_chain_interface.clone(),396 transaction_pool,397 network,398 params.keystore_container.sync_keystore(),399 force_authoring,400 )?;401402 let spawner = task_manager.spawn_handle();403404 let params = StartCollatorParams {405 para_id: id,406 block_status: client.clone(),407 announce_block,408 client: client.clone(),409 task_manager: &mut task_manager,410 spawner,411 parachain_consensus,412 import_queue,413 collator_key,414 relay_chain_interface,415 relay_chain_slot_duration,416 };417418 start_collator(params).await?;419 } else {420 let params = StartFullNodeParams {421 client: client.clone(),422 announce_block,423 task_manager: &mut task_manager,424 para_id: id,425 import_queue,426 relay_chain_interface,427 relay_chain_slot_duration,428 };429430 start_full_node(params)?;431 }432433 start_network.start_network();434435 Ok((task_manager, client))436}437438439pub fn parachain_build_import_queue(440 client: Arc<FullClient>,441 config: &Configuration,442 telemetry: Option<TelemetryHandle>,443 task_manager: &TaskManager,444) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {445 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;446447 cumulus_client_consensus_aura::import_queue::<448 sp_consensus_aura::sr25519::AuthorityPair,449 _,450 _,451 _,452 _,453 _,454 _,455 >(cumulus_client_consensus_aura::ImportQueueParams {456 block_import: client.clone(),457 client: client.clone(),458 create_inherent_data_providers: move |_, _| async move {459 let time = sp_timestamp::InherentDataProvider::from_system_time();460461 let slot =462 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(463 *time,464 slot_duration.slot_duration(),465 );466467 Ok((time, slot))468 },469 registry: config.prometheus_registry(),470 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),471 spawner: &task_manager.spawn_essential_handle(),472 telemetry,473 })474 .map_err(Into::into)475}476477478pub async fn start_node(479 parachain_config: Configuration,480 polkadot_config: Configuration,481 id: ParaId,482) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {483 start_node_impl::<_, _>(484 parachain_config,485 polkadot_config,486 id,487 parachain_build_import_queue,488 |client,489 prometheus_registry,490 telemetry,491 task_manager,492 relay_chain_interface,493 transaction_pool,494 sync_oracle,495 keystore,496 force_authoring| {497 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;498499 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(500 task_manager.spawn_handle(),501 client.clone(),502 transaction_pool,503 prometheus_registry,504 telemetry.clone(),505 );506507 Ok(AuraConsensus::build::<508 sp_consensus_aura::sr25519::AuthorityPair,509 _,510 _,511 _,512 _,513 _,514 _,515 >(BuildAuraConsensusParams {516 proposer_factory,517 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {518 let relay_chain_interface = relay_chain_interface.clone();519 async move {520 let parachain_inherent =521 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at(522 relay_parent,523 &relay_chain_interface,524 &validation_data,525 id,526 ).await;527528 let time = sp_timestamp::InherentDataProvider::from_system_time();529530 let slot =531 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(532 *time,533 slot_duration.slot_duration(),534 );535536 let parachain_inherent = parachain_inherent.ok_or_else(|| {537 Box::<dyn std::error::Error + Send + Sync>::from(538 "Failed to create parachain inherent",539 )540 })?;541 Ok((time, slot, parachain_inherent))542 }543 },544 block_import: client.clone(),545 para_client: client,546 backoff_authoring_blocks: Option::<()>::None,547 sync_oracle,548 keystore,549 force_authoring,550 slot_duration: *slot_duration,551 552 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),553 telemetry,554 max_block_proposal_slot_portion: None,555 }))556 },557 )558 .await559}