mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-05-30 22:11:02 +00:00
use tokio reactor to execute jsonrpsee futures (#1061)
This commit is contained in:
committed by
Bastian Köcher
parent
08fd53adef
commit
63d6fc436a
@@ -16,4 +16,5 @@ jsonrpsee-ws-client = "0.2"
|
||||
libsecp256k1 = { version = "0.3.4", default-features = false, features = ["hmac"] }
|
||||
log = "0.4.11"
|
||||
relay-utils = { path = "../utils" }
|
||||
tokio = "1.8"
|
||||
web3 = { git = "https://github.com/svyatonik/rust-web3.git", branch = "bump-deps" }
|
||||
|
||||
@@ -23,7 +23,7 @@ use crate::{ConnectionParams, Error, Result};
|
||||
|
||||
use jsonrpsee_ws_client::{WsClient as RpcClient, WsClientBuilder as RpcClientBuilder};
|
||||
use relay_utils::relay_loop::RECONNECT_DELAY;
|
||||
use std::sync::Arc;
|
||||
use std::{future::Future, sync::Arc};
|
||||
|
||||
/// Number of headers missing from the Ethereum node for us to consider node not synced.
|
||||
const MAJOR_SYNC_BLOCKS: u64 = 5;
|
||||
@@ -31,6 +31,7 @@ const MAJOR_SYNC_BLOCKS: u64 = 5;
|
||||
/// The client used to interact with an Ethereum node through RPC.
|
||||
#[derive(Clone)]
|
||||
pub struct Client {
|
||||
tokio: Arc<tokio::runtime::Runtime>,
|
||||
params: ConnectionParams,
|
||||
client: Arc<RpcClient>,
|
||||
}
|
||||
@@ -59,22 +60,23 @@ impl Client {
|
||||
/// Try to connect to Ethereum node. Returns Ethereum RPC client if connection has been established
|
||||
/// or error otherwise.
|
||||
pub async fn try_connect(params: ConnectionParams) -> Result<Self> {
|
||||
Ok(Self {
|
||||
client: Self::build_client(¶ms).await?,
|
||||
params,
|
||||
})
|
||||
let (tokio, client) = Self::build_client(¶ms).await?;
|
||||
Ok(Self { tokio, client, params })
|
||||
}
|
||||
|
||||
/// Build client to use in connection.
|
||||
async fn build_client(params: &ConnectionParams) -> Result<Arc<RpcClient>> {
|
||||
async fn build_client(params: &ConnectionParams) -> Result<(Arc<tokio::runtime::Runtime>, Arc<RpcClient>)> {
|
||||
let tokio = tokio::runtime::Runtime::new()?;
|
||||
let uri = format!("ws://{}:{}", params.host, params.port);
|
||||
let client = RpcClientBuilder::default().build(&uri).await?;
|
||||
Ok(Arc::new(client))
|
||||
let client = tokio
|
||||
.spawn(async move { RpcClientBuilder::default().build(&uri).await })
|
||||
.await??;
|
||||
Ok((Arc::new(tokio), Arc::new(client)))
|
||||
}
|
||||
|
||||
/// Reopen client connection.
|
||||
pub async fn reconnect(&mut self) -> Result<()> {
|
||||
self.client = Self::build_client(&self.params).await?;
|
||||
self.client = Self::build_client(&self.params).await?.1;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -82,113 +84,153 @@ impl Client {
|
||||
impl Client {
|
||||
/// Returns true if client is connected to at least one peer and is in synced state.
|
||||
pub async fn ensure_synced(&self) -> Result<()> {
|
||||
match Ethereum::syncing(&*self.client).await? {
|
||||
SyncState::NotSyncing => Ok(()),
|
||||
SyncState::Syncing(syncing) => {
|
||||
let missing_headers = syncing.highest_block.saturating_sub(syncing.current_block);
|
||||
if missing_headers > MAJOR_SYNC_BLOCKS.into() {
|
||||
return Err(Error::ClientNotSynced(missing_headers));
|
||||
}
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
match Ethereum::syncing(&*client).await? {
|
||||
SyncState::NotSyncing => Ok(()),
|
||||
SyncState::Syncing(syncing) => {
|
||||
let missing_headers = syncing.highest_block.saturating_sub(syncing.current_block);
|
||||
if missing_headers > MAJOR_SYNC_BLOCKS.into() {
|
||||
return Err(Error::ClientNotSynced(missing_headers));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Estimate gas usage for the given call.
|
||||
pub async fn estimate_gas(&self, call_request: CallRequest) -> Result<U256> {
|
||||
Ok(Ethereum::estimate_gas(&*self.client, call_request).await?)
|
||||
self.jsonrpsee_execute(move |client| async move { Ok(Ethereum::estimate_gas(&*client, call_request).await?) })
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve number of the best known block from the Ethereum node.
|
||||
pub async fn best_block_number(&self) -> Result<u64> {
|
||||
Ok(Ethereum::block_number(&*self.client).await?.as_u64())
|
||||
self.jsonrpsee_execute(move |client| async move { Ok(Ethereum::block_number(&*client).await?.as_u64()) })
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve number of the best known block from the Ethereum node.
|
||||
pub async fn header_by_number(&self, block_number: u64) -> Result<Header> {
|
||||
let get_full_tx_objects = false;
|
||||
let header = Ethereum::get_block_by_number(&*self.client, block_number, get_full_tx_objects).await?;
|
||||
match header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some() {
|
||||
true => Ok(header),
|
||||
false => Err(Error::IncompleteHeader),
|
||||
}
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
let get_full_tx_objects = false;
|
||||
let header = Ethereum::get_block_by_number(&*client, block_number, get_full_tx_objects).await?;
|
||||
match header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some() {
|
||||
true => Ok(header),
|
||||
false => Err(Error::IncompleteHeader),
|
||||
}
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve block header by its hash from Ethereum node.
|
||||
pub async fn header_by_hash(&self, hash: H256) -> Result<Header> {
|
||||
let get_full_tx_objects = false;
|
||||
let header = Ethereum::get_block_by_hash(&*self.client, hash, get_full_tx_objects).await?;
|
||||
match header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some() {
|
||||
true => Ok(header),
|
||||
false => Err(Error::IncompleteHeader),
|
||||
}
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
let get_full_tx_objects = false;
|
||||
let header = Ethereum::get_block_by_hash(&*client, hash, get_full_tx_objects).await?;
|
||||
match header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some() {
|
||||
true => Ok(header),
|
||||
false => Err(Error::IncompleteHeader),
|
||||
}
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve block header and its transactions by its number from Ethereum node.
|
||||
pub async fn header_by_number_with_transactions(&self, number: u64) -> Result<HeaderWithTransactions> {
|
||||
let get_full_tx_objects = true;
|
||||
let header =
|
||||
Ethereum::get_block_by_number_with_transactions(&*self.client, number, get_full_tx_objects).await?;
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
let get_full_tx_objects = true;
|
||||
let header = Ethereum::get_block_by_number_with_transactions(&*client, number, get_full_tx_objects).await?;
|
||||
|
||||
let is_complete_header = header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some();
|
||||
if !is_complete_header {
|
||||
return Err(Error::IncompleteHeader);
|
||||
}
|
||||
let is_complete_header = header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some();
|
||||
if !is_complete_header {
|
||||
return Err(Error::IncompleteHeader);
|
||||
}
|
||||
|
||||
let is_complete_transactions = header.transactions.iter().all(|tx| tx.raw.is_some());
|
||||
if !is_complete_transactions {
|
||||
return Err(Error::IncompleteTransaction);
|
||||
}
|
||||
let is_complete_transactions = header.transactions.iter().all(|tx| tx.raw.is_some());
|
||||
if !is_complete_transactions {
|
||||
return Err(Error::IncompleteTransaction);
|
||||
}
|
||||
|
||||
Ok(header)
|
||||
Ok(header)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve block header and its transactions by its hash from Ethereum node.
|
||||
pub async fn header_by_hash_with_transactions(&self, hash: H256) -> Result<HeaderWithTransactions> {
|
||||
let get_full_tx_objects = true;
|
||||
let header = Ethereum::get_block_by_hash_with_transactions(&*self.client, hash, get_full_tx_objects).await?;
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
let get_full_tx_objects = true;
|
||||
let header = Ethereum::get_block_by_hash_with_transactions(&*client, hash, get_full_tx_objects).await?;
|
||||
|
||||
let is_complete_header = header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some();
|
||||
if !is_complete_header {
|
||||
return Err(Error::IncompleteHeader);
|
||||
}
|
||||
let is_complete_header = header.number.is_some() && header.hash.is_some() && header.logs_bloom.is_some();
|
||||
if !is_complete_header {
|
||||
return Err(Error::IncompleteHeader);
|
||||
}
|
||||
|
||||
let is_complete_transactions = header.transactions.iter().all(|tx| tx.raw.is_some());
|
||||
if !is_complete_transactions {
|
||||
return Err(Error::IncompleteTransaction);
|
||||
}
|
||||
let is_complete_transactions = header.transactions.iter().all(|tx| tx.raw.is_some());
|
||||
if !is_complete_transactions {
|
||||
return Err(Error::IncompleteTransaction);
|
||||
}
|
||||
|
||||
Ok(header)
|
||||
Ok(header)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve transaction by its hash from Ethereum node.
|
||||
pub async fn transaction_by_hash(&self, hash: H256) -> Result<Option<Transaction>> {
|
||||
Ok(Ethereum::transaction_by_hash(&*self.client, hash).await?)
|
||||
self.jsonrpsee_execute(move |client| async move { Ok(Ethereum::transaction_by_hash(&*client, hash).await?) })
|
||||
.await
|
||||
}
|
||||
|
||||
/// Retrieve transaction receipt by transaction hash.
|
||||
pub async fn transaction_receipt(&self, transaction_hash: H256) -> Result<Receipt> {
|
||||
Ok(Ethereum::get_transaction_receipt(&*self.client, transaction_hash).await?)
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
Ok(Ethereum::get_transaction_receipt(&*client, transaction_hash).await?)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Get the nonce of the given account.
|
||||
pub async fn account_nonce(&self, address: Address) -> Result<U256> {
|
||||
Ok(Ethereum::get_transaction_count(&*self.client, address).await?)
|
||||
self.jsonrpsee_execute(
|
||||
move |client| async move { Ok(Ethereum::get_transaction_count(&*client, address).await?) },
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Submit an Ethereum transaction.
|
||||
///
|
||||
/// The transaction must already be signed before sending it through this method.
|
||||
pub async fn submit_transaction(&self, signed_raw_tx: SignedRawTx) -> Result<TransactionHash> {
|
||||
let transaction = Bytes(signed_raw_tx);
|
||||
let tx_hash = Ethereum::submit_transaction(&*self.client, transaction).await?;
|
||||
log::trace!(target: "bridge", "Sent transaction to Ethereum node: {:?}", tx_hash);
|
||||
Ok(tx_hash)
|
||||
self.jsonrpsee_execute(move |client| async move {
|
||||
let transaction = Bytes(signed_raw_tx);
|
||||
let tx_hash = Ethereum::submit_transaction(&*client, transaction).await?;
|
||||
log::trace!(target: "bridge", "Sent transaction to Ethereum node: {:?}", tx_hash);
|
||||
Ok(tx_hash)
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
/// Call Ethereum smart contract.
|
||||
pub async fn eth_call(&self, call_transaction: CallRequest) -> Result<Bytes> {
|
||||
Ok(Ethereum::call(&*self.client, call_transaction).await?)
|
||||
self.jsonrpsee_execute(move |client| async move { Ok(Ethereum::call(&*client, call_transaction).await?) })
|
||||
.await
|
||||
}
|
||||
|
||||
/// Execute jsonrpsee future in tokio context.
|
||||
async fn jsonrpsee_execute<MF, F, T>(&self, make_jsonrpsee_future: MF) -> Result<T>
|
||||
where
|
||||
MF: FnOnce(Arc<RpcClient>) -> F + Send + 'static,
|
||||
F: Future<Output = Result<T>> + Send,
|
||||
T: Send + 'static,
|
||||
{
|
||||
let client = self.client.clone();
|
||||
self.tokio
|
||||
.spawn(async move { make_jsonrpsee_future(client).await })
|
||||
.await?
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,8 @@ pub type Result<T> = std::result::Result<T, Error>;
|
||||
/// an Ethereum node through RPC.
|
||||
#[derive(Debug)]
|
||||
pub enum Error {
|
||||
/// IO error.
|
||||
Io(std::io::Error),
|
||||
/// An error that can occur when making an HTTP request to
|
||||
/// an JSON-RPC client.
|
||||
RpcError(RpcError),
|
||||
@@ -45,6 +47,8 @@ pub enum Error {
|
||||
/// The client we're connected to is not synced, so we can't rely on its state. Contains
|
||||
/// number of unsynced headers.
|
||||
ClientNotSynced(U256),
|
||||
/// Custom logic error.
|
||||
Custom(String),
|
||||
}
|
||||
|
||||
impl From<RpcError> for Error {
|
||||
@@ -53,6 +57,18 @@ impl From<RpcError> for Error {
|
||||
}
|
||||
}
|
||||
|
||||
impl From<std::io::Error> for Error {
|
||||
fn from(error: std::io::Error) -> Self {
|
||||
Error::Io(error)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<tokio::task::JoinError> for Error {
|
||||
fn from(error: tokio::task::JoinError) -> Self {
|
||||
Error::Custom(format!("Failed to wait tokio task: {}", error))
|
||||
}
|
||||
}
|
||||
|
||||
impl MaybeConnectionError for Error {
|
||||
fn is_connection_error(&self) -> bool {
|
||||
matches!(
|
||||
@@ -68,6 +84,7 @@ impl MaybeConnectionError for Error {
|
||||
impl ToString for Error {
|
||||
fn to_string(&self) -> String {
|
||||
match self {
|
||||
Self::Io(e) => e.to_string(),
|
||||
Self::RpcError(e) => e.to_string(),
|
||||
Self::ResponseParseFailed(e) => e.to_string(),
|
||||
Self::IncompleteHeader => {
|
||||
@@ -80,6 +97,7 @@ impl ToString for Error {
|
||||
Self::ClientNotSynced(missing_headers) => {
|
||||
format!("Ethereum client is not synced: syncing {} headers", missing_headers)
|
||||
}
|
||||
Self::Custom(ref e) => e.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user