use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde_json::Value;
use spectra_backend_remote_common::RemoteMetricsBackend;
use spectra_core::{
MetricPoint, MetricWriteRow, MetricsQueryRange, MetricsStorageBackend, Result,
StorageEngineType,
};
pub struct TensorBaseMetricsBackend(RemoteMetricsBackend);
impl TensorBaseMetricsBackend {
pub async fn connect(url: &str) -> Result<Self> {
Ok(Self(
RemoteMetricsBackend::connect(
url,
StorageEngineType::TensorBase,
&crate::ddl::metrics_ddl(),
)
.await?,
))
}
pub async fn connect_host(host: &str) -> Result<Self> {
Self::connect(&crate::ddl::default_url(host)).await
}
pub fn in_memory_stub() -> Self {
Self(RemoteMetricsBackend::in_memory_for_test(
StorageEngineType::TensorBase,
))
}
#[cfg(test)]
pub fn in_memory_for_test() -> Self {
Self::in_memory_stub()
}
}
#[async_trait]
impl MetricsStorageBackend for TensorBaseMetricsBackend {
fn engine_type(&self) -> StorageEngineType {
self.0.engine_type()
}
async fn record_counter(
&self,
name: &str,
labels: &Value,
delta: i64,
ts: DateTime<Utc>,
) -> Result<()> {
self.0.record_counter(name, labels, delta, ts).await
}
async fn record_gauge(
&self,
name: &str,
labels: &Value,
value: f64,
ts: DateTime<Utc>,
) -> Result<()> {
self.0.record_gauge(name, labels, value, ts).await
}
async fn record_metrics_batch(&self, rows: &[MetricWriteRow]) -> Result<()> {
self.0.record_metrics_batch(rows).await
}
async fn query_range(&self, query: MetricsQueryRange) -> Result<Vec<MetricPoint>> {
self.0.query_range(query).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Duration;
use serde_json::json;
#[tokio::test]
async fn tensorbase_metrics_roundtrip_in_memory() {
let backend = TensorBaseMetricsBackend::in_memory_for_test();
let ts = Utc::now();
backend
.record_counter("hits", &json!({}), 2, ts)
.await
.expect("write");
let points = backend
.query_range(MetricsQueryRange {
metric_name: "hits".into(),
start: ts - Duration::seconds(1),
end: ts + Duration::seconds(1),
label_matchers: vec![],
})
.await
.expect("query");
assert_eq!(points.len(), 1);
}
#[tokio::test]
#[ignore = "requires SPECTRA_TENSORBASE_URL"]
async fn tensorbase_metrics_integration() {
let url = std::env::var("SPECTRA_TENSORBASE_URL").expect("SPECTRA_TENSORBASE_URL");
let backend = TensorBaseMetricsBackend::connect(&url)
.await
.expect("connect");
let ts = Utc::now();
backend
.record_counter("integration_hits", &json!({}), 1, ts)
.await
.expect("write");
let points = backend
.query_range(MetricsQueryRange {
metric_name: "integration_hits".into(),
start: ts - Duration::seconds(5),
end: ts + Duration::seconds(5),
label_matchers: vec![],
})
.await
.expect("query");
assert!(!points.is_empty());
}
}