123456789use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::time::Duration;13use futures::StreamExt;141516use nft_runtime::RuntimeApi;171819use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};20use cumulus_client_consensus_common::ParachainConsensus;21use cumulus_client_network::build_block_announce_validator;22use cumulus_client_service::{23 prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,24};25use cumulus_primitives_core::ParaId;262728use sc_client_api::ExecutorProvider;29use sc_executor::NativeElseWasmExecutor;30use sc_executor::NativeExecutionDispatch;31use sc_network::NetworkService;32use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};33use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};34use sp_consensus::SlotData;35use sp_keystore::SyncCryptoStorePtr;36use sp_runtime::traits::BlakeTwo256;37use substrate_prometheus_endpoint::Registry;38use sc_client_api::BlockchainEvents;394041use fc_rpc_core::types::FilterPool;42use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};434445type BlockNumber = u32;46type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;47pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;48type Hash = sp_core::H256;495051pub struct ParachainRuntimeExecutor;5253impl NativeExecutionDispatch for ParachainRuntimeExecutor {54 type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;5556 fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {57 nft_runtime::api::dispatch(method, data)58 }5960 fn native_version() -> sc_executor::NativeVersion {61 nft_runtime::native_version()62 }63}6465pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {66 let config_dir = config67 .base_path68 .as_ref()69 .map(|base_path| base_path.config_dir(config.chain_spec.id()))70 .unwrap_or_else(|| {71 BasePath::from_project("", "", "nft").config_dir(config.chain_spec.id())72 });73 let database_dir = config_dir.join("frontier").join("db");7475 Ok(Arc::new(fc_db::Backend::<Block>::new(76 &fc_db::DatabaseSettings {77 source: fc_db::DatabaseSettingsSrc::RocksDb {78 path: database_dir,79 cache_size: 0,80 },81 },82 )?))83}8485type ExecutorDispatch = ParachainRuntimeExecutor;8687type FullClient =88 sc_service::TFullClient<Block, RuntimeApi, NativeElseWasmExecutor<ExecutorDispatch>>;89type FullBackend = sc_service::TFullBackend<Block>;90type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;919293949596#[allow(clippy::type_complexity)]97pub fn new_partial<BIQ>(98 config: &Configuration,99 build_import_queue: BIQ,100) -> Result<101 PartialComponents<102 FullClient,103 FullBackend,104 FullSelectChain,105 sc_consensus::DefaultImportQueue<Block, FullClient>,106 sc_transaction_pool::FullPool<Block, FullClient>,107 (108 Option<Telemetry>,109 Option<FilterPool>,110 Arc<fc_db::Backend<Block>>,111 Option<TelemetryWorkerHandle>,112 ),113 >,114 sc_service::Error,115>116where117 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,118 ExecutorDispatch: NativeExecutionDispatch + 'static,119 BIQ: FnOnce(120 Arc<FullClient>,121 &Configuration,122 Option<TelemetryHandle>,123 &TaskManager,124 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,125{126 let _telemetry = config127 .telemetry_endpoints128 .clone()129 .filter(|x| !x.is_empty())130 .map(|endpoints| -> Result<_, sc_telemetry::Error> {131 let worker = TelemetryWorker::new(16)?;132 let telemetry = worker.handle().new_telemetry(endpoints);133 Ok((worker, telemetry))134 })135 .transpose()?;136137 let telemetry = config138 .telemetry_endpoints139 .clone()140 .filter(|x| !x.is_empty())141 .map(|endpoints| -> Result<_, sc_telemetry::Error> {142 let worker = TelemetryWorker::new(16)?;143 let telemetry = worker.handle().new_telemetry(endpoints);144 Ok((worker, telemetry))145 })146 .transpose()?;147148 let executor = NativeElseWasmExecutor::<ExecutorDispatch>::new(149 config.wasm_method,150 config.default_heap_pages,151 config.max_runtime_instances,152 );153154 let (client, backend, keystore_container, task_manager) =155 sc_service::new_full_parts::<Block, RuntimeApi, _>(156 config,157 telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),158 executor,159 )?;160 let client = Arc::new(client);161162 let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());163164 let telemetry = telemetry.map(|(worker, telemetry)| {165 task_manager.spawn_handle().spawn("telemetry", worker.run());166 telemetry167 });168169 let select_chain = sc_consensus::LongestChain::new(backend.clone());170171 let transaction_pool = sc_transaction_pool::BasicPool::new_full(172 config.transaction_pool.clone(),173 config.role.is_authority().into(),174 config.prometheus_registry(),175 task_manager.spawn_essential_handle(),176 client.clone(),177 );178179 let filter_pool: Option<FilterPool> = Some(Arc::new(Mutex::new(BTreeMap::new())));180181 let frontier_backend = open_frontier_backend(config)?;182183 let import_queue = build_import_queue(184 client.clone(),185 config,186 telemetry.as_ref().map(|telemetry| telemetry.handle()),187 &task_manager,188 )?;189190 let params = PartialComponents {191 backend,192 client,193 import_queue,194 keystore_container,195 task_manager,196 transaction_pool,197 select_chain,198 other: (199 telemetry,200 filter_pool,201 frontier_backend,202 telemetry_worker_handle,203 ),204 };205206 Ok(params)207}208209210211212#[sc_tracing::logging::prefix_logs_with("Parachain")]213async fn start_node_impl<BIQ, BIC>(214 parachain_config: Configuration,215 polkadot_config: Configuration,216 id: ParaId,217 build_import_queue: BIQ,218 build_consensus: BIC,219) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>220where221 sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,222 ExecutorDispatch: NativeExecutionDispatch + 'static,223 BIQ: FnOnce(224 Arc<FullClient>,225 &Configuration,226 Option<TelemetryHandle>,227 &TaskManager,228 ) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,229 BIC: FnOnce(230 Arc<FullClient>,231 Option<&Registry>,232 Option<TelemetryHandle>,233 &TaskManager,234 &polkadot_service::NewFull<polkadot_service::Client>,235 Arc<sc_transaction_pool::FullPool<Block, FullClient>>,236 Arc<NetworkService<Block, Hash>>,237 SyncCryptoStorePtr,238 bool,239 ) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,240{241 if matches!(parachain_config.role, Role::Light) {242 return Err("Light client not supported!".into());243 }244245 let parachain_config = prepare_node_config(parachain_config);246247 let params = new_partial::<BIQ>(¶chain_config, build_import_queue)?;248 let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle) = params.other;249250 let relay_chain_full_node =251 cumulus_client_service::build_polkadot_full_node(polkadot_config, telemetry_worker_handle)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);272273 let (network, system_rpc_tx, start_network) =274 sc_service::build_network(sc_service::BuildNetworkParams {275 config: ¶chain_config,276 client: client.clone(),277 transaction_pool: transaction_pool.clone(),278 spawn_handle: task_manager.spawn_handle(),279 import_queue: import_queue.clone(),280 on_demand: None,281 block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),282 warp_sync: None,283 })?;284285 let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());286 let rpc_client = client.clone();287 let rpc_pool = transaction_pool.clone();288 let select_chain = params.select_chain.clone();289 let is_authority = parachain_config.role.clone().is_authority();290 let rpc_network = network.clone();291292 let rpc_frontier_backend = frontier_backend.clone();293 let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {294 let full_deps = nft_rpc::FullDeps {295 backend: rpc_frontier_backend.clone(),296 deny_unsafe,297 client: rpc_client.clone(),298 pool: rpc_pool.clone(),299 graph: rpc_pool.pool().clone(),300 301 enable_dev_signer: false,302 filter_pool: filter_pool.clone(),303 network: rpc_network.clone(),304 select_chain: select_chain.clone(),305 is_authority,306 307 max_past_logs: 10000,308 };309310 Ok(nft_rpc::create_full::<_, _, _, _, RuntimeApi, _>(311 full_deps,312 subscription_executor.clone(),313 ))314 });315316 task_manager.spawn_essential_handle().spawn(317 "frontier-mapping-sync-worker",318 MappingSyncWorker::new(319 client.import_notification_stream(),320 Duration::new(6, 0),321 client.clone(),322 backend.clone(),323 frontier_backend.clone(),324 SyncStrategy::Normal,325 )326 .for_each(|()| futures::future::ready(())),327 );328329 sc_service::spawn_tasks(sc_service::SpawnTasksParams {330 on_demand: None,331 remote_blockchain: None,332 rpc_extensions_builder,333 client: client.clone(),334 transaction_pool: transaction_pool.clone(),335 task_manager: &mut task_manager,336 config: parachain_config,337 keystore: params.keystore_container.sync_keystore(),338 backend: backend.clone(),339 network: network.clone(),340 system_rpc_tx,341 telemetry: telemetry.as_mut(),342 })?;343344 let announce_block = {345 let network = network.clone();346 Arc::new(move |hash, data| network.announce_block(hash, data))347 };348349 if validator {350 let parachain_consensus = build_consensus(351 client.clone(),352 prometheus_registry.as_ref(),353 telemetry.as_ref().map(|t| t.handle()),354 &task_manager,355 &relay_chain_full_node,356 transaction_pool,357 network,358 params.keystore_container.sync_keystore(),359 force_authoring,360 )?;361362 let spawner = task_manager.spawn_handle();363364 let params = StartCollatorParams {365 para_id: id,366 block_status: client.clone(),367 announce_block,368 client: client.clone(),369 task_manager: &mut task_manager,370 relay_chain_full_node,371 spawner,372 parachain_consensus,373 import_queue,374 };375376 start_collator(params).await?;377 } else {378 let params = StartFullNodeParams {379 client: client.clone(),380 announce_block,381 task_manager: &mut task_manager,382 para_id: id,383 relay_chain_full_node,384 };385386 start_full_node(params)?;387 }388389 start_network.start_network();390391 Ok((task_manager, client))392}393394395pub fn parachain_build_import_queue(396 client: Arc<FullClient>,397 config: &Configuration,398 telemetry: Option<TelemetryHandle>,399 task_manager: &TaskManager,400) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {401 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;402403 cumulus_client_consensus_aura::import_queue::<404 sp_consensus_aura::sr25519::AuthorityPair,405 _,406 _,407 _,408 _,409 _,410 _,411 >(cumulus_client_consensus_aura::ImportQueueParams {412 block_import: client.clone(),413 client: client.clone(),414 create_inherent_data_providers: move |_, _| async move {415 let time = sp_timestamp::InherentDataProvider::from_system_time();416417 let slot =418 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(419 *time,420 slot_duration.slot_duration(),421 );422423 Ok((time, slot))424 },425 registry: config.prometheus_registry(),426 can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),427 spawner: &task_manager.spawn_essential_handle(),428 telemetry,429 })430 .map_err(Into::into)431}432433434pub async fn start_node(435 parachain_config: Configuration,436 polkadot_config: Configuration,437 id: ParaId,438) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {439 start_node_impl::<_, _>(440 parachain_config,441 polkadot_config,442 id,443 parachain_build_import_queue,444 |client,445 prometheus_registry,446 telemetry,447 task_manager,448 relay_chain_node,449 transaction_pool,450 sync_oracle,451 keystore,452 force_authoring| {453 let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;454455 let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(456 task_manager.spawn_handle(),457 client.clone(),458 transaction_pool,459 prometheus_registry,460 telemetry.clone(),461 );462463 let relay_chain_backend = relay_chain_node.backend.clone();464 let relay_chain_client = relay_chain_node.client.clone();465 Ok(build_aura_consensus::<466 sp_consensus_aura::sr25519::AuthorityPair,467 _,468 _,469 _,470 _,471 _,472 _,473 _,474 _,475 _,476 >(BuildAuraConsensusParams {477 proposer_factory,478 create_inherent_data_providers: move |_, (relay_parent, validation_data)| {479 let parachain_inherent =480 cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(481 relay_parent,482 &relay_chain_client,483 &*relay_chain_backend,484 &validation_data,485 id,486 );487 async move {488 let time = sp_timestamp::InherentDataProvider::from_system_time();489490 let slot =491 sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(492 *time,493 slot_duration.slot_duration(),494 );495496 let parachain_inherent = parachain_inherent.ok_or_else(|| {497 Box::<dyn std::error::Error + Send + Sync>::from(498 "Failed to create parachain inherent",499 )500 })?;501 Ok((time, slot, parachain_inherent))502 }503 },504 block_import: client.clone(),505 relay_chain_client: relay_chain_node.client.clone(),506 relay_chain_backend: relay_chain_node.backend.clone(),507 para_client: client,508 backoff_authoring_blocks: Option::<()>::None,509 sync_oracle,510 keystore,511 force_authoring,512 slot_duration,513 514 block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),515 telemetry,516 max_block_proposal_slot_portion: None,517 }))518 },519 )520 .await521}