mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-07-22 06:45:41 +00:00
Shared reference to conversion rate metric value (#1034)
* shared conversion rate metric value * clippy
This commit is contained in:
committed by
Bastian Köcher
parent
db0216dabb
commit
ecd20d9d24
@@ -7,6 +7,7 @@ license = "GPL-3.0-or-later WITH Classpath-exception-2.0"
|
|||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
ansi_term = "0.12"
|
ansi_term = "0.12"
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = "1.9.0"
|
async-std = "1.9.0"
|
||||||
async-trait = "0.1.42"
|
async-trait = "0.1.42"
|
||||||
clap = { version = "2.33.3", features = ["yaml"] }
|
clap = { version = "2.33.3", features = ["yaml"] }
|
||||||
|
|||||||
@@ -355,7 +355,7 @@ async fn run_single_transaction_relay(params: EthereumExchangeParams, eth_tx_has
|
|||||||
async fn run_auto_transactions_relay_loop(
|
async fn run_auto_transactions_relay_loop(
|
||||||
params: EthereumExchangeParams,
|
params: EthereumExchangeParams,
|
||||||
eth_start_with_block_number: Option<u64>,
|
eth_start_with_block_number: Option<u64>,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let EthereumExchangeParams {
|
let EthereumExchangeParams {
|
||||||
eth_params,
|
eth_params,
|
||||||
sub_params,
|
sub_params,
|
||||||
@@ -375,7 +375,7 @@ async fn run_auto_transactions_relay_loop(
|
|||||||
.best_ethereum_finalized_block()
|
.best_ethereum_finalized_block()
|
||||||
.await
|
.await
|
||||||
.map_err(|err| {
|
.map_err(|err| {
|
||||||
format!(
|
anyhow::format_err!(
|
||||||
"Error retrieving best finalized Ethereum block from Substrate node: {:?}",
|
"Error retrieving best finalized Ethereum block from Substrate node: {:?}",
|
||||||
err
|
err
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -292,7 +292,7 @@ pub async fn run(params: EthereumSyncParams) -> Result<(), RpcError> {
|
|||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(RpcError::SyncLoop)?;
|
.map_err(|e| RpcError::SyncLoop(e.to_string()))?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -194,7 +194,7 @@ pub async fn run(params: SubstrateSyncParams) -> Result<(), RpcError> {
|
|||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(RpcError::SyncLoop)?;
|
.map_err(|e| RpcError::SyncLoop(e.to_string()))?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -17,7 +17,8 @@
|
|||||||
//! Millau-to-Rialto messages sync entrypoint.
|
//! Millau-to-Rialto messages sync entrypoint.
|
||||||
|
|
||||||
use crate::messages_lane::{
|
use crate::messages_lane::{
|
||||||
select_delivery_transaction_limits, MessagesRelayParams, SubstrateMessageLane, SubstrateMessageLaneToSubstrate,
|
select_delivery_transaction_limits, MessagesRelayParams, StandaloneMessagesMetrics, SubstrateMessageLane,
|
||||||
|
SubstrateMessageLaneToSubstrate,
|
||||||
};
|
};
|
||||||
use crate::messages_source::SubstrateMessagesSource;
|
use crate::messages_source::SubstrateMessagesSource;
|
||||||
use crate::messages_target::SubstrateMessagesTarget;
|
use crate::messages_target::SubstrateMessagesTarget;
|
||||||
@@ -30,10 +31,8 @@ use frame_support::dispatch::GetDispatchInfo;
|
|||||||
use messages_relay::message_lane::MessageLane;
|
use messages_relay::message_lane::MessageLane;
|
||||||
use relay_millau_client::{HeaderId as MillauHeaderId, Millau, SigningParams as MillauSigningParams};
|
use relay_millau_client::{HeaderId as MillauHeaderId, Millau, SigningParams as MillauSigningParams};
|
||||||
use relay_rialto_client::{HeaderId as RialtoHeaderId, Rialto, SigningParams as RialtoSigningParams};
|
use relay_rialto_client::{HeaderId as RialtoHeaderId, Rialto, SigningParams as RialtoSigningParams};
|
||||||
use relay_substrate_client::{
|
use relay_substrate_client::{Chain, Client, TransactionSignScheme};
|
||||||
metrics::{FloatStorageValueMetric, StorageProofOverheadMetric},
|
use relay_utils::metrics::MetricsParams;
|
||||||
Chain, TransactionSignScheme,
|
|
||||||
};
|
|
||||||
use sp_core::{Bytes, Pair};
|
use sp_core::{Bytes, Pair};
|
||||||
use std::{ops::RangeInclusive, time::Duration};
|
use std::{ops::RangeInclusive, time::Duration};
|
||||||
|
|
||||||
@@ -136,7 +135,7 @@ type RialtoTargetClient =
|
|||||||
/// Run Millau-to-Rialto messages sync.
|
/// Run Millau-to-Rialto messages sync.
|
||||||
pub async fn run(
|
pub async fn run(
|
||||||
params: MessagesRelayParams<Millau, MillauSigningParams, Rialto, RialtoSigningParams>,
|
params: MessagesRelayParams<Millau, MillauSigningParams, Rialto, RialtoSigningParams>,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let stall_timeout = Duration::from_secs(5 * 60);
|
let stall_timeout = Duration::from_secs(5 * 60);
|
||||||
let relayer_id_at_millau = (*params.source_sign.public().as_array_ref()).into();
|
let relayer_id_at_millau = (*params.source_sign.public().as_array_ref()).into();
|
||||||
|
|
||||||
@@ -172,6 +171,7 @@ pub async fn run(
|
|||||||
max_messages_weight_in_single_batch,
|
max_messages_weight_in_single_batch,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
let (metrics_params, _) = add_standalone_metrics(params.metrics_params, source_client.clone())?;
|
||||||
messages_relay::message_lane_loop::run(
|
messages_relay::message_lane_loop::run(
|
||||||
messages_relay::message_lane_loop::Params {
|
messages_relay::message_lane_loop::Params {
|
||||||
lane: lane_id,
|
lane: lane_id,
|
||||||
@@ -202,36 +202,25 @@ pub async fn run(
|
|||||||
MILLAU_CHAIN_ID,
|
MILLAU_CHAIN_ID,
|
||||||
params.source_to_target_headers_relay,
|
params.source_to_target_headers_relay,
|
||||||
),
|
),
|
||||||
relay_utils::relay_metrics(
|
metrics_params,
|
||||||
Some(messages_relay::message_lane_loop::metrics_prefix::<
|
|
||||||
MillauMessagesToRialto,
|
|
||||||
>(&lane_id)),
|
|
||||||
params.metrics_params,
|
|
||||||
)
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
StorageProofOverheadMetric::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client.clone(),
|
|
||||||
"millau_storage_proof_overhead".into(),
|
|
||||||
"Millau storage proof overhead".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
FloatStorageValueMetric::<_, sp_runtime::FixedU128>::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client,
|
|
||||||
sp_core::storage::StorageKey(
|
|
||||||
millau_runtime::rialto_messages::RialtoToMillauConversionRate::key().to_vec(),
|
|
||||||
),
|
|
||||||
Some(millau_runtime::rialto_messages::INITIAL_RIALTO_TO_MILLAU_CONVERSION_RATE),
|
|
||||||
"millau_rialto_to_millau_conversion_rate".into(),
|
|
||||||
"Rialto to Millau tokens conversion rate (used by Rialto)".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.into_params(),
|
|
||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Add standalone metrics for the Millau -> Rialto messages loop.
|
||||||
|
pub(crate) fn add_standalone_metrics(
|
||||||
|
metrics_params: MetricsParams,
|
||||||
|
source_client: Client<Millau>,
|
||||||
|
) -> anyhow::Result<(MetricsParams, StandaloneMessagesMetrics)> {
|
||||||
|
crate::messages_lane::add_standalone_metrics::<MillauMessagesToRialto>(
|
||||||
|
metrics_params,
|
||||||
|
source_client,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
Some((
|
||||||
|
sp_core::storage::StorageKey(millau_runtime::rialto_messages::RialtoToMillauConversionRate::key().to_vec()),
|
||||||
|
millau_runtime::rialto_messages::INITIAL_RIALTO_TO_MILLAU_CONVERSION_RATE,
|
||||||
|
)),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -32,38 +32,38 @@ mod rococo;
|
|||||||
mod westend;
|
mod westend;
|
||||||
mod wococo;
|
mod wococo;
|
||||||
|
|
||||||
use relay_utils::metrics::{FloatJsonValueMetric, MetricsParams};
|
use relay_utils::metrics::{FloatJsonValueMetric, MetricsParams, PrometheusError, Registry};
|
||||||
|
|
||||||
pub(crate) fn add_polkadot_kusama_price_metrics<T: finality_relay::FinalitySyncPipeline>(
|
pub(crate) fn add_polkadot_kusama_price_metrics<T: finality_relay::FinalitySyncPipeline>(
|
||||||
params: MetricsParams,
|
params: MetricsParams,
|
||||||
) -> anyhow::Result<MetricsParams> {
|
) -> anyhow::Result<MetricsParams> {
|
||||||
Ok(
|
// Polkadot/Kusama prices are added as metrics here, because atm we don't have Polkadot <-> Kusama
|
||||||
relay_utils::relay_metrics(Some(finality_relay::metrics_prefix::<T>()), params)
|
// relays, but we want to test metrics/dashboards in advance
|
||||||
// Polkadot/Kusama prices are added as metrics here, because atm we don't have Polkadot <-> Kusama
|
Ok(relay_utils::relay_metrics(None, params)
|
||||||
// relays, but we want to test metrics/dashboards in advance
|
.standalone_metric(|registry, prefix| token_price_metric(registry, prefix, "polkadot"))?
|
||||||
.standalone_metric(|registry, prefix| {
|
.standalone_metric(|registry, prefix| token_price_metric(registry, prefix, "kusama"))?
|
||||||
FloatJsonValueMetric::new(
|
.into_params())
|
||||||
registry,
|
}
|
||||||
prefix,
|
|
||||||
"https://api.coingecko.com/api/v3/simple/price?ids=Polkadot&vs_currencies=btc".into(),
|
/// Creates standalone token price metric.
|
||||||
"$.polkadot.btc".into(),
|
pub(crate) fn token_price_metric(
|
||||||
"polkadot_to_base_conversion_rate".into(),
|
registry: &Registry,
|
||||||
"Rate used to convert from DOT to some BASE tokens".into(),
|
prefix: Option<&str>,
|
||||||
)
|
token_id: &str,
|
||||||
})
|
) -> Result<FloatJsonValueMetric, PrometheusError> {
|
||||||
.map_err(|e| anyhow::format_err!("{}", e))?
|
FloatJsonValueMetric::new(
|
||||||
.standalone_metric(|registry, prefix| {
|
registry,
|
||||||
FloatJsonValueMetric::new(
|
prefix,
|
||||||
registry,
|
format!(
|
||||||
prefix,
|
"https://api.coingecko.com/api/v3/simple/price?ids={}&vs_currencies=btc",
|
||||||
"https://api.coingecko.com/api/v3/simple/price?ids=Kusama&vs_currencies=btc".into(),
|
token_id
|
||||||
"$.kusama.btc".into(),
|
),
|
||||||
"kusama_to_base_conversion_rate".into(),
|
format!("$.{}.btc", token_id),
|
||||||
"Rate used to convert from KSM to some BASE tokens".into(),
|
format!("{}_to_base_conversion_rate", token_id.replace("-", "_")),
|
||||||
)
|
format!(
|
||||||
})
|
"Rate used to convert from {} to some BASE tokens",
|
||||||
.map_err(|e| anyhow::format_err!("{}", e))?
|
token_id.to_uppercase()
|
||||||
.into_params(),
|
),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -17,7 +17,8 @@
|
|||||||
//! Rialto-to-Millau messages sync entrypoint.
|
//! Rialto-to-Millau messages sync entrypoint.
|
||||||
|
|
||||||
use crate::messages_lane::{
|
use crate::messages_lane::{
|
||||||
select_delivery_transaction_limits, MessagesRelayParams, SubstrateMessageLane, SubstrateMessageLaneToSubstrate,
|
select_delivery_transaction_limits, MessagesRelayParams, StandaloneMessagesMetrics, SubstrateMessageLane,
|
||||||
|
SubstrateMessageLaneToSubstrate,
|
||||||
};
|
};
|
||||||
use crate::messages_source::SubstrateMessagesSource;
|
use crate::messages_source::SubstrateMessagesSource;
|
||||||
use crate::messages_target::SubstrateMessagesTarget;
|
use crate::messages_target::SubstrateMessagesTarget;
|
||||||
@@ -30,10 +31,8 @@ use frame_support::dispatch::GetDispatchInfo;
|
|||||||
use messages_relay::message_lane::MessageLane;
|
use messages_relay::message_lane::MessageLane;
|
||||||
use relay_millau_client::{HeaderId as MillauHeaderId, Millau, SigningParams as MillauSigningParams};
|
use relay_millau_client::{HeaderId as MillauHeaderId, Millau, SigningParams as MillauSigningParams};
|
||||||
use relay_rialto_client::{HeaderId as RialtoHeaderId, Rialto, SigningParams as RialtoSigningParams};
|
use relay_rialto_client::{HeaderId as RialtoHeaderId, Rialto, SigningParams as RialtoSigningParams};
|
||||||
use relay_substrate_client::{
|
use relay_substrate_client::{Chain, Client, TransactionSignScheme};
|
||||||
metrics::{FloatStorageValueMetric, StorageProofOverheadMetric},
|
use relay_utils::metrics::MetricsParams;
|
||||||
Chain, TransactionSignScheme,
|
|
||||||
};
|
|
||||||
use sp_core::{Bytes, Pair};
|
use sp_core::{Bytes, Pair};
|
||||||
use std::{ops::RangeInclusive, time::Duration};
|
use std::{ops::RangeInclusive, time::Duration};
|
||||||
|
|
||||||
@@ -136,7 +135,7 @@ type MillauTargetClient =
|
|||||||
/// Run Rialto-to-Millau messages sync.
|
/// Run Rialto-to-Millau messages sync.
|
||||||
pub async fn run(
|
pub async fn run(
|
||||||
params: MessagesRelayParams<Rialto, RialtoSigningParams, Millau, MillauSigningParams>,
|
params: MessagesRelayParams<Rialto, RialtoSigningParams, Millau, MillauSigningParams>,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let stall_timeout = Duration::from_secs(5 * 60);
|
let stall_timeout = Duration::from_secs(5 * 60);
|
||||||
let relayer_id_at_rialto = (*params.source_sign.public().as_array_ref()).into();
|
let relayer_id_at_rialto = (*params.source_sign.public().as_array_ref()).into();
|
||||||
|
|
||||||
@@ -171,6 +170,7 @@ pub async fn run(
|
|||||||
max_messages_weight_in_single_batch,
|
max_messages_weight_in_single_batch,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
let (metrics_params, _) = add_standalone_metrics(params.metrics_params, source_client.clone())?;
|
||||||
messages_relay::message_lane_loop::run(
|
messages_relay::message_lane_loop::run(
|
||||||
messages_relay::message_lane_loop::Params {
|
messages_relay::message_lane_loop::Params {
|
||||||
lane: lane_id,
|
lane: lane_id,
|
||||||
@@ -201,36 +201,25 @@ pub async fn run(
|
|||||||
RIALTO_CHAIN_ID,
|
RIALTO_CHAIN_ID,
|
||||||
params.source_to_target_headers_relay,
|
params.source_to_target_headers_relay,
|
||||||
),
|
),
|
||||||
relay_utils::relay_metrics(
|
metrics_params,
|
||||||
Some(messages_relay::message_lane_loop::metrics_prefix::<
|
|
||||||
RialtoMessagesToMillau,
|
|
||||||
>(&lane_id)),
|
|
||||||
params.metrics_params,
|
|
||||||
)
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
StorageProofOverheadMetric::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client.clone(),
|
|
||||||
"rialto_storage_proof_overhead".into(),
|
|
||||||
"Rialto storage proof overhead".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
FloatStorageValueMetric::<_, sp_runtime::FixedU128>::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client,
|
|
||||||
sp_core::storage::StorageKey(
|
|
||||||
rialto_runtime::millau_messages::MillauToRialtoConversionRate::key().to_vec(),
|
|
||||||
),
|
|
||||||
Some(rialto_runtime::millau_messages::INITIAL_MILLAU_TO_RIALTO_CONVERSION_RATE),
|
|
||||||
"rialto_millau_to_rialto_conversion_rate".into(),
|
|
||||||
"Millau to Rialto tokens conversion rate (used by Millau)".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.into_params(),
|
|
||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Add standalone metrics for the Rialto -> Millau messages loop.
|
||||||
|
pub(crate) fn add_standalone_metrics(
|
||||||
|
metrics_params: MetricsParams,
|
||||||
|
source_client: Client<Rialto>,
|
||||||
|
) -> anyhow::Result<(MetricsParams, StandaloneMessagesMetrics)> {
|
||||||
|
crate::messages_lane::add_standalone_metrics::<RialtoMessagesToMillau>(
|
||||||
|
metrics_params,
|
||||||
|
source_client,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
Some((
|
||||||
|
sp_core::storage::StorageKey(rialto_runtime::millau_messages::MillauToRialtoConversionRate::key().to_vec()),
|
||||||
|
rialto_runtime::millau_messages::INITIAL_MILLAU_TO_RIALTO_CONVERSION_RATE,
|
||||||
|
)),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -17,7 +17,8 @@
|
|||||||
//! Rococo-to-Wococo messages sync entrypoint.
|
//! Rococo-to-Wococo messages sync entrypoint.
|
||||||
|
|
||||||
use crate::messages_lane::{
|
use crate::messages_lane::{
|
||||||
select_delivery_transaction_limits, MessagesRelayParams, SubstrateMessageLane, SubstrateMessageLaneToSubstrate,
|
select_delivery_transaction_limits, MessagesRelayParams, StandaloneMessagesMetrics, SubstrateMessageLane,
|
||||||
|
SubstrateMessageLaneToSubstrate,
|
||||||
};
|
};
|
||||||
use crate::messages_source::SubstrateMessagesSource;
|
use crate::messages_source::SubstrateMessagesSource;
|
||||||
use crate::messages_target::SubstrateMessagesTarget;
|
use crate::messages_target::SubstrateMessagesTarget;
|
||||||
@@ -28,7 +29,8 @@ use bridge_runtime_common::messages::target::FromBridgedChainMessagesProof;
|
|||||||
use codec::Encode;
|
use codec::Encode;
|
||||||
use messages_relay::message_lane::MessageLane;
|
use messages_relay::message_lane::MessageLane;
|
||||||
use relay_rococo_client::{HeaderId as RococoHeaderId, Rococo, SigningParams as RococoSigningParams};
|
use relay_rococo_client::{HeaderId as RococoHeaderId, Rococo, SigningParams as RococoSigningParams};
|
||||||
use relay_substrate_client::{metrics::StorageProofOverheadMetric, Chain, TransactionSignScheme};
|
use relay_substrate_client::{Chain, Client, TransactionSignScheme};
|
||||||
|
use relay_utils::metrics::MetricsParams;
|
||||||
use relay_wococo_client::{HeaderId as WococoHeaderId, SigningParams as WococoSigningParams, Wococo};
|
use relay_wococo_client::{HeaderId as WococoHeaderId, SigningParams as WococoSigningParams, Wococo};
|
||||||
use sp_core::{Bytes, Pair};
|
use sp_core::{Bytes, Pair};
|
||||||
use std::{ops::RangeInclusive, time::Duration};
|
use std::{ops::RangeInclusive, time::Duration};
|
||||||
@@ -142,7 +144,7 @@ type WococoTargetClient = SubstrateMessagesTarget<
|
|||||||
/// Run Rococo-to-Wococo messages sync.
|
/// Run Rococo-to-Wococo messages sync.
|
||||||
pub async fn run(
|
pub async fn run(
|
||||||
params: MessagesRelayParams<Rococo, RococoSigningParams, Wococo, WococoSigningParams>,
|
params: MessagesRelayParams<Rococo, RococoSigningParams, Wococo, WococoSigningParams>,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let stall_timeout = Duration::from_secs(5 * 60);
|
let stall_timeout = Duration::from_secs(5 * 60);
|
||||||
let relayer_id_at_rococo = (*params.source_sign.public().as_array_ref()).into();
|
let relayer_id_at_rococo = (*params.source_sign.public().as_array_ref()).into();
|
||||||
|
|
||||||
@@ -183,6 +185,7 @@ pub async fn run(
|
|||||||
max_messages_weight_in_single_batch,
|
max_messages_weight_in_single_batch,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
let (metrics_params, _) = add_standalone_metrics(params.metrics_params, source_client.clone())?;
|
||||||
messages_relay::message_lane_loop::run(
|
messages_relay::message_lane_loop::run(
|
||||||
messages_relay::message_lane_loop::Params {
|
messages_relay::message_lane_loop::Params {
|
||||||
lane: lane_id,
|
lane: lane_id,
|
||||||
@@ -213,23 +216,22 @@ pub async fn run(
|
|||||||
ROCOCO_CHAIN_ID,
|
ROCOCO_CHAIN_ID,
|
||||||
params.source_to_target_headers_relay,
|
params.source_to_target_headers_relay,
|
||||||
),
|
),
|
||||||
relay_utils::relay_metrics(
|
metrics_params,
|
||||||
Some(messages_relay::message_lane_loop::metrics_prefix::<
|
|
||||||
RococoMessagesToWococo,
|
|
||||||
>(&lane_id)),
|
|
||||||
params.metrics_params,
|
|
||||||
)
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
StorageProofOverheadMetric::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client.clone(),
|
|
||||||
"rococo_storage_proof_overhead".into(),
|
|
||||||
"Rococo storage proof overhead".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.into_params(),
|
|
||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Add standalone metrics for the Rococo -> Wococo messages loop.
|
||||||
|
pub(crate) fn add_standalone_metrics(
|
||||||
|
metrics_params: MetricsParams,
|
||||||
|
source_client: Client<Rococo>,
|
||||||
|
) -> anyhow::Result<(MetricsParams, StandaloneMessagesMetrics)> {
|
||||||
|
crate::messages_lane::add_standalone_metrics::<RococoMessagesToWococo>(
|
||||||
|
metrics_params,
|
||||||
|
source_client,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -17,7 +17,8 @@
|
|||||||
//! Wococo-to-Rococo messages sync entrypoint.
|
//! Wococo-to-Rococo messages sync entrypoint.
|
||||||
|
|
||||||
use crate::messages_lane::{
|
use crate::messages_lane::{
|
||||||
select_delivery_transaction_limits, MessagesRelayParams, SubstrateMessageLane, SubstrateMessageLaneToSubstrate,
|
select_delivery_transaction_limits, MessagesRelayParams, StandaloneMessagesMetrics, SubstrateMessageLane,
|
||||||
|
SubstrateMessageLaneToSubstrate,
|
||||||
};
|
};
|
||||||
use crate::messages_source::SubstrateMessagesSource;
|
use crate::messages_source::SubstrateMessagesSource;
|
||||||
use crate::messages_target::SubstrateMessagesTarget;
|
use crate::messages_target::SubstrateMessagesTarget;
|
||||||
@@ -28,7 +29,8 @@ use bridge_runtime_common::messages::target::FromBridgedChainMessagesProof;
|
|||||||
use codec::Encode;
|
use codec::Encode;
|
||||||
use messages_relay::message_lane::MessageLane;
|
use messages_relay::message_lane::MessageLane;
|
||||||
use relay_rococo_client::{HeaderId as RococoHeaderId, Rococo, SigningParams as RococoSigningParams};
|
use relay_rococo_client::{HeaderId as RococoHeaderId, Rococo, SigningParams as RococoSigningParams};
|
||||||
use relay_substrate_client::{metrics::StorageProofOverheadMetric, Chain, TransactionSignScheme};
|
use relay_substrate_client::{Chain, Client, TransactionSignScheme};
|
||||||
|
use relay_utils::metrics::MetricsParams;
|
||||||
use relay_wococo_client::{HeaderId as WococoHeaderId, SigningParams as WococoSigningParams, Wococo};
|
use relay_wococo_client::{HeaderId as WococoHeaderId, SigningParams as WococoSigningParams, Wococo};
|
||||||
use sp_core::{Bytes, Pair};
|
use sp_core::{Bytes, Pair};
|
||||||
use std::{ops::RangeInclusive, time::Duration};
|
use std::{ops::RangeInclusive, time::Duration};
|
||||||
@@ -142,7 +144,7 @@ type RococoTargetClient = SubstrateMessagesTarget<
|
|||||||
/// Run Wococo-to-Rococo messages sync.
|
/// Run Wococo-to-Rococo messages sync.
|
||||||
pub async fn run(
|
pub async fn run(
|
||||||
params: MessagesRelayParams<Wococo, WococoSigningParams, Rococo, RococoSigningParams>,
|
params: MessagesRelayParams<Wococo, WococoSigningParams, Rococo, RococoSigningParams>,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let stall_timeout = Duration::from_secs(5 * 60);
|
let stall_timeout = Duration::from_secs(5 * 60);
|
||||||
let relayer_id_at_wococo = (*params.source_sign.public().as_array_ref()).into();
|
let relayer_id_at_wococo = (*params.source_sign.public().as_array_ref()).into();
|
||||||
|
|
||||||
@@ -183,6 +185,7 @@ pub async fn run(
|
|||||||
max_messages_weight_in_single_batch,
|
max_messages_weight_in_single_batch,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
let (metrics_params, _) = add_standalone_metrics(params.metrics_params, source_client.clone())?;
|
||||||
messages_relay::message_lane_loop::run(
|
messages_relay::message_lane_loop::run(
|
||||||
messages_relay::message_lane_loop::Params {
|
messages_relay::message_lane_loop::Params {
|
||||||
lane: lane_id,
|
lane: lane_id,
|
||||||
@@ -213,23 +216,22 @@ pub async fn run(
|
|||||||
WOCOCO_CHAIN_ID,
|
WOCOCO_CHAIN_ID,
|
||||||
params.source_to_target_headers_relay,
|
params.source_to_target_headers_relay,
|
||||||
),
|
),
|
||||||
relay_utils::relay_metrics(
|
metrics_params,
|
||||||
Some(messages_relay::message_lane_loop::metrics_prefix::<
|
|
||||||
WococoMessagesToRococo,
|
|
||||||
>(&lane_id)),
|
|
||||||
params.metrics_params,
|
|
||||||
)
|
|
||||||
.standalone_metric(|registry, prefix| {
|
|
||||||
StorageProofOverheadMetric::new(
|
|
||||||
registry,
|
|
||||||
prefix,
|
|
||||||
source_client.clone(),
|
|
||||||
"wococo_storage_proof_overhead".into(),
|
|
||||||
"Wococo storage proof overhead".into(),
|
|
||||||
)
|
|
||||||
})?
|
|
||||||
.into_params(),
|
|
||||||
futures::future::pending(),
|
futures::future::pending(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Add standalone metrics for the Wococo -> Rococo messages loop.
|
||||||
|
pub(crate) fn add_standalone_metrics(
|
||||||
|
metrics_params: MetricsParams,
|
||||||
|
source_client: Client<Wococo>,
|
||||||
|
) -> anyhow::Result<(MetricsParams, StandaloneMessagesMetrics)> {
|
||||||
|
crate::messages_lane::add_standalone_metrics::<WococoMessagesToRococo>(
|
||||||
|
metrics_params,
|
||||||
|
source_client,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|||||||
@@ -100,8 +100,12 @@ macro_rules! select_bridge {
|
|||||||
const MAX_MISSING_LEFT_HEADERS_AT_RIGHT: bp_millau::BlockNumber = bp_millau::SESSION_LENGTH;
|
const MAX_MISSING_LEFT_HEADERS_AT_RIGHT: bp_millau::BlockNumber = bp_millau::SESSION_LENGTH;
|
||||||
const MAX_MISSING_RIGHT_HEADERS_AT_LEFT: bp_rialto::BlockNumber = bp_rialto::SESSION_LENGTH;
|
const MAX_MISSING_RIGHT_HEADERS_AT_LEFT: bp_rialto::BlockNumber = bp_rialto::SESSION_LENGTH;
|
||||||
|
|
||||||
use crate::chains::millau_messages_to_rialto::run as left_to_right_messages;
|
use crate::chains::millau_messages_to_rialto::{
|
||||||
use crate::chains::rialto_messages_to_millau::run as right_to_left_messages;
|
add_standalone_metrics as add_left_to_right_standalone_metrics, run as left_to_right_messages,
|
||||||
|
};
|
||||||
|
use crate::chains::rialto_messages_to_millau::{
|
||||||
|
add_standalone_metrics as add_right_to_left_standalone_metrics, run as right_to_left_messages,
|
||||||
|
};
|
||||||
|
|
||||||
$generic
|
$generic
|
||||||
}
|
}
|
||||||
@@ -120,8 +124,12 @@ macro_rules! select_bridge {
|
|||||||
const MAX_MISSING_LEFT_HEADERS_AT_RIGHT: bp_rococo::BlockNumber = bp_rococo::SESSION_LENGTH;
|
const MAX_MISSING_LEFT_HEADERS_AT_RIGHT: bp_rococo::BlockNumber = bp_rococo::SESSION_LENGTH;
|
||||||
const MAX_MISSING_RIGHT_HEADERS_AT_LEFT: bp_wococo::BlockNumber = bp_wococo::SESSION_LENGTH;
|
const MAX_MISSING_RIGHT_HEADERS_AT_LEFT: bp_wococo::BlockNumber = bp_wococo::SESSION_LENGTH;
|
||||||
|
|
||||||
use crate::chains::rococo_messages_to_wococo::run as left_to_right_messages;
|
use crate::chains::rococo_messages_to_wococo::{
|
||||||
use crate::chains::wococo_messages_to_rococo::run as right_to_left_messages;
|
add_standalone_metrics as add_left_to_right_standalone_metrics, run as left_to_right_messages,
|
||||||
|
};
|
||||||
|
use crate::chains::wococo_messages_to_rococo::{
|
||||||
|
add_standalone_metrics as add_right_to_left_standalone_metrics, run as right_to_left_messages,
|
||||||
|
};
|
||||||
|
|
||||||
$generic
|
$generic
|
||||||
}
|
}
|
||||||
@@ -153,6 +161,8 @@ impl RelayHeadersAndMessages {
|
|||||||
|
|
||||||
let metrics_params: MetricsParams = params.shared.prometheus_params.into();
|
let metrics_params: MetricsParams = params.shared.prometheus_params.into();
|
||||||
let metrics_params = relay_utils::relay_metrics(None, metrics_params).into_params();
|
let metrics_params = relay_utils::relay_metrics(None, metrics_params).into_params();
|
||||||
|
let (metrics_params, _) = add_left_to_right_standalone_metrics(metrics_params, left_client.clone())?;
|
||||||
|
let (metrics_params, _) = add_right_to_left_standalone_metrics(metrics_params, right_client.clone())?;
|
||||||
|
|
||||||
let left_to_right_on_demand_headers = OnDemandHeadersRelay::new(
|
let left_to_right_on_demand_headers = OnDemandHeadersRelay::new(
|
||||||
left_client.clone(),
|
left_client.clone(),
|
||||||
|
|||||||
@@ -21,9 +21,16 @@ use crate::on_demand_headers::OnDemandHeadersRelay;
|
|||||||
use bp_messages::{LaneId, MessageNonce};
|
use bp_messages::{LaneId, MessageNonce};
|
||||||
use frame_support::weights::Weight;
|
use frame_support::weights::Weight;
|
||||||
use messages_relay::message_lane::{MessageLane, SourceHeaderIdOf, TargetHeaderIdOf};
|
use messages_relay::message_lane::{MessageLane, SourceHeaderIdOf, TargetHeaderIdOf};
|
||||||
use relay_substrate_client::{BlockNumberOf, Chain, Client, HashOf};
|
use relay_substrate_client::{
|
||||||
use relay_utils::{metrics::MetricsParams, BlockNumberBase};
|
metrics::{FloatStorageValueMetric, StorageProofOverheadMetric},
|
||||||
use sp_core::Bytes;
|
BlockNumberOf, Chain, Client, HashOf,
|
||||||
|
};
|
||||||
|
use relay_utils::{
|
||||||
|
metrics::{F64SharedRef, MetricsParams},
|
||||||
|
BlockNumberBase,
|
||||||
|
};
|
||||||
|
use sp_core::{storage::StorageKey, Bytes};
|
||||||
|
use sp_runtime::FixedU128;
|
||||||
use std::ops::RangeInclusive;
|
use std::ops::RangeInclusive;
|
||||||
|
|
||||||
/// Substrate <-> Substrate messages relay parameters.
|
/// Substrate <-> Substrate messages relay parameters.
|
||||||
@@ -185,6 +192,84 @@ pub fn select_delivery_transaction_limits<W: pallet_bridge_messages::WeightInfoE
|
|||||||
(max_number_of_messages, weight_for_messages_dispatch)
|
(max_number_of_messages, weight_for_messages_dispatch)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Shared references to the values of standalone metrics of the message lane relay loop.
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
|
pub struct StandaloneMessagesMetrics {
|
||||||
|
/// Shared reference to the actual target -> <base> chain token conversion rate.
|
||||||
|
pub target_to_base_conversion_rate: Option<F64SharedRef>,
|
||||||
|
/// Shared reference to the actual source -> <base> chain token conversion rate.
|
||||||
|
pub source_to_base_conversion_rate: Option<F64SharedRef>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Add general standalone metrics for the message lane relay loop.
|
||||||
|
pub fn add_standalone_metrics<P: SubstrateMessageLane>(
|
||||||
|
metrics_params: MetricsParams,
|
||||||
|
source_client: Client<P::SourceChain>,
|
||||||
|
source_chain_token_id: Option<&str>,
|
||||||
|
target_chain_token_id: Option<&str>,
|
||||||
|
target_to_source_conversion_rate_params: Option<(StorageKey, FixedU128)>,
|
||||||
|
) -> anyhow::Result<(MetricsParams, StandaloneMessagesMetrics)> {
|
||||||
|
let mut source_to_base_conversion_rate = None;
|
||||||
|
let mut target_to_base_conversion_rate = None;
|
||||||
|
let mut metrics_params =
|
||||||
|
relay_utils::relay_metrics(None, metrics_params).standalone_metric(|registry, prefix| {
|
||||||
|
StorageProofOverheadMetric::new(
|
||||||
|
registry,
|
||||||
|
prefix,
|
||||||
|
source_client.clone(),
|
||||||
|
format!("{}_storage_proof_overhead", P::SourceChain::NAME.to_lowercase()),
|
||||||
|
format!("{} storage proof overhead", P::SourceChain::NAME),
|
||||||
|
)
|
||||||
|
})?;
|
||||||
|
if let Some((target_to_source_conversion_rate_storage_key, initial_target_to_source_conversion_rate)) =
|
||||||
|
target_to_source_conversion_rate_params
|
||||||
|
{
|
||||||
|
metrics_params = metrics_params.standalone_metric(|registry, prefix| {
|
||||||
|
let metric = FloatStorageValueMetric::<_, sp_runtime::FixedU128>::new(
|
||||||
|
registry,
|
||||||
|
prefix,
|
||||||
|
source_client,
|
||||||
|
target_to_source_conversion_rate_storage_key,
|
||||||
|
Some(initial_target_to_source_conversion_rate),
|
||||||
|
format!(
|
||||||
|
"{}_{}_to_{}_conversion_rate",
|
||||||
|
P::SourceChain::NAME,
|
||||||
|
P::TargetChain::NAME,
|
||||||
|
P::SourceChain::NAME
|
||||||
|
),
|
||||||
|
format!(
|
||||||
|
"{} to {} tokens conversion rate (used by {})",
|
||||||
|
P::TargetChain::NAME,
|
||||||
|
P::SourceChain::NAME,
|
||||||
|
P::SourceChain::NAME
|
||||||
|
),
|
||||||
|
)?;
|
||||||
|
Ok(metric)
|
||||||
|
})?;
|
||||||
|
}
|
||||||
|
if let Some(source_chain_token_id) = source_chain_token_id {
|
||||||
|
metrics_params = metrics_params.standalone_metric(|registry, prefix| {
|
||||||
|
let metric = crate::chains::token_price_metric(registry, prefix, source_chain_token_id)?;
|
||||||
|
source_to_base_conversion_rate = Some(metric.shared_value_ref());
|
||||||
|
Ok(metric)
|
||||||
|
})?;
|
||||||
|
}
|
||||||
|
if let Some(target_chain_token_id) = target_chain_token_id {
|
||||||
|
metrics_params = metrics_params.standalone_metric(|registry, prefix| {
|
||||||
|
let metric = crate::chains::token_price_metric(registry, prefix, target_chain_token_id)?;
|
||||||
|
target_to_base_conversion_rate = Some(metric.shared_value_ref());
|
||||||
|
Ok(metric)
|
||||||
|
})?;
|
||||||
|
}
|
||||||
|
Ok((
|
||||||
|
metrics_params.into_params(),
|
||||||
|
StandaloneMessagesMetrics {
|
||||||
|
source_to_base_conversion_rate,
|
||||||
|
target_to_base_conversion_rate,
|
||||||
|
},
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|||||||
@@ -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]
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = "1.6.5"
|
async-std = "1.6.5"
|
||||||
async-trait = "0.1.40"
|
async-trait = "0.1.40"
|
||||||
backoff = "0.2"
|
backoff = "0.2"
|
||||||
|
|||||||
@@ -90,7 +90,7 @@ pub async fn run<P: TransactionProofPipeline>(
|
|||||||
target_client: impl TargetClient<P>,
|
target_client: impl TargetClient<P>,
|
||||||
metrics_params: MetricsParams,
|
metrics_params: MetricsParams,
|
||||||
exit_signal: impl Future<Output = ()> + 'static + Send,
|
exit_signal: impl Future<Output = ()> + 'static + Send,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let exit_signal = exit_signal.shared();
|
let exit_signal = exit_signal.shared();
|
||||||
|
|
||||||
relay_utils::relay_loop(source_client, target_client)
|
relay_utils::relay_loop(source_client, target_client)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ license = "GPL-3.0-or-later WITH Classpath-exception-2.0"
|
|||||||
description = "Finality proofs relay"
|
description = "Finality proofs relay"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = "1.6.5"
|
async-std = "1.6.5"
|
||||||
async-trait = "0.1.40"
|
async-trait = "0.1.40"
|
||||||
backoff = "0.2"
|
backoff = "0.2"
|
||||||
|
|||||||
@@ -104,7 +104,7 @@ pub async fn run<P: FinalitySyncPipeline>(
|
|||||||
sync_params: FinalitySyncParams,
|
sync_params: FinalitySyncParams,
|
||||||
metrics_params: MetricsParams,
|
metrics_params: MetricsParams,
|
||||||
exit_signal: impl Future<Output = ()> + 'static + Send,
|
exit_signal: impl Future<Output = ()> + 'static + Send,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let exit_signal = exit_signal.shared();
|
let exit_signal = exit_signal.shared();
|
||||||
relay_utils::relay_loop(source_client, target_client)
|
relay_utils::relay_loop(source_client, target_client)
|
||||||
.with_metrics(Some(metrics_prefix::<P>()), metrics_params)
|
.with_metrics(Some(metrics_prefix::<P>()), metrics_params)
|
||||||
|
|||||||
@@ -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]
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = "1.6.5"
|
async-std = "1.6.5"
|
||||||
async-trait = "0.1.40"
|
async-trait = "0.1.40"
|
||||||
backoff = "0.2"
|
backoff = "0.2"
|
||||||
|
|||||||
@@ -126,7 +126,7 @@ pub async fn run<P: HeadersSyncPipeline, TC: TargetClient<P>>(
|
|||||||
sync_params: HeadersSyncParams,
|
sync_params: HeadersSyncParams,
|
||||||
metrics_params: MetricsParams,
|
metrics_params: MetricsParams,
|
||||||
exit_signal: impl Future<Output = ()> + 'static + Send,
|
exit_signal: impl Future<Output = ()> + 'static + Send,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let exit_signal = exit_signal.shared();
|
let exit_signal = exit_signal.shared();
|
||||||
relay_utils::relay_loop(source_client, target_client)
|
relay_utils::relay_loop(source_client, target_client)
|
||||||
.with_metrics(Some(metrics_prefix::<P>()), metrics_params)
|
.with_metrics(Some(metrics_prefix::<P>()), metrics_params)
|
||||||
|
|||||||
@@ -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]
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = { version = "1.6.5", features = ["attributes"] }
|
async-std = { version = "1.6.5", features = ["attributes"] }
|
||||||
async-trait = "0.1.40"
|
async-trait = "0.1.40"
|
||||||
futures = "0.3.5"
|
futures = "0.3.5"
|
||||||
|
|||||||
@@ -258,7 +258,7 @@ pub async fn run<P: MessageLane>(
|
|||||||
target_client: impl TargetClient<P>,
|
target_client: impl TargetClient<P>,
|
||||||
metrics_params: MetricsParams,
|
metrics_params: MetricsParams,
|
||||||
exit_signal: impl Future<Output = ()> + Send + 'static,
|
exit_signal: impl Future<Output = ()> + Send + 'static,
|
||||||
) -> Result<(), String> {
|
) -> anyhow::Result<()> {
|
||||||
let exit_signal = exit_signal.shared();
|
let exit_signal = exit_signal.shared();
|
||||||
relay_utils::relay_loop(source_client, target_client)
|
relay_utils::relay_loop(source_client, target_client)
|
||||||
.reconnect_delay(params.reconnect_delay)
|
.reconnect_delay(params.reconnect_delay)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ license = "GPL-3.0-or-later WITH Classpath-exception-2.0"
|
|||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
ansi_term = "0.12"
|
ansi_term = "0.12"
|
||||||
|
anyhow = "1.0"
|
||||||
async-std = "1.6.5"
|
async-std = "1.6.5"
|
||||||
async-trait = "0.1.40"
|
async-trait = "0.1.40"
|
||||||
backoff = "0.2"
|
backoff = "0.2"
|
||||||
|
|||||||
@@ -21,12 +21,16 @@ pub use substrate_prometheus_endpoint::{
|
|||||||
register, Counter, CounterVec, Gauge, GaugeVec, Opts, PrometheusError, Registry, F64, U64,
|
register, Counter, CounterVec, Gauge, GaugeVec, Opts, PrometheusError, Registry, F64, U64,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
use async_std::sync::{Arc, RwLock};
|
||||||
use async_trait::async_trait;
|
use async_trait::async_trait;
|
||||||
use std::{fmt::Debug, time::Duration};
|
use std::{fmt::Debug, time::Duration};
|
||||||
|
|
||||||
mod float_json_value;
|
mod float_json_value;
|
||||||
mod global;
|
mod global;
|
||||||
|
|
||||||
|
/// Shared reference to `f64` value that is updated by the metric.
|
||||||
|
pub type F64SharedRef = Arc<RwLock<Option<f64>>>;
|
||||||
|
|
||||||
/// Unparsed address that needs to be used to expose Prometheus metrics.
|
/// Unparsed address that needs to be used to expose Prometheus metrics.
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct MetricsAddress {
|
pub struct MetricsAddress {
|
||||||
|
|||||||
@@ -14,8 +14,9 @@
|
|||||||
// You should have received a copy of the GNU General Public License
|
// You should have received a copy of the GNU General Public License
|
||||||
// along with Parity Bridges Common. If not, see <http://www.gnu.org/licenses/>.
|
// along with Parity Bridges Common. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
use crate::metrics::{metric_name, register, Gauge, PrometheusError, Registry, StandaloneMetrics, F64};
|
use crate::metrics::{metric_name, register, F64SharedRef, Gauge, PrometheusError, Registry, StandaloneMetrics, F64};
|
||||||
|
|
||||||
|
use async_std::sync::{Arc, RwLock};
|
||||||
use async_trait::async_trait;
|
use async_trait::async_trait;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
@@ -28,6 +29,7 @@ pub struct FloatJsonValueMetric {
|
|||||||
url: String,
|
url: String,
|
||||||
json_path: String,
|
json_path: String,
|
||||||
metric: Gauge<F64>,
|
metric: Gauge<F64>,
|
||||||
|
shared_value_ref: F64SharedRef,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl FloatJsonValueMetric {
|
impl FloatJsonValueMetric {
|
||||||
@@ -40,13 +42,20 @@ impl FloatJsonValueMetric {
|
|||||||
name: String,
|
name: String,
|
||||||
help: String,
|
help: String,
|
||||||
) -> Result<Self, PrometheusError> {
|
) -> Result<Self, PrometheusError> {
|
||||||
|
let shared_value_ref = Arc::new(RwLock::new(None));
|
||||||
Ok(FloatJsonValueMetric {
|
Ok(FloatJsonValueMetric {
|
||||||
url,
|
url,
|
||||||
json_path,
|
json_path,
|
||||||
metric: register(Gauge::new(metric_name(prefix, &name), help)?, registry)?,
|
metric: register(Gauge::new(metric_name(prefix, &name), help)?, registry)?,
|
||||||
|
shared_value_ref,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Get shared reference to metric value.
|
||||||
|
pub fn shared_value_ref(&self) -> F64SharedRef {
|
||||||
|
self.shared_value_ref.clone()
|
||||||
|
}
|
||||||
|
|
||||||
/// Read value from HTTP service.
|
/// Read value from HTTP service.
|
||||||
async fn read_value(&self) -> Result<f64, String> {
|
async fn read_value(&self) -> Result<f64, String> {
|
||||||
use isahc::{AsyncReadResponseExt, HttpClient, Request};
|
use isahc::{AsyncReadResponseExt, HttpClient, Request};
|
||||||
@@ -79,7 +88,9 @@ impl StandaloneMetrics for FloatJsonValueMetric {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn update(&self) {
|
async fn update(&self) {
|
||||||
crate::metrics::set_gauge_value(&self.metric, self.read_value().await.map(Some));
|
let value = self.read_value().await;
|
||||||
|
crate::metrics::set_gauge_value(&self.metric, value.clone().map(Some));
|
||||||
|
*self.shared_value_ref.write().await = value.ok();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -105,7 +105,7 @@ impl<SC, TC, LM> Loop<SC, TC, LM> {
|
|||||||
/// This function represents an outer loop, which in turn calls provided `run_loop` function to do
|
/// This function represents an outer loop, which in turn calls provided `run_loop` function to do
|
||||||
/// actual job. When `run_loop` returns, this outer loop reconnects to failed client (source,
|
/// actual job. When `run_loop` returns, this outer loop reconnects to failed client (source,
|
||||||
/// target or both) and calls `run_loop` again.
|
/// target or both) and calls `run_loop` again.
|
||||||
pub async fn run<R, F>(mut self, loop_name: String, run_loop: R) -> Result<(), String>
|
pub async fn run<R, F>(mut self, loop_name: String, run_loop: R) -> anyhow::Result<()>
|
||||||
where
|
where
|
||||||
R: 'static + Send + Fn(SC, TC, Option<LM>) -> F,
|
R: 'static + Send + Fn(SC, TC, Option<LM>) -> F,
|
||||||
F: 'static + Send + Future<Output = Result<(), FailedClient>>,
|
F: 'static + Send + Future<Output = Result<(), FailedClient>>,
|
||||||
@@ -151,8 +151,8 @@ impl<SC, TC, LM> LoopMetrics<SC, TC, LM> {
|
|||||||
pub fn loop_metric<NewLM: Metrics>(
|
pub fn loop_metric<NewLM: Metrics>(
|
||||||
self,
|
self,
|
||||||
create_metric: impl FnOnce(&Registry, Option<&str>) -> Result<NewLM, PrometheusError>,
|
create_metric: impl FnOnce(&Registry, Option<&str>) -> Result<NewLM, PrometheusError>,
|
||||||
) -> Result<LoopMetrics<SC, TC, NewLM>, String> {
|
) -> anyhow::Result<LoopMetrics<SC, TC, NewLM>> {
|
||||||
let loop_metric = create_metric(&self.registry, self.metrics_prefix.as_deref()).map_err(|e| e.to_string())?;
|
let loop_metric = create_metric(&self.registry, self.metrics_prefix.as_deref())?;
|
||||||
|
|
||||||
Ok(LoopMetrics {
|
Ok(LoopMetrics {
|
||||||
relay_loop: self.relay_loop,
|
relay_loop: self.relay_loop,
|
||||||
@@ -167,13 +167,13 @@ impl<SC, TC, LM> LoopMetrics<SC, TC, LM> {
|
|||||||
pub fn standalone_metric<M: StandaloneMetrics>(
|
pub fn standalone_metric<M: StandaloneMetrics>(
|
||||||
self,
|
self,
|
||||||
create_metric: impl FnOnce(&Registry, Option<&str>) -> Result<M, PrometheusError>,
|
create_metric: impl FnOnce(&Registry, Option<&str>) -> Result<M, PrometheusError>,
|
||||||
) -> Result<Self, String> {
|
) -> anyhow::Result<Self> {
|
||||||
// since standalone metrics are updating themselves, we may just ignore the fact that the same
|
// since standalone metrics are updating themselves, we may just ignore the fact that the same
|
||||||
// standalone metric is exposed by several loops && only spawn single metric
|
// standalone metric is exposed by several loops && only spawn single metric
|
||||||
match create_metric(&self.registry, self.metrics_prefix.as_deref()) {
|
match create_metric(&self.registry, self.metrics_prefix.as_deref()) {
|
||||||
Ok(standalone_metrics) => standalone_metrics.spawn(),
|
Ok(standalone_metrics) => standalone_metrics.spawn(),
|
||||||
Err(PrometheusError::AlreadyReg) => (),
|
Err(PrometheusError::AlreadyReg) => (),
|
||||||
Err(e) => return Err(e.to_string()),
|
Err(e) => anyhow::bail!(e),
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(self)
|
Ok(self)
|
||||||
@@ -191,13 +191,14 @@ impl<SC, TC, LM> LoopMetrics<SC, TC, LM> {
|
|||||||
/// Expose metrics using address passed at creation.
|
/// Expose metrics using address passed at creation.
|
||||||
///
|
///
|
||||||
/// If passed `address` is `None`, metrics are not exposed.
|
/// If passed `address` is `None`, metrics are not exposed.
|
||||||
pub async fn expose(self) -> Result<Loop<SC, TC, LM>, String> {
|
pub async fn expose(self) -> anyhow::Result<Loop<SC, TC, LM>> {
|
||||||
if let Some(address) = self.address {
|
if let Some(address) = self.address {
|
||||||
let socket_addr = SocketAddr::new(
|
let socket_addr = SocketAddr::new(
|
||||||
address.host.parse().map_err(|err| {
|
address.host.parse().map_err(|err| {
|
||||||
format!(
|
anyhow::format_err!(
|
||||||
"Invalid host {} is used to expose Prometheus metrics: {}",
|
"Invalid host {} is used to expose Prometheus metrics: {}",
|
||||||
address.host, err,
|
address.host,
|
||||||
|
err,
|
||||||
)
|
)
|
||||||
})?,
|
})?,
|
||||||
address.port,
|
address.port,
|
||||||
|
|||||||
Reference in New Issue
Block a user