From 8cc39ba6d3e1e043cdea8b7d16f84b2754403a45 Mon Sep 17 00:00:00 2001 From: Daniel Shiposha Date: Tue, 15 Mar 2022 14:03:05 +0000 Subject: [PATCH] Implement autoseal in dev mode --- --- a/node/cli/Cargo.toml +++ b/node/cli/Cargo.toml @@ -168,6 +168,9 @@ [dependencies.serde_json] version = '1.0.68' +[dependencies.sc-consensus-manual-seal] +git = 'https://github.com/paritytech/substrate.git' +branch = 'polkadot-v0.9.17' ################################################################################ # Cumulus dependencies --- a/node/cli/src/chain_spec.rs +++ b/node/cli/src/chain_spec.rs @@ -68,6 +68,25 @@ } } +pub enum ServiceId { + Prod, + Dev +} + +pub trait ServiceIdentification { + fn service_id(&self) -> ServiceId; +} + +impl ServiceIdentification for Box { + fn service_id(&self) -> ServiceId { + if self.id().ends_with("dev") { + ServiceId::Dev + } else { + ServiceId::Prod + } + } +} + /// Helper function to generate a crypto pair from seed pub fn get_from_seed(seed: &str) -> ::Public { TPublic::Pair::from_string(&format!("//{}", seed), None) --- a/node/cli/src/command.rs +++ b/node/cli/src/command.rs @@ -33,9 +33,9 @@ // limitations under the License. use crate::{ - chain_spec::{self, RuntimeId, RuntimeIdentification}, + chain_spec::{self, RuntimeId, RuntimeIdentification, ServiceId, ServiceIdentification}, cli::{Cli, RelayChainCli, Subcommand}, - service::new_partial, + service::{new_partial, start_node, start_dev_node}, }; #[cfg(feature = "unique-runtime")] @@ -210,6 +210,7 @@ >( &$config, crate::service::parachain_build_import_queue, + ServiceId::Prod, )?; let task_manager = $components.task_manager; @@ -245,6 +246,34 @@ }} } +macro_rules! start_node_using_chain_runtime { + ($start_node_fn:ident($config:expr $(, $($args:expr),+)?) $($code:tt)*) => { + match $config.chain_spec.runtime_id() { + #[cfg(feature = "unique-runtime")] + RuntimeId::Unique => $start_node_fn::< + unique_runtime::Runtime, + unique_runtime::RuntimeApi, + UniqueRuntimeExecutor, + >($config $(, $($args),+)?) $($code)*, + + #[cfg(feature = "quartz-runtime")] + RuntimeId::Quartz => $start_node_fn::< + quartz_runtime::Runtime, + quartz_runtime::RuntimeApi, + QuartzRuntimeExecutor, + >($config $(, $($args),+)?) $($code)*, + + RuntimeId::Opal => $start_node_fn::< + opal_runtime::Runtime, + opal_runtime::RuntimeApi, + OpalRuntimeExecutor, + >($config $(, $($args),+)?) $($code)*, + + RuntimeId::Unknown(chain) => Err(no_runtime_err!(chain).into()), + } + }; +} + /// Parse command line arguments into service configuration. pub fn run() -> Result<()> { let cli = Cli::from_args(); @@ -365,7 +394,20 @@ let runner = cli.create_runner(&cli.run.normalize())?; runner.run_node_until_exit(|config| async move { - let para_id = chain_spec::Extensions::try_get(&*config.chain_spec) + let extensions = chain_spec::Extensions::try_get(&*config.chain_spec); + + let service_id = config.chain_spec.service_id(); + let relay_chain_id = extensions.map(|e| e.relay_chain.clone()); + let is_dev_service = matches![service_id, ServiceId::Dev] + || relay_chain_id == Some("dev-service".into()); + + if is_dev_service { + return start_node_using_chain_runtime! { + start_dev_node(config).map_err(Into::into) + }; + }; + + let para_id = extensions .map(|e| e.para_id) .ok_or("Could not find parachain ID in chain-spec.")?; @@ -376,10 +418,10 @@ .chain(cli.relaychain_args.iter()), ); - let id = ParaId::from(para_id); + let para_id = ParaId::from(para_id); let parachain_account = - AccountIdConversion::::into_account(&id); + AccountIdConversion::::into_account(¶_id); let state_version = RelayChainCli::native_runtime_version(&config.chain_spec).state_version(); @@ -395,7 +437,7 @@ ) .map_err(|err| format!("Relay chain argument error: {}", err))?; - info!("Parachain id: {:?}", id); + info!("Parachain id: {:?}", para_id); info!("Parachain Account: {}", parachain_account); info!("Parachain genesis state: {}", genesis_state); info!("Parachain genesis hash: {}", genesis_hash); @@ -408,37 +450,11 @@ } ); - match config.chain_spec.runtime_id() { - #[cfg(feature = "unique-runtime")] - RuntimeId::Unique => crate::service::start_node::< - unique_runtime::Runtime, - unique_runtime::RuntimeApi, - UniqueRuntimeExecutor, - >(config, polkadot_config, id) - .await - .map(|r| r.0) - .map_err(Into::into), - - #[cfg(feature = "quartz-runtime")] - RuntimeId::Quartz => crate::service::start_node::< - quartz_runtime::Runtime, - quartz_runtime::RuntimeApi, - QuartzRuntimeExecutor, - >(config, polkadot_config, id) - .await - .map(|r| r.0) - .map_err(Into::into), - - RuntimeId::Opal => crate::service::start_node::< - opal_runtime::Runtime, - opal_runtime::RuntimeApi, - OpalRuntimeExecutor, - >(config, polkadot_config, id) - .await - .map(|r| r.0) - .map_err(Into::into), - - RuntimeId::Unknown(chain) => Err(no_runtime_err!(chain).into()), + start_node_using_chain_runtime! { + start_node(config, polkadot_config, para_id) + .await + .map(|r| r.0) + .map_err(Into::into) } }) } --- a/node/cli/src/service.rs +++ b/node/cli/src/service.rs @@ -57,6 +57,7 @@ use fc_mapping_sync::{MappingSyncWorker, SyncStrategy}; use unique_runtime_common::types::{AuraId, RuntimeInstance, AccountId, Balance, Index, Hash, Block}; +use crate::chain_spec::ServiceId; /// Native executor instance. pub struct UniqueRuntimeExecutor; @@ -125,6 +126,7 @@ sc_service::TFullClient>; type FullBackend = sc_service::TFullBackend; type FullSelectChain = sc_consensus::LongestChain; +type MaybeSelectChain = Option; /// Starts a `ServiceBuilder` for a full service. /// @@ -134,11 +136,12 @@ pub fn new_partial( config: &Configuration, build_import_queue: BIQ, + service_id: ServiceId, ) -> Result< PartialComponents< FullClient, FullBackend, - FullSelectChain, + MaybeSelectChain, sc_consensus::DefaultImportQueue>, sc_transaction_pool::FullPool>, ( @@ -215,7 +218,10 @@ telemetry }); - let select_chain = sc_consensus::LongestChain::new(backend.clone()); + let select_chain = match service_id { + ServiceId::Prod => Some(sc_consensus::LongestChain::new(backend.clone())), + ServiceId::Dev => None + }; let transaction_pool = sc_transaction_pool::BasicPool::new_full( config.transaction_pool.clone(), @@ -317,7 +323,9 @@ let parachain_config = prepare_node_config(parachain_config); let params = - new_partial::(¶chain_config, build_import_queue)?; + new_partial::( + ¶chain_config, build_import_queue, ServiceId::Prod + )?; let (mut telemetry, filter_pool, frontier_backend, telemetry_worker_handle, fee_history_cache) = params.other; @@ -356,7 +364,9 @@ let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle()); let rpc_client = client.clone(); let rpc_pool = transaction_pool.clone(); - let select_chain = params.select_chain.clone(); + let select_chain = params.select_chain + .expect("select_chain always exists when running Prod service; qed") + .clone(); let rpc_network = network.clone(); let rpc_frontier_backend = frontier_backend.clone(); @@ -638,3 +648,255 @@ ) .await } + +fn dev_build_import_queue( + client: Arc>, + config: &Configuration, + _: Option, + task_manager: &TaskManager, +) -> Result>, sc_service::Error> +where + RuntimeApi: sp_api::ConstructRuntimeApi> + + Send + + Sync + + 'static, + RuntimeApi::RuntimeApi: sp_transaction_pool::runtime_api::TaggedTransactionQueue + + sp_api::ApiExt>, + ExecutorDispatch: NativeExecutionDispatch + 'static, +{ + Ok(sc_consensus_manual_seal::import_queue( + Box::new(client.clone()), + &task_manager.spawn_essential_handle(), + config.prometheus_registry(), + )) +} + +/// Builds a new development service. This service uses instant seal, and mocks +/// the parachain inherent +pub fn start_dev_node(config: Configuration) + -> sc_service::error::Result +where + Runtime: RuntimeInstance + Send + Sync + 'static, + ::CrossAccountId: Serialize, + for<'de> ::CrossAccountId: Deserialize<'de>, + RuntimeApi: sp_api::ConstructRuntimeApi> + + Send + + Sync + + 'static, + RuntimeApi::RuntimeApi: sp_transaction_pool::runtime_api::TaggedTransactionQueue + + fp_rpc::EthereumRuntimeRPCApi + + sp_session::SessionKeys + + sp_block_builder::BlockBuilder + + pallet_transaction_payment_rpc_runtime_api::TransactionPaymentApi + + sp_api::ApiExt> + + up_rpc::UniqueApi + + substrate_frame_rpc_system::AccountNonceApi + + sp_api::Metadata + + sp_offchain::OffchainWorkerApi + + cumulus_primitives_core::CollectCollationInfo + + sp_consensus_aura::AuraApi, + ExecutorDispatch: NativeExecutionDispatch + 'static, +{ + use futures::Stream; + use sc_consensus_manual_seal::{run_manual_seal, EngineCommand, ManualSealParams}; + use fc_consensus::FrontierBlockImport; + use sc_client_api::HeaderBackend; + + let sc_service::PartialComponents { + client, + backend, + mut task_manager, + import_queue, + keystore_container, + select_chain: maybe_select_chain, + transaction_pool, + other: + ( + telemetry, + filter_pool, + frontier_backend, + _telemetry_worker_handle, + fee_history_cache, + ), + } = new_partial::( + &config, + dev_build_import_queue::, + ServiceId::Dev + )?; + + let block_data_cache = Arc::new(fc_rpc::EthBlockDataCache::new( + task_manager.spawn_handle(), + overrides_handle::<_, _, Runtime>(client.clone()), + 50, + 50, + )); + + let (network, system_rpc_tx, network_starter) = + sc_service::build_network(sc_service::BuildNetworkParams { + config: &config, + client: client.clone(), + transaction_pool: transaction_pool.clone(), + spawn_handle: task_manager.spawn_handle(), + import_queue, + block_announce_validator_builder: None, + warp_sync: None, + })?; + + if config.offchain_worker.enabled { + sc_service::build_offchain_workers( + &config, + task_manager.spawn_handle(), + client.clone(), + network.clone(), + ); + } + + let prometheus_registry = config.prometheus_registry().cloned(); + let collator = config.role.is_authority(); + + let select_chain = maybe_select_chain.clone().expect( + "`new_partial` builds a `LongestChainRule` when building dev service.\ + We specified the dev service when calling `new_partial`.\ + Therefore, a `LongestChainRule` is present. qed.", + ); + + if collator { + let block_import = + FrontierBlockImport::new(client.clone(), client.clone(), frontier_backend.clone()); + + let env = sc_basic_authorship::ProposerFactory::new( + task_manager.spawn_handle(), + client.clone(), + transaction_pool.clone(), + prometheus_registry.as_ref(), + telemetry.as_ref().map(|x| x.handle()), + ); + + let commands_stream: Box> + Send + Sync + Unpin> = + Box::new( + // This bit cribbed from the implementation of instant seal. + transaction_pool + .pool() + .validated_pool() + .import_notification_stream() + .map(|_| EngineCommand::SealNewBlock { + create_empty: true, // was false in Moonbeam + finalize: false, + parent_hash: None, + sender: None, + }), + ); + + let slot_duration = cumulus_client_consensus_aura::slot_duration(&*client)?; + let client_set_aside_for_cidp = client.clone(); + + task_manager.spawn_essential_handle().spawn_blocking( + "authorship_task", + Some("block-authoring"), + run_manual_seal(ManualSealParams { + block_import, + env, + client: client.clone(), + pool: transaction_pool.clone(), + commands_stream, + select_chain: select_chain.clone(), + consensus_data_provider: None, + create_inherent_data_providers: move |block: Hash, ()| { + let current_para_block = client_set_aside_for_cidp + .number(block) + .expect("Header lookup should succeed") + .expect("Header passed in as parent should be present in backend."); + + let client_for_xcm = client_set_aside_for_cidp.clone(); + async move { + let time = sp_timestamp::InherentDataProvider::from_system_time(); + + let mocked_parachain = cumulus_primitives_parachain_inherent::MockValidationDataInherentDataProvider { + current_para_block, + relay_offset: 1000, + relay_blocks_per_para_block: 2, + xcm_config: cumulus_primitives_parachain_inherent::MockXcmConfig::new( + &*client_for_xcm, + block, + Default::default(), + Default::default(), + ), + raw_downward_messages: vec![], + raw_horizontal_messages: vec![], + }; + + let slot = + sp_consensus_aura::inherents::InherentDataProvider::from_timestamp_and_duration( + *time, + slot_duration.slot_duration(), + ); + + Ok((time, slot, mocked_parachain)) + } + }, + }), + ); + } + + task_manager.spawn_essential_handle().spawn( + "frontier-mapping-sync-worker", + Some("block-authoring"), + MappingSyncWorker::new( + client.import_notification_stream(), + Duration::new(6, 0), + client.clone(), + backend.clone(), + frontier_backend.clone(), + SyncStrategy::Normal, + ) + .for_each(|()| futures::future::ready(())), + ); + + let subscription_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle()); + let rpc_client = client.clone(); + let rpc_pool = transaction_pool.clone(); + let rpc_network = network.clone(); + let rpc_frontier_backend = frontier_backend.clone(); + let rpc_extensions_builder = Box::new(move |deny_unsafe, _| { + let full_deps = unique_rpc::FullDeps { + backend: rpc_frontier_backend.clone(), + deny_unsafe, + client: rpc_client.clone(), + pool: rpc_pool.clone(), + graph: rpc_pool.pool().clone(), + // TODO: Unhardcode + enable_dev_signer: false, + filter_pool: filter_pool.clone(), + network: rpc_network.clone(), + select_chain: select_chain.clone(), + is_authority: collator, + // TODO: Unhardcode + max_past_logs: 10000, + block_data_cache: block_data_cache.clone(), + fee_history_cache: fee_history_cache.clone(), + // TODO: Unhardcode + fee_history_limit: 2048, + }; + + Ok(unique_rpc::create_full::<_, _, _, _, Runtime, RuntimeApi, _>( + full_deps, + subscription_executor.clone(), + )) + }); + + sc_service::spawn_tasks(sc_service::SpawnTasksParams { + network, + client, + keystore: keystore_container.sync_keystore(), + task_manager: &mut task_manager, + transaction_pool, + rpc_extensions_builder, + backend, + system_rpc_tx, + config, + telemetry: None, + })?; + + network_starter.start_network(); + Ok(task_manager) +} -- gitstuff