use crate::server::AppState;
use axum::{
extract::{Query, State},
http::StatusCode,
response::Json,
};
use otelite_core::api::{
AgentRolesResponse, AgentRollup, AgentRollupResponse, ConversationCostRow,
ConversationDepthStats, CostSeriesPoint, ErrorResponse, GenAiCapabilityResponse,
LatencyPercentilesResponse, ProjectRollupResponse, ProviderMixResponse, ReasoningShareResponse,
RequestParamProfile, RetrievalStats, RetryStats, SessionCostRow, TokenUsageResponse,
ToolApprovalStats, TopSpan, TopSpanSort,
};
use otelite_core::filters::{GenAiFilters, FILTER_DIMENSIONS};
use otelite_core::pricing::{PricingDatabase, TokenUsage};
use serde::{Deserialize, Serialize};
macro_rules! genai_filter_impl {
($t:ident) => {
impl $t {
fn filters(&self) -> GenAiFilters {
GenAiFilters {
agent: self.agent.clone(),
model: self.model.clone(),
models: None,
provider: self.provider.clone(),
project: self.project.clone(),
session: self.session.clone(),
}
}
}
};
}
#[derive(Debug, Serialize, utoipa::ToSchema)]
pub struct GenAiItemsResponse {
pub items: serde_json::Value,
pub filters_applied: Vec<String>,
}
fn enrich_top_spans(rows: &mut [TopSpan], db: &PricingDatabase) {
for row in rows {
let usage = TokenUsage {
input: row.input_tokens,
output: row.output_tokens,
cache_creation: row.cache_creation_tokens,
cache_read: row.cache_read_tokens,
};
let result = db.compute_cost(row.model.as_deref(), usage, row.system.as_deref());
row.cost = result.cost;
row.cost_source = Some(result.source.as_str().to_string());
row.cost_reason = result.reason;
let duration_ms = row.duration / 1_000_000;
if row.output_tokens > 0 && duration_ms > 0 {
row.derived_output_tokens_per_sec =
Some(row.output_tokens as f64 / (duration_ms as f64 / 1000.0));
}
}
}
fn enrich_cost_series(rows: &mut [CostSeriesPoint], db: &PricingDatabase) {
for row in rows {
let usage = TokenUsage {
input: row.input_tokens,
output: row.output_tokens,
cache_creation: row.cache_creation_tokens,
cache_read: row.cache_read_tokens,
};
let result = db.compute_cost(row.model.as_deref(), usage, None);
row.cost = result.cost;
row.cost_source = Some(result.source.as_str().to_string());
}
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct TokenUsageQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/usage",
params(TokenUsageQuery),
responses(
(status = 200, description = "Token usage summary", body = TokenUsageResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_token_usage(
State(state): State<AppState>,
Query(query): Query<TokenUsageQuery>,
) -> Result<Json<TokenUsageResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let (summary, by_model, by_system) = state
.storage
.query_token_usage(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query token usage: {}",
e
))),
)
})?;
Ok(Json(TokenUsageResponse {
summary,
by_model,
by_system,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct CostSeriesQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub bucket: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/cost_series",
params(CostSeriesQuery),
responses(
(status = 200, description = "Cost series points", body = GenAiItemsResponse),
(status = 400, description = "Invalid bucket parameter", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_cost_series(
State(state): State<AppState>,
Query(query): Query<CostSeriesQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let bucket_seconds = query.bucket.unwrap_or(3600);
if bucket_seconds <= 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"bucket must be a positive number of seconds",
)),
));
}
let bucket_ns = bucket_seconds.saturating_mul(1_000_000_000);
let filters = query.filters();
let mut series = state
.storage
.query_cost_series(query.start_time, query.end_time, bucket_ns, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query cost series: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
enrich_cost_series(&mut series, &pricing.db);
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&series).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize cost series: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct TopSpansQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub limit: Option<usize>,
#[serde(default)]
pub sort_by: TopSpanSort,
#[serde(default)]
pub truncated_only: bool,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct TopGroupQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub limit: Option<usize>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/top_spans",
params(TopSpansQuery),
responses(
(status = 200, description = "Top spans", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_top_spans(
State(state): State<AppState>,
Query(query): Query<TopSpansQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let limit = query.limit.unwrap_or(20).clamp(1, 100);
let filters = query.filters();
let mut spans = state
.storage
.query_top_spans(
query.start_time,
query.end_time,
&filters,
limit,
query.sort_by,
query.truncated_only,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query top spans: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
enrich_top_spans(&mut spans, &pricing.db);
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&spans).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize top spans: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
fn enrich_session_rows(rows: &mut [SessionCostRow], db: &PricingDatabase) {
for row in rows {
let usage = TokenUsage {
input: row.input_tokens,
output: row.output_tokens,
..Default::default()
};
let result = db.compute_cost(None, usage, None);
row.cost = result.cost;
row.cost_source = Some(result.source.as_str().to_string());
}
}
fn enrich_conversation_rows(rows: &mut [ConversationCostRow], db: &PricingDatabase) {
for row in rows {
let usage = TokenUsage {
input: row.input_tokens,
output: row.output_tokens,
..Default::default()
};
let result = db.compute_cost(None, usage, None);
row.cost = result.cost;
row.cost_source = Some(result.source.as_str().to_string());
}
}
#[utoipa::path(
get,
path = "/api/genai/top_sessions",
params(TopGroupQuery),
responses(
(status = 200, description = "Top sessions", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_top_sessions(
State(state): State<AppState>,
Query(query): Query<TopGroupQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let limit = query.limit.unwrap_or(20).clamp(1, 100);
let filters = query.filters();
let mut rows = state
.storage
.query_top_sessions(query.start_time, query.end_time, &filters, limit)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query top sessions: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
enrich_session_rows(&mut rows, &pricing.db);
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize top sessions: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/top_conversations",
params(TopGroupQuery),
responses(
(status = 200, description = "Top conversations", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_top_conversations(
State(state): State<AppState>,
Query(query): Query<TopGroupQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let limit = query.limit.unwrap_or(20).clamp(1, 100);
let filters = query.filters();
let mut rows = state
.storage
.query_top_conversations(query.start_time, query.end_time, &filters, limit)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query top conversations: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
enrich_conversation_rows(&mut rows, &pricing.db);
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize top conversations: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct FinishReasonsQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/finish_reasons",
params(FinishReasonsQuery),
responses(
(status = 200, description = "Finish reason counts", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_finish_reasons(
State(state): State<AppState>,
Query(query): Query<FinishReasonsQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_finish_reasons(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query finish reasons: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize finish_reasons: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct LatencyQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct LatencyPercentileQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub bucket_secs: Option<u64>,
pub metrics: Option<String>,
pub calendar_day: Option<String>,
pub timezone: Option<String>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/latency_stats",
params(LatencyQuery),
responses(
(status = 200, description = "Latency statistics per model", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_latency_stats(
State(state): State<AppState>,
Query(query): Query<LatencyQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_latency_stats(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query latency stats: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize latency_stats: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/latency_percentiles",
params(LatencyPercentileQuery),
responses(
(status = 200, description = "Percentile series per metric and model", body = LatencyPercentilesResponse),
(status = 400, description = "Invalid bucket_secs or metrics", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_latency_percentiles(
State(state): State<AppState>,
Query(query): Query<LatencyPercentileQuery>,
) -> Result<Json<LatencyPercentilesResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let bucket_secs = query.bucket_secs.unwrap_or(3600);
if bucket_secs == 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"bucket_secs must be a positive number of seconds",
)),
));
}
let metrics: Vec<&str> = query
.metrics
.as_deref()
.unwrap_or("duration,ttft")
.split(',')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.collect();
for m in &metrics {
if *m != "duration" && *m != "ttft" {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"unknown metric '{m}' — expected \"duration\" and/or \"ttft\""
))),
));
}
}
if metrics.is_empty() {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"metrics must include \"duration\" and/or \"ttft\"",
)),
));
}
let calendar_day = match query.calendar_day.as_deref() {
None => false,
Some("1" | "true" | "yes") => true,
Some("0" | "false" | "no") => false,
Some(v) => {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"calendar_day must be 1/0, true/false or yes/no, got '{v}'"
))),
))
},
};
let timezone: Option<String> = match (query.timezone.as_deref(), calendar_day) {
(Some(tz), true) => {
let tz = tz.trim();
let _parsed: chrono_tz::Tz = std::str::FromStr::from_str(tz).map_err(|e| {
(
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"unknown IANA timezone '{tz}': {e} (e.g. Europe/London, America/New_York, UTC)"
))),
)
})?;
Some(tz.to_string())
},
(None, true) => Some("UTC".to_string()),
(Some(tz), false) => {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"timezone '{tz}' is only used with calendar_day=1"
))),
))
},
(None, false) => None,
};
if calendar_day && (query.start_time.is_none() || query.end_time.is_none()) {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"calendar_day=1 requires explicit start_time and end_time",
)),
));
}
let timezone = timezone.as_deref();
let mut rows = state
.storage
.query_latency_percentiles(
query.start_time,
query.end_time,
&filters,
bucket_secs,
&metrics,
timezone,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query latency percentiles: {}",
e
))),
)
})?;
rows.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(rows))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct DistributionQuery {
pub metric: String,
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub buckets: Option<usize>,
pub scale: Option<String>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/distributions",
params(DistributionQuery),
responses(
(status = 200, description = "Distribution buckets + summary stats", body = otelite_core::api::DistributionResponse),
(status = 400, description = "Invalid metric, scale or buckets", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_distributions(
State(state): State<AppState>,
Query(query): Query<DistributionQuery>,
) -> Result<Json<otelite_core::api::DistributionResponse>, (StatusCode, Json<ErrorResponse>)> {
use otelite_core::distribution;
use otelite_core::session_cost;
const KNOWN: &[&str] = &[
"session_cost",
"tool_duration",
"llm_duration",
"ttft",
"output_tokens",
];
if !KNOWN.contains(&query.metric.as_str()) {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"unknown metric '{}' — expected one of: {}",
query.metric,
KNOWN.join(" | ")
))),
));
}
let scale = query.scale.as_deref().unwrap_or("linear");
if scale != "linear" && scale != "log" {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"unknown scale '{scale}' — expected \"linear\" or \"log\""
))),
));
}
let buckets = query.buckets.unwrap_or(20);
if buckets == 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"buckets must be a positive integer",
)),
));
}
let filters = query.filters();
let resp = if query.metric == "session_cost" {
let rows = state
.storage
.query_session_costs(query.start_time, query.end_time)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query session costs: {e}"
))),
)
})?;
let pricing = state.pricing.snapshot().await;
let sessions = session_cost::build_session_costs(rows, &pricing.db);
let values: Vec<f64> = sessions.iter().filter_map(|s| s.cost_usd).collect();
distribution::build("session_cost", "usd", scale, buckets, values)
} else {
state
.storage
.query_distribution(
&query.metric,
query.start_time,
query.end_time,
&filters,
buckets,
scale,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query distribution: {e}"
))),
)
})?
};
Ok(Json(resp))
}
#[utoipa::path(
get,
path = "/api/genai/capabilities",
params(ModelAnalyticsQuery),
responses(
(status = 200, description = "GenAI telemetry capability report", body = GenAiCapabilityResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_genai_capabilities(
State(state): State<AppState>,
Query(query): Query<ModelAnalyticsQuery>,
) -> Result<Json<GenAiCapabilityResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut report = state
.storage
.query_genai_capabilities(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query GenAI capabilities: {e}"
))),
)
})?;
report.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(report))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct ModelPerformanceQueryParams {
pub start_time: i64,
pub end_time: i64,
pub rolling_ns: Option<i64>,
pub model: Option<String>,
pub provider: Option<String>,
pub timezone: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/model-performance",
params(ModelPerformanceQueryParams),
responses(
(status = 200, description = "Model-performance diagnosis", body = otelite_core::api::ModelPerformanceDiagnosis),
(status = 400, description = "Invalid interval or baseline", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_model_performance(
State(state): State<AppState>,
Query(query): Query<ModelPerformanceQueryParams>,
) -> Result<Json<otelite_core::api::ModelPerformanceDiagnosis>, (StatusCode, Json<ErrorResponse>)> {
if query.end_time <= query.start_time {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"end_time must be after start_time (got start={} end={})",
query.start_time, query.end_time
))),
));
}
if let Some(len) = query.rolling_ns {
if len <= 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"rolling_ns must be a positive length in nanoseconds (got {len})"
))),
));
}
}
if let Some(tz) = query.timezone.as_deref() {
let tz = tz.trim();
let _parsed: chrono_tz::Tz = std::str::FromStr::from_str(tz).map_err(|e| {
(
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(format!(
"unknown IANA timezone '{tz}': {e} (e.g. Europe/London, America/New_York, UTC)"
))),
)
})?;
}
let current_len = query.end_time - query.start_time;
let preceding_start = query.start_time - current_len;
let query_core = otelite_core::api::ModelPerformanceQuery {
current: otelite_core::api::ModelPerformanceWindow {
start_time: query.start_time,
end_time: query.end_time,
},
rolling: query
.rolling_ns
.map(|len| otelite_core::api::ModelPerformanceWindow {
start_time: preceding_start - len,
end_time: preceding_start,
}),
model: query.model.clone(),
provider: query.provider.clone(),
};
let response = state
.storage
.query_model_performance(&query_core)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query model performance: {e}"
))),
)
})?;
let capability = state
.storage
.query_genai_capabilities(
Some(query_core.current.start_time),
Some(query_core.current.end_time),
&GenAiFilters {
model: query.model.clone(),
provider: query.provider.clone(),
..Default::default()
},
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query GenAI capabilities for TTFT trust: {e}"
))),
)
})?;
Ok(Json(otelite_core::model_performance::build_diagnosis(
&response,
&capability,
query.timezone.clone(),
)))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct ErrorRateQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/error_rate",
params(ErrorRateQuery),
responses(
(status = 200, description = "Error rate per model", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_error_rate(
State(state): State<AppState>,
Query(query): Query<ErrorRateQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_error_rate(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query error rate: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize error_rate: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct ToolUsageQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub limit: Option<usize>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/tool_usage",
params(ToolUsageQuery),
responses(
(status = 200, description = "Tool usage aggregates", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_tool_usage(
State(state): State<AppState>,
Query(query): Query<ToolUsageQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let limit = query.limit.unwrap_or(20).clamp(1, 100);
let filters = query.filters();
let rows = state
.storage
.query_tool_usage(query.start_time, query.end_time, &filters, limit)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query tool usage: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize tool_usage: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct RetryStatsQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/retry_stats",
params(RetryStatsQuery),
responses(
(status = 200, description = "Retry statistics", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_retry_stats(
State(state): State<AppState>,
Query(query): Query<RetryStatsQuery>,
) -> Result<Json<RetryStats>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut stats = state
.storage
.query_retry_stats(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query retry stats: {}",
e
))),
)
})?;
stats.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(stats))
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct RetrievalStatsQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub limit: Option<usize>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/retrieval_stats",
params(RetrievalStatsQuery),
responses(
(status = 200, description = "Retrieval statistics", body = RetrievalStats),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_retrieval_stats(
State(state): State<AppState>,
Query(query): Query<RetrievalStatsQuery>,
) -> Result<Json<RetrievalStats>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let limit = query.limit.unwrap_or(5).clamp(1, 20);
let mut stats = state
.storage
.query_retrieval_stats(query.start_time, query.end_time, &filters, limit)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query retrieval stats: {}",
e
))),
)
})?;
stats.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(stats))
}
#[derive(Debug, Clone, Serialize, utoipa::ToSchema)]
pub struct PricingMetadata {
pub source: &'static str,
pub entry_count: usize,
pub last_fetched_unix_ms: Option<i64>,
pub last_failed_unix_ms: Option<i64>,
pub fallback_last_verified: &'static str,
pub source_url: &'static str,
pub license: &'static str,
pub disclaimer: &'static str,
pub filters_applied: Vec<String>,
}
#[utoipa::path(
get,
path = "/api/genai/agent_framework_defs",
responses(
(status = 200, description = "Agent framework recognizers"),
),
tag = "genai"
)]
pub async fn get_agent_framework_defs(
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let items =
serde_json::to_value(otelite_core::agent_frameworks::AGENT_FRAMEWORKS).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize agent framework defs: {e}"
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items,
filters_applied: Vec::new(),
}))
}
const PRICING_DISCLAIMER: &str =
"Cost figures are best-effort estimates. Per-token rates sourced from the LiteLLM \
community pricing database (MIT-licensed, © 2023 Berri AI). When the upstream \
fetch is unavailable, a small hand-curated Claude 4.x fallback table is used.";
#[utoipa::path(
get,
path = "/api/genai/pricing_metadata",
responses(
(status = 200, description = "Pricing metadata", body = PricingMetadata),
),
tag = "genai"
)]
pub async fn get_pricing_metadata(State(state): State<AppState>) -> Json<PricingMetadata> {
let snapshot = state.pricing.snapshot().await;
Json(PricingMetadata {
source: if snapshot.db.is_litellm() {
"litellm"
} else {
"fallback"
},
entry_count: snapshot.db.len(),
last_fetched_unix_ms: snapshot.last_fetched_unix_ms,
last_failed_unix_ms: snapshot.last_failed_unix_ms,
fallback_last_verified: otelite_core::pricing::FALLBACK_LAST_VERIFIED,
source_url: otelite_core::pricing::LITELLM_SOURCE_URL,
license: otelite_core::pricing::LITELLM_LICENSE,
disclaimer: PRICING_DISCLAIMER,
filters_applied: Vec::new(),
})
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct ModelAnalyticsQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct CacheQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
pub by_model: Option<String>,
pub bucket_secs: Option<u64>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct TimeSeriesQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub bucket_secs: Option<u64>,
pub span_filter: Option<String>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct ModelTimeSeriesQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub bucket_secs: Option<u64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
pub span_filter: Option<String>,
}
#[derive(Debug, Deserialize, Serialize, utoipa::IntoParams, utoipa::ToSchema)]
pub struct TimeRangeQuery {
pub start_time: Option<i64>,
pub end_time: Option<i64>,
pub agent: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub project: Option<String>,
pub session: Option<String>,
}
#[utoipa::path(
get,
path = "/api/genai/truncation_rate",
params(ModelAnalyticsQuery),
responses(
(status = 200, description = "Truncation rate by model", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_truncation_rate(
State(state): State<AppState>,
Query(query): Query<ModelAnalyticsQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_truncation_rate(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query truncation rate: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize truncation_rate: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
fn by_model_enabled(v: Option<&str>) -> bool {
matches!(v, Some("1") | Some("true"))
}
#[utoipa::path(
get,
path = "/api/genai/cache_hit_rate",
params(CacheQuery),
responses(
(status = 200, description = "Per-model cache hit rate (default) or cache economics with by_model=1", body = serde_json::Value),
(status = 400, description = "Invalid bucket_secs", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_cache_hit_rate(
State(state): State<AppState>,
Query(query): Query<CacheQuery>,
) -> Result<Json<serde_json::Value>, (StatusCode, Json<ErrorResponse>)> {
if by_model_enabled(query.by_model.as_deref()) {
let bucket_secs = query.bucket_secs.unwrap_or(3600);
if bucket_secs == 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"bucket_secs must be a positive number of seconds",
)),
));
}
let mut response = state
.storage
.query_cache_economics(
query.start_time,
query.end_time,
(bucket_secs as i64) * 1_000_000_000,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query cache economics: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
for m in &mut response.models {
let r = pricing
.db
.compute_cache_savings(Some(&m.model), m.cache_read_tokens, None);
m.est_savings_usd = r.cost;
m.savings_known = r.cost.is_some();
}
let value = serde_json::to_value(response).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize cache economics: {e}"
))),
)
})?;
let obj = value.as_object().ok_or_else(|| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(
"cache economics payload is not a JSON object".to_string(),
)),
)
})?;
let mut obj = obj.clone();
obj.insert(
"filters_applied".to_string(),
serde_json::Value::Array(Vec::new()),
);
return Ok(Json(serde_json::Value::Object(obj)));
}
let filters = query.filters();
let rows = state
.storage
.query_cache_hit_rate(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query cache hit rate: {}",
e
))),
)
})?;
let items = serde_json::to_value(rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize cache hit rate: {e}"
))),
)
})?;
Ok(Json(serde_json::json!({
"items": items,
"filters_applied": filters.applied(&FILTER_DIMENSIONS),
})))
}
#[utoipa::path(
get,
path = "/api/genai/reasoning_share",
params(TimeRangeQuery),
responses(
(status = 200, description = "Reasoning token share by model and effort", body = ReasoningShareResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_reasoning_share(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<ReasoningShareResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut response = state
.storage
.query_reasoning_share(query.start_time, query.end_time)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query reasoning share: {e}"
))),
)
})?;
let pricing = state.pricing.snapshot().await;
for m in &mut response.models {
let usage = TokenUsage {
input: 0,
output: m.reasoning_tokens,
cache_creation: 0,
cache_read: 0,
};
m.cost_usd = pricing
.db
.compute_cost(Some(m.model.as_str()), usage, None)
.cost;
}
response.filters_applied = filters.applied(&[]);
Ok(Json(response))
}
#[utoipa::path(
get,
path = "/api/genai/agents",
params(TimeSeriesQuery),
responses(
(status = 200, description = "Per-harness sessions, cost, tokens, tool calls and retries", body = AgentRollupResponse),
(status = 400, description = "Invalid bucket_secs", body = ErrorResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_agents(
State(state): State<AppState>,
Query(query): Query<TimeSeriesQuery>,
) -> Result<Json<AgentRollupResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let bucket_secs = query.bucket_secs.unwrap_or(3600);
if bucket_secs == 0 {
return Err((
StatusCode::BAD_REQUEST,
Json(ErrorResponse::bad_request(
"bucket_secs must be a positive number of seconds",
)),
));
}
let rollups = state
.storage
.query_agent_rollup(query.start_time, query.end_time, bucket_secs)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query agent rollup: {e}"
))),
)
})?;
let pricing = state.pricing.snapshot().await;
let mut agents: Vec<AgentRollup> = rollups.into_iter().map(|r| r.enrich(&pricing.db)).collect();
agents.sort_by(|a, b| {
b.cost_usd
.unwrap_or(0.0)
.partial_cmp(&a.cost_usd.unwrap_or(0.0))
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| a.agent.cmp(&b.agent))
});
Ok(Json(AgentRollupResponse {
agents,
filters_applied: filters.applied(&[]),
}))
}
#[utoipa::path(
get,
path = "/api/genai/projects",
params(TimeRangeQuery),
responses(
(status = 200, description = "Cost, sessions and tokens per project", body = ProjectRollupResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_projects(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<ProjectRollupResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rollups = state
.storage
.query_project_rollup(query.start_time, query.end_time)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query project rollup: {e}"
))),
)
})?;
let pricing = state.pricing.snapshot().await;
let mut projects: Vec<_> = rollups.into_iter().map(|r| r.enrich(&pricing.db)).collect();
projects.sort_by(|a, b| {
b.cost_usd
.unwrap_or(0.0)
.partial_cmp(&a.cost_usd.unwrap_or(0.0))
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| a.project_id.cmp(&b.project_id))
});
Ok(Json(ProjectRollupResponse {
projects,
filters_applied: filters.applied(&[]),
}))
}
#[utoipa::path(
get,
path = "/api/genai/agent_roles",
params(TimeRangeQuery),
responses(
(status = 200, description = "Cost and token attribution per sub-agent role", body = AgentRolesResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_agent_roles(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<AgentRolesResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut response = state
.storage
.query_agent_roles(query.start_time, query.end_time)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query agent roles: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
for role in &mut response.roles {
let mut total: f64 = 0.0;
let mut all_priced = true;
for m in &mut role.top_models {
let usage = TokenUsage {
input: m.tokens.input,
output: m.tokens.output,
cache_creation: m.tokens.cache_write,
cache_read: m.tokens.cache_read,
};
let result = pricing.db.compute_cost(Some(m.model.as_str()), usage, None);
m.cost = result.cost;
m.cost_source = Some(result.source.as_str().to_string());
m.cost_reason = result.reason;
match result.cost {
Some(c) => total += c,
None => all_priced = false,
}
}
role.cost = if all_priced && !role.top_models.is_empty() {
Some(total)
} else {
None
};
}
response.filters_applied = filters.applied(&[]);
Ok(Json(response))
}
#[utoipa::path(
get,
path = "/api/genai/provider_mix",
params(TimeRangeQuery),
responses(
(status = 200, description = "Provider x model token and cost mix", body = ProviderMixResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_provider_mix(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<ProviderMixResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut response = state
.storage
.query_provider_mix(query.start_time, query.end_time)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query provider mix: {}",
e
))),
)
})?;
let pricing = state.pricing.snapshot().await;
for provider in &mut response.providers {
let mut total: f64 = 0.0;
let mut any_priced = false;
for m in &mut provider.models {
let usage = TokenUsage {
input: m.tokens.input,
output: m.tokens.output,
cache_creation: m.tokens.cache_write,
cache_read: m.tokens.cache_read,
};
let result = pricing.db.compute_cost(Some(m.model.as_str()), usage, None);
if let Some(c) = result.cost {
total += c;
any_priced = true;
}
m.cost_usd = result.cost;
m.cost_source = Some(result.source.as_str().to_string());
}
provider.cost_usd = if any_priced { Some(total) } else { None };
}
response.filters_applied = filters.applied(&[]);
Ok(Json(response))
}
#[utoipa::path(
get,
path = "/api/genai/request_param_profile",
params(TimeRangeQuery),
responses(
(status = 200, description = "Request parameter profile", body = RequestParamProfile),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_request_param_profile(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<RequestParamProfile>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut profile = state
.storage
.query_request_param_profile(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query request param profile: {}",
e
))),
)
})?;
profile.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(profile))
}
#[utoipa::path(
get,
path = "/api/genai/conversation_depth",
params(TimeRangeQuery),
responses(
(status = 200, description = "Conversation depth statistics", body = ConversationDepthStats),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_conversation_depth(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<ConversationDepthStats>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut stats = state
.storage
.query_conversation_depth(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query conversation depth: {}",
e
))),
)
})?;
stats.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(stats))
}
#[utoipa::path(
get,
path = "/api/genai/latency_series",
params(ModelTimeSeriesQuery),
responses(
(status = 200, description = "Latency stats per time bucket", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_latency_series(
State(state): State<AppState>,
Query(query): Query<ModelTimeSeriesQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let bucket_secs = query.bucket_secs.unwrap_or(3600).clamp(60, 86400);
let all_spans = query.span_filter.as_deref() == Some("all");
let rows = state
.storage
.query_latency_series(
query.start_time,
query.end_time,
bucket_secs,
&filters,
all_spans,
None,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query latency series: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize latency_series: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/calls_series",
params(TimeSeriesQuery),
responses(
(status = 200, description = "Calls per time bucket", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_calls_series(
State(state): State<AppState>,
Query(query): Query<TimeSeriesQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let bucket_secs = query.bucket_secs.unwrap_or(3600).clamp(60, 86400);
let all_spans = query.span_filter.as_deref() == Some("all");
let rows = state
.storage
.query_calls_series(
query.start_time,
query.end_time,
&filters,
bucket_secs,
all_spans,
)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query calls series: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize calls_series: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/latency_by_context",
params(ModelAnalyticsQuery),
responses(
(status = 200, description = "Latency per context size bin", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_latency_by_context(
State(state): State<AppState>,
Query(query): Query<ModelAnalyticsQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_latency_by_context(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query latency by context: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize latency_by_context: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/error_types",
params(ModelAnalyticsQuery),
responses(
(status = 200, description = "Error type breakdown per model", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_error_types(
State(state): State<AppState>,
Query(query): Query<ModelAnalyticsQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_error_types(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query error types: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize error_types: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/model_drift",
params(TimeRangeQuery),
responses(
(status = 200, description = "Request→response model pairs", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_model_drift(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_model_drift(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query model drift: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize model_drift: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/tool_approvals",
params(TimeRangeQuery),
responses(
(status = 200, description = "Tool approval statistics", body = ToolApprovalStats),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_tool_approvals(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<ToolApprovalStats>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let mut stats = state
.storage
.query_tool_approvals(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query tool approvals: {}",
e
))),
)
})?;
stats.filters_applied = filters.applied(&FILTER_DIMENSIONS);
Ok(Json(stats))
}
#[utoipa::path(
get,
path = "/api/genai/stop_reasons",
params(TimeRangeQuery),
responses(
(status = 200, description = "Stop reason distribution", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_stop_reasons(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_stop_reasons(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query stop reasons: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize stop_reasons: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/context_type_split",
params(TimeRangeQuery),
responses(
(status = 200, description = "Context type token split", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_context_type_split(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_context_type_split(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query context type split: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize context_type_split: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/tool_errors",
params(ToolUsageQuery),
responses(
(status = 200, description = "Tool error messages", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_tool_errors(
State(state): State<AppState>,
Query(query): Query<ToolUsageQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let limit = query.limit.unwrap_or(30).clamp(1, 100);
let filters = query.filters();
let rows = state
.storage
.query_tool_errors(query.start_time, query.end_time, &filters, limit)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query tool errors: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize tool_errors: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
#[utoipa::path(
get,
path = "/api/genai/hour_of_day",
params(TimeRangeQuery),
responses(
(status = 200, description = "Hour-of-day buckets", body = GenAiItemsResponse),
(status = 500, description = "Internal server error", body = ErrorResponse)
),
tag = "genai"
)]
pub async fn get_hour_of_day(
State(state): State<AppState>,
Query(query): Query<TimeRangeQuery>,
) -> Result<Json<GenAiItemsResponse>, (StatusCode, Json<ErrorResponse>)> {
let filters = query.filters();
let rows = state
.storage
.query_hour_of_day(query.start_time, query.end_time, &filters)
.await
.map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"query hour of day: {}",
e
))),
)
})?;
Ok(Json(GenAiItemsResponse {
items: serde_json::to_value(&rows).map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse::storage_error(format!(
"serialize hour_of_day: {e}"
))),
)
})?,
filters_applied: filters.applied(&FILTER_DIMENSIONS),
}))
}
genai_filter_impl!(TokenUsageQuery);
genai_filter_impl!(CostSeriesQuery);
genai_filter_impl!(FinishReasonsQuery);
genai_filter_impl!(LatencyQuery);
genai_filter_impl!(ErrorRateQuery);
genai_filter_impl!(ModelAnalyticsQuery);
genai_filter_impl!(CacheQuery);
genai_filter_impl!(ModelTimeSeriesQuery);
genai_filter_impl!(TopSpansQuery);
genai_filter_impl!(TopGroupQuery);
genai_filter_impl!(LatencyPercentileQuery);
genai_filter_impl!(DistributionQuery);
genai_filter_impl!(ToolUsageQuery);
genai_filter_impl!(RetryStatsQuery);
genai_filter_impl!(RetrievalStatsQuery);
genai_filter_impl!(TimeRangeQuery);
genai_filter_impl!(TimeSeriesQuery);
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn by_model_flag_accepts_one_and_true_only() {
assert!(by_model_enabled(Some("1")));
assert!(by_model_enabled(Some("true")));
assert!(!by_model_enabled(Some("0")));
assert!(!by_model_enabled(Some("yes")));
assert!(!by_model_enabled(Some("2")));
assert!(!by_model_enabled(None));
}
}