git.delta.rocks / unique-network / refs/commits / 90fcef986d5a

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::{19	build_aura_consensus, BuildAuraConsensusParams, SlotProportion,20};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;2728// Polkadot Imports29use polkadot_primitives::v1::CollatorPair;3031// Substrate Imports32use sc_client_api::ExecutorProvider;33pub use sc_executor::NativeExecutor;34use sc_executor::native_executor_instance;35use sc_network::NetworkService;36use sc_service::{BasePath, Configuration, PartialComponents, Role, TaskManager};37use sc_telemetry::{Telemetry, TelemetryHandle, TelemetryWorker, TelemetryWorkerHandle};38use sp_consensus::SlotData;39use sp_keystore::SyncCryptoStorePtr;40use sp_runtime::traits::BlakeTwo256;41use substrate_prometheus_endpoint::Registry;4243// Frontier Imports44use fc_rpc_core::types::FilterPool;45use fc_rpc_core::types::PendingTransactions;4647// Runtime type overrides48type BlockNumber = u32;49type Header = sp_runtime::generic::Header<BlockNumber, sp_runtime::traits::BlakeTwo256>;50pub type Block = sp_runtime::generic::Block<Header, sp_runtime::OpaqueExtrinsic>;51type Hash = sp_core::H256;5253// Native executor instance.54native_executor_instance!(55	pub ParachainRuntimeExecutor,56	nft_runtime::api::dispatch,57	nft_runtime::native_version,58	frame_benchmarking::benchmarking::HostFunctions,59);6061pub fn open_frontier_backend(config: &Configuration) -> Result<Arc<fc_db::Backend<Block>>, String> {62	let config_dir = config.base_path.as_ref()63		.map(|base_path| base_path.config_dir(config.chain_spec.id()))64		.unwrap_or_else(|| {65			BasePath::from_project("", "", "nft")66				.config_dir(config.chain_spec.id())67		});68	let database_dir = config_dir.join("frontier").join("db");6970	Ok(Arc::new(fc_db::Backend::<Block>::new(&fc_db::DatabaseSettings {71		source: fc_db::DatabaseSettingsSrc::RocksDb {72			path: database_dir,73			cache_size: 0,74		}75	})?))76}7778type Executor = ParachainRuntimeExecutor;7980type FullClient = sc_service::TFullClient<Block, RuntimeApi, Executor>;81type FullBackend = sc_service::TFullBackend<Block>;82type FullSelectChain = sc_consensus::LongestChain<FullBackend, Block>;8384/// Starts a `ServiceBuilder` for a full service.85///86/// Use this macro if you don't actually need the full service, but just the builder in order to87/// be able to perform chain operations.88pub fn new_partial<BIQ>(89	config: &Configuration,90	build_import_queue: BIQ,91) -> Result<92	PartialComponents<93		FullClient,94		FullBackend,95		FullSelectChain,96		sp_consensus::DefaultImportQueue<Block, FullClient>,97		sc_transaction_pool::FullPool<Block, FullClient>,98		(Option<Telemetry>, PendingTransactions, Option<FilterPool>, Arc<fc_db::Backend<Block>>, Option<TelemetryWorkerHandle>),99	>,100	sc_service::Error,101>102where103	sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,104	Executor: sc_executor::NativeExecutionDispatch + 'static,105	BIQ: FnOnce(106		Arc<FullClient>,107		&Configuration,108		Option<TelemetryHandle>,109		&TaskManager,110	) -> Result<111		sp_consensus::DefaultImportQueue<Block, FullClient>,112		sc_service::Error,113	>,114{115	let telemetry = config116		.telemetry_endpoints117		.clone()118		.filter(|x| !x.is_empty())119		.map(|endpoints| -> Result<_, sc_telemetry::Error> {120			let worker = TelemetryWorker::new(16)?;121			let telemetry = worker.handle().new_telemetry(endpoints);122			Ok((worker, telemetry))123		})124		.transpose()?;125126	let telemetry = config.telemetry_endpoints.clone()127		.filter(|x| !x.is_empty())128		.map(|endpoints| -> Result<_, sc_telemetry::Error> {129			let worker = TelemetryWorker::new(16)?;130			let telemetry = worker.handle().new_telemetry(endpoints);131			Ok((worker, telemetry))132		})133		.transpose()?;134135	let (client, backend, keystore_container, task_manager) =136		sc_service::new_full_parts::<Block, RuntimeApi, Executor>(137			&config,138			telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()),139		)?;140	let client = Arc::new(client);141142	let telemetry_worker_handle = telemetry.as_ref().map(|(worker, _)| worker.handle());143144	let telemetry = telemetry.map(|(worker, telemetry)| {145		task_manager.spawn_handle().spawn("telemetry", worker.run());146		telemetry147	});148149	let select_chain = sc_consensus::LongestChain::new(backend.clone());150151	let transaction_pool = sc_transaction_pool::BasicPool::new_full(152		config.transaction_pool.clone(),153		config.role.is_authority().into(),154		config.prometheus_registry(),155		task_manager.spawn_handle(),156		client.clone(),157	);158159	let pending_transactions: PendingTransactions = Some(Arc::new(Mutex::new(HashMap::new())));160161	let filter_pool: Option<FilterPool>162		= Some(Arc::new(Mutex::new(BTreeMap::new())));163164	let frontier_backend = open_frontier_backend(config)?;165166	let import_queue = build_import_queue(167		client.clone(),168		config,169		telemetry.as_ref().map(|telemetry| telemetry.handle()),170		&task_manager,171	)?;172173	let params = PartialComponents {174		backend,175		client,176		import_queue,177		keystore_container,178		task_manager,179		transaction_pool,180		select_chain,181		other: (telemetry, pending_transactions, filter_pool, frontier_backend, telemetry_worker_handle),182	};183184	Ok(params)185}186187/// Start a node with the given parachain `Configuration` and relay chain `Configuration`.188///189/// This is the actual implementation that is abstract over the executor and the runtime api.190#[sc_tracing::logging::prefix_logs_with("Parachain")]191async fn start_node_impl<BIQ, BIC>(192	parachain_config: Configuration,193	collator_key: CollatorPair,194	polkadot_config: Configuration,195	id: ParaId,196	build_import_queue: BIQ,197	build_consensus: BIC,198) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>199where200	sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,201	Executor: sc_executor::NativeExecutionDispatch + 'static,202	BIQ: FnOnce(203		Arc<FullClient>,204		&Configuration,205		Option<TelemetryHandle>,206		&TaskManager,207	) -> Result<208		sp_consensus::DefaultImportQueue<Block, FullClient>,209		sc_service::Error,210	>,211	BIC: FnOnce(212		Arc<FullClient>,213		Option<&Registry>,214		Option<TelemetryHandle>,215		&TaskManager,216		&polkadot_service::NewFull<polkadot_service::Client>,217		Arc<sc_transaction_pool::FullPool<Block, FullClient>>,218		Arc<NetworkService<Block, Hash>>,219		SyncCryptoStorePtr,220		bool,221	) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,222{223	if matches!(parachain_config.role, Role::Light) {224		return Err("Light client not supported!".into());225	}226227	let parachain_config = prepare_node_config(parachain_config);228229	let params = new_partial::<BIQ>(&parachain_config, build_import_queue)?;230	let (mut telemetry, pending_transactions, filter_pool, frontier_backend, telemetry_worker_handle) = params.other;231232	let relay_chain_full_node = cumulus_client_service::build_polkadot_full_node(233		polkadot_config,234		collator_key.clone(),235		telemetry_worker_handle,236	)237	.map_err(|e| match e {238		polkadot_service::Error::Sub(x) => x,239		s => format!("{}", s).into(),240	})?;241242	let client = params.client.clone();243	let backend = params.backend.clone();244	let block_announce_validator = build_block_announce_validator(245		relay_chain_full_node.client.clone(),246		id,247		Box::new(relay_chain_full_node.network.clone()),248		relay_chain_full_node.backend.clone(),249	);250251	let force_authoring = parachain_config.force_authoring;252	let validator = parachain_config.role.is_authority();253	let prometheus_registry = parachain_config.prometheus_registry().cloned();254	let transaction_pool = params.transaction_pool.clone();255	let mut task_manager = params.task_manager;256	let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);257	let (network, network_status_sinks, system_rpc_tx, start_network) =258		sc_service::build_network(sc_service::BuildNetworkParams {259			config: &parachain_config,260			client: client.clone(),261			transaction_pool: transaction_pool.clone(),262			spawn_handle: task_manager.spawn_handle(),263			import_queue: import_queue.clone(),264			on_demand: None,265			block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),266		})?;267268	let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());269	let rpc_client = client.clone();270	let rpc_pool = transaction_pool.clone();271	let select_chain = params.select_chain.clone();272	let is_authority = parachain_config.role.clone().is_authority();273	let rpc_network = network.clone();274	let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {275		let full_deps = nft_rpc::FullDeps {276			backend: frontier_backend.clone(),277			deny_unsafe,278			client: rpc_client.clone(),279			pool: rpc_pool.clone(),280			// TODO: Unhardcode281			enable_dev_signer: false,282			filter_pool: filter_pool.clone(),283			network: rpc_network.clone(),284			pending_transactions: pending_transactions.clone(),285			select_chain: select_chain.clone(),286			is_authority: is_authority.clone(),287			// TODO: Unhardcode288			max_past_logs: 10000,289		};290291		nft_rpc::create_full::<_, _, _, RuntimeApi, _>(full_deps, subscription_executor.clone())292	});293294	sc_service::spawn_tasks(sc_service::SpawnTasksParams {295		on_demand: None,296		remote_blockchain: None,297		rpc_extensions_builder,298		client: client.clone(),299		transaction_pool: transaction_pool.clone(),300		task_manager: &mut task_manager,301		config: parachain_config,302		keystore: params.keystore_container.sync_keystore(),303		backend: backend.clone(),304		network: network.clone(),305		network_status_sinks,306		system_rpc_tx,307		telemetry: telemetry.as_mut(),308	})?;309310	let announce_block = {311		let network = network.clone();312		Arc::new(move |hash, data| network.announce_block(hash, data))313	};314315	if validator {316		let parachain_consensus = build_consensus(317			client.clone(),318			prometheus_registry.as_ref(),319			telemetry.as_ref().map(|t| t.handle()),320			&task_manager,321			&relay_chain_full_node,322			transaction_pool,323			network,324			params.keystore_container.sync_keystore(),325			force_authoring,326		)?;327328		let spawner = task_manager.spawn_handle();329330		let params = StartCollatorParams {331			para_id: id,332			block_status: client.clone(),333			announce_block,334			client: client.clone(),335			task_manager: &mut task_manager,336			collator_key,337			relay_chain_full_node,338			spawner,339			parachain_consensus,340			import_queue,341		};342343		start_collator(params).await?;344	} else {345		let params = StartFullNodeParams {346			client: client.clone(),347			announce_block,348			task_manager: &mut task_manager,349			para_id: id,350			relay_chain_full_node,351		};352353		start_full_node(params)?;354	}355356	start_network.start_network();357358	Ok((task_manager, client))359}360361/// Build the import queue for the the parachain runtime.362pub fn parachain_build_import_queue(363	client: Arc<FullClient>,364	config: &Configuration,365	telemetry: Option<TelemetryHandle>,366	task_manager: &TaskManager,367) -> Result<368	sp_consensus::DefaultImportQueue<369		Block,370		FullClient,371	>,372	sc_service::Error,373> {374	let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;375376377	cumulus_client_consensus_aura::import_queue::<378		sp_consensus_aura::sr25519::AuthorityPair,379		_,380		_,381		_,382		_,383		_,384		_,385	>(cumulus_client_consensus_aura::ImportQueueParams {386		block_import:  client.clone(),387		client: client.clone(),388		create_inherent_data_providers: move |_, _| async move {389			let time = sp_timestamp::InherentDataProvider::from_system_time();390391			let slot =392				sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(393					*time,394					slot_duration.slot_duration(),395				);396397			Ok((time, slot))398		},399		registry: config.prometheus_registry().clone(),400		can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),401		spawner: &task_manager.spawn_essential_handle(),402		telemetry,403	})404	.map_err(Into::into)405}406407/// Start a normal parachain node.408pub async fn start_node(409	parachain_config: Configuration,410	collator_key: CollatorPair,411	polkadot_config: Configuration,412	id: ParaId,413) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {414	start_node_impl::<_, _>(415		parachain_config,416		collator_key,417		polkadot_config,418		id,419		parachain_build_import_queue,420		|client,421		 prometheus_registry,422		 telemetry,423		 task_manager,424		 relay_chain_node,425		 transaction_pool,426		 sync_oracle,427		 keystore,428		 force_authoring| {429			let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;430431			let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(432				task_manager.spawn_handle(),433				client.clone(),434				transaction_pool,435				prometheus_registry.clone(),436				telemetry.clone(),437			);438439			let relay_chain_backend = relay_chain_node.backend.clone();440			let relay_chain_client = relay_chain_node.client.clone();441			Ok(build_aura_consensus::<442				sp_consensus_aura::sr25519::AuthorityPair,443				_,444				_,445				_,446				_,447				_,448				_,449				_,450				_,451				_,452			>(BuildAuraConsensusParams {453				proposer_factory,454				create_inherent_data_providers: move |_, (relay_parent, validation_data)| {455					let parachain_inherent =456					cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(457						relay_parent,458						&relay_chain_client,459						&*relay_chain_backend,460						&validation_data,461						id,462					);463					async move {464						let time = sp_timestamp::InherentDataProvider::from_system_time();465466						let slot =467						sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(468							*time,469							slot_duration.slot_duration(),470						);471472						let parachain_inherent = parachain_inherent.ok_or_else(|| {473							Box::<dyn std::error::Error + Send + Sync>::from(474								"Failed to create parachain inherent",475							)476						})?;477						Ok((time, slot, parachain_inherent))478					}479				},480				block_import: client.clone(),481				relay_chain_client: relay_chain_node.client.clone(),482				relay_chain_backend: relay_chain_node.backend.clone(),483				para_client: client.clone(),484				backoff_authoring_blocks: Option::<()>::None,485				sync_oracle,486				keystore,487				force_authoring,488				slot_duration,489				// We got around 500ms for proposing490				block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),491				telemetry,492			}))493		},494	)495	.await496}