git.delta.rocks / unique-network / refs/commits / ef0b4c2a9b8b

difftreelog

source

node/cli/src/service.rs14.2 KiBsourcehistory
1//! Service and ServiceFactory implementation. Specialized wrapper over substrate service.23//4// This file is subject to the terms and conditions defined in5// file 'LICENSE', which is part of this source code package.6//78// std9use std::sync::Arc;10use std::sync::Mutex;11use std::collections::BTreeMap;12use std::collections::HashMap;1314// Local Runtime Types15use nft_runtime::RuntimeApi;1617// Cumulus Imports18use 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;2526// Polkadot Imports27use polkadot_primitives::v1::CollatorPair;2829// Substrate Imports30use 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;4041// Frontier Imports42use fc_rpc_core::types::FilterPool;43use fc_rpc_core::types::PendingTransactions;4445// Runtime type overrides46type 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;5051// Native executor instance.52native_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>;8485/// Starts a `ServiceBuilder` for a full service.86///87/// Use this macro if you don't actually need the full service, but just the builder in order to88/// be able to perform chain operations.89#[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}198199/// Start a node with the given parachain `Configuration` and relay chain `Configuration`.200///201/// This is the actual implementation that is abstract over the executor and the runtime api.202#[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>(&parachain_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: &parachain_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			// TODO: Unhardcode296			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			// TODO: Unhardcode303			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}375376/// Build the import queue for the the parachain runtime.377pub 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}414415/// Start a normal parachain node.416pub 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				// We got around 500ms for proposing498				block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),499				telemetry,500			}))501		},502	)503	.await504}