123456789use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::collections::HashMap;13use std::time::Duration;14use futures::StreamExt;151617use nft_runtime::RuntimeApi;181920use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};21use cumulus_client_consensus_common::ParachainConsensus;22use cumulus_client_network::build_block_announce_validator;23use cumulus_client_service::{24 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,25};26use cumulus_primitives_core::ParaId;272829use sc_client_api::ExecutorProvider;30use sc_executor::NativeElseWasmExecutor;31use sc_executor::NativeExecutionDispatch;32use sc_network::NetworkService;33use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};34use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};35use sp_consensus::SlotData;36use sp_keystore::SyncCryptoStorePtr;37use sp_runtime::traits::BlakeTwo256;38use substrate_prometheus_endpoint::Registry;39use sc_client_api::BlockchainEvents;404142use fc_rpc_core::types::FilterPool;43use fc_rpc_core::types::PendingTransactions;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 nft_runtime::api::dispatch(method, data)60 }6162 fn native_version() -> sc_executor::NativeVersion {63 nft_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("", "", "nft").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 PendingTransactions,112 Option<FilterPool>,113 Arc<fc_db::Backend<Block>>,114 Option<TelemetryWorkerHandle>,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_manager.spawn_handle().spawn("telemetry", worker.run());169 telemetry170 });171172 let select_chain = sc_consensus::LongestChain::new(backend.clone());173174 let transaction_pool = sc_transaction_pool::BasicPool::new_full(175 config.transaction_pool.clone(),176 config.role.is_authority().into(),177 config.prometheus_registry(),178 task_manager.spawn_essential_handle(),179 client.clone(),180 );181182 let pending_transactions: PendingTransactions = Some(Arc::new(Mutex::new(HashMap::new())));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 )?;194195 let params = PartialComponents {196 backend,197 client,198 import_queue,199 keystore_container,200 task_manager,201 transaction_pool,202 select_chain,203 other: (204 telemetry,205 pending_transactions,206 filter_pool,207 frontier_backend,208 telemetry_worker_handle,209 ),210 };211212 Ok(params)213}214215216217218#[sc_tracing::logging::prefix_logs_with("Parachain")]219async fn start_node_impl<BIQ, BIC>(220 parachain_config: Configuration,221 polkadot_config: Configuration,222 id: ParaId,223 build_import_queue: BIQ,224 build_consensus: BIC,225) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>226where227 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,228 ExecutorDispatch: NativeExecutionDispatch + 'static,229 BIQ: FnOnce(230 Arc<FullClient>,231 &Configuration,232 Option<TelemetryHandle>,233 &TaskManager,234 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,235 BIC: FnOnce(236 Arc<FullClient>,237 Option<&Registry>,238 Option<TelemetryHandle>,239 &TaskManager,240 &polkadot_service::NewFull<polkadot_service::Client>,241 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,242 Arc<NetworkService<Block, Hash>>,243 SyncCryptoStorePtr,244 bool,245 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,246{247 if matches!(parachain_config.role, Role::Light) {248 return Err("Light client not supported!".into());249 }250251 let parachain_config = prepare_node_config(parachain_config);252253 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;254 let (255 mut telemetry,256 pending_transactions,257 filter_pool,258 frontier_backend,259 telemetry_worker_handle,260 ) = params.other;261262 let relay_chain_full_node =263 cumulus_client_service::build_polkadot_full_node(polkadot_config, telemetry_worker_handle)264 .map_err(|e| match e {265 polkadot_service::Error::Sub(x) => x,266 s => format!("{}", s).into(),267 })?;268269 let client = params.client.clone();270 let backend = params.backend.clone();271 let block_announce_validator = build_block_announce_validator(272 relay_chain_full_node.client.clone(),273 id,274 Box::new(relay_chain_full_node.network.clone()),275 relay_chain_full_node.backend.clone(),276 );277278 let force_authoring = parachain_config.force_authoring;279 let validator = parachain_config.role.is_authority();280 let prometheus_registry = parachain_config.prometheus_registry().cloned();281 let transaction_pool = params.transaction_pool.clone();282 let mut task_manager = params.task_manager;283 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);284285 let (network, system_rpc_tx, start_network) =286 sc_service::build_network(sc_service::BuildNetworkParams {287 config: ¶chain_config,288 client: client.clone(),289 transaction_pool: transaction_pool.clone(),290 spawn_handle: task_manager.spawn_handle(),291 import_queue: import_queue.clone(),292 on_demand: None,293 block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),294 warp_sync: None,295 })?;296297 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());298 let rpc_client = client.clone();299 let rpc_pool = transaction_pool.clone();300 let select_chain = params.select_chain.clone();301 let is_authority = parachain_config.role.clone().is_authority();302 let rpc_network = network.clone();303304 let rpc_frontier_backend = frontier_backend.clone();305 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {306 let full_deps = nft_rpc::FullDeps {307 backend: rpc_frontier_backend.clone(),308 deny_unsafe,309 client: rpc_client.clone(),310 pool: rpc_pool.clone(),311 312 enable_dev_signer: false,313 filter_pool: filter_pool.clone(),314 network: rpc_network.clone(),315 pending_transactions: pending_transactions.clone(),316 select_chain: select_chain.clone(),317 is_authority,318 319 max_past_logs: 10000,320 };321322 Ok(nft_rpc::create_full::<_, _, _, RuntimeApi, _>(323 full_deps,324 subscription_executor.clone(),325 ))326 });327328 task_manager.spawn_essential_handle().spawn(329 "frontier-mapping-sync-worker",330 MappingSyncWorker::new(331 client.import_notification_stream(),332 Duration::new(6, 0),333 client.clone(),334 backend.clone(),335 frontier_backend.clone(),336 SyncStrategy::Normal,337 )338 .for_each(|()| futures::future::ready(())),339 );340341 sc_service::spawn_tasks(sc_service::SpawnTasksParams {342 on_demand: None,343 remote_blockchain: None,344 rpc_extensions_builder,345 client: client.clone(),346 transaction_pool: transaction_pool.clone(),347 task_manager: &mut task_manager,348 config: parachain_config,349 keystore: params.keystore_container.sync_keystore(),350 backend: backend.clone(),351 network: network.clone(),352 system_rpc_tx,353 telemetry: telemetry.as_mut(),354 })?;355356 let announce_block = {357 let network = network.clone();358 Arc::new(move |hash, data| network.announce_block(hash, data))359 };360361 if validator {362 let parachain_consensus = build_consensus(363 client.clone(),364 prometheus_registry.as_ref(),365 telemetry.as_ref().map(|t| t.handle()),366 &task_manager,367 &relay_chain_full_node,368 transaction_pool,369 network,370 params.keystore_container.sync_keystore(),371 force_authoring,372 )?;373374 let spawner = task_manager.spawn_handle();375376 let params = StartCollatorParams {377 para_id: id,378 block_status: client.clone(),379 announce_block,380 client: client.clone(),381 task_manager: &mut task_manager,382 relay_chain_full_node,383 spawner,384 parachain_consensus,385 import_queue,386 };387388 start_collator(params).await?;389 } else {390 let params = StartFullNodeParams {391 client: client.clone(),392 announce_block,393 task_manager: &mut task_manager,394 para_id: id,395 relay_chain_full_node,396 };397398 start_full_node(params)?;399 }400401 start_network.start_network();402403 Ok((task_manager, client))404}405406407pub fn parachain_build_import_queue(408 client: Arc<FullClient>,409 config: &Configuration,410 telemetry: Option<TelemetryHandle>,411 task_manager: &TaskManager,412) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {413 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;414415 cumulus_client_consensus_aura::import_queue::<416 sp_consensus_aura::sr25519::AuthorityPair,417 _,418 _,419 _,420 _,421 _,422 _,423 >(cumulus_client_consensus_aura::ImportQueueParams {424 block_import: client.clone(),425 client: client.clone(),426 create_inherent_data_providers: move |_, _| async move {427 let time = sp_timestamp::InherentDataProvider::from_system_time();428429 let slot =430 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(431 *time,432 slot_duration.slot_duration(),433 );434435 Ok((time, slot))436 },437 registry: config.prometheus_registry(),438 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),439 spawner: &task_manager.spawn_essential_handle(),440 telemetry,441 })442 .map_err(Into::into)443}444445446pub async fn start_node(447 parachain_config: Configuration,448 polkadot_config: Configuration,449 id: ParaId,450) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {451 start_node_impl::<_, _>(452 parachain_config,453 polkadot_config,454 id,455 parachain_build_import_queue,456 |client,457 prometheus_registry,458 telemetry,459 task_manager,460 relay_chain_node,461 transaction_pool,462 sync_oracle,463 keystore,464 force_authoring| {465 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;466467 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(468 task_manager.spawn_handle(),469 client.clone(),470 transaction_pool,471 prometheus_registry,472 telemetry.clone(),473 );474475 let relay_chain_backend = relay_chain_node.backend.clone();476 let relay_chain_client = relay_chain_node.client.clone();477 Ok(build_aura_consensus::<478 sp_consensus_aura::sr25519::AuthorityPair,479 _,480 _,481 _,482 _,483 _,484 _,485 _,486 _,487 _,488 >(BuildAuraConsensusParams {489 proposer_factory,490 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {491 let parachain_inherent =492 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(493 relay_parent,494 &relay_chain_client,495 &*relay_chain_backend,496 &validation_data,497 id,498 );499 async move {500 let time = sp_timestamp::InherentDataProvider::from_system_time();501502 let slot =503 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(504 *time,505 slot_duration.slot_duration(),506 );507508 let parachain_inherent = parachain_inherent.ok_or_else(|| {509 Box::<dyn std::error::Error + Send + Sync>::from(510 "Failed to create parachain inherent",511 )512 })?;513 Ok((time, slot, parachain_inherent))514 }515 },516 block_import: client.clone(),517 relay_chain_client: relay_chain_node.client.clone(),518 relay_chain_backend: relay_chain_node.backend.clone(),519 para_client: client,520 backoff_authoring_blocks: Option::<()>::None,521 sync_oracle,522 keystore,523 force_authoring,524 slot_duration,525 526 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),527 telemetry,528 max_block_proposal_slot_portion: None,529 }))530 },531 )532 .await533}