use ecat_circuit_breaker::Breaker;
use ecat_data::{BackendKind, timeout_counter};
use std::sync::Arc;
use std::sync::atomic::Ordering;
pub fn register_outbound_metrics(breaker: Arc<Breaker>) {
register_as("iotdb", breaker);
}
fn register_as(backend: &'static str, breaker: Arc<Breaker>) {
let opened = Arc::clone(&breaker);
ecat_metrics::register_outbound_metrics(
backend,
Box::new(|| timeout_counter(BackendKind::Tsdb).load(Ordering::Relaxed)),
Box::new(move || opened.opened_total()),
Box::new(move || breaker.state().code()),
);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{IotdbClient, IotdbConfig};
use ecat_circuit_breaker::{BreakerConfig, BreakerState};
fn sample(text: &str, prefix: &str) -> Option<f64> {
text.lines()
.find(|l| l.starts_with(prefix) && !l.starts_with('#'))
.and_then(|l| l.rsplit(' ').next())
.and_then(|v| v.parse().ok())
}
#[tokio::test]
async fn outbound_metrics_appear_with_live_values() {
let breaker = Arc::new(Breaker::new(BreakerConfig::default()));
register_as("iotdb-live-test", Arc::clone(&breaker));
timeout_counter(BackendKind::Tsdb).fetch_add(1000, Ordering::Relaxed);
let text = ecat_metrics::metrics_text();
assert!(
sample(
&text,
r#"ecat_outbound_timeouts_total{backend="iotdb-live-test"}"#
)
.is_some_and(|v| v >= 1000.0),
"超时样本应现读静态量(先 +1000 再抓取),实际输出:\n{text}"
);
assert_eq!(
sample(
&text,
r#"ecat_outbound_breaker_state{backend="iotdb-live-test"}"#
),
Some(0.0),
"未打开时状态应为 0,实际输出:\n{text}"
);
let fail = || async { Err::<(), &str>("backend down") };
for _ in 0..5 {
let _ = breaker.call(fail).await;
}
assert_eq!(breaker.state(), BreakerState::Open);
let text = ecat_metrics::metrics_text();
assert_eq!(
sample(
&text,
r#"ecat_outbound_breaker_open_total{backend="iotdb-live-test"}"#
),
Some(1.0),
"打开次数应为 1,实际输出:\n{text}"
);
assert_eq!(
sample(
&text,
r#"ecat_outbound_breaker_state{backend="iotdb-live-test"}"#
),
Some(1.0),
"打开后状态应为 1(现读),实际输出:\n{text}"
);
}
#[tokio::test]
async fn from_config_registers_outbound_metrics() {
let cfg: IotdbConfig = serde_json::from_str(
r#"{"base_url":"http://127.0.0.1:1","username":"u","password":"p"}"#,
)
.unwrap();
let _c = IotdbClient::from_config(cfg).unwrap();
let text = ecat_metrics::metrics_text();
assert!(
text.contains(r#"ecat_outbound_timeouts_total{backend="iotdb"}"#),
"from_config 没自动注册?实际输出:\n{text}"
);
}
}