git.delta.rocks / unique-network / refs/commits / 4e00e0c4cd4e

difftreelog

source

node/cli/src/service.rs14.8 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 futures::StreamExt;1415// Local Runtime Types16use unique_runtime::RuntimeApi;1718// Cumulus Imports19use 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;2627// Substrate Imports28use 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;3940// Frontier Imports41use fc_rpc_core::types::FilterPool;42use fc_mapping_sync::{MappingSyncWorker, SyncStrategy};4344// Runtime type overrides45type 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;4950/// Native executor instance.51pub struct ParachainRuntimeExecutor;5253impl NativeExecutionDispatch for ParachainRuntimeExecutor {54	type ExtendHostFunctions = frame_benchmarking::benchmarking::HostFunctions;5556	fn dispatch(method: &str, data: &[u8]) -> Option<Vec<u8>> {57		unique_runtime::api::dispatch(method, data)58	}5960	fn native_version() -> sc_executor::NativeVersion {61		unique_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("", "", "unique").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>;9192/// Starts a `ServiceBuilder` for a full service.93///94/// Use this macro if you don't actually need the full service, but just the builder in order to95/// be able to perform chain operations.96#[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}208209/// Start a node with the given parachain `Configuration` and relay chain `Configuration`.210///211/// This is the actual implementation that is abstract over the executor and the runtime api.212#[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>(&parachain_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: &parachain_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 = unique_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			// TODO: Unhardcode301			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			// TODO: Unhardcode307			max_past_logs: 10000,308		};309310		Ok(unique_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}393394/// Build the import queue for the the parachain runtime.395pub 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}432433/// Start a normal parachain node.434pub 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				// We got around 500ms for proposing514				block_proposal_slot_portion: SlotProportion::new(1f32 / 24f32),515				telemetry,516				max_block_proposal_slot_portion: None,517			}))518		},519	)520	.await521}