--- a/node/src/chain_spec.rs +++ b/node/src/chain_spec.rs @@ -3,19 +3,18 @@ // file 'LICENSE', which is part of this source code package. // -// use nft_runtime::{ -// AccountId, AuraConfig, BalancesConfig, GenesisConfig, GrandpaConfig, Signature, SudoConfig, -// SystemConfig, WASM_BINARY, -// }; -// use nft_runtime::{ContractsConfig, ContractsSchedule, NftConfig, CollectionType}; use nft_runtime::*; -use sc_service::ChainType; +use sp_core::{Pair, Public, sr25519}; +use nft_runtime::{ + AccountId, AuraConfig, BalancesConfig, EVMConfig, EthereumConfig, GenesisConfig, GrandpaConfig, + SudoConfig, SystemConfig, WASM_BINARY, Signature +}; use sp_consensus_aura::sr25519::AuthorityId as AuraId; -use sp_core::{sr25519, Pair, Public}; use sp_finality_grandpa::AuthorityId as GrandpaId; -use sp_runtime::traits::{IdentifyAccount, Verify}; +use sp_runtime::traits::{Verify, IdentifyAccount}; +use sc_service::ChainType; use serde_json::map::Map; -// use crate::chain_spec::api::chain_extension::*; +use std::collections::BTreeMap; // Note this is the URL for the telemetry server //const STAGING_TELEMETRY_URL: &str = "wss://telemetry.polkadot.io/submit/"; @@ -25,9 +24,9 @@ /// Helper function to generate a crypto pair from seed pub fn get_from_seed(seed: &str) -> ::Public { - TPublic::Pair::from_string(&format!("//{}", seed), None) - .expect("static values are valid; qed") - .public() + TPublic::Pair::from_string(&format!("//{}", seed), None) + .expect("static values are valid; qed") + .public() } type AccountPublic = ::Signer; @@ -35,14 +34,14 @@ /// Helper function to generate an account ID from seed pub fn get_account_id_from_seed(seed: &str) -> AccountId where - AccountPublic: From<::Public>, + AccountPublic: From<::Public>, { - AccountPublic::from(get_from_seed::(seed)).into_account() + AccountPublic::from(get_from_seed::(seed)).into_account() } /// Helper function to generate an authority key for Aura pub fn authority_keys_from_seed(s: &str) -> (AuraId, GrandpaId) { - (get_from_seed::(s), get_from_seed::(s)) + (get_from_seed::(s), get_from_seed::(s)) } pub fn development_config() -> Result { @@ -138,88 +137,92 @@ } fn testnet_genesis( - wasm_binary: &[u8], - initial_authorities: Vec<(AuraId, GrandpaId)>, - root_key: AccountId, - endowed_accounts: Vec, - enable_println: bool, + wasm_binary: &[u8], + initial_authorities: Vec<(AuraId, GrandpaId)>, + root_key: AccountId, + endowed_accounts: Vec, + enable_println: bool, ) -> GenesisConfig { let vested_accounts = vec![ get_account_id_from_seed::("Bob"), ]; - GenesisConfig { - system: Some(SystemConfig { - code: wasm_binary.to_vec(), - changes_trie_config: Default::default(), - }), - pallet_balances: Some(BalancesConfig { - balances: endowed_accounts - .iter() - .cloned() - .map(|k| (k, 1 << 100)) - .collect(), - }), - pallet_aura: Some(AuraConfig { - authorities: initial_authorities.iter().map(|x| (x.0.clone())).collect(), - }), - pallet_grandpa: Some(GrandpaConfig { - authorities: initial_authorities - .iter() - .map(|x| (x.1.clone(), 1)) - .collect(), - }), - pallet_treasury: Some(Default::default()), - pallet_sudo: Some(SudoConfig { key: root_key }), - pallet_vesting: Some(VestingConfig { - vesting: vested_accounts - .iter() - .cloned() - .map(|k| (k, 1000, 100, 1 << 98)) - .collect(), - }), - pallet_nft: Some(NftConfig { - collection_id: vec![( - 1, - Collection { - owner: get_account_id_from_seed::("Alice"), - mode: CollectionMode::NFT, - access: AccessMode::Normal, - decimal_points: 0, - name: vec![], - description: vec![], - token_prefix: vec![], - mint_mode: false, + GenesisConfig { + system: SystemConfig { + code: wasm_binary.to_vec(), + changes_trie_config: Default::default(), + }, + pallet_balances: BalancesConfig { + balances: endowed_accounts + .iter() + .cloned() + .map(|k| (k, 1 << 100)) + .collect(), + }, + pallet_aura: AuraConfig { + authorities: initial_authorities.iter().map(|x| (x.0.clone())).collect(), + }, + pallet_grandpa: GrandpaConfig { + authorities: initial_authorities + .iter() + .map(|x| (x.1.clone(), 1)) + .collect(), + }, + pallet_treasury: Default::default(), + pallet_sudo: SudoConfig { key: root_key }, + pallet_vesting: VestingConfig { + vesting: vested_accounts + .iter() + .cloned() + .map(|k| (k, 1000, 100, 1 << 98)) + .collect(), + }, + pallet_nft: NftConfig { + collection_id: vec![( + 1, + Collection { + owner: get_account_id_from_seed::("Alice"), + mode: CollectionMode::NFT, + access: AccessMode::Normal, + decimal_points: 0, + name: vec![], + description: vec![], + token_prefix: vec![], + mint_mode: false, offchain_schema: vec![], schema_version: SchemaVersion::default(), - sponsorship: SponsorshipState::Confirmed(get_account_id_from_seed::("Alice")), - const_on_chain_schema: vec![], + sponsorship: SponsorshipState::Confirmed(get_account_id_from_seed::("Alice")), + const_on_chain_schema: vec![], variable_on_chain_schema: vec![], limits: CollectionLimits::default() - }, - )], - nft_item_id: vec![], - fungible_item_id: vec![], - refungible_item_id: vec![], - chain_limit: ChainLimits { - collection_numbers_limit: 100000, - account_token_ownership_limit: 1000000, - collections_admins_limit: 5, - custom_data_limit: 2048, - nft_sponsor_transfer_timeout: 15, - fungible_sponsor_transfer_timeout: 15, + }, + )], + nft_item_id: vec![], + fungible_item_id: vec![], + refungible_item_id: vec![], + chain_limit: ChainLimits { + collection_numbers_limit: 100000, + account_token_ownership_limit: 1000000, + collections_admins_limit: 5, + custom_data_limit: 2048, + nft_sponsor_transfer_timeout: 15, + fungible_sponsor_transfer_timeout: 15, refungible_sponsor_transfer_timeout: 15, offchain_schema_limit: 1024, variable_on_chain_schema_limit: 1024, const_on_chain_schema_limit: 1024, - }, - }), - pallet_contracts: Some(ContractsConfig { - current_schedule: ContractsSchedule { - enable_println, - ..Default::default() - }, - }), - } + }, + }, + pallet_contracts: ContractsConfig { + current_schedule: ContractsSchedule { + enable_println, + ..Default::default() + }, + }, + pallet_evm: EVMConfig { + accounts: BTreeMap::new(), + }, + pallet_ethereum: EthereumConfig {}, + } } --- a/node/src/command.rs +++ b/node/src/command.rs @@ -75,7 +75,7 @@ let runner = cli.create_runner(cmd)?; runner.async_run(|config| { let PartialComponents { client, task_manager, import_queue, ..} - = service::new_partial(&config)?; + = service::new_partial(&config, &cli)?; Ok((cmd.run(client, import_queue), task_manager)) }) }, @@ -83,7 +83,7 @@ let runner = cli.create_runner(cmd)?; runner.async_run(|config| { let PartialComponents { client, task_manager, ..} - = service::new_partial(&config)?; + = service::new_partial(&config, &cli)?; Ok((cmd.run(client, config.database), task_manager)) }) }, @@ -91,7 +91,7 @@ let runner = cli.create_runner(cmd)?; runner.async_run(|config| { let PartialComponents { client, task_manager, ..} - = service::new_partial(&config)?; + = service::new_partial(&config, &cli)?; Ok((cmd.run(client, config.chain_spec), task_manager)) }) }, @@ -99,7 +99,7 @@ let runner = cli.create_runner(cmd)?; runner.async_run(|config| { let PartialComponents { client, task_manager, import_queue, ..} - = service::new_partial(&config)?; + = service::new_partial(&config, &cli)?; Ok((cmd.run(client, import_queue), task_manager)) }) }, @@ -111,7 +111,7 @@ let runner = cli.create_runner(cmd)?; runner.async_run(|config| { let PartialComponents { client, task_manager, backend, ..} - = service::new_partial(&config)?; + = service::new_partial(&config, &cli)?; Ok((cmd.run(client, backend), task_manager)) }) }, @@ -130,7 +130,7 @@ runner.run_node_until_exit(|config| async move { match config.role { Role::Light => service::new_light(config), - _ => service::new_full(config), + _ => service::new_full(config, &cli), }.map_err(sc_cli::Error::Service) }) } --- a/node/src/lib.rs +++ /dev/null @@ -1,3 +0,0 @@ -pub mod chain_spec; -pub mod service; -pub mod rpc; --- a/node/src/main.rs +++ b/node/src/main.rs @@ -14,5 +14,5 @@ mod rpc; fn main() -> sc_cli::Result<()> { - command::run() + command::run() } --- a/node/src/rpc.rs +++ b/node/src/rpc.rs @@ -3,17 +3,38 @@ //! used by Substrate nodes. This file extends those RPC definitions with //! capabilities that are specific to this project's runtime configuration. -#![warn(missing_docs)] - use std::sync::Arc; -use nft_runtime::{opaque::Block, AccountId, Balance, Index, BlockNumber}; +use std::collections::BTreeMap; +use fc_rpc_core::types::{PendingTransactions, FilterPool}; +use nft_runtime::{Hash, AccountId, Index, opaque::Block, Balance}; use sp_api::ProvideRuntimeApi; use sp_blockchain::{Error as BlockChainError, HeaderMetadata, HeaderBackend}; +use sc_client_api::{ + backend::{StorageProvider, Backend, StateBackend, AuxStore}, + client::BlockchainEvents +}; +use sc_rpc::SubscriptionTaskExecutor; +use sp_runtime::traits::BlakeTwo256; use sp_block_builder::BlockBuilder; -pub use sc_rpc_api::DenyUnsafe; +use sc_rpc_api::DenyUnsafe; use sp_transaction_pool::TransactionPool; -use pallet_contracts_rpc::{Contracts, ContractsApi}; +use sc_network::NetworkService; +use jsonrpc_pubsub::manager::SubscriptionManager; +use pallet_ethereum::EthereumStorageSchema; +use fc_rpc::{StorageOverride, SchemaV1Override}; + +/// Light client extra dependencies. +pub struct LightDeps { + /// The client instance to use. + pub client: Arc, + /// Transaction pool instance. + pub pool: Arc

, + /// Remote access to the blockchain (async). + pub remote_blockchain: Arc>, + /// Fetcher instance. + pub fetcher: Arc, +} /// Full client dependencies. pub struct FullDeps { @@ -23,47 +44,155 @@ pub pool: Arc

, /// Whether to deny unsafe calls pub deny_unsafe: DenyUnsafe, + /// The Node authority flag + pub is_authority: bool, + /// Whether to enable dev signer + pub enable_dev_signer: bool, + /// Network service + pub network: Arc>, + /// Ethereum pending transactions. + pub pending_transactions: PendingTransactions, + /// EthFilterApi pool. + pub filter_pool: Option, + /// Backend. + pub backend: Arc>, } /// Instantiate all full RPC extensions. -pub fn create_full( +pub fn create_full( deps: FullDeps, + subscription_task_executor: SubscriptionTaskExecutor ) -> jsonrpc_core::IoHandler where - C: ProvideRuntimeApi, - C: HeaderBackend + HeaderMetadata + 'static, + BE: Backend + 'static, + BE::State: StateBackend, + C: ProvideRuntimeApi + StorageProvider + AuxStore, + C: BlockchainEvents, + C: HeaderBackend + HeaderMetadata, C: Send + Sync + 'static, C::Api: substrate_frame_rpc_system::AccountNonceApi, - C::Api: pallet_transaction_payment_rpc::TransactionPaymentRuntimeApi, C::Api: BlockBuilder, - C::Api: pallet_contracts_rpc::ContractsRuntimeApi, - P: TransactionPool + 'static, + C::Api: pallet_transaction_payment_rpc::TransactionPaymentRuntimeApi, + C::Api: fp_rpc::EthereumRuntimeRPCApi, + P: TransactionPool + 'static, { use substrate_frame_rpc_system::{FullSystem, SystemApi}; use pallet_transaction_payment_rpc::{TransactionPayment, TransactionPaymentApi}; + use fc_rpc::{ + EthApi, EthApiServer, EthFilterApi, EthFilterApiServer, NetApi, NetApiServer, + EthPubSubApi, EthPubSubApiServer, Web3Api, Web3ApiServer, EthDevSigner, EthSigner, + HexEncodedIdProvider, + }; let mut io = jsonrpc_core::IoHandler::default(); let FullDeps { client, pool, deny_unsafe, + is_authority, + network, + pending_transactions, + filter_pool, + backend, + enable_dev_signer, } = deps; io.extend_with( - SystemApi::to_delegate(FullSystem::new(client.clone(), pool, deny_unsafe)) + SystemApi::to_delegate(FullSystem::new(client.clone(), pool.clone(), deny_unsafe)) ); io.extend_with( TransactionPaymentApi::to_delegate(TransactionPayment::new(client.clone())) ); - io.extend_with( - ContractsApi::to_delegate(Contracts::new(client.clone())) + // io.extend_with( + // ContractsApi::to_delegate(Contracts::new(client.clone())) + // ); + + let mut signers = Vec::new(); + if enable_dev_signer { + signers.push(Box::new(EthDevSigner::new()) as Box); + } + let mut overrides = BTreeMap::new(); + overrides.insert( + EthereumStorageSchema::V1, + Box::new(SchemaV1Override::new(client.clone())) as Box + Send + Sync> + ); + io.extend_with( + EthApiServer::to_delegate(EthApi::new( + client.clone(), + pool.clone(), + nft_runtime::TransactionConverter, + network.clone(), + pending_transactions.clone(), + signers, + overrides, + backend, + is_authority, + )) + ); + + if let Some(filter_pool) = filter_pool { + io.extend_with( + EthFilterApiServer::to_delegate(EthFilterApi::new( + client.clone(), + filter_pool.clone(), + 500 as usize, // max stored filters + )) + ); + } + + io.extend_with( + NetApiServer::to_delegate(NetApi::new( + client.clone(), + network.clone(), + )) ); - - // Extend this RPC with a custom API by using the following syntax. - // `YourRpcStruct` should have a reference to a client, which is needed - // to call into the runtime. - // `io.extend_with(YourRpcTrait::to_delegate(YourRpcStruct::new(ReferenceToClient, ...)));` + io.extend_with( + Web3ApiServer::to_delegate(Web3Api::new( + client.clone(), + )) + ); + + io.extend_with( + EthPubSubApiServer::to_delegate(EthPubSubApi::new( + pool.clone(), + client.clone(), + network.clone(), + SubscriptionManager::::with_id_provider( + HexEncodedIdProvider::default(), + Arc::new(subscription_task_executor) + ), + )) + ); + + io +} + +/// Instantiate all Light RPC extensions. +pub fn create_light( + deps: LightDeps, +) -> jsonrpc_core::IoHandler where + C: sp_blockchain::HeaderBackend, + C: Send + Sync + 'static, + F: sc_client_api::light::Fetcher + 'static, + P: TransactionPool + 'static, + M: jsonrpc_core::Metadata + Default, +{ + use substrate_frame_rpc_system::{LightSystem, SystemApi}; + + let LightDeps { + client, + pool, + remote_blockchain, + fetcher + } = deps; + let mut io = jsonrpc_core::IoHandler::default(); + io.extend_with( + SystemApi::::to_delegate( + LightSystem::new(client, remote_blockchain, fetcher, pool) + ) + ); + io } --- a/node/src/service.rs +++ b/node/src/service.rs @@ -5,17 +5,26 @@ // file 'LICENSE', which is part of this source code package. // -use std::sync::Arc; -use std::time::Duration; -use sc_client_api::{ExecutorProvider, RemoteBackend}; -use nft_runtime::{self, opaque::Block, RuntimeApi}; -use sc_service::{error::Error as ServiceError, Configuration, TaskManager}; -use sp_inherents::InherentDataProviders; +use std::{sync::{Arc, Mutex}, cell::RefCell, time::Duration, collections::{HashMap, BTreeMap}}; +use fc_rpc::EthTask; +use fc_rpc_core::types::{FilterPool, PendingTransactions}; +use sc_client_api::{ExecutorProvider, RemoteBackend, BlockchainEvents}; +use fc_consensus::FrontierBlockImport; +use fc_mapping_sync::MappingSyncWorker; +use nft_runtime::{self, opaque::Block, RuntimeApi, SLOT_DURATION}; +use sc_service::{error::Error as ServiceError, Configuration, TaskManager, BasePath}; +use sp_inherents::{InherentDataProviders, ProvideInherentData, InherentIdentifier, InherentData}; use sc_executor::native_executor_instance; pub use sc_executor::NativeExecutor; +use sc_consensus_aura::{ImportQueueParams, StartAuraParams, SlotProportion}; use sp_consensus_aura::sr25519::{AuthorityPair as AuraPair}; use sc_finality_grandpa::SharedVoterState; -use sc_keystore::LocalKeystore; +use sp_timestamp::InherentError; +use sc_telemetry::{Telemetry, TelemetryWorker}; +use sc_cli::SubstrateCli; +use futures::StreamExt; + +use crate::cli::Cli; // Our native executor instance. native_executor_instance!( @@ -29,30 +38,96 @@ type FullBackend = sc_service::TFullBackend; type FullSelectChain = sc_consensus::LongestChain; -pub fn new_partial(config: &Configuration) -> Result, - sc_transaction_pool::FullPool, - ( - sc_consensus_aura::AuraBlockImport< +pub type ConsensusResult = ( + sc_consensus_aura::AuraBlockImport< + Block, + FullClient, + FrontierBlockImport< Block, - FullClient, sc_finality_grandpa::GrandpaBlockImport, - AuraPair + FullClient >, - sc_finality_grandpa::LinkHalf, - ) ->, ServiceError> { - if config.keystore_remote.is_some() { - return Err(ServiceError::Other( - format!("Remote Keystores are not supported."))) + AuraPair + >, + sc_finality_grandpa::LinkHalf +); + +/// Provide a mock duration starting at 0 in millisecond for timestamp inherent. +/// Each call will increment timestamp by slot_duration making Aura think time has passed. +pub struct MockTimestampInherentDataProvider; + +pub const INHERENT_IDENTIFIER: InherentIdentifier = *b"timstap0"; + +thread_local!(static TIMESTAMP: RefCell = RefCell::new(0)); + +impl ProvideInherentData for MockTimestampInherentDataProvider { + fn inherent_identifier(&self) -> &'static InherentIdentifier { + &INHERENT_IDENTIFIER } + + fn provide_inherent_data( + &self, + inherent_data: &mut InherentData, + ) -> Result<(), sp_inherents::Error> { + TIMESTAMP.with(|x| { + *x.borrow_mut() += SLOT_DURATION; + inherent_data.put_data(INHERENT_IDENTIFIER, &*x.borrow()) + }) + } + + fn error_to_string(&self, error: &[u8]) -> Option { + InherentError::try_from(&INHERENT_IDENTIFIER, error).map(|e| format!("{:?}", e)) + } +} + +pub fn open_frontier_backend(config: &Configuration) -> Result>, String> { + let config_dir = config.base_path.as_ref() + .map(|base_path| base_path.config_dir(config.chain_spec.id())) + .unwrap_or_else(|| { + BasePath::from_project("", "", &crate::cli::Cli::executable_name()) + .config_dir(config.chain_spec.id()) + }); + let database_dir = config_dir.join("frontier").join("db"); + + Ok(Arc::new(fc_db::Backend::::new(&fc_db::DatabaseSettings { + source: fc_db::DatabaseSettingsSrc::RocksDb { + path: database_dir, + cache_size: 0, + } + })?)) +} + +pub fn new_partial(config: &Configuration, #[allow(unused_variables)] cli: &Cli) -> Result< + sc_service::PartialComponents< + FullClient, FullBackend, FullSelectChain, + sp_consensus::import_queue::BasicQueue>, + sc_transaction_pool::FullPool, + (ConsensusResult, PendingTransactions, Option, Arc>, Option), +>, ServiceError> { let inherent_data_providers = sp_inherents::InherentDataProviders::new(); + let telemetry = config.telemetry_endpoints.clone() + .filter(|x| !x.is_empty()) + .map(|endpoints| -> Result<_, sc_telemetry::Error> { + let worker = TelemetryWorker::new(16)?; + let telemetry = worker.handle().new_telemetry(endpoints); + Ok((worker, telemetry)) + }) + .transpose()?; + let (client, backend, keystore_container, task_manager) = - sc_service::new_full_parts::(&config)?; + sc_service::new_full_parts::( + &config, + telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()), + )?; let client = Arc::new(client); + let telemetry = telemetry + .map(|(worker, telemetry)| { + task_manager.spawn_handle().spawn("telemetry", worker.run()); + telemetry + }); + let select_chain = sc_consensus::LongestChain::new(backend.clone()); let transaction_pool = sc_transaction_pool::BasicPool::new_full( @@ -63,64 +138,77 @@ client.clone(), ); + let pending_transactions: PendingTransactions + = Some(Arc::new(Mutex::new(HashMap::new()))); + + let filter_pool: Option + = Some(Arc::new(Mutex::new(BTreeMap::new()))); + + let frontier_backend = open_frontier_backend(config)?; + let (grandpa_block_import, grandpa_link) = sc_finality_grandpa::block_import( - client.clone(), &(client.clone() as Arc<_>), select_chain.clone(), + client.clone(), + &(client.clone() as Arc<_>), + select_chain.clone(), + telemetry.as_ref().map(|x| x.handle()), )?; + let frontier_block_import = FrontierBlockImport::new( + grandpa_block_import.clone(), + client.clone(), + frontier_backend.clone(), + ); + let aura_block_import = sc_consensus_aura::AuraBlockImport::<_, _, _, AuraPair>::new( - grandpa_block_import.clone(), client.clone(), + frontier_block_import, client.clone(), ); - let import_queue = sc_consensus_aura::import_queue::<_, _, _, AuraPair, _, _>( - sc_consensus_aura::slot_duration(&*client)?, - aura_block_import.clone(), - Some(Box::new(grandpa_block_import.clone())), - client.clone(), - inherent_data_providers.clone(), - &task_manager.spawn_handle(), - config.prometheus_registry(), - sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()), + let import_queue = sc_consensus_aura::import_queue::( + ImportQueueParams { + slot_duration: sc_consensus_aura::slot_duration(&*client)?, + block_import: aura_block_import.clone(), + justification_import: Some(Box::new(grandpa_block_import.clone())), + client: client.clone(), + inherent_data_providers: inherent_data_providers.clone(), + spawner: &task_manager.spawn_essential_handle(), + registry: config.prometheus_registry(), + can_author_with: sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()), + check_for_equivocation: Default::default(), + telemetry: telemetry.as_ref().map(|x| x.handle()), + } )?; Ok(sc_service::PartialComponents { - client, backend, task_manager, import_queue, keystore_container, select_chain, transaction_pool, - inherent_data_providers, - other: (aura_block_import, grandpa_link), + client, backend, task_manager, import_queue, keystore_container, + select_chain, transaction_pool, inherent_data_providers, + other: ( + (aura_block_import, grandpa_link), + pending_transactions, + filter_pool, + frontier_backend, + telemetry, + ) }) } -fn remote_keystore(_url: &String) -> Result, &'static str> { - // FIXME: here would the concrete keystore be built, - // must return a concrete type (NOT `LocalKeystore`) that - // implements `CryptoStore` and `SyncCryptoStore` - Err("Remote Keystore not supported.") -} - /// Builds a new service for a full client. -pub fn new_full(mut config: Configuration) -> Result { - let sc_service::PartialComponents { - client, backend, mut task_manager, import_queue, mut keystore_container, select_chain, transaction_pool, - inherent_data_providers, - other: (block_import, grandpa_link), - } = new_partial(&config)?; +pub fn new_full( + config: Configuration, + cli: &Cli, +) -> Result { + let enable_dev_signer = cli.run.shared_params.dev; - if let Some(url) = &config.keystore_remote { - match remote_keystore(url) { - Ok(k) => keystore_container.set_remote_keystore(k), - Err(e) => { - return Err(ServiceError::Other( - format!("Error hooking up remote keystore for {}: {}", url, e))) - } - }; - } - - config.network.extra_sets.push(sc_finality_grandpa::grandpa_peers_set_config()); + let sc_service::PartialComponents { + client, backend, mut task_manager, import_queue, keystore_container, + select_chain, transaction_pool, inherent_data_providers, + other: (consensus_result, pending_transactions, filter_pool, frontier_backend, mut telemetry), + } = new_partial(&config, cli)?; let (network, network_status_sinks, system_rpc_tx, network_starter) = sc_service::build_network(sc_service::BuildNetworkParams { config: &config, - client: client.clone(), - transaction_pool: transaction_pool.clone(), + client: Arc::clone(&client), + transaction_pool: Arc::clone(&transaction_pool), spawn_handle: task_manager.spawn_handle(), import_queue, on_demand: None, @@ -129,7 +217,7 @@ if config.offchain_worker.enabled { sc_service::build_offchain_workers( - &config, backend.clone(), task_manager.spawn_handle(), client.clone(), network.clone(), + &config, task_manager.spawn_handle(), Arc::clone(&client), network.clone(), ); } @@ -139,121 +227,191 @@ let name = config.network.node_name.clone(); let enable_grandpa = !config.disable_grandpa; let prometheus_registry = config.prometheus_registry().cloned(); + let is_authority = role.is_authority(); + let subscription_task_executor = sc_rpc::SubscriptionTaskExecutor::new(task_manager.spawn_handle()); let rpc_extensions_builder = { - let client = client.clone(); - let pool = transaction_pool.clone(); + let client = Arc::clone(&client); + let pool = Arc::clone(&transaction_pool); + let network = network.clone(); + let pending = pending_transactions.clone(); + let filter_pool = filter_pool.clone(); + let frontier_backend = frontier_backend.clone(); Box::new(move |deny_unsafe, _| { let deps = crate::rpc::FullDeps { - client: client.clone(), - pool: pool.clone(), + client: Arc::clone(&client), + pool: Arc::clone(&pool), deny_unsafe, + is_authority, + enable_dev_signer, + network: network.clone(), + pending_transactions: pending.clone(), + filter_pool: filter_pool.clone(), + backend: frontier_backend.clone(), }; - - crate::rpc::create_full(deps) + crate::rpc::create_full( + deps, + subscription_task_executor.clone() + ) }) }; - let (_rpc_handlers, telemetry_connection_notifier) = sc_service::spawn_tasks( - sc_service::SpawnTasksParams { - network: network.clone(), - client: client.clone(), - keystore: keystore_container.sync_keystore(), - task_manager: &mut task_manager, - transaction_pool: transaction_pool.clone(), - rpc_extensions_builder, - on_demand: None, - remote_blockchain: None, - backend, - network_status_sinks, - system_rpc_tx, - config, - }, - )?; + task_manager.spawn_essential_handle().spawn( + "frontier-mapping-sync-worker", + MappingSyncWorker::new( + client.import_notification_stream(), + Duration::new(6, 0), + Arc::clone(&client), + Arc::clone(&backend), + frontier_backend.clone(), + ).for_each(|()| futures::future::ready(())) + ); + + let _rpc_handlers = sc_service::spawn_tasks(sc_service::SpawnTasksParams { + network: network.clone(), + client: Arc::clone(&client), + keystore: keystore_container.sync_keystore(), + task_manager: &mut task_manager, + transaction_pool: Arc::clone(&transaction_pool), + rpc_extensions_builder: rpc_extensions_builder, + on_demand: None, + remote_blockchain: None, + backend, network_status_sinks, system_rpc_tx, config, telemetry: telemetry.as_mut(), + })?; + + // Spawn Frontier EthFilterApi maintenance task. + if let Some(filter_pool) = filter_pool { + // Each filter is allowed to stay in the pool for 100 blocks. + const FILTER_RETAIN_THRESHOLD: u64 = 100; + task_manager.spawn_essential_handle().spawn( + "frontier-filter-pool", + EthTask::filter_pool_task( + Arc::clone(&client), + filter_pool, + FILTER_RETAIN_THRESHOLD, + ) + ); + } + // Spawn Frontier pending transactions maintenance task (as essential, otherwise we leak). + if let Some(pending_transactions) = pending_transactions { + const TRANSACTION_RETAIN_THRESHOLD: u64 = 5; + task_manager.spawn_essential_handle().spawn( + "frontier-pending-transactions", + EthTask::pending_transaction_task( + Arc::clone(&client), + pending_transactions, + TRANSACTION_RETAIN_THRESHOLD, + ) + ); + } + + let (aura_block_import, grandpa_link) = consensus_result; + if role.is_authority() { let proposer = sc_basic_authorship::ProposerFactory::new( task_manager.spawn_handle(), - client.clone(), - transaction_pool, + Arc::clone(&client), + Arc::clone(&transaction_pool), prometheus_registry.as_ref(), + telemetry.as_ref().map(|x| x.handle()), ); let can_author_with = sp_consensus::CanAuthorWithNativeVersion::new(client.executor().clone()); - - let aura = sc_consensus_aura::start_aura::<_, _, _, _, _, AuraPair, _, _, _, _>( - sc_consensus_aura::slot_duration(&*client)?, - client.clone(), - select_chain, - block_import, - proposer, - network.clone(), - inherent_data_providers.clone(), - force_authoring, - backoff_authoring_blocks, - keystore_container.sync_keystore(), - can_author_with, + let aura = sc_consensus_aura::start_aura::( + StartAuraParams { + slot_duration: sc_consensus_aura::slot_duration(&*client)?, + client: Arc::clone(&client), + select_chain, + block_import: aura_block_import, + proposer_factory: proposer, + sync_oracle: network.clone(), + inherent_data_providers: inherent_data_providers.clone(), + force_authoring, + backoff_authoring_blocks, + keystore: keystore_container.sync_keystore(), + can_author_with, + block_proposal_slot_portion: SlotProportion::new(2f32 / 3f32), + telemetry: telemetry.as_ref().map(|x| x.handle()), + } )?; // the AURA authoring task is considered essential, i.e. if it // fails we take down the service with it. task_manager.spawn_essential_handle().spawn_blocking("aura", aura); - } - // if the node isn't actively participating in consensus then it doesn't - // need a keystore, regardless of which protocol we use below. - let keystore = if role.is_authority() { - Some(keystore_container.sync_keystore()) - } else { - None - }; + // if the node isn't actively participating in consensus then it doesn't + // need a keystore, regardless of which protocol we use below. + let keystore = if role.is_authority() { + Some(keystore_container.sync_keystore()) + } else { + None + }; - let grandpa_config = sc_finality_grandpa::Config { - // FIXME #1578 make this available through chainspec - gossip_duration: Duration::from_millis(333), - justification_period: 512, - name: Some(name), - observer_enabled: false, - keystore, - is_authority: role.is_network_authority(), - }; - - if enable_grandpa { - // start the full GRANDPA voter - // NOTE: non-authorities could run the GRANDPA observer protocol, but at - // this point the full voter should provide better guarantees of block - // and vote data availability than the observer. The observer has not - // been tested extensively yet and having most nodes in a network run it - // could lead to finality stalls. - let grandpa_config = sc_finality_grandpa::GrandpaParams { - config: grandpa_config, - link: grandpa_link, - network, - telemetry_on_connect: telemetry_connection_notifier.map(|x| x.on_connect_stream()), - voting_rule: sc_finality_grandpa::VotingRulesBuilder::default().build(), - prometheus_registry, - shared_voter_state: SharedVoterState::empty(), + let grandpa_config = sc_finality_grandpa::Config { + // FIXME #1578 make this available through chainspec + gossip_duration: Duration::from_millis(333), + justification_period: 512, + name: Some(name), + observer_enabled: false, + keystore, + is_authority: role.is_authority(), + telemetry: telemetry.as_ref().map(|x| x.handle()), }; - // the GRANDPA voter task is considered infallible, i.e. - // if it fails we take down the service with it. - task_manager.spawn_essential_handle().spawn_blocking( - "grandpa-voter", - sc_finality_grandpa::run_grandpa_voter(grandpa_config)? - ); + if enable_grandpa { + // start the full GRANDPA voter + // NOTE: non-authorities could run the GRANDPA observer protocol, but at + // this point the full voter should provide better guarantees of block + // and vote data availability than the observer. The observer has not + // been tested extensively yet and having most nodes in a network run it + // could lead to finality stalls. + let grandpa_config = sc_finality_grandpa::GrandpaParams { + config: grandpa_config, + link: grandpa_link, + network, + telemetry: telemetry.as_ref().map(|x| x.handle()), + voting_rule: sc_finality_grandpa::VotingRulesBuilder::default().build(), + prometheus_registry, + shared_voter_state: SharedVoterState::empty(), + }; + + // the GRANDPA voter task is considered infallible, i.e. + // if it fails we take down the service with it. + task_manager.spawn_essential_handle().spawn_blocking( + "grandpa-voter", + sc_finality_grandpa::run_grandpa_voter(grandpa_config)? + ); + } } network_starter.start_network(); Ok(task_manager) } /// Builds a new service for a light client. -pub fn new_light(mut config: Configuration) -> Result { +pub fn new_light(config: Configuration) -> Result { + let telemetry = config.telemetry_endpoints.clone() + .filter(|x| !x.is_empty()) + .map(|endpoints| -> Result<_, sc_telemetry::Error> { + let worker = TelemetryWorker::new(16)?; + let telemetry = worker.handle().new_telemetry(endpoints); + Ok((worker, telemetry)) + }) + .transpose()?; + let (client, backend, keystore_container, mut task_manager, on_demand) = - sc_service::new_light_parts::(&config)?; + sc_service::new_light_parts::( + &config, + telemetry.as_ref().map(|(_, telemetry)| telemetry.handle()), + )?; - config.network.extra_sets.push(sc_finality_grandpa::grandpa_peers_set_config()); + let mut telemetry = telemetry + .map(|(worker, telemetry)| { + task_manager.spawn_handle().spawn("telemetry", worker.run()); + telemetry + }); let select_chain = sc_consensus::LongestChain::new(backend.clone()); @@ -269,23 +427,32 @@ client.clone(), &(client.clone() as Arc<_>), select_chain.clone(), + telemetry.as_ref().map(|x| x.handle()), + )?; + + let import_queue = sc_consensus_aura::import_queue::( + ImportQueueParams { + slot_duration: sc_consensus_aura::slot_duration(&*client)?, + block_import: grandpa_block_import.clone(), + justification_import: Some(Box::new(grandpa_block_import)), + client: client.clone(), + inherent_data_providers: InherentDataProviders::new(), + spawner: &task_manager.spawn_essential_handle(), + registry: config.prometheus_registry(), + can_author_with: sp_consensus::NeverCanAuthor, + check_for_equivocation: Default::default(), + telemetry: telemetry.as_ref().map(|x| x.handle()), + } )?; - let aura_block_import = sc_consensus_aura::AuraBlockImport::<_, _, _, AuraPair>::new( - grandpa_block_import.clone(), - client.clone(), - ); + let light_deps = crate::rpc::LightDeps { + remote_blockchain: backend.remote_blockchain(), + fetcher: on_demand.clone(), + client: client.clone(), + pool: transaction_pool.clone(), + }; - let import_queue = sc_consensus_aura::import_queue::<_, _, _, AuraPair, _, _>( - sc_consensus_aura::slot_duration(&*client)?, - aura_block_import, - Some(Box::new(grandpa_block_import)), - client.clone(), - InherentDataProviders::new(), - &task_manager.spawn_handle(), - config.prometheus_registry(), - sp_consensus::NeverCanAuthor, - )?; + let rpc_extensions = crate::rpc::create_light(light_deps); let (network, network_status_sinks, system_rpc_tx, network_starter) = sc_service::build_network(sc_service::BuildNetworkParams { @@ -300,7 +467,7 @@ if config.offchain_worker.enabled { sc_service::build_offchain_workers( - &config, backend.clone(), task_manager.spawn_handle(), client.clone(), network.clone(), + &config, task_manager.spawn_handle(), client.clone(), network.clone(), ); } @@ -309,7 +476,7 @@ transaction_pool, task_manager: &mut task_manager, on_demand: Some(on_demand), - rpc_extensions_builder: Box::new(|_, _| ()), + rpc_extensions_builder: Box::new(sc_service::NoopRpcExtensionBuilder(rpc_extensions)), config, client, keystore: keystore_container.sync_keystore(), @@ -317,6 +484,7 @@ network, network_status_sinks, system_rpc_tx, + telemetry: telemetry.as_mut(), })?; network_starter.start_network();