git.delta.rocks / unique-network / refs/commits / 54cf3dbad6fe

difftreelog

source

node/cli/src/service.rs15.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::time::Duration;13use fc_rpc_core::types::FeeHistoryCache;14use futures::StreamExt;1516use unique_rpc::overrides_handle;17// Local Runtime Types18use unique_runtime::RuntimeApi;1920// Cumulus Imports21use cumulus_client_consensus_aura::{build_aura_consensus, BuildAuraConsensusParams, SlotProportion};22use cumulus_client_consensus_common::ParachainConsensus;23use cumulus_client_network::build_block_announce_validator;24use cumulus_client_service::{25	prepare_node_config, start_collator, start_full_node, StartCollatorParams, StartFullNodeParams,26};27use cumulus_primitives_core::ParaId;2829// Substrate Imports30use sc_client_api::ExecutorProvider;31use sc_executor::NativeElseWasmExecutor;32use sc_executor::NativeExecutionDispatch;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;40use sc_client_api::BlockchainEvents;4142// Frontier Imports43use fc_rpc_core::types::FilterPool;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		unique_runtime::api::dispatch(method, data)60	}6162	fn native_version() -> sc_executor::NativeVersion {63		unique_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("", "", "unique").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			Option<FilterPool>,112			Arc<fc_db::Backend<Block>>,113			Option<TelemetryWorkerHandle>,114			FeeHistoryCache,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_manager169			.spawn_handle()170			.spawn("telemetry", None, worker.run());171		telemetry172	});173174	let select_chain = sc_consensus::LongestChain::new(backend.clone());175176	let transaction_pool = sc_transaction_pool::BasicPool::new_full(177		config.transaction_pool.clone(),178		config.role.is_authority().into(),179		config.prometheus_registry(),180		task_manager.spawn_essential_handle(),181		client.clone(),182	);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	)?;194	let fee_history_cache: FeeHistoryCache = Arc::new(Mutex::new(BTreeMap::new()));195196	let params = PartialComponents {197		backend,198		client,199		import_queue,200		keystore_container,201		task_manager,202		transaction_pool,203		select_chain,204		other: (205			telemetry,206			filter_pool,207			frontier_backend,208			telemetry_worker_handle,209			fee_history_cache,210		),211	};212213	Ok(params)214}215216/// Start a node with the given parachain `Configuration` and relay chain `Configuration`.217///218/// This is the actual implementation that is abstract over the executor and the runtime api.219#[sc_tracing::logging::prefix_logs_with("Parachain")]220async fn start_node_impl<BIQ, BIC>(221	parachain_config: Configuration,222	polkadot_config: Configuration,223	id: ParaId,224	build_import_queue: BIQ,225	build_consensus: BIC,226) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)>227where228	sc_client_api::StateBackendFor<FullBackend, Block>: sp_api::StateBackend<BlakeTwo256>,229	ExecutorDispatch: NativeExecutionDispatch + 'static,230	BIQ: FnOnce(231		Arc<FullClient>,232		&Configuration,233		Option<TelemetryHandle>,234		&TaskManager,235	) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error>,236	BIC: FnOnce(237		Arc<FullClient>,238		Option<&Registry>,239		Option<TelemetryHandle>,240		&TaskManager,241		&polkadot_service::NewFull<polkadot_service::Client>,242		Arc<sc_transaction_pool::FullPool<Block, FullClient>>,243		Arc<NetworkService<Block, Hash>>,244		SyncCryptoStorePtr,245		bool,246	) -> Result<Box<dyn ParachainConsensus<Block>>, sc_service::Error>,247{248	if matches!(parachain_config.role, Role::Light) {249		return Err("Light client not supported!".into());250	}251252	let parachain_config = prepare_node_config(parachain_config);253254	let params = new_partial::<BIQ>(&parachain_config, build_import_queue)?;255	let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle, fee_history_cache) =256		params.other;257258	let relay_chain_full_node =259		cumulus_client_service::build_polkadot_full_node(polkadot_config, telemetry_worker_handle)260			.map_err(|e| match e {261				polkadot_service::Error::Sub(x) => x,262				s => format!("{}", s).into(),263			})?;264265	let client = params.client.clone();266	let backend = params.backend.clone();267	let block_announce_validator = build_block_announce_validator(268		relay_chain_full_node.client.clone(),269		id,270		Box::new(relay_chain_full_node.network.clone()),271		relay_chain_full_node.backend.clone(),272	);273274	let force_authoring = parachain_config.force_authoring;275	let validator = parachain_config.role.is_authority();276	let prometheus_registry = parachain_config.prometheus_registry().cloned();277	let transaction_pool = params.transaction_pool.clone();278	let mut task_manager = params.task_manager;279	let import_queue = cumulus_client_service::SharedImportQueue::new(params.import_queue);280281	let (network, system_rpc_tx, start_network) =282		sc_service::build_network(sc_service::BuildNetworkParams {283			config: &parachain_config,284			client: client.clone(),285			transaction_pool: transaction_pool.clone(),286			spawn_handle: task_manager.spawn_handle(),287			import_queue: import_queue.clone(),288			block_announce_validator_builder: Some(Box::new(|_| block_announce_validator)),289			warp_sync: None,290		})?;291292	let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle());293	let rpc_client = client.clone();294	let rpc_pool = transaction_pool.clone();295	let select_chain = params.select_chain.clone();296	let is_authority = parachain_config.role.clone().is_authority();297	let rpc_network = network.clone();298299	let rpc_frontier_backend = frontier_backend.clone();300301	let block_data_cache = Arc::new(fc_rpc::EthBlockDataCache::new(302		task_manager.spawn_handle(),303		overrides_handle(client.clone()),304		50,305		50,306	));307308	let rpc_extensions_builder = Box::new(move |deny_unsafe, _| {309		let full_deps = unique_rpc::FullDeps {310			backend: rpc_frontier_backend.clone(),311			deny_unsafe,312			client: rpc_client.clone(),313			pool: rpc_pool.clone(),314			graph: rpc_pool.pool().clone(),315			// TODO: Unhardcode316			enable_dev_signer: false,317			filter_pool: filter_pool.clone(),318			network: rpc_network.clone(),319			select_chain: select_chain.clone(),320			is_authority,321			// TODO: Unhardcode322			max_past_logs: 10000,323			block_data_cache: block_data_cache.clone(),324			fee_history_cache: fee_history_cache.clone(),325			// TODO: Unhardcode326			fee_history_limit: 2048,327		};328329		Ok(unique_rpc::create_full::<_, _, _, _, RuntimeApi, _>(330			full_deps,331			subscription_executor.clone(),332		))333	});334335	task_manager.spawn_essential_handle().spawn(336		"frontier-mapping-sync-worker",337		None,338		MappingSyncWorker::new(339			client.import_notification_stream(),340			Duration::new(6, 0),341			client.clone(),342			backend.clone(),343			frontier_backend.clone(),344			SyncStrategy::Normal,345		)346		.for_each(|()| futures::future::ready(())),347	);348349	sc_service::spawn_tasks(sc_service::SpawnTasksParams {350		rpc_extensions_builder,351		client: client.clone(),352		transaction_pool: transaction_pool.clone(),353		task_manager: &mut task_manager,354		config: parachain_config,355		keystore: params.keystore_container.sync_keystore(),356		backend: backend.clone(),357		network: network.clone(),358		system_rpc_tx,359		telemetry: telemetry.as_mut(),360	})?;361362	let announce_block = {363		let network = network.clone();364		Arc::new(move |hash, data| network.announce_block(hash, data))365	};366367	if validator {368		let parachain_consensus = build_consensus(369			client.clone(),370			prometheus_registry.as_ref(),371			telemetry.as_ref().map(|t| t.handle()),372			&task_manager,373			&relay_chain_full_node,374			transaction_pool,375			network,376			params.keystore_container.sync_keystore(),377			force_authoring,378		)?;379380		let spawner = task_manager.spawn_handle();381382		let params = StartCollatorParams {383			para_id: id,384			block_status: client.clone(),385			announce_block,386			client: client.clone(),387			task_manager: &mut task_manager,388			relay_chain_full_node,389			spawner,390			parachain_consensus,391			import_queue,392		};393394		start_collator(params).await?;395	} else {396		let params = StartFullNodeParams {397			client: client.clone(),398			announce_block,399			task_manager: &mut task_manager,400			para_id: id,401			relay_chain_full_node,402		};403404		start_full_node(params)?;405	}406407	start_network.start_network();408409	Ok((task_manager, client))410}411412/// Build the import queue for the the parachain runtime.413pub fn parachain_build_import_queue(414	client: Arc<FullClient>,415	config: &Configuration,416	telemetry: Option<TelemetryHandle>,417	task_manager: &TaskManager,418) -> Result<sc_consensus::DefaultImportQueue<Block, FullClient>, sc_service::Error> {419	let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;420421	cumulus_client_consensus_aura::import_queue::<422		sp_consensus_aura::sr25519::AuthorityPair,423		_,424		_,425		_,426		_,427		_,428		_,429	>(cumulus_client_consensus_aura::ImportQueueParams {430		block_import: client.clone(),431		client: client.clone(),432		create_inherent_data_providers: move |_, _| async move {433			let time = sp_timestamp::InherentDataProvider::from_system_time();434435			let slot =436				sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(437					*time,438					slot_duration.slot_duration(),439				);440441			Ok((time, slot))442		},443		registry: config.prometheus_registry(),444		can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()),445		spawner: &task_manager.spawn_essential_handle(),446		telemetry,447	})448	.map_err(Into::into)449}450451/// Start a normal parachain node.452pub async fn start_node(453	parachain_config: Configuration,454	polkadot_config: Configuration,455	id: ParaId,456) -> sc_service::error::Result<(TaskManager, Arc<FullClient>)> {457	start_node_impl::<_, _>(458		parachain_config,459		polkadot_config,460		id,461		parachain_build_import_queue,462		|client,463		 prometheus_registry,464		 telemetry,465		 task_manager,466		 relay_chain_node,467		 transaction_pool,468		 sync_oracle,469		 keystore,470		 force_authoring| {471			let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?;472473			let proposer_factory = sc_basic_authorship::ProposerFactory::with_proof_recording(474				task_manager.spawn_handle(),475				client.clone(),476				transaction_pool,477				prometheus_registry,478				telemetry.clone(),479			);480481			let relay_chain_backend = relay_chain_node.backend.clone();482			let relay_chain_client = relay_chain_node.client.clone();483			Ok(build_aura_consensus::<484				sp_consensus_aura::sr25519::AuthorityPair,485				_,486				_,487				_,488				_,489				_,490				_,491				_,492				_,493				_,494			>(BuildAuraConsensusParams {495				proposer_factory,496				create_inherent_data_providers: move |_, (relay_parent, validation_data)| {497					let parachain_inherent =498					cumulus_primitives_parachain_inherent::ParachainInherentData::create_at_with_client(499						relay_parent,500						&relay_chain_client,501						&*relay_chain_backend,502						&validation_data,503						id,504					);505					async move {506						let time = sp_timestamp::InherentDataProvider::from_system_time();507508						let slot =509						sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration(510							*time,511							slot_duration.slot_duration(),512						);513514						let parachain_inherent = parachain_inherent.ok_or_else(|| {515							Box::<dyn std::error::Error + Send + Sync>::from(516								"Failed to create parachain inherent",517							)518						})?;519						Ok((time, slot, parachain_inherent))520					}521				},522				block_import: client.clone(),523				relay_chain_client: relay_chain_node.client.clone(),524				relay_chain_backend: relay_chain_node.backend.clone(),525				para_client: client,526				backoff_authoring_blocks: Option::<()>::None,527				sync_oracle,528				keystore,529				force_authoring,530				slot_duration,531				// We got around 500ms for proposing532				block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),533				telemetry,534				max_block_proposal_slot_portion: None,535			}))536		},537	)538	.await539}