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;2829use unique_runtime::RuntimeApi;303132use cumulus_client_consensus_aura::{AuraConsensus, BuildAuraConsensusParams, SlotProportion};33use cumulus_client_consensus_common::ParachainConsensus;34use cumulus_client_service::{35 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,36};37use cumulus_client_network::BlockAnnounceValidator;38use cumulus_primitives_core::ParaId;39use cumulus_relay_chain_interface::RelayChainInterface;40use cumulus_relay_chain_local::build_relay_chain_interface;414243use sc_client_api::ExecutorProvider;44use sc_executor::NativeElseWasmExecutor;45use sc_executor::NativeExecutionDispatch;46use sc_network::NetworkService;47use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};48use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};49use sp_consensus::SlotData;50use sp_keystore::SyncCryptoStorePtr;51use sp_runtime::traits::BlakeTwo256;52use substrate_prometheus_endpoint::Registry;53use sc_client_api::BlockchainEvents;545556use fc_rpc_core::types::FilterPool;57use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};585960type BlockNumber = u32;61type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;62pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;63type Hash = sp_core::H256;646566pub struct ParachainRuntimeExecutor;6768impl NativeExecutionDispatch for ParachainRuntimeExecutor {69 type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;7071 fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {72 unique_runtime::api::dispatch(method, data)73 }7475 fn native_version() -> sc_executor::NativeVersion {76 unique_runtime::native_version()77 }78}7980pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {81 let config_dir = config82 .base_path83 .as_ref()84 .map(|base_path| base_path.config_dir(config.chain_spec.id()))85 .unwrap_or_else(|| {86 BasePath::from_project("", "", "unique").config_dir(config.chain_spec.id())87 });88 let database_dir = config_dir.join("frontier").join("db");8990 Ok(Arc::new(fc_db::Backend::<Block>::new(91 &fc_db::DatabaseSettings {92 source: fc_db::DatabaseSettingsSrc::RocksDb {93 path: database_dir,94 cache_size: 0,95 },96 },97 )?))98}99100type ExecutorDispatch = ParachainRuntimeExecutor;101102type FullClient =103 sc_service::TFullClient<Block, RuntimeApi, NativeElseWasmExecutor<ExecutorDispatch>>;104type FullBackend = sc_service::TFullBackend<Block>;105type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;106107108109110111#[allow(clippy::type_complexity)]112pub fn new_partial<BIQ>(113 config: &Configuration,114 build_import_queue: BIQ,115) -> Result<116 PartialComponents<117 FullClient,118 FullBackend,119 FullSelectChain,120 sc_consensus::DefaultImportQueue<Block, FullClient>,121 sc_transaction_pool::FullPool<Block, FullClient>,122 (123 Option<Telemetry>,124 Option<FilterPool>,125 Arc<fc_db::Backend<Block>>,126 Option<TelemetryWorkerHandle>,127 FeeHistoryCache,128 ),129 >,130 sc_service::Error,131>132where133 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,134 ExecutorDispatch: NativeExecutionDispatch + 'static,135 BIQ: FnOnce(136 Arc<FullClient>,137 &Configuration,138 Option<TelemetryHandle>,139 &TaskManager,140 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,141{142 let _telemetry = config143 .telemetry_endpoints144 .clone()145 .filter(|x| !x.is_empty())146 .map(|endpoints| -> Result<_, sc_telemetry::Error> {147 let worker = TelemetryWorker::new(16)?;148 let telemetry = worker.handle().new_telemetry(endpoints);149 Ok((worker, telemetry))150 })151 .transpose()?;152153 let telemetry = config154 .telemetry_endpoints155 .clone()156 .filter(|x| !x.is_empty())157 .map(|endpoints| -> Result<_, sc_telemetry::Error> {158 let worker = TelemetryWorker::new(16)?;159 let telemetry = worker.handle().new_telemetry(endpoints);160 Ok((worker, telemetry))161 })162 .transpose()?;163164 let executor = NativeElseWasmExecutor::<ExecutorDispatch>::new(165 config.wasm_method,166 config.default_heap_pages,167 config.max_runtime_instances,168 config.runtime_cache_size,169 );170171 let (client, backend, keystore_container, task_manager) =172 sc_service::new_full_parts::<Block, RuntimeApi, _>(173 config,174 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),175 executor,176 )?;177 let client = Arc::new(client);178179 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());180181 let telemetry = telemetry.map(|(worker, telemetry)| {182 task_manager183 .spawn_handle()184 .spawn("telemetry", None, worker.run());185 telemetry186 });187188 let select_chain = sc_consensus::LongestChain::new(backend.clone());189190 let transaction_pool = sc_transaction_pool::BasicPool::new_full(191 config.transaction_pool.clone(),192 config.role.is_authority().into(),193 config.prometheus_registry(),194 task_manager.spawn_essential_handle(),195 client.clone(),196 );197198 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));199200 let frontier_backend = open_frontier_backend(config)?;201202 let import_queue = build_import_queue(203 client.clone(),204 config,205 telemetry.as_ref().map(|telemetry| telemetry.handle()),206 &task_manager,207 )?;208 let fee_history_cache: FeeHistoryCache = Arc::new(Mutex::new(BTreeMap::new()));209210 let params = PartialComponents {211 backend,212 client,213 import_queue,214 keystore_container,215 task_manager,216 transaction_pool,217 select_chain,218 other: (219 telemetry,220 filter_pool,221 frontier_backend,222 telemetry_worker_handle,223 fee_history_cache,224 ),225 };226227 Ok(params)228}229230231232233#[sc_tracing::logging::prefix_logs_with("Parachain")]234async fn start_node_impl<BIQ, BIC>(235 parachain_config: Configuration,236 polkadot_config: Configuration,237 id: ParaId,238 build_import_queue: BIQ,239 build_consensus: BIC,240) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>241where242 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,243 ExecutorDispatch: NativeExecutionDispatch + 'static,244 BIQ: FnOnce(245 Arc<FullClient>,246 &Configuration,247 Option<TelemetryHandle>,248 &TaskManager,249 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,250 BIC: FnOnce(251 Arc<FullClient>,252 Option<&Registry>,253 Option<TelemetryHandle>,254 &TaskManager,255 Arc<dyn RelayChainInterface>,256 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,257 Arc<NetworkService<Block, Hash>>,258 SyncCryptoStorePtr,259 bool,260 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,261{262 if matches!(parachain_config.role, Role::Light) {263 return Err("Light client not supported!".into());264 }265266 let parachain_config = prepare_node_config(parachain_config);267268 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;269 let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle, fee_history_cache) =270 params.other;271272 let client = params.client.clone();273 let backend = params.backend.clone();274 let mut task_manager = params.task_manager;275276 let (relay_chain_interface, collator_key) =277 build_relay_chain_interface(polkadot_config, telemetry_worker_handle, &mut task_manager)278 .map_err(|e| match e {279 polkadot_service::Error::Sub(x) => x,280 s => format!("{}", s).into(),281 })?;282283 let block_announce_validator = BlockAnnounceValidator::new(relay_chain_interface.clone(), id);284285 let force_authoring = parachain_config.force_authoring;286 let validator = parachain_config.role.is_authority();287 let prometheus_registry = parachain_config.prometheus_registry().cloned();288 let transaction_pool = params.transaction_pool.clone();289 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);290291 let (network, system_rpc_tx, start_network) =292 sc_service::build_network(sc_service::BuildNetworkParams {293 config: ¶chain_config,294 client: client.clone(),295 transaction_pool: transaction_pool.clone(),296 spawn_handle: task_manager.spawn_handle(),297 import_queue: import_queue.clone(),298 block_announce_validator_builder: Some(Box::new(|_| {299 Box::new(block_announce_validator)300 })),301 warp_sync: None,302 })?;303304 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());305 let rpc_client = client.clone();306 let rpc_pool = transaction_pool.clone();307 let select_chain = params.select_chain.clone();308 let rpc_network = network.clone();309310 let rpc_frontier_backend = frontier_backend.clone();311312 let block_data_cache = Arc::new(fc_rpc::EthBlockDataCache::new(313 task_manager.spawn_handle(),314 overrides_handle(client.clone()),315 50,316 50,317 ));318319 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {320 let full_deps = unique_rpc::FullDeps {321 backend: rpc_frontier_backend.clone(),322 deny_unsafe,323 client: rpc_client.clone(),324 pool: rpc_pool.clone(),325 graph: rpc_pool.pool().clone(),326 327 enable_dev_signer: false,328 filter_pool: filter_pool.clone(),329 network: rpc_network.clone(),330 select_chain: select_chain.clone(),331 is_authority: validator,332 333 max_past_logs: 10000,334 block_data_cache: block_data_cache.clone(),335 fee_history_cache: fee_history_cache.clone(),336 337 fee_history_limit: 2048,338 };339340 Ok(unique_rpc::create_full::<_, _, _, _, RuntimeApi, _>(341 full_deps,342 subscription_executor.clone(),343 ))344 });345346 task_manager.spawn_essential_handle().spawn(347 "frontier-mapping-sync-worker",348 None,349 MappingSyncWorker::new(350 client.import_notification_stream(),351 Duration::new(6, 0),352 client.clone(),353 backend.clone(),354 frontier_backend.clone(),355 SyncStrategy::Normal,356 )357 .for_each(|()| futures::future::ready(())),358 );359360 sc_service::spawn_tasks(sc_service::SpawnTasksParams {361 rpc_extensions_builder,362 client: client.clone(),363 transaction_pool: transaction_pool.clone(),364 task_manager: &mut task_manager,365 config: parachain_config,366 keystore: params.keystore_container.sync_keystore(),367 backend: backend.clone(),368 network: network.clone(),369 system_rpc_tx,370 telemetry: telemetry.as_mut(),371 })?;372373 let announce_block = {374 let network = network.clone();375 Arc::new(move |hash, data| network.announce_block(hash, data))376 };377378 let relay_chain_slot_duration = Duration::from_secs(6);379380 if validator {381 let parachain_consensus = build_consensus(382 client.clone(),383 prometheus_registry.as_ref(),384 telemetry.as_ref().map(|t| t.handle()),385 &task_manager,386 relay_chain_interface.clone(),387 transaction_pool,388 network,389 params.keystore_container.sync_keystore(),390 force_authoring,391 )?;392393 let spawner = task_manager.spawn_handle();394395 let params = StartCollatorParams {396 para_id: id,397 block_status: client.clone(),398 announce_block,399 client: client.clone(),400 task_manager: &mut task_manager,401 spawner,402 parachain_consensus,403 import_queue,404 collator_key,405 relay_chain_interface,406 relay_chain_slot_duration,407 };408409 start_collator(params).await?;410 } else {411 let params = StartFullNodeParams {412 client: client.clone(),413 announce_block,414 task_manager: &mut task_manager,415 para_id: id,416 import_queue,417 relay_chain_interface,418 relay_chain_slot_duration,419 };420421 start_full_node(params)?;422 }423424 start_network.start_network();425426 Ok((task_manager, client))427}428429430pub fn parachain_build_import_queue(431 client: Arc<FullClient>,432 config: &Configuration,433 telemetry: Option<TelemetryHandle>,434 task_manager: &TaskManager,435) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {436 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;437438 cumulus_client_consensus_aura::import_queue::<439 sp_consensus_aura::sr25519::AuthorityPair,440 _,441 _,442 _,443 _,444 _,445 _,446 >(cumulus_client_consensus_aura::ImportQueueParams {447 block_import: client.clone(),448 client: client.clone(),449 create_inherent_data_providers: move |_, _| async move {450 let time = sp_timestamp::InherentDataProvider::from_system_time();451452 let slot =453 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(454 *time,455 slot_duration.slot_duration(),456 );457458 Ok((time, slot))459 },460 registry: config.prometheus_registry(),461 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),462 spawner: &task_manager.spawn_essential_handle(),463 telemetry,464 })465 .map_err(Into::into)466}467468469pub async fn start_node(470 parachain_config: Configuration,471 polkadot_config: Configuration,472 id: ParaId,473) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {474 start_node_impl::<_, _>(475 parachain_config,476 polkadot_config,477 id,478 parachain_build_import_queue,479 |client,480 prometheus_registry,481 telemetry,482 task_manager,483 relay_chain_interface,484 transaction_pool,485 sync_oracle,486 keystore,487 force_authoring| {488 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;489490 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(491 task_manager.spawn_handle(),492 client.clone(),493 transaction_pool,494 prometheus_registry,495 telemetry.clone(),496 );497498 Ok(AuraConsensus::build::<499 sp_consensus_aura::sr25519::AuthorityPair,500 _,501 _,502 _,503 _,504 _,505 _,506 >(BuildAuraConsensusParams {507 proposer_factory,508 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {509 let relay_chain_interface = relay_chain_interface.clone();510 async move {511 let parachain_inherent =512 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at(513 relay_parent,514 &relay_chain_interface,515 &validation_data,516 id,517 ).await;518519 let time = sp_timestamp::InherentDataProvider::from_system_time();520521 let slot =522 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(523 *time,524 slot_duration.slot_duration(),525 );526527 let parachain_inherent = parachain_inherent.ok_or_else(|| {528 Box::<dyn std::error::Error + Send + Sync>::from(529 "Failed to create parachain inherent",530 )531 })?;532 Ok((time, slot, parachain_inherent))533 }534 },535 block_import: client.clone(),536 para_client: client,537 backoff_authoring_blocks: Option::<()>::None,538 sync_oracle,539 keystore,540 force_authoring,541 slot_duration: *slot_duration,542 543 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),544 telemetry,545 max_block_proposal_slot_portion: None,546 }))547 },548 )549 .await550}