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

difftreelog

source

node/cli/src/service.rs15.0 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;13use std::time::Duration;14use futures::StreamExt;1516// Local Runtime Types17use nft_runtime::RuntimeApi;1819// Cumulus Imports20use 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;2728// Substrate Imports29use 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;4041// Frontier Imports42use fc_rpc_core::types::FilterPool;43use fc_rpc_core::types::PendingTransactions;44use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};4546// Runtime type overrides47type 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;5152/// Native executor instance.53pub 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>;9394/// Starts a `ServiceBuilder` for a full service.95///96/// Use this macro if you don't actually need the full service, but just the builder in order to97/// be able to perform chain operations.98#[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}214215/// Start a node with the given parachain `Configuration` and relay chain `Configuration`.216///217/// This is the actual implementation that is abstract over the executor and the runtime api.218#[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>(&parachain_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: &parachain_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			// TODO: Unhardcode312			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			// TODO: Unhardcode319			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}405406/// Build the import queue for the the parachain runtime.407pub 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}444445/// Start a normal parachain node.446pub 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				// We got around 500ms for proposing526				block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),527				telemetry,528				max_block_proposal_slot_portion: None,529			}))530		},531	)532	.await533}