use std::sync::Arc;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde_json::Value;
use crate::error::Result;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StorageEngineType {
NoOp,
Mem,
Sqlite,
TensorBase,
ClickHouse,
}
#[derive(Debug, Clone)]
pub struct MetricsQueryRange {
pub metric_name: String,
pub start: DateTime<Utc>,
pub end: DateTime<Utc>,
pub label_matchers: Vec<crate::query::LabelMatcher>,
}
#[derive(Debug, Clone, Default)]
pub struct EventsQueryFilter {
pub table: String,
pub start: Option<DateTime<Utc>>,
pub end: Option<DateTime<Utc>>,
pub partition: Option<String>,
pub limit: Option<u32>,
pub offset: Option<u32>,
pub sort_field: Option<String>,
pub sort_desc: bool,
pub filter: crate::query::GridFilterModel,
}
#[derive(Debug, Clone)]
pub struct EventsAggregateFilter {
pub table: String,
pub start: DateTime<Utc>,
pub end: DateTime<Utc>,
pub partition: Option<String>,
pub filter: crate::query::GridFilterModel,
pub measure: crate::query::EventMeasure,
pub measure_field: Option<String>,
pub time_bucket_secs: Option<u64>,
pub group_by_field: Option<String>,
}
#[derive(Debug, Clone)]
pub struct MetricPoint {
pub ts: DateTime<Utc>,
pub value: f64,
pub labels: Value,
}
#[derive(Debug, Clone)]
pub struct EventRow {
pub ts: DateTime<Utc>,
pub fields: Value,
}
#[derive(Debug, Clone)]
pub struct MetricWriteRow {
pub name: String,
pub kind: &'static str,
pub value: Value,
pub labels: Value,
pub ts: DateTime<Utc>,
pub correlation_id: Option<String>,
}
#[derive(Debug, Clone)]
pub struct EventWriteRow {
pub table: String,
pub fields: Value,
pub ts: DateTime<Utc>,
pub correlation_id: Option<String>,
}
#[async_trait]
pub trait MetricsStorageBackend: Send + Sync {
fn engine_type(&self) -> StorageEngineType;
async fn record_counter(
&self,
name: &str,
labels: &Value,
delta: i64,
ts: DateTime<Utc>,
) -> Result<()>;
async fn record_gauge(
&self,
name: &str,
labels: &Value,
value: f64,
ts: DateTime<Utc>,
) -> Result<()>;
async fn query_range(&self, _query: MetricsQueryRange) -> Result<Vec<MetricPoint>> {
Ok(Vec::new())
}
async fn record_metrics_batch(&self, rows: &[MetricWriteRow]) -> Result<()> {
for row in rows {
if row.kind == "counter" {
let delta = row.value.as_i64().unwrap_or(0);
self.record_counter(&row.name, &row.labels, delta, row.ts)
.await?;
} else {
let value = row.value.as_f64().unwrap_or(0.0);
self.record_gauge(&row.name, &row.labels, value, row.ts)
.await?;
}
}
Ok(())
}
}
#[async_trait]
pub trait EventStorageBackend: Send + Sync {
fn engine_type(&self) -> StorageEngineType;
async fn append_row(
&self,
table: &str,
fields: &Value,
ts: DateTime<Utc>,
correlation_id: Option<&str>,
) -> Result<()>;
async fn query_rows(&self, _filter: EventsQueryFilter) -> Result<Vec<EventRow>> {
Ok(Vec::new())
}
async fn query_aggregate(
&self,
_filter: EventsAggregateFilter,
) -> Result<crate::query::EventAggregateResult> {
Ok(crate::query::EventAggregateResult::TimeSeries {
series: Vec::new(),
headline: Vec::new(),
})
}
async fn append_rows_batch(&self, rows: &[EventWriteRow]) -> Result<()> {
for row in rows {
self.append_row(
&row.table,
&row.fields,
row.ts,
row.correlation_id.as_deref(),
)
.await?;
}
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoOpMetricsBackend;
#[async_trait]
impl MetricsStorageBackend for NoOpMetricsBackend {
fn engine_type(&self) -> StorageEngineType {
StorageEngineType::NoOp
}
async fn record_counter(&self, _: &str, _: &Value, _: i64, _: DateTime<Utc>) -> Result<()> {
Ok(())
}
async fn record_gauge(&self, _: &str, _: &Value, _: f64, _: DateTime<Utc>) -> Result<()> {
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NoOpEventBackend;
#[async_trait]
impl EventStorageBackend for NoOpEventBackend {
fn engine_type(&self) -> StorageEngineType {
StorageEngineType::NoOp
}
async fn append_row(
&self,
_: &str,
_: &Value,
_: DateTime<Utc>,
_: Option<&str>,
) -> Result<()> {
Ok(())
}
}
pub type SharedMetricsBackend = Arc<dyn MetricsStorageBackend>;
pub type SharedEventBackend = Arc<dyn EventStorageBackend>;