mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
fix(admin): 修复 usage 观测全表扫描并下推聚合查询
- 收紧 admin usage/stats 默认时间范围和 active 轮询查询 - 将 summary/records/aggregation/time-series/leaderboard 下推到 SQL 侧 - 统一前端时间参数并补齐 admin usage/stats 回归测试
This commit is contained in:
@@ -2,9 +2,13 @@ use std::collections::BTreeMap;
|
||||
use std::sync::RwLock;
|
||||
|
||||
use aether_data_contracts::repository::usage::{
|
||||
parse_usage_body_ref, usage_body_ref, UsageBodyField,
|
||||
parse_usage_body_ref, usage_body_ref, StoredUsageAuditAggregation, StoredUsageAuditSummary,
|
||||
StoredUsageLeaderboardSummary, StoredUsageTimeSeriesBucket, UsageAuditAggregationGroupBy,
|
||||
UsageAuditAggregationQuery, UsageAuditSummaryQuery, UsageBodyField, UsageLeaderboardGroupBy,
|
||||
UsageLeaderboardQuery, UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use chrono::Utc;
|
||||
use serde_json::Value;
|
||||
|
||||
use super::{
|
||||
@@ -109,6 +113,309 @@ fn usage_status_is_lifecycle(status: &str) -> bool {
|
||||
matches!(status, "pending" | "streaming")
|
||||
}
|
||||
|
||||
fn usage_matches_list_query(item: &StoredRequestUsageAudit, query: &UsageAuditListQuery) -> bool {
|
||||
// The field is historically named `created_at_unix_ms`, but usage audit rows
|
||||
// across gateway handlers, SQL repositories and tests are stored as epoch seconds.
|
||||
if let Some(created_from_unix_secs) = query.created_from_unix_secs {
|
||||
if item.created_at_unix_ms < created_from_unix_secs {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(created_until_unix_secs) = query.created_until_unix_secs {
|
||||
if item.created_at_unix_ms >= created_until_unix_secs {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
if item.user_id.as_deref() != Some(user_id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
if item.provider_name != provider_name {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
if item.model != model {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(api_format) = query.api_format.as_deref() {
|
||||
if item.api_format.as_deref() != Some(api_format) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(statuses) = query.statuses.as_ref() {
|
||||
if !statuses.iter().any(|status| status == &item.status) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(is_stream) = query.is_stream {
|
||||
if item.is_stream != is_stream {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if query.error_only
|
||||
&& item.status != "failed"
|
||||
&& item.status_code.unwrap_or_default() < 400
|
||||
&& item
|
||||
.error_message
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.unwrap_or_default()
|
||||
.is_empty()
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
true
|
||||
}
|
||||
|
||||
fn usage_matches_summary_query(
|
||||
item: &StoredRequestUsageAudit,
|
||||
query: &UsageAuditSummaryQuery,
|
||||
) -> bool {
|
||||
if item.created_at_unix_ms < query.created_from_unix_secs
|
||||
|| item.created_at_unix_ms >= query.created_until_unix_secs
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
if item.user_id.as_deref() != Some(user_id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
if item.provider_name != provider_name {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
if item.model != model {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
fn usage_matches_time_series_query(
|
||||
item: &StoredRequestUsageAudit,
|
||||
query: &UsageTimeSeriesQuery,
|
||||
) -> bool {
|
||||
if item.created_at_unix_ms < query.created_from_unix_secs
|
||||
|| item.created_at_unix_ms >= query.created_until_unix_secs
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
if item.user_id.as_deref() != Some(user_id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
if item.provider_name != provider_name {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
if item.model != model {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
fn usage_time_series_bucket_key(
|
||||
item: &StoredRequestUsageAudit,
|
||||
granularity: UsageTimeSeriesGranularity,
|
||||
tz_offset_minutes: i32,
|
||||
) -> Option<String> {
|
||||
let timestamp =
|
||||
chrono::DateTime::<Utc>::from_timestamp(i64::try_from(item.created_at_unix_ms).ok()?, 0)?;
|
||||
let local =
|
||||
timestamp.checked_add_signed(chrono::Duration::minutes(i64::from(tz_offset_minutes)))?;
|
||||
Some(match granularity {
|
||||
UsageTimeSeriesGranularity::Day => local.date_naive().to_string(),
|
||||
UsageTimeSeriesGranularity::Hour => local.format("%Y-%m-%dT%H:00:00+00:00").to_string(),
|
||||
})
|
||||
}
|
||||
|
||||
fn usage_matches_leaderboard_query(
|
||||
item: &StoredRequestUsageAudit,
|
||||
query: &UsageLeaderboardQuery,
|
||||
) -> bool {
|
||||
if item.created_at_unix_ms < query.created_from_unix_secs
|
||||
|| item.created_at_unix_ms >= query.created_until_unix_secs
|
||||
|| matches!(item.status.as_str(), "pending" | "streaming")
|
||||
|| matches!(item.provider_name.as_str(), "unknown" | "pending")
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
if item.user_id.as_deref() != Some(user_id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
if item.provider_name != provider_name {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
if item.model != model {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
fn sort_usage_items(items: &mut [StoredRequestUsageAudit], newest_first: bool) {
|
||||
items.sort_by(|left, right| {
|
||||
let created_order = if newest_first {
|
||||
right.created_at_unix_ms.cmp(&left.created_at_unix_ms)
|
||||
} else {
|
||||
left.created_at_unix_ms.cmp(&right.created_at_unix_ms)
|
||||
};
|
||||
if newest_first {
|
||||
created_order.then_with(|| left.id.cmp(&right.id))
|
||||
} else {
|
||||
created_order.then_with(|| left.request_id.cmp(&right.request_id))
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
fn usage_cache_creation_tokens(item: &StoredRequestUsageAudit) -> u64 {
|
||||
let classified = item
|
||||
.cache_creation_ephemeral_5m_input_tokens
|
||||
.saturating_add(item.cache_creation_ephemeral_1h_input_tokens);
|
||||
if item.cache_creation_input_tokens == 0 && classified > 0 {
|
||||
classified
|
||||
} else {
|
||||
item.cache_creation_input_tokens
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
||||
enum UsageApiFamily {
|
||||
OpenAi,
|
||||
Claude,
|
||||
Gemini,
|
||||
Unknown,
|
||||
}
|
||||
|
||||
fn usage_api_family(api_format: Option<&str>) -> UsageApiFamily {
|
||||
let Some(api_format) = api_format else {
|
||||
return UsageApiFamily::Unknown;
|
||||
};
|
||||
let family = api_format
|
||||
.split(':')
|
||||
.next()
|
||||
.unwrap_or_default()
|
||||
.trim()
|
||||
.to_ascii_lowercase();
|
||||
match family.as_str() {
|
||||
"openai" => UsageApiFamily::OpenAi,
|
||||
"claude" | "anthropic" => UsageApiFamily::Claude,
|
||||
"gemini" | "google" => UsageApiFamily::Gemini,
|
||||
_ => UsageApiFamily::Unknown,
|
||||
}
|
||||
}
|
||||
|
||||
fn normalize_usage_input_tokens(
|
||||
api_format: Option<&str>,
|
||||
input_tokens: i64,
|
||||
cache_read_tokens: i64,
|
||||
) -> i64 {
|
||||
if input_tokens <= 0 {
|
||||
return input_tokens.max(0);
|
||||
}
|
||||
if cache_read_tokens <= 0 {
|
||||
return input_tokens;
|
||||
}
|
||||
|
||||
match usage_api_family(api_format) {
|
||||
UsageApiFamily::OpenAi | UsageApiFamily::Gemini => {
|
||||
(input_tokens - cache_read_tokens).max(0)
|
||||
}
|
||||
UsageApiFamily::Claude | UsageApiFamily::Unknown => input_tokens,
|
||||
}
|
||||
}
|
||||
|
||||
fn normalize_usage_total_input_context(
|
||||
api_format: Option<&str>,
|
||||
input_tokens: i64,
|
||||
cache_creation_tokens: i64,
|
||||
cache_read_tokens: i64,
|
||||
) -> i64 {
|
||||
let normalized_input_tokens = input_tokens.max(0);
|
||||
let normalized_cache_creation_tokens = cache_creation_tokens.max(0);
|
||||
let normalized_cache_read_tokens = cache_read_tokens.max(0);
|
||||
|
||||
let fresh_input_tokens = match usage_api_family(api_format) {
|
||||
UsageApiFamily::Claude => {
|
||||
normalized_input_tokens.saturating_add(normalized_cache_creation_tokens)
|
||||
}
|
||||
UsageApiFamily::OpenAi | UsageApiFamily::Gemini => normalize_usage_input_tokens(
|
||||
api_format,
|
||||
normalized_input_tokens,
|
||||
normalized_cache_read_tokens,
|
||||
),
|
||||
UsageApiFamily::Unknown => {
|
||||
if normalized_cache_creation_tokens > 0 {
|
||||
normalized_input_tokens.saturating_add(normalized_cache_creation_tokens)
|
||||
} else {
|
||||
normalized_input_tokens
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
fresh_input_tokens.saturating_add(normalized_cache_read_tokens)
|
||||
}
|
||||
|
||||
fn usage_total_input_context(item: &StoredRequestUsageAudit) -> u64 {
|
||||
let api_format = item
|
||||
.endpoint_api_format
|
||||
.as_deref()
|
||||
.or(item.api_format.as_deref());
|
||||
let input_tokens = i64::try_from(item.input_tokens).unwrap_or(i64::MAX);
|
||||
let cache_creation_tokens =
|
||||
i64::try_from(usage_cache_creation_tokens(item)).unwrap_or(i64::MAX);
|
||||
let cache_read_tokens = i64::try_from(item.cache_read_input_tokens).unwrap_or(i64::MAX);
|
||||
normalize_usage_total_input_context(
|
||||
api_format,
|
||||
input_tokens,
|
||||
cache_creation_tokens,
|
||||
cache_read_tokens,
|
||||
) as u64
|
||||
}
|
||||
|
||||
fn usage_effective_input_tokens(item: &StoredRequestUsageAudit) -> u64 {
|
||||
let api_format = item
|
||||
.endpoint_api_format
|
||||
.as_deref()
|
||||
.or(item.api_format.as_deref());
|
||||
let input_tokens = i64::try_from(item.input_tokens).unwrap_or(i64::MAX);
|
||||
let cache_read_tokens = i64::try_from(item.cache_read_input_tokens).unwrap_or(i64::MAX);
|
||||
normalize_usage_input_tokens(api_format, input_tokens, cache_read_tokens) as u64
|
||||
}
|
||||
|
||||
fn usage_is_success(item: &StoredRequestUsageAudit) -> bool {
|
||||
matches!(
|
||||
item.status.as_str(),
|
||||
"completed" | "success" | "ok" | "billed" | "settled"
|
||||
) && item.status_code.is_none_or(|code| code < 400)
|
||||
}
|
||||
|
||||
fn usage_provider_display_name(item: &StoredRequestUsageAudit) -> Option<String> {
|
||||
let provider_name = item.provider_name.trim();
|
||||
if provider_name.is_empty() || matches!(provider_name, "unknown" | "pending") {
|
||||
None
|
||||
} else {
|
||||
Some(item.provider_name.clone())
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl UsageReadRepository for InMemoryUsageReadRepository {
|
||||
async fn find_by_id(
|
||||
@@ -172,54 +479,302 @@ impl UsageReadRepository for InMemoryUsageReadRepository {
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
.filter(|item| {
|
||||
// The field is historically named `created_at_unix_ms`, but usage audit rows
|
||||
// across gateway handlers, SQL repositories and tests are stored as epoch seconds.
|
||||
if let Some(created_from_unix_secs) = query.created_from_unix_secs {
|
||||
if item.created_at_unix_ms < created_from_unix_secs {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(created_until_unix_secs) = query.created_until_unix_secs {
|
||||
if item.created_at_unix_ms >= created_until_unix_secs {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
if item.user_id.as_deref() != Some(user_id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
if item.provider_name != provider_name {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
if item.model != model {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
if let Some(statuses) = query.statuses.as_ref() {
|
||||
if !statuses.iter().any(|s| s == &item.status) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
true
|
||||
})
|
||||
.filter(|item| usage_matches_list_query(item, query))
|
||||
.cloned()
|
||||
.collect();
|
||||
items.sort_by(|left, right| {
|
||||
left.created_at_unix_ms
|
||||
.cmp(&right.created_at_unix_ms)
|
||||
.then_with(|| left.request_id.cmp(&right.request_id))
|
||||
});
|
||||
sort_usage_items(&mut items, query.newest_first);
|
||||
if let Some(offset) = query.offset {
|
||||
if offset >= items.len() {
|
||||
items.clear();
|
||||
} else {
|
||||
items.drain(..offset);
|
||||
}
|
||||
}
|
||||
if let Some(limit) = query.limit {
|
||||
items.truncate(limit);
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
async fn count_usage_audits(&self, query: &UsageAuditListQuery) -> Result<u64, DataLayerError> {
|
||||
Ok(self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
.filter(|item| usage_matches_list_query(item, query))
|
||||
.count() as u64)
|
||||
}
|
||||
|
||||
async fn aggregate_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditAggregationQuery,
|
||||
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
|
||||
#[derive(Default)]
|
||||
struct AggregateBucket {
|
||||
display_name: Option<String>,
|
||||
secondary_name: Option<String>,
|
||||
request_count: u64,
|
||||
total_tokens: u64,
|
||||
output_tokens: u64,
|
||||
effective_input_tokens: u64,
|
||||
total_input_context: u64,
|
||||
cache_creation_tokens: u64,
|
||||
cache_creation_ephemeral_5m_tokens: u64,
|
||||
cache_creation_ephemeral_1h_tokens: u64,
|
||||
cache_read_tokens: u64,
|
||||
total_cost_usd: f64,
|
||||
actual_total_cost_usd: f64,
|
||||
response_time_ms_sum: u64,
|
||||
success_count: u64,
|
||||
}
|
||||
|
||||
let mut grouped: BTreeMap<String, AggregateBucket> = BTreeMap::new();
|
||||
for item in self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
{
|
||||
if item.created_at_unix_ms < query.created_from_unix_secs
|
||||
|| item.created_at_unix_ms >= query.created_until_unix_secs
|
||||
|| matches!(item.status.as_str(), "pending" | "streaming")
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
let group_key = match query.group_by {
|
||||
UsageAuditAggregationGroupBy::Model => item.model.clone(),
|
||||
UsageAuditAggregationGroupBy::Provider => item
|
||||
.provider_id
|
||||
.clone()
|
||||
.unwrap_or_else(|| "unknown".to_string()),
|
||||
UsageAuditAggregationGroupBy::ApiFormat => item
|
||||
.api_format
|
||||
.clone()
|
||||
.unwrap_or_else(|| "unknown".to_string()),
|
||||
UsageAuditAggregationGroupBy::User => match item.user_id.clone() {
|
||||
Some(value) => value,
|
||||
None => continue,
|
||||
},
|
||||
};
|
||||
let bucket = grouped.entry(group_key).or_default();
|
||||
if matches!(query.group_by, UsageAuditAggregationGroupBy::Provider)
|
||||
&& (bucket.display_name.is_none()
|
||||
|| bucket.display_name.as_deref() == Some("Unknown"))
|
||||
{
|
||||
bucket.display_name =
|
||||
usage_provider_display_name(item).or(Some("Unknown".to_string()));
|
||||
}
|
||||
bucket.request_count = bucket.request_count.saturating_add(1);
|
||||
bucket.total_tokens = bucket.total_tokens.saturating_add(item.total_tokens);
|
||||
bucket.output_tokens = bucket.output_tokens.saturating_add(item.output_tokens);
|
||||
bucket.effective_input_tokens = bucket
|
||||
.effective_input_tokens
|
||||
.saturating_add(usage_effective_input_tokens(item));
|
||||
bucket.total_input_context = bucket
|
||||
.total_input_context
|
||||
.saturating_add(usage_total_input_context(item));
|
||||
bucket.cache_creation_tokens = bucket
|
||||
.cache_creation_tokens
|
||||
.saturating_add(usage_cache_creation_tokens(item));
|
||||
bucket.cache_creation_ephemeral_5m_tokens = bucket
|
||||
.cache_creation_ephemeral_5m_tokens
|
||||
.saturating_add(item.cache_creation_ephemeral_5m_input_tokens);
|
||||
bucket.cache_creation_ephemeral_1h_tokens = bucket
|
||||
.cache_creation_ephemeral_1h_tokens
|
||||
.saturating_add(item.cache_creation_ephemeral_1h_input_tokens);
|
||||
bucket.cache_read_tokens = bucket
|
||||
.cache_read_tokens
|
||||
.saturating_add(item.cache_read_input_tokens);
|
||||
bucket.total_cost_usd += item.total_cost_usd;
|
||||
bucket.actual_total_cost_usd += item.actual_total_cost_usd;
|
||||
bucket.response_time_ms_sum = bucket
|
||||
.response_time_ms_sum
|
||||
.saturating_add(item.response_time_ms.unwrap_or_default());
|
||||
bucket.success_count = bucket
|
||||
.success_count
|
||||
.saturating_add(if usage_is_success(item) { 1 } else { 0 });
|
||||
}
|
||||
|
||||
let mut items = grouped
|
||||
.into_iter()
|
||||
.map(|(group_key, bucket)| StoredUsageAuditAggregation {
|
||||
group_key,
|
||||
display_name: bucket.display_name,
|
||||
secondary_name: bucket.secondary_name,
|
||||
request_count: bucket.request_count,
|
||||
total_tokens: bucket.total_tokens,
|
||||
output_tokens: bucket.output_tokens,
|
||||
effective_input_tokens: bucket.effective_input_tokens,
|
||||
total_input_context: bucket.total_input_context,
|
||||
cache_creation_tokens: bucket.cache_creation_tokens,
|
||||
cache_creation_ephemeral_5m_tokens: bucket.cache_creation_ephemeral_5m_tokens,
|
||||
cache_creation_ephemeral_1h_tokens: bucket.cache_creation_ephemeral_1h_tokens,
|
||||
cache_read_tokens: bucket.cache_read_tokens,
|
||||
total_cost_usd: bucket.total_cost_usd,
|
||||
actual_total_cost_usd: bucket.actual_total_cost_usd,
|
||||
avg_response_time_ms: match query.group_by {
|
||||
UsageAuditAggregationGroupBy::Provider
|
||||
| UsageAuditAggregationGroupBy::ApiFormat => {
|
||||
Some(if bucket.request_count == 0 {
|
||||
0.0
|
||||
} else {
|
||||
bucket.response_time_ms_sum as f64 / bucket.request_count as f64
|
||||
})
|
||||
}
|
||||
_ => None,
|
||||
},
|
||||
success_count: match query.group_by {
|
||||
UsageAuditAggregationGroupBy::Provider => Some(bucket.success_count),
|
||||
_ => None,
|
||||
},
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
items.sort_by(|left, right| {
|
||||
right
|
||||
.request_count
|
||||
.cmp(&left.request_count)
|
||||
.then_with(|| left.group_key.cmp(&right.group_key))
|
||||
});
|
||||
items.truncate(query.limit);
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
async fn summarize_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditSummaryQuery,
|
||||
) -> Result<StoredUsageAuditSummary, DataLayerError> {
|
||||
let mut summary = StoredUsageAuditSummary::default();
|
||||
for item in self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
{
|
||||
if !usage_matches_summary_query(item, query) {
|
||||
continue;
|
||||
}
|
||||
summary.total_requests = summary.total_requests.saturating_add(1);
|
||||
summary.input_tokens = summary.input_tokens.saturating_add(item.input_tokens);
|
||||
summary.output_tokens = summary.output_tokens.saturating_add(item.output_tokens);
|
||||
summary.recorded_total_tokens = summary
|
||||
.recorded_total_tokens
|
||||
.saturating_add(item.total_tokens);
|
||||
summary.cache_creation_tokens = summary
|
||||
.cache_creation_tokens
|
||||
.saturating_add(usage_cache_creation_tokens(item));
|
||||
summary.cache_creation_ephemeral_5m_tokens = summary
|
||||
.cache_creation_ephemeral_5m_tokens
|
||||
.saturating_add(item.cache_creation_ephemeral_5m_input_tokens);
|
||||
summary.cache_creation_ephemeral_1h_tokens = summary
|
||||
.cache_creation_ephemeral_1h_tokens
|
||||
.saturating_add(item.cache_creation_ephemeral_1h_input_tokens);
|
||||
summary.cache_read_tokens = summary
|
||||
.cache_read_tokens
|
||||
.saturating_add(item.cache_read_input_tokens);
|
||||
summary.total_cost_usd += item.total_cost_usd;
|
||||
summary.actual_total_cost_usd += item.actual_total_cost_usd;
|
||||
summary.cache_creation_cost_usd += item.cache_creation_cost_usd;
|
||||
summary.cache_read_cost_usd += item.cache_read_cost_usd;
|
||||
summary.total_response_time_ms += item.response_time_ms.unwrap_or(0) as f64;
|
||||
if item.status_code.is_some_and(|value| value >= 400) || item.error_message.is_some() {
|
||||
summary.error_requests = summary.error_requests.saturating_add(1);
|
||||
}
|
||||
}
|
||||
Ok(summary)
|
||||
}
|
||||
|
||||
async fn summarize_usage_time_series(
|
||||
&self,
|
||||
query: &UsageTimeSeriesQuery,
|
||||
) -> Result<Vec<StoredUsageTimeSeriesBucket>, DataLayerError> {
|
||||
let mut buckets = BTreeMap::<String, StoredUsageTimeSeriesBucket>::new();
|
||||
for item in self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
{
|
||||
if !usage_matches_time_series_query(item, query) {
|
||||
continue;
|
||||
}
|
||||
let Some(bucket_key) =
|
||||
usage_time_series_bucket_key(item, query.granularity, query.tz_offset_minutes)
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
let bucket =
|
||||
buckets
|
||||
.entry(bucket_key.clone())
|
||||
.or_insert_with(|| StoredUsageTimeSeriesBucket {
|
||||
bucket_key,
|
||||
..Default::default()
|
||||
});
|
||||
bucket.total_requests = bucket.total_requests.saturating_add(1);
|
||||
bucket.input_tokens = bucket.input_tokens.saturating_add(item.input_tokens);
|
||||
bucket.output_tokens = bucket.output_tokens.saturating_add(item.output_tokens);
|
||||
bucket.cache_creation_tokens = bucket
|
||||
.cache_creation_tokens
|
||||
.saturating_add(item.cache_creation_input_tokens);
|
||||
bucket.cache_read_tokens = bucket
|
||||
.cache_read_tokens
|
||||
.saturating_add(item.cache_read_input_tokens);
|
||||
bucket.total_cost_usd += item.total_cost_usd;
|
||||
bucket.total_response_time_ms += item.response_time_ms.unwrap_or(0) as f64;
|
||||
}
|
||||
Ok(buckets.into_values().collect())
|
||||
}
|
||||
|
||||
async fn summarize_usage_leaderboard(
|
||||
&self,
|
||||
query: &UsageLeaderboardQuery,
|
||||
) -> Result<Vec<StoredUsageLeaderboardSummary>, DataLayerError> {
|
||||
let mut grouped = BTreeMap::<String, StoredUsageLeaderboardSummary>::new();
|
||||
for item in self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
{
|
||||
if !usage_matches_leaderboard_query(item, query) {
|
||||
continue;
|
||||
}
|
||||
let (group_key, legacy_name) = match query.group_by {
|
||||
UsageLeaderboardGroupBy::Model => (item.model.clone(), None),
|
||||
UsageLeaderboardGroupBy::User => match item.user_id.clone() {
|
||||
Some(user_id) => (user_id, item.username.clone()),
|
||||
None => continue,
|
||||
},
|
||||
UsageLeaderboardGroupBy::ApiKey => match item.api_key_id.clone() {
|
||||
Some(api_key_id) => (api_key_id, item.api_key_name.clone()),
|
||||
None => continue,
|
||||
},
|
||||
};
|
||||
let entry =
|
||||
grouped
|
||||
.entry(group_key.clone())
|
||||
.or_insert_with(|| StoredUsageLeaderboardSummary {
|
||||
group_key,
|
||||
legacy_name: legacy_name.clone(),
|
||||
..Default::default()
|
||||
});
|
||||
if entry.legacy_name.is_none() {
|
||||
entry.legacy_name = legacy_name;
|
||||
}
|
||||
entry.request_count = entry.request_count.saturating_add(1);
|
||||
entry.total_tokens = entry.total_tokens.saturating_add(
|
||||
item.input_tokens
|
||||
.saturating_add(item.output_tokens)
|
||||
.saturating_add(item.cache_creation_input_tokens)
|
||||
.saturating_add(item.cache_read_input_tokens),
|
||||
);
|
||||
entry.total_cost_usd += item.total_cost_usd;
|
||||
}
|
||||
Ok(grouped.into_values().collect())
|
||||
}
|
||||
|
||||
async fn list_recent_usage_audits(
|
||||
&self,
|
||||
user_id: Option<&str>,
|
||||
|
||||
@@ -4,8 +4,12 @@ mod sql;
|
||||
#[allow(unused_imports)]
|
||||
pub(crate) use aether_data_contracts::repository::usage::{
|
||||
StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
|
||||
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
|
||||
UsageDailyHeatmapQuery, UsageReadRepository, UsageRepository, UsageWriteRepository,
|
||||
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
|
||||
StoredUsageDailySummary, StoredUsageLeaderboardSummary, StoredUsageTimeSeriesBucket,
|
||||
UpsertUsageRecord, UsageAuditAggregationGroupBy, UsageAuditAggregationQuery,
|
||||
UsageAuditListQuery, UsageAuditSummaryQuery, UsageDailyHeatmapQuery, UsageLeaderboardGroupBy,
|
||||
UsageLeaderboardQuery, UsageReadRepository, UsageRepository, UsageTimeSeriesGranularity,
|
||||
UsageTimeSeriesQuery, UsageWriteRepository,
|
||||
};
|
||||
pub use memory::InMemoryUsageReadRepository;
|
||||
pub use sql::SqlxUsageReadRepository;
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
use aether_data_contracts::repository::usage::{
|
||||
parse_usage_body_ref, usage_body_ref, UsageBodyField,
|
||||
parse_usage_body_ref, usage_body_ref, StoredUsageAuditAggregation, StoredUsageAuditSummary,
|
||||
StoredUsageLeaderboardSummary, StoredUsageTimeSeriesBucket, UsageAuditAggregationGroupBy,
|
||||
UsageAuditAggregationQuery, UsageAuditSummaryQuery, UsageBodyField, UsageLeaderboardGroupBy,
|
||||
UsageLeaderboardQuery, UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use flate2::{read::GzDecoder, write::GzEncoder, Compression};
|
||||
@@ -620,6 +623,93 @@ LEFT JOIN usage_settlement_snapshots
|
||||
ON usage_settlement_snapshots.request_id = "usage".request_id
|
||||
"#;
|
||||
|
||||
struct UsageAuditAggregationSqlFragments {
|
||||
filtered_extra_where: &'static str,
|
||||
group_key_expr: &'static str,
|
||||
display_name_expr: &'static str,
|
||||
secondary_name_expr: &'static str,
|
||||
aggregate_display_name_expr: &'static str,
|
||||
aggregate_secondary_name_expr: &'static str,
|
||||
avg_response_time_expr: &'static str,
|
||||
success_count_expr: &'static str,
|
||||
}
|
||||
|
||||
fn usage_audit_aggregation_sql_fragments(
|
||||
group_by: UsageAuditAggregationGroupBy,
|
||||
) -> UsageAuditAggregationSqlFragments {
|
||||
match group_by {
|
||||
UsageAuditAggregationGroupBy::Model => UsageAuditAggregationSqlFragments {
|
||||
filtered_extra_where: "",
|
||||
group_key_expr: "model",
|
||||
display_name_expr: "NULL::varchar",
|
||||
secondary_name_expr: "NULL::varchar",
|
||||
aggregate_display_name_expr: "NULL::varchar",
|
||||
aggregate_secondary_name_expr: "NULL::varchar",
|
||||
avg_response_time_expr: "NULL::DOUBLE PRECISION",
|
||||
success_count_expr: "NULL::BIGINT",
|
||||
},
|
||||
UsageAuditAggregationGroupBy::Provider => UsageAuditAggregationSqlFragments {
|
||||
filtered_extra_where: "",
|
||||
group_key_expr: "provider_group_key",
|
||||
display_name_expr: "provider_display_name",
|
||||
secondary_name_expr: "NULL::varchar",
|
||||
aggregate_display_name_expr:
|
||||
"COALESCE(MAX(NULLIF(display_name, 'Unknown')), 'Unknown')",
|
||||
aggregate_secondary_name_expr: "NULL::varchar",
|
||||
avg_response_time_expr: "AVG(response_time_ms::DOUBLE PRECISION)",
|
||||
success_count_expr: "COALESCE(SUM(success_flag), 0)::BIGINT",
|
||||
},
|
||||
UsageAuditAggregationGroupBy::ApiFormat => UsageAuditAggregationSqlFragments {
|
||||
filtered_extra_where: "",
|
||||
group_key_expr: "api_format_group_key",
|
||||
display_name_expr: "NULL::varchar",
|
||||
secondary_name_expr: "NULL::varchar",
|
||||
aggregate_display_name_expr: "NULL::varchar",
|
||||
aggregate_secondary_name_expr: "NULL::varchar",
|
||||
avg_response_time_expr: "AVG(response_time_ms::DOUBLE PRECISION)",
|
||||
success_count_expr: "NULL::BIGINT",
|
||||
},
|
||||
UsageAuditAggregationGroupBy::User => UsageAuditAggregationSqlFragments {
|
||||
filtered_extra_where: " AND \"usage\".user_id IS NOT NULL",
|
||||
group_key_expr: "user_id",
|
||||
display_name_expr: "NULL::varchar",
|
||||
secondary_name_expr: "NULL::varchar",
|
||||
aggregate_display_name_expr: "NULL::varchar",
|
||||
aggregate_secondary_name_expr: "NULL::varchar",
|
||||
avg_response_time_expr: "NULL::DOUBLE PRECISION",
|
||||
success_count_expr: "NULL::BIGINT",
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
struct UsageLeaderboardSqlFragments {
|
||||
filtered_extra_where: &'static str,
|
||||
group_key_expr: &'static str,
|
||||
legacy_name_expr: &'static str,
|
||||
}
|
||||
|
||||
fn usage_leaderboard_sql_fragments(
|
||||
group_by: UsageLeaderboardGroupBy,
|
||||
) -> UsageLeaderboardSqlFragments {
|
||||
match group_by {
|
||||
UsageLeaderboardGroupBy::Model => UsageLeaderboardSqlFragments {
|
||||
filtered_extra_where: "",
|
||||
group_key_expr: "\"usage\".model",
|
||||
legacy_name_expr: "NULL::varchar",
|
||||
},
|
||||
UsageLeaderboardGroupBy::User => UsageLeaderboardSqlFragments {
|
||||
filtered_extra_where: " AND \"usage\".user_id IS NOT NULL",
|
||||
group_key_expr: "\"usage\".user_id",
|
||||
legacy_name_expr: "NULLIF(BTRIM(\"usage\".username), '')",
|
||||
},
|
||||
UsageLeaderboardGroupBy::ApiKey => UsageLeaderboardSqlFragments {
|
||||
filtered_extra_where: " AND \"usage\".api_key_id IS NOT NULL",
|
||||
group_key_expr: "\"usage\".api_key_id",
|
||||
legacy_name_expr: "NULLIF(BTRIM(\"usage\".api_key_name), '')",
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
const LIST_RECENT_USAGE_AUDITS_PREFIX: &str = r#"
|
||||
SELECT
|
||||
"usage".id,
|
||||
@@ -1206,9 +1296,17 @@ impl SqlxUsageReadRepository {
|
||||
.push("\"usage\".model = ")
|
||||
.push_bind(model.to_string());
|
||||
}
|
||||
if let Some(api_format) = query.api_format.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".api_format = ")
|
||||
.push_bind(api_format.to_string());
|
||||
}
|
||||
if let Some(statuses) = query.statuses.as_deref() {
|
||||
if !statuses.is_empty() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder.push("\"usage\".status IN (");
|
||||
let mut separated = builder.separated(", ");
|
||||
for status in statuses {
|
||||
@@ -1217,11 +1315,31 @@ impl SqlxUsageReadRepository {
|
||||
separated.push_unseparated(")");
|
||||
}
|
||||
}
|
||||
if let Some(is_stream) = query.is_stream {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder.push("\"usage\".is_stream = ").push_bind(is_stream);
|
||||
}
|
||||
if query.error_only {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder.push(
|
||||
"(\"usage\".status = 'failed' \
|
||||
OR COALESCE(\"usage\".status_code, 0) >= 400 \
|
||||
OR (\"usage\".error_message IS NOT NULL AND BTRIM(\"usage\".error_message) <> ''))",
|
||||
);
|
||||
}
|
||||
|
||||
builder.push(" ORDER BY \"usage\".created_at ASC, \"usage\".request_id ASC");
|
||||
if query.newest_first {
|
||||
builder.push(" ORDER BY \"usage\".created_at DESC, \"usage\".id ASC");
|
||||
} else {
|
||||
builder.push(" ORDER BY \"usage\".created_at ASC, \"usage\".request_id ASC");
|
||||
}
|
||||
if let Some(limit) = query.limit {
|
||||
builder.push(" LIMIT ").push_bind(limit as i64);
|
||||
}
|
||||
if let Some(offset) = query.offset {
|
||||
builder.push(" OFFSET ").push_bind(offset as i64);
|
||||
}
|
||||
let query = builder.build();
|
||||
let mut rows = query.fetch(&self.pool);
|
||||
let mut items = Vec::new();
|
||||
@@ -1231,6 +1349,626 @@ impl SqlxUsageReadRepository {
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
pub async fn count_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditListQuery,
|
||||
) -> Result<u64, DataLayerError> {
|
||||
let mut builder =
|
||||
QueryBuilder::<Postgres>::new(r#"SELECT COUNT(*)::BIGINT AS total FROM "usage""#);
|
||||
let mut has_where = false;
|
||||
|
||||
if let Some(created_from_unix_secs) = query.created_from_unix_secs {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".created_at >= TO_TIMESTAMP(")
|
||||
.push_bind(created_from_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
}
|
||||
if let Some(created_until_unix_secs) = query.created_until_unix_secs {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".created_at < TO_TIMESTAMP(")
|
||||
.push_bind(created_until_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
}
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".user_id = ")
|
||||
.push_bind(user_id.to_string());
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".provider_name = ")
|
||||
.push_bind(provider_name.to_string());
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".model = ")
|
||||
.push_bind(model.to_string());
|
||||
}
|
||||
if let Some(api_format) = query.api_format.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".api_format = ")
|
||||
.push_bind(api_format.to_string());
|
||||
}
|
||||
if let Some(statuses) = query.statuses.as_deref() {
|
||||
if !statuses.is_empty() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder.push("\"usage\".status IN (");
|
||||
let mut separated = builder.separated(", ");
|
||||
for status in statuses {
|
||||
separated.push_bind(status.to_string());
|
||||
}
|
||||
separated.push_unseparated(")");
|
||||
}
|
||||
}
|
||||
if let Some(is_stream) = query.is_stream {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder.push("\"usage\".is_stream = ").push_bind(is_stream);
|
||||
}
|
||||
if query.error_only {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder.push(
|
||||
"(\"usage\".status = 'failed' \
|
||||
OR COALESCE(\"usage\".status_code, 0) >= 400 \
|
||||
OR (\"usage\".error_message IS NOT NULL AND BTRIM(\"usage\".error_message) <> ''))",
|
||||
);
|
||||
}
|
||||
|
||||
let row = builder
|
||||
.build()
|
||||
.fetch_one(&self.pool)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
Ok(row.try_get::<i64, _>("total").map_postgres_err()?.max(0) as u64)
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditSummaryQuery,
|
||||
) -> Result<StoredUsageAuditSummary, DataLayerError> {
|
||||
let mut builder = QueryBuilder::<Postgres>::new(
|
||||
r#"
|
||||
SELECT
|
||||
COUNT(*)::BIGINT AS total_requests,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".input_tokens, 0), 0)), 0)::BIGINT AS input_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".output_tokens, 0), 0)), 0)::BIGINT AS output_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".total_tokens, 0), 0)), 0)::BIGINT AS recorded_total_tokens,
|
||||
COALESCE(SUM(
|
||||
CASE
|
||||
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE("usage".cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE("usage".cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE("usage".cache_creation_input_tokens, 0)
|
||||
END
|
||||
), 0)::BIGINT AS cache_creation_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".cache_creation_input_tokens_5m, 0), 0)), 0)::BIGINT
|
||||
AS cache_creation_ephemeral_5m_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".cache_creation_input_tokens_1h, 0), 0)), 0)::BIGINT
|
||||
AS cache_creation_ephemeral_1h_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)), 0)::BIGINT
|
||||
AS cache_read_tokens,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS total_cost_usd,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".actual_total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS actual_total_cost_usd,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".cache_creation_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS cache_creation_cost_usd,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".cache_read_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS cache_read_cost_usd,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".response_time_ms, 0), 0)::DOUBLE PRECISION), 0)
|
||||
AS total_response_time_ms,
|
||||
COALESCE(SUM(
|
||||
CASE
|
||||
WHEN COALESCE("usage".status_code, 0) >= 400 OR "usage".error_message IS NOT NULL THEN 1
|
||||
ELSE 0
|
||||
END
|
||||
), 0)::BIGINT AS error_requests
|
||||
FROM "usage"
|
||||
"#,
|
||||
);
|
||||
let mut has_where = false;
|
||||
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".created_at >= TO_TIMESTAMP(")
|
||||
.push_bind(query.created_from_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder
|
||||
.push("\"usage\".created_at < TO_TIMESTAMP(")
|
||||
.push_bind(query.created_until_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".user_id = ")
|
||||
.push_bind(user_id.to_string());
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".provider_name = ")
|
||||
.push_bind(provider_name.to_string());
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder
|
||||
.push("\"usage\".model = ")
|
||||
.push_bind(model.to_string());
|
||||
}
|
||||
|
||||
let row = builder
|
||||
.build()
|
||||
.fetch_one(&self.pool)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
Ok(StoredUsageAuditSummary {
|
||||
total_requests: row
|
||||
.try_get::<i64, _>("total_requests")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
input_tokens: row
|
||||
.try_get::<i64, _>("input_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
output_tokens: row
|
||||
.try_get::<i64, _>("output_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
recorded_total_tokens: row
|
||||
.try_get::<i64, _>("recorded_total_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_ephemeral_5m_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_ephemeral_1h_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_read_tokens: row
|
||||
.try_get::<i64, _>("cache_read_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
|
||||
actual_total_cost_usd: row
|
||||
.try_get::<f64, _>("actual_total_cost_usd")
|
||||
.map_postgres_err()?,
|
||||
cache_creation_cost_usd: row
|
||||
.try_get::<f64, _>("cache_creation_cost_usd")
|
||||
.map_postgres_err()?,
|
||||
cache_read_cost_usd: row
|
||||
.try_get::<f64, _>("cache_read_cost_usd")
|
||||
.map_postgres_err()?,
|
||||
total_response_time_ms: row
|
||||
.try_get::<f64, _>("total_response_time_ms")
|
||||
.map_postgres_err()?,
|
||||
error_requests: row
|
||||
.try_get::<i64, _>("error_requests")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_time_series(
|
||||
&self,
|
||||
query: &UsageTimeSeriesQuery,
|
||||
) -> Result<Vec<StoredUsageTimeSeriesBucket>, DataLayerError> {
|
||||
let mut builder = QueryBuilder::<Postgres>::new("SELECT ");
|
||||
match query.granularity {
|
||||
UsageTimeSeriesGranularity::Day => {
|
||||
builder
|
||||
.push("TO_CHAR(date_trunc('day', \"usage\".created_at + (")
|
||||
.push_bind(query.tz_offset_minutes)
|
||||
.push("::integer * INTERVAL '1 minute')), 'YYYY-MM-DD') AS bucket_key");
|
||||
}
|
||||
UsageTimeSeriesGranularity::Hour => {
|
||||
builder
|
||||
.push("TO_CHAR(date_trunc('hour', \"usage\".created_at + (")
|
||||
.push_bind(query.tz_offset_minutes)
|
||||
.push("::integer * INTERVAL '1 minute')), 'YYYY-MM-DD\"T\"HH24:00:00+00:00') AS bucket_key");
|
||||
}
|
||||
}
|
||||
builder.push(
|
||||
r#",
|
||||
COUNT(*)::BIGINT AS total_requests,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".input_tokens, 0), 0)), 0)::BIGINT AS input_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".output_tokens, 0), 0)), 0)::BIGINT AS output_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".cache_creation_input_tokens, 0), 0)), 0)::BIGINT
|
||||
AS cache_creation_tokens,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)), 0)::BIGINT
|
||||
AS cache_read_tokens,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS total_cost_usd,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".response_time_ms, 0), 0)::DOUBLE PRECISION), 0)
|
||||
AS total_response_time_ms
|
||||
FROM "usage"
|
||||
"#,
|
||||
);
|
||||
let mut has_where = false;
|
||||
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".created_at >= TO_TIMESTAMP(")
|
||||
.push_bind(query.created_from_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder
|
||||
.push("\"usage\".created_at < TO_TIMESTAMP(")
|
||||
.push_bind(query.created_until_unix_secs as f64)
|
||||
.push("::double precision)");
|
||||
if let Some(user_id) = query.user_id.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".user_id = ")
|
||||
.push_bind(user_id.to_string());
|
||||
}
|
||||
if let Some(provider_name) = query.provider_name.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
has_where = true;
|
||||
builder
|
||||
.push("\"usage\".provider_name = ")
|
||||
.push_bind(provider_name.to_string());
|
||||
}
|
||||
if let Some(model) = query.model.as_deref() {
|
||||
builder.push(if has_where { " AND " } else { " WHERE " });
|
||||
builder
|
||||
.push("\"usage\".model = ")
|
||||
.push_bind(model.to_string());
|
||||
}
|
||||
builder.push(" GROUP BY bucket_key ORDER BY bucket_key ASC");
|
||||
|
||||
let mut rows = builder.build().fetch(&self.pool);
|
||||
let mut items = Vec::new();
|
||||
while let Some(row) = rows.try_next().await.map_postgres_err()? {
|
||||
items.push(StoredUsageTimeSeriesBucket {
|
||||
bucket_key: row.try_get::<String, _>("bucket_key").map_postgres_err()?,
|
||||
total_requests: row
|
||||
.try_get::<i64, _>("total_requests")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
input_tokens: row
|
||||
.try_get::<i64, _>("input_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
output_tokens: row
|
||||
.try_get::<i64, _>("output_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_read_tokens: row
|
||||
.try_get::<i64, _>("cache_read_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
|
||||
total_response_time_ms: row
|
||||
.try_get::<f64, _>("total_response_time_ms")
|
||||
.map_postgres_err()?,
|
||||
});
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_leaderboard(
|
||||
&self,
|
||||
query: &UsageLeaderboardQuery,
|
||||
) -> Result<Vec<StoredUsageLeaderboardSummary>, DataLayerError> {
|
||||
let fragments = usage_leaderboard_sql_fragments(query.group_by);
|
||||
let sql = format!(
|
||||
r#"
|
||||
SELECT
|
||||
{group_key_expr} AS group_key,
|
||||
MAX({legacy_name_expr}) AS legacy_name,
|
||||
COUNT(*)::BIGINT AS request_count,
|
||||
COALESCE(SUM(
|
||||
GREATEST(COALESCE("usage".input_tokens, 0), 0)
|
||||
+ GREATEST(COALESCE("usage".output_tokens, 0), 0)
|
||||
+ GREATEST(COALESCE("usage".cache_creation_input_tokens, 0), 0)
|
||||
+ GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)
|
||||
), 0)::BIGINT AS total_tokens,
|
||||
COALESCE(SUM(COALESCE(CAST("usage".total_cost_usd AS DOUBLE PRECISION), 0)), 0)
|
||||
AS total_cost_usd
|
||||
FROM "usage"
|
||||
WHERE "usage".created_at >= TO_TIMESTAMP($1::double precision)
|
||||
AND "usage".created_at < TO_TIMESTAMP($2::double precision)
|
||||
AND "usage".status NOT IN ('pending', 'streaming')
|
||||
AND "usage".provider_name NOT IN ('unknown', 'pending')
|
||||
{filtered_extra_where}
|
||||
AND ($3::varchar IS NULL OR "usage".user_id = $3)
|
||||
AND ($4::varchar IS NULL OR "usage".provider_name = $4)
|
||||
AND ($5::varchar IS NULL OR "usage".model = $5)
|
||||
GROUP BY group_key
|
||||
ORDER BY group_key ASC
|
||||
"#,
|
||||
group_key_expr = fragments.group_key_expr,
|
||||
legacy_name_expr = fragments.legacy_name_expr,
|
||||
filtered_extra_where = fragments.filtered_extra_where,
|
||||
);
|
||||
let mut rows = sqlx::query(&sql)
|
||||
.bind(query.created_from_unix_secs as f64)
|
||||
.bind(query.created_until_unix_secs as f64)
|
||||
.bind(query.user_id.as_deref())
|
||||
.bind(query.provider_name.as_deref())
|
||||
.bind(query.model.as_deref())
|
||||
.fetch(&self.pool);
|
||||
let mut items = Vec::new();
|
||||
while let Some(row) = rows.try_next().await.map_postgres_err()? {
|
||||
items.push(StoredUsageLeaderboardSummary {
|
||||
group_key: row.try_get::<String, _>("group_key").map_postgres_err()?,
|
||||
legacy_name: row
|
||||
.try_get::<Option<String>, _>("legacy_name")
|
||||
.map_postgres_err()?,
|
||||
request_count: row
|
||||
.try_get::<i64, _>("request_count")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_tokens: row
|
||||
.try_get::<i64, _>("total_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
|
||||
});
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
pub async fn aggregate_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditAggregationQuery,
|
||||
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
|
||||
let fragments = usage_audit_aggregation_sql_fragments(query.group_by);
|
||||
let sql = format!(
|
||||
r#"
|
||||
WITH filtered_usage AS (
|
||||
SELECT
|
||||
"usage".model AS model,
|
||||
"usage".user_id AS user_id,
|
||||
COALESCE("usage".provider_id, 'unknown') AS provider_group_key,
|
||||
CASE
|
||||
WHEN BTRIM(COALESCE("usage".provider_name, '')) = ''
|
||||
OR "usage".provider_name IN ('unknown', 'pending')
|
||||
THEN 'Unknown'
|
||||
ELSE "usage".provider_name
|
||||
END AS provider_display_name,
|
||||
COALESCE("usage".api_format, 'unknown') AS api_format_group_key,
|
||||
GREATEST(COALESCE("usage".input_tokens, 0), 0) AS input_tokens,
|
||||
GREATEST(COALESCE("usage".output_tokens, 0), 0) AS output_tokens,
|
||||
GREATEST(COALESCE("usage".total_tokens, 0), 0) AS total_tokens,
|
||||
CASE
|
||||
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE("usage".cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE("usage".cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE("usage".cache_creation_input_tokens, 0)
|
||||
END AS cache_creation_tokens,
|
||||
GREATEST(COALESCE("usage".cache_creation_input_tokens_5m, 0), 0)
|
||||
AS cache_creation_ephemeral_5m_tokens,
|
||||
GREATEST(COALESCE("usage".cache_creation_input_tokens_1h, 0), 0)
|
||||
AS cache_creation_ephemeral_1h_tokens,
|
||||
GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0) AS cache_read_tokens,
|
||||
COALESCE("usage".endpoint_api_format, "usage".api_format) AS normalized_api_format,
|
||||
COALESCE(CAST("usage".total_cost_usd AS DOUBLE PRECISION), 0) AS total_cost_usd,
|
||||
COALESCE(CAST("usage".actual_total_cost_usd AS DOUBLE PRECISION), 0) AS actual_total_cost_usd,
|
||||
GREATEST(COALESCE("usage".response_time_ms, 0), 0) AS response_time_ms,
|
||||
CASE
|
||||
WHEN "usage".status IN ('completed', 'success', 'ok', 'billed', 'settled')
|
||||
AND ("usage".status_code IS NULL OR "usage".status_code < 400)
|
||||
THEN 1
|
||||
ELSE 0
|
||||
END AS success_flag
|
||||
FROM "usage"
|
||||
WHERE "usage".created_at >= TO_TIMESTAMP($1::double precision)
|
||||
AND "usage".created_at < TO_TIMESTAMP($2::double precision)
|
||||
AND "usage".status NOT IN ('pending', 'streaming')
|
||||
{filtered_extra_where}
|
||||
),
|
||||
normalized_usage AS (
|
||||
SELECT
|
||||
{group_key_expr} AS group_key,
|
||||
{display_name_expr} AS display_name,
|
||||
{secondary_name_expr} AS secondary_name,
|
||||
total_tokens,
|
||||
output_tokens,
|
||||
cache_creation_tokens,
|
||||
cache_creation_ephemeral_5m_tokens,
|
||||
cache_creation_ephemeral_1h_tokens,
|
||||
cache_read_tokens,
|
||||
total_cost_usd,
|
||||
actual_total_cost_usd,
|
||||
response_time_ms,
|
||||
success_flag,
|
||||
CASE
|
||||
WHEN input_tokens <= 0 THEN 0
|
||||
WHEN cache_read_tokens <= 0 THEN input_tokens
|
||||
WHEN split_part(lower(COALESCE(normalized_api_format, '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN GREATEST(input_tokens - cache_read_tokens, 0)
|
||||
ELSE input_tokens
|
||||
END AS effective_input_tokens,
|
||||
CASE
|
||||
WHEN split_part(lower(COALESCE(normalized_api_format, '')), ':', 1)
|
||||
IN ('claude', 'anthropic')
|
||||
THEN input_tokens + cache_creation_tokens + cache_read_tokens
|
||||
WHEN split_part(lower(COALESCE(normalized_api_format, '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN (
|
||||
CASE
|
||||
WHEN input_tokens <= 0 THEN 0
|
||||
WHEN cache_read_tokens <= 0 THEN input_tokens
|
||||
WHEN split_part(lower(COALESCE(normalized_api_format, '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN GREATEST(input_tokens - cache_read_tokens, 0)
|
||||
ELSE input_tokens
|
||||
END
|
||||
) + cache_read_tokens
|
||||
ELSE CASE
|
||||
WHEN cache_creation_tokens > 0
|
||||
THEN input_tokens + cache_creation_tokens + cache_read_tokens
|
||||
ELSE input_tokens + cache_read_tokens
|
||||
END
|
||||
END AS total_input_context
|
||||
FROM filtered_usage
|
||||
),
|
||||
aggregated_usage AS (
|
||||
SELECT
|
||||
group_key,
|
||||
{aggregate_display_name_expr} AS display_name,
|
||||
{aggregate_secondary_name_expr} AS secondary_name,
|
||||
COUNT(*)::BIGINT AS request_count,
|
||||
COALESCE(SUM(total_tokens), 0)::BIGINT AS total_tokens,
|
||||
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
|
||||
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
|
||||
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
|
||||
COALESCE(SUM(cache_creation_tokens), 0)::BIGINT AS cache_creation_tokens,
|
||||
COALESCE(SUM(cache_creation_ephemeral_5m_tokens), 0)::BIGINT
|
||||
AS cache_creation_ephemeral_5m_tokens,
|
||||
COALESCE(SUM(cache_creation_ephemeral_1h_tokens), 0)::BIGINT
|
||||
AS cache_creation_ephemeral_1h_tokens,
|
||||
COALESCE(SUM(cache_read_tokens), 0)::BIGINT AS cache_read_tokens,
|
||||
COALESCE(SUM(total_cost_usd), 0) AS total_cost_usd,
|
||||
COALESCE(SUM(actual_total_cost_usd), 0) AS actual_total_cost_usd,
|
||||
{avg_response_time_expr} AS avg_response_time_ms,
|
||||
{success_count_expr} AS success_count
|
||||
FROM normalized_usage
|
||||
GROUP BY group_key
|
||||
)
|
||||
SELECT
|
||||
group_key,
|
||||
display_name,
|
||||
secondary_name,
|
||||
request_count,
|
||||
total_tokens,
|
||||
output_tokens,
|
||||
effective_input_tokens,
|
||||
total_input_context,
|
||||
cache_creation_tokens,
|
||||
cache_creation_ephemeral_5m_tokens,
|
||||
cache_creation_ephemeral_1h_tokens,
|
||||
cache_read_tokens,
|
||||
total_cost_usd,
|
||||
actual_total_cost_usd,
|
||||
avg_response_time_ms,
|
||||
success_count
|
||||
FROM aggregated_usage
|
||||
ORDER BY request_count DESC, group_key ASC
|
||||
LIMIT $3
|
||||
"#,
|
||||
filtered_extra_where = fragments.filtered_extra_where,
|
||||
group_key_expr = fragments.group_key_expr,
|
||||
display_name_expr = fragments.display_name_expr,
|
||||
secondary_name_expr = fragments.secondary_name_expr,
|
||||
aggregate_display_name_expr = fragments.aggregate_display_name_expr,
|
||||
aggregate_secondary_name_expr = fragments.aggregate_secondary_name_expr,
|
||||
avg_response_time_expr = fragments.avg_response_time_expr,
|
||||
success_count_expr = fragments.success_count_expr,
|
||||
);
|
||||
|
||||
let mut rows = sqlx::query(&sql)
|
||||
.bind(query.created_from_unix_secs as f64)
|
||||
.bind(query.created_until_unix_secs as f64)
|
||||
.bind(i64::try_from(query.limit).map_err(|_| {
|
||||
DataLayerError::InvalidInput(format!(
|
||||
"invalid usage aggregation limit: {}",
|
||||
query.limit
|
||||
))
|
||||
})?)
|
||||
.fetch(&self.pool);
|
||||
|
||||
let mut items = Vec::new();
|
||||
while let Some(row) = rows.try_next().await.map_postgres_err()? {
|
||||
items.push(StoredUsageAuditAggregation {
|
||||
group_key: row.try_get::<String, _>("group_key").map_postgres_err()?,
|
||||
display_name: row
|
||||
.try_get::<Option<String>, _>("display_name")
|
||||
.map_postgres_err()?,
|
||||
secondary_name: row
|
||||
.try_get::<Option<String>, _>("secondary_name")
|
||||
.map_postgres_err()?,
|
||||
request_count: row
|
||||
.try_get::<i64, _>("request_count")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_tokens: row
|
||||
.try_get::<i64, _>("total_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
output_tokens: row
|
||||
.try_get::<i64, _>("output_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
effective_input_tokens: row
|
||||
.try_get::<i64, _>("effective_input_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_input_context: row
|
||||
.try_get::<i64, _>("total_input_context")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_ephemeral_5m_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_creation_ephemeral_1h_tokens: row
|
||||
.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
cache_read_tokens: row
|
||||
.try_get::<i64, _>("cache_read_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
|
||||
actual_total_cost_usd: row
|
||||
.try_get::<f64, _>("actual_total_cost_usd")
|
||||
.map_postgres_err()?,
|
||||
avg_response_time_ms: row
|
||||
.try_get::<Option<f64>, _>("avg_response_time_ms")
|
||||
.map_postgres_err()?,
|
||||
success_count: row
|
||||
.try_get::<Option<i64>, _>("success_count")
|
||||
.map_postgres_err()?
|
||||
.map(|value| value.max(0) as u64),
|
||||
});
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_daily_heatmap(
|
||||
&self,
|
||||
query: &UsageDailyHeatmapQuery,
|
||||
@@ -1756,6 +2494,38 @@ impl UsageReadRepository for SqlxUsageReadRepository {
|
||||
Self::list_usage_audits(self, query).await
|
||||
}
|
||||
|
||||
async fn count_usage_audits(&self, query: &UsageAuditListQuery) -> Result<u64, DataLayerError> {
|
||||
Self::count_usage_audits(self, query).await
|
||||
}
|
||||
|
||||
async fn aggregate_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditAggregationQuery,
|
||||
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
|
||||
Self::aggregate_usage_audits(self, query).await
|
||||
}
|
||||
|
||||
async fn summarize_usage_audits(
|
||||
&self,
|
||||
query: &UsageAuditSummaryQuery,
|
||||
) -> Result<StoredUsageAuditSummary, DataLayerError> {
|
||||
Self::summarize_usage_audits(self, query).await
|
||||
}
|
||||
|
||||
async fn summarize_usage_time_series(
|
||||
&self,
|
||||
query: &UsageTimeSeriesQuery,
|
||||
) -> Result<Vec<StoredUsageTimeSeriesBucket>, DataLayerError> {
|
||||
Self::summarize_usage_time_series(self, query).await
|
||||
}
|
||||
|
||||
async fn summarize_usage_leaderboard(
|
||||
&self,
|
||||
query: &UsageLeaderboardQuery,
|
||||
) -> Result<Vec<StoredUsageLeaderboardSummary>, DataLayerError> {
|
||||
Self::summarize_usage_leaderboard(self, query).await
|
||||
}
|
||||
|
||||
async fn list_recent_usage_audits(
|
||||
&self,
|
||||
user_id: Option<&str>,
|
||||
|
||||
Reference in New Issue
Block a user