123456789use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::collections::HashMap;131415use nft_runtime::RuntimeApi;161718use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};19use cumulus_client_consensus_common::ParachainConsensus;20use cumulus_client_network::build_block_announce_validator;21use cumulus_client_service::{22 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,23};24use cumulus_primitives_core::ParaId;252627use polkadot_primitives::v1::CollatorPair;282930use sc_client_api::ExecutorProvider;31pub use sc_executor::NativeExecutor;32use sc_executor::native_executor_instance;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;404142use fc_rpc_core::types::FilterPool;43use fc_rpc_core::types::PendingTransactions;444546type BlockNumber = u32;47type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;48pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;49type Hash = sp_core::H256;505152native_executor_instance!(53 pub ParachainRuntimeExecutor,54 nft_runtime::api::dispatch,55 nft_runtime::native_version,56 frame_benchmarking::benchmarking::HostFunctions,57);5859pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {60 let config_dir = config61 .base_path62 .as_ref()63 .map(|base_path| base_path.config_dir(config.chain_spec.id()))64 .unwrap_or_else(|| {65 BasePath::from_project("", "", "nft").config_dir(config.chain_spec.id())66 });67 let database_dir = config_dir.join("frontier").join("db");6869 Ok(Arc::new(fc_db::Backend::<Block>::new(70 &fc_db::DatabaseSettings {71 source: fc_db::DatabaseSettingsSrc::RocksDb {72 path: database_dir,73 cache_size: 0,74 },75 },76 )?))77}7879type Executor = ParachainRuntimeExecutor;8081type FullClient = sc_service::TFullClient<Block, RuntimeApi, Executor>;82type FullBackend = sc_service::TFullBackend<Block>;83type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;848586878889#[allow(clippy::type_complexity)]90pub fn new_partial<BIQ>(91 config: &Configuration,92 build_import_queue: BIQ,93) -> Result<94 PartialComponents<95 FullClient,96 FullBackend,97 FullSelectChain,98 sp_consensus::DefaultImportQueue<Block, FullClient>,99 sc_transaction_pool::FullPool<Block, FullClient>,100 (101 Option<Telemetry>,102 PendingTransactions,103 Option<FilterPool>,104 Arc<fc_db::Backend<Block>>,105 Option<TelemetryWorkerHandle>,106 ),107 >,108 sc_service::Error,109>110where111 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,112 Executor: sc_executor::NativeExecutionDispatch + 'static,113 BIQ: FnOnce(114 Arc<FullClient>,115 &Configuration,116 Option<TelemetryHandle>,117 &TaskManager,118 ) -> Result<sp_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,119{120 let _telemetry = config121 .telemetry_endpoints122 .clone()123 .filter(|x| !x.is_empty())124 .map(|endpoints| -> Result<_, sc_telemetry::Error> {125 let worker = TelemetryWorker::new(16)?;126 let telemetry = worker.handle().new_telemetry(endpoints);127 Ok((worker, telemetry))128 })129 .transpose()?;130131 let telemetry = config132 .telemetry_endpoints133 .clone()134 .filter(|x| !x.is_empty())135 .map(|endpoints| -> Result<_, sc_telemetry::Error> {136 let worker = TelemetryWorker::new(16)?;137 let telemetry = worker.handle().new_telemetry(endpoints);138 Ok((worker, telemetry))139 })140 .transpose()?;141142 let (client, backend, keystore_container, task_manager) =143 sc_service::new_full_parts::<Block, RuntimeApi, Executor>(144 &config,145 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),146 )?;147 let client = Arc::new(client);148149 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());150151 let telemetry = telemetry.map(|(worker, telemetry)| {152 task_manager.spawn_handle().spawn("telemetry", worker.run());153 telemetry154 });155156 let select_chain = sc_consensus::LongestChain::new(backend.clone());157158 let transaction_pool = sc_transaction_pool::BasicPool::new_full(159 config.transaction_pool.clone(),160 config.role.is_authority().into(),161 config.prometheus_registry(),162 task_manager.spawn_handle(),163 client.clone(),164 );165166 let pending_transactions: PendingTransactions = Some(Arc::new(Mutex::new(HashMap::new())));167168 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));169170 let frontier_backend = open_frontier_backend(config)?;171172 let import_queue = build_import_queue(173 client.clone(),174 config,175 telemetry.as_ref().map(|telemetry| telemetry.handle()),176 &task_manager,177 )?;178179 let params = PartialComponents {180 backend,181 client,182 import_queue,183 keystore_container,184 task_manager,185 transaction_pool,186 select_chain,187 other: (188 telemetry,189 pending_transactions,190 filter_pool,191 frontier_backend,192 telemetry_worker_handle,193 ),194 };195196 Ok(params)197}198199200201202#[sc_tracing::logging::prefix_logs_with("Parachain")]203async fn start_node_impl<BIQ, BIC>(204 parachain_config: Configuration,205 collator_key: CollatorPair,206 polkadot_config: Configuration,207 id: ParaId,208 build_import_queue: BIQ,209 build_consensus: BIC,210) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>211where212 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,213 Executor: sc_executor::NativeExecutionDispatch + 'static,214 BIQ: FnOnce(215 Arc<FullClient>,216 &Configuration,217 Option<TelemetryHandle>,218 &TaskManager,219 ) -> Result<sp_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,220 BIC: FnOnce(221 Arc<FullClient>,222 Option<&Registry>,223 Option<TelemetryHandle>,224 &TaskManager,225 &polkadot_service::NewFull<polkadot_service::Client>,226 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,227 Arc<NetworkService<Block, Hash>>,228 SyncCryptoStorePtr,229 bool,230 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,231{232 if matches!(parachain_config.role, Role::Light) {233 return Err("Light client not supported!".into());234 }235236 let parachain_config = prepare_node_config(parachain_config);237238 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;239 let (240 mut telemetry,241 pending_transactions,242 filter_pool,243 frontier_backend,244 telemetry_worker_handle,245 ) = params.other;246247 let relay_chain_full_node = cumulus_client_service::build_polkadot_full_node(248 polkadot_config,249 collator_key.clone(),250 telemetry_worker_handle,251 )252 .map_err(|e| match e {253 polkadot_service::Error::Sub(x) => x,254 s => format!("{}", s).into(),255 })?;256257 let client = params.client.clone();258 let backend = params.backend.clone();259 let block_announce_validator = build_block_announce_validator(260 relay_chain_full_node.client.clone(),261 id,262 Box::new(relay_chain_full_node.network.clone()),263 relay_chain_full_node.backend.clone(),264 );265266 let force_authoring = parachain_config.force_authoring;267 let validator = parachain_config.role.is_authority();268 let prometheus_registry = parachain_config.prometheus_registry().cloned();269 let transaction_pool = params.transaction_pool.clone();270 let mut task_manager = params.task_manager;271 let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);272 let (network, network_status_sinks, system_rpc_tx, start_network) =273 sc_service::build_network(sc_service::BuildNetworkParams {274 config: ¶chain_config,275 client: client.clone(),276 transaction_pool: transaction_pool.clone(),277 spawn_handle: task_manager.spawn_handle(),278 import_queue: import_queue.clone(),279 on_demand: None,280 block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),281 })?;282283 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());284 let rpc_client = client.clone();285 let rpc_pool = transaction_pool.clone();286 let select_chain = params.select_chain.clone();287 let is_authority = parachain_config.role.clone().is_authority();288 let rpc_network = network.clone();289 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {290 let full_deps = nft_rpc::FullDeps {291 backend: frontier_backend.clone(),292 deny_unsafe,293 client: rpc_client.clone(),294 pool: rpc_pool.clone(),295 296 enable_dev_signer: false,297 filter_pool: filter_pool.clone(),298 network: rpc_network.clone(),299 pending_transactions: pending_transactions.clone(),300 select_chain: select_chain.clone(),301 is_authority,302 303 max_past_logs: 10000,304 };305306 nft_rpc::create_full::<_, _, _, RuntimeApi, _>(full_deps, subscription_executor.clone())307 });308309 sc_service::spawn_tasks(sc_service::SpawnTasksParams {310 on_demand: None,311 remote_blockchain: None,312 rpc_extensions_builder,313 client: client.clone(),314 transaction_pool: transaction_pool.clone(),315 task_manager: &mut task_manager,316 config: parachain_config,317 keystore: params.keystore_container.sync_keystore(),318 backend: backend.clone(),319 network: network.clone(),320 network_status_sinks,321 system_rpc_tx,322 telemetry: telemetry.as_mut(),323 })?;324325 let announce_block = {326 let network = network.clone();327 Arc::new(move |hash, data| network.announce_block(hash, data))328 };329330 if validator {331 let parachain_consensus = build_consensus(332 client.clone(),333 prometheus_registry.as_ref(),334 telemetry.as_ref().map(|t| t.handle()),335 &task_manager,336 &relay_chain_full_node,337 transaction_pool,338 network,339 params.keystore_container.sync_keystore(),340 force_authoring,341 )?;342343 let spawner = task_manager.spawn_handle();344345 let params = StartCollatorParams {346 para_id: id,347 block_status: client.clone(),348 announce_block,349 client: client.clone(),350 task_manager: &mut task_manager,351 collator_key,352 relay_chain_full_node,353 spawner,354 parachain_consensus,355 import_queue,356 };357358 start_collator(params).await?;359 } else {360 let params = StartFullNodeParams {361 client: client.clone(),362 announce_block,363 task_manager: &mut task_manager,364 para_id: id,365 relay_chain_full_node,366 };367368 start_full_node(params)?;369 }370371 start_network.start_network();372373 Ok((task_manager, client))374}375376377pub fn parachain_build_import_queue(378 client: Arc<FullClient>,379 config: &Configuration,380 telemetry: Option<TelemetryHandle>,381 task_manager: &TaskManager,382) -> Result<sp_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {383 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;384385 cumulus_client_consensus_aura::import_queue::<386 sp_consensus_aura::sr25519::AuthorityPair,387 _,388 _,389 _,390 _,391 _,392 _,393 >(cumulus_client_consensus_aura::ImportQueueParams {394 block_import: client.clone(),395 client: client.clone(),396 create_inherent_data_providers: move |_, _| async move {397 let time = sp_timestamp::InherentDataProvider::from_system_time();398399 let slot =400 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(401 *time,402 slot_duration.slot_duration(),403 );404405 Ok((time, slot))406 },407 registry: config.prometheus_registry(),408 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),409 spawner: &task_manager.spawn_essential_handle(),410 telemetry,411 })412 .map_err(Into::into)413}414415416pub async fn start_node(417 parachain_config: Configuration,418 collator_key: CollatorPair,419 polkadot_config: Configuration,420 id: ParaId,421) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {422 start_node_impl::<_, _>(423 parachain_config,424 collator_key,425 polkadot_config,426 id,427 parachain_build_import_queue,428 |client,429 prometheus_registry,430 telemetry,431 task_manager,432 relay_chain_node,433 transaction_pool,434 sync_oracle,435 keystore,436 force_authoring| {437 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;438439 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(440 task_manager.spawn_handle(),441 client.clone(),442 transaction_pool,443 prometheus_registry,444 telemetry.clone(),445 );446447 let relay_chain_backend = relay_chain_node.backend.clone();448 let relay_chain_client = relay_chain_node.client.clone();449 Ok(build_aura_consensus::<450 sp_consensus_aura::sr25519::AuthorityPair,451 _,452 _,453 _,454 _,455 _,456 _,457 _,458 _,459 _,460 >(BuildAuraConsensusParams {461 proposer_factory,462 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {463 let parachain_inherent =464 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(465 relay_parent,466 &relay_chain_client,467 &*relay_chain_backend,468 &validation_data,469 id,470 );471 async move {472 let time = sp_timestamp::InherentDataProvider::from_system_time();473474 let slot =475 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(476 *time,477 slot_duration.slot_duration(),478 );479480 let parachain_inherent = parachain_inherent.ok_or_else(|| {481 Box::<dyn std::error::Error + Send + Sync>::from(482 "Failed to create parachain inherent",483 )484 })?;485 Ok((time, slot, parachain_inherent))486 }487 },488 block_import: client.clone(),489 relay_chain_client: relay_chain_node.client.clone(),490 relay_chain_backend: relay_chain_node.backend.clone(),491 para_client: client,492 backoff_authoring_blocks: Option::<()>::None,493 sync_oracle,494 keystore,495 force_authoring,496 slot_duration,497 498 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),499 telemetry,500 }))501 },502 )503 .await504}