mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-05-30 23:21:02 +00:00
in auto-relays keep trying to connect to nodes until connection is established (#971)
This commit is contained in:
committed by
Bastian Köcher
parent
ea82ad67cf
commit
f8f8f42a89
@@ -60,8 +60,8 @@ pub async fn run(params: EthereumDeployContractParams) {
|
|||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let result = async move {
|
let result = async move {
|
||||||
let eth_client = EthereumClient::new(eth_params).await.map_err(RpcError::Ethereum)?;
|
let eth_client = EthereumClient::try_connect(eth_params).await.map_err(RpcError::Ethereum)?;
|
||||||
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await.map_err(RpcError::Substrate)?;
|
let sub_client = SubstrateClient::<Rialto>::try_connect(sub_params).await.map_err(RpcError::Substrate)?;
|
||||||
|
|
||||||
let (initial_header_id, initial_header) = prepare_initial_header(&sub_client, sub_initial_header).await?;
|
let (initial_header_id, initial_header) = prepare_initial_header(&sub_client, sub_initial_header).await?;
|
||||||
let initial_set_id = sub_initial_authorities_set_id.unwrap_or(0);
|
let initial_set_id = sub_initial_authorities_set_id.unwrap_or(0);
|
||||||
|
|||||||
@@ -335,8 +335,10 @@ async fn run_single_transaction_relay(params: EthereumExchangeParams, eth_tx_has
|
|||||||
..
|
..
|
||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let eth_client = EthereumClient::new(eth_params).await.map_err(RpcError::Ethereum)?;
|
let eth_client = EthereumClient::try_connect(eth_params)
|
||||||
let sub_client = SubstrateClient::<Rialto>::new(sub_params)
|
.await
|
||||||
|
.map_err(RpcError::Ethereum)?;
|
||||||
|
let sub_client = SubstrateClient::<Rialto>::try_connect(sub_params)
|
||||||
.await
|
.await
|
||||||
.map_err(RpcError::Substrate)?;
|
.map_err(RpcError::Substrate)?;
|
||||||
|
|
||||||
@@ -363,12 +365,8 @@ async fn run_auto_transactions_relay_loop(
|
|||||||
..
|
..
|
||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let eth_client = EthereumClient::new(eth_params)
|
let eth_client = EthereumClient::new(eth_params).await;
|
||||||
.await
|
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await;
|
||||||
.map_err(|err| format!("Error starting Ethereum client: {:?}", err))?;
|
|
||||||
let sub_client = SubstrateClient::<Rialto>::new(sub_params)
|
|
||||||
.await
|
|
||||||
.map_err(|err| format!("Error starting Substrate client: {:?}", err))?;
|
|
||||||
|
|
||||||
let eth_start_with_block_number = match eth_start_with_block_number {
|
let eth_start_with_block_number = match eth_start_with_block_number {
|
||||||
Some(eth_start_with_block_number) => eth_start_with_block_number,
|
Some(eth_start_with_block_number) => eth_start_with_block_number,
|
||||||
|
|||||||
@@ -52,7 +52,7 @@ pub async fn run(params: EthereumExchangeSubmitParams) {
|
|||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let result: Result<_, String> = async move {
|
let result: Result<_, String> = async move {
|
||||||
let eth_client = EthereumClient::new(eth_params)
|
let eth_client = EthereumClient::try_connect(eth_params)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| format!("error connecting to Ethereum node: {:?}", err))?;
|
.map_err(|err| format!("error connecting to Ethereum node: {:?}", err))?;
|
||||||
|
|
||||||
|
|||||||
@@ -270,8 +270,8 @@ pub async fn run(params: EthereumSyncParams) -> Result<(), RpcError> {
|
|||||||
instance,
|
instance,
|
||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let eth_client = EthereumClient::new(eth_params).await?;
|
let eth_client = EthereumClient::new(eth_params).await;
|
||||||
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await?;
|
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await;
|
||||||
|
|
||||||
let sign_sub_transactions = match sync_params.target_tx_mode {
|
let sign_sub_transactions = match sync_params.target_tx_mode {
|
||||||
TargetTransactionMode::Signed | TargetTransactionMode::Backup => true,
|
TargetTransactionMode::Signed | TargetTransactionMode::Backup => true,
|
||||||
|
|||||||
@@ -177,8 +177,8 @@ pub async fn run(params: SubstrateSyncParams) -> Result<(), RpcError> {
|
|||||||
metrics_params,
|
metrics_params,
|
||||||
} = params;
|
} = params;
|
||||||
|
|
||||||
let eth_client = EthereumClient::new(eth_params).await?;
|
let eth_client = EthereumClient::new(eth_params).await;
|
||||||
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await?;
|
let sub_client = SubstrateClient::<Rialto>::new(sub_params).await;
|
||||||
|
|
||||||
let target = EthereumHeadersTarget::new(eth_client, eth_contract_address, eth_sign);
|
let target = EthereumHeadersTarget::new(eth_client, eth_contract_address, eth_sign);
|
||||||
let source = SubstrateHeadersSource::new(sub_client);
|
let source = SubstrateHeadersSource::new(sub_client);
|
||||||
|
|||||||
@@ -406,7 +406,7 @@ macro_rules! declare_chain_options {
|
|||||||
port: self.[<$chain_prefix _port>],
|
port: self.[<$chain_prefix _port>],
|
||||||
secure: self.[<$chain_prefix _secure>],
|
secure: self.[<$chain_prefix _secure>],
|
||||||
})
|
})
|
||||||
.await?
|
.await
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ edition = "2018"
|
|||||||
license = "GPL-3.0-or-later WITH Classpath-exception-2.0"
|
license = "GPL-3.0-or-later WITH Classpath-exception-2.0"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
async-std = "1.6.5"
|
||||||
bp-eth-poa = { path = "../../primitives/ethereum-poa" }
|
bp-eth-poa = { path = "../../primitives/ethereum-poa" }
|
||||||
codec = { package = "parity-scale-codec", version = "2.0.0" }
|
codec = { package = "parity-scale-codec", version = "2.0.0" }
|
||||||
headers-relay = { path = "../headers" }
|
headers-relay = { path = "../headers" }
|
||||||
|
|||||||
@@ -22,6 +22,7 @@ use crate::types::{
|
|||||||
use crate::{ConnectionParams, Error, Result};
|
use crate::{ConnectionParams, Error, Result};
|
||||||
|
|
||||||
use jsonrpsee_ws_client::{WsClient as RpcClient, WsClientBuilder as RpcClientBuilder};
|
use jsonrpsee_ws_client::{WsClient as RpcClient, WsClientBuilder as RpcClientBuilder};
|
||||||
|
use relay_utils::relay_loop::RECONNECT_DELAY;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
/// Number of headers missing from the Ethereum node for us to consider node not synced.
|
/// Number of headers missing from the Ethereum node for us to consider node not synced.
|
||||||
@@ -36,7 +37,28 @@ pub struct Client {
|
|||||||
|
|
||||||
impl Client {
|
impl Client {
|
||||||
/// Create a new Ethereum RPC Client.
|
/// Create a new Ethereum RPC Client.
|
||||||
pub async fn new(params: ConnectionParams) -> Result<Self> {
|
///
|
||||||
|
/// This function will keep connecting to given Ethereum node until connection is established
|
||||||
|
/// and is functional. If attempt fail, it will wait for `RECONNECT_DELAY` and retry again.
|
||||||
|
pub async fn new(params: ConnectionParams) -> Self {
|
||||||
|
loop {
|
||||||
|
match Self::try_connect(params.clone()).await {
|
||||||
|
Ok(client) => return client,
|
||||||
|
Err(error) => log::error!(
|
||||||
|
target: "bridge",
|
||||||
|
"Failed to connect to Ethereum node: {:?}. Going to retry in {}s",
|
||||||
|
error,
|
||||||
|
RECONNECT_DELAY.as_secs(),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
|
||||||
|
async_std::task::sleep(RECONNECT_DELAY).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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 {
|
Ok(Self {
|
||||||
client: Self::build_client(¶ms).await?,
|
client: Self::build_client(¶ms).await?,
|
||||||
params,
|
params,
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ use jsonrpsee_ws_client::{traits::SubscriptionClient, v2::params::JsonRpcParams,
|
|||||||
use jsonrpsee_ws_client::{Subscription, WsClient as RpcClient, WsClientBuilder as RpcClientBuilder};
|
use jsonrpsee_ws_client::{Subscription, WsClient as RpcClient, WsClientBuilder as RpcClientBuilder};
|
||||||
use num_traits::Zero;
|
use num_traits::Zero;
|
||||||
use pallet_balances::AccountData;
|
use pallet_balances::AccountData;
|
||||||
|
use relay_utils::relay_loop::RECONNECT_DELAY;
|
||||||
use sp_core::{storage::StorageKey, Bytes};
|
use sp_core::{storage::StorageKey, Bytes};
|
||||||
use sp_trie::StorageProof;
|
use sp_trie::StorageProof;
|
||||||
use sp_version::RuntimeVersion;
|
use sp_version::RuntimeVersion;
|
||||||
@@ -77,7 +78,29 @@ impl<C: Chain> std::fmt::Debug for Client<C> {
|
|||||||
|
|
||||||
impl<C: Chain> Client<C> {
|
impl<C: Chain> Client<C> {
|
||||||
/// Returns client that is able to call RPCs on Substrate node over websocket connection.
|
/// Returns client that is able to call RPCs on Substrate node over websocket connection.
|
||||||
pub async fn new(params: ConnectionParams) -> Result<Self> {
|
///
|
||||||
|
/// This function will keep connecting to given Sustrate node until connection is established
|
||||||
|
/// and is functional. If attempt fail, it will wait for `RECONNECT_DELAY` and retry again.
|
||||||
|
pub async fn new(params: ConnectionParams) -> Self {
|
||||||
|
loop {
|
||||||
|
match Self::try_connect(params.clone()).await {
|
||||||
|
Ok(client) => return client,
|
||||||
|
Err(error) => log::error!(
|
||||||
|
target: "bridge",
|
||||||
|
"Failed to connect to {} node: {:?}. Going to retry in {}s",
|
||||||
|
C::NAME,
|
||||||
|
error,
|
||||||
|
RECONNECT_DELAY.as_secs(),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
|
||||||
|
async_std::task::sleep(RECONNECT_DELAY).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Try to connect to Substrate node over websocket. Returns Substrate RPC client if connection
|
||||||
|
/// has been established or error otherwise.
|
||||||
|
pub async fn try_connect(params: ConnectionParams) -> Result<Self> {
|
||||||
let client = Self::build_client(params.clone()).await?;
|
let client = Self::build_client(params.clone()).await?;
|
||||||
|
|
||||||
let number: C::BlockNumber = Zero::zero();
|
let number: C::BlockNumber = Zero::zero();
|
||||||
|
|||||||
Reference in New Issue
Block a user