// This file is part of Substrate. // Copyright (C) 2017-2022 Parity Technologies (UK) Ltd. // SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0 // This program is free software: you can redistribute it and/or modify // it under the terms of the GNU General Public License as published by // the Free Software Foundation, either version 3 of the License, or // (at your option) any later version. // This program is distributed in the hope that it will be useful, // but WITHOUT ANY WARRANTY; without even the implied warranty of // MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the // GNU General Public License for more details. // You should have received a copy of the GNU General Public License // along with this program. If not, see . //! Substrate blockchain API. mod chain_full; #[cfg(test)] mod tests; use futures::{future, StreamExt, TryStreamExt}; use log::warn; use rpc::{ futures::{stream, FutureExt, SinkExt, Stream}, Result as RpcResult, }; use std::sync::Arc; use jsonrpc_pubsub::{manager::SubscriptionManager, typed::Subscriber, SubscriptionId}; use sc_client_api::BlockchainEvents; use sp_rpc::{list::ListOrValue, number::NumberOrHex}; use sp_runtime::{ generic::{BlockId, SignedBlock}, traits::{Block as BlockT, Header, NumberFor}, }; use self::error::{Error, FutureResult, Result}; use sc_client_api::BlockBackend; pub use sc_rpc_api::chain::*; use sp_blockchain::HeaderBackend; /// Blockchain backend API trait ChainBackend: Send + Sync + 'static where Block: BlockT + 'static, Block::Header: Unpin, Client: HeaderBackend + BlockchainEvents + 'static, { /// Get client reference. fn client(&self) -> &Arc; /// Get subscriptions reference. fn subscriptions(&self) -> &SubscriptionManager; /// Tries to unwrap passed block hash, or uses best block hash otherwise. fn unwrap_or_best(&self, hash: Option) -> Block::Hash { match hash { None => self.client().info().best_hash, Some(hash) => hash, } } /// Get header of a relay chain block. fn header(&self, hash: Option) -> FutureResult>; /// Get header and body of a relay chain block. fn block(&self, hash: Option) -> FutureResult>>; /// Get hash of the n-th block in the canon chain. /// /// By default returns latest block hash. fn block_hash(&self, number: Option) -> Result> { match number { None => Ok(Some(self.client().info().best_hash)), Some(num_or_hex) => { // FIXME <2329>: Database seems to limit the block number to u32 for no reason let block_num: u32 = num_or_hex.try_into().map_err(|_| { Error::Other(format!( "`{:?}` > u32::MAX, the max block number is u32.", num_or_hex )) })?; let block_num = >::from(block_num); Ok(self .client() .header(BlockId::number(block_num)) .map_err(client_err)? .map(|h| h.hash())) }, } } /// Get hash of the last finalized block in the canon chain. fn finalized_head(&self) -> Result { Ok(self.client().info().finalized_hash) } /// All new head subscription fn subscribe_all_heads( &self, _metadata: crate::Metadata, subscriber: Subscriber, ) { subscribe_headers( self.client(), self.subscriptions(), subscriber, || self.client().info().best_hash, || { self.client() .import_notification_stream() .map(|notification| Ok::<_, rpc::Error>(notification.header)) }, ) } /// Unsubscribe from all head subscription. fn unsubscribe_all_heads( &self, _metadata: Option, id: SubscriptionId, ) -> RpcResult { Ok(self.subscriptions().cancel(id)) } /// New best head subscription fn subscribe_new_heads( &self, _metadata: crate::Metadata, subscriber: Subscriber, ) { subscribe_headers( self.client(), self.subscriptions(), subscriber, || self.client().info().best_hash, || { self.client() .import_notification_stream() .filter(|notification| future::ready(notification.is_new_best)) .map(|notification| Ok::<_, rpc::Error>(notification.header)) }, ) } /// Unsubscribe from new best head subscription. fn unsubscribe_new_heads( &self, _metadata: Option, id: SubscriptionId, ) -> RpcResult { Ok(self.subscriptions().cancel(id)) } /// Finalized head subscription fn subscribe_finalized_heads( &self, _metadata: crate::Metadata, subscriber: Subscriber, ) { subscribe_headers( self.client(), self.subscriptions(), subscriber, || self.client().info().finalized_hash, || { self.client() .finality_notification_stream() .map(|notification| Ok::<_, rpc::Error>(notification.header)) }, ) } /// Unsubscribe from finalized head subscription. fn unsubscribe_finalized_heads( &self, _metadata: Option, id: SubscriptionId, ) -> RpcResult { Ok(self.subscriptions().cancel(id)) } } /// Create new state API that works on full node. pub fn new_full( client: Arc, subscriptions: SubscriptionManager, ) -> Chain where Block: BlockT + 'static, Block::Header: Unpin, Client: BlockBackend + HeaderBackend + BlockchainEvents + 'static, { Chain { backend: Box::new(self::chain_full::FullChain::new(client, subscriptions)) } } /// Chain API with subscriptions support. pub struct Chain { backend: Box>, } impl ChainApi, Block::Hash, Block::Header, SignedBlock> for Chain where Block: BlockT + 'static, Block::Header: Unpin, Client: HeaderBackend + BlockchainEvents + 'static, { type Metadata = crate::Metadata; fn header(&self, hash: Option) -> FutureResult> { self.backend.header(hash) } fn block(&self, hash: Option) -> FutureResult>> { self.backend.block(hash) } fn block_hash( &self, number: Option>, ) -> Result>> { match number { None => self.backend.block_hash(None).map(ListOrValue::Value), Some(ListOrValue::Value(number)) => self.backend.block_hash(Some(number)).map(ListOrValue::Value), Some(ListOrValue::List(list)) => Ok(ListOrValue::List( list.into_iter() .map(|number| self.backend.block_hash(Some(number))) .collect::>()?, )), } } fn finalized_head(&self) -> Result { self.backend.finalized_head() } fn subscribe_all_heads(&self, metadata: Self::Metadata, subscriber: Subscriber) { self.backend.subscribe_all_heads(metadata, subscriber) } fn unsubscribe_all_heads( &self, metadata: Option, id: SubscriptionId, ) -> RpcResult { self.backend.unsubscribe_all_heads(metadata, id) } fn subscribe_new_heads(&self, metadata: Self::Metadata, subscriber: Subscriber) { self.backend.subscribe_new_heads(metadata, subscriber) } fn unsubscribe_new_heads( &self, metadata: Option, id: SubscriptionId, ) -> RpcResult { self.backend.unsubscribe_new_heads(metadata, id) } fn subscribe_finalized_heads( &self, metadata: Self::Metadata, subscriber: Subscriber, ) { self.backend.subscribe_finalized_heads(metadata, subscriber) } fn unsubscribe_finalized_heads( &self, metadata: Option, id: SubscriptionId, ) -> RpcResult { self.backend.unsubscribe_finalized_heads(metadata, id) } } /// Subscribe to new headers. fn subscribe_headers( client: &Arc, subscriptions: &SubscriptionManager, subscriber: Subscriber, best_block_hash: G, stream: F, ) where Block: BlockT + 'static, Block::Header: Unpin, Client: HeaderBackend + 'static, F: FnOnce() -> S, G: FnOnce() -> Block::Hash, S: Stream> + Send + 'static, { subscriptions.add(subscriber, |sink| { // send current head right at the start. let header = client .header(BlockId::Hash(best_block_hash())) .map_err(client_err) .and_then(|header| { header.ok_or_else(|| Error::Other("Best header missing.".to_string())) }) .map_err(Into::into); // send further subscriptions let stream = stream() .inspect_err(|e| warn!("Block notification stream error: {:?}", e)) .map(Ok); stream::iter(vec![Ok(header)]) .chain(stream) .forward(sink.sink_map_err(|e| warn!("Error sending notifications: {:?}", e))) // we ignore the resulting Stream (if the first stream is over we are unsubscribed) .map(|_| ()) }); } fn client_err(err: sp_blockchain::Error) -> Error { Error::Client(Box::new(err)) }