fix(usage): repair aggregate usage statistics

This commit is contained in:
fawney19
2026-05-06 01:13:23 +08:00
parent f1358dd845
commit 3264857e4a
10 changed files with 736 additions and 121 deletions

View File

@@ -1429,7 +1429,7 @@ pub fn build_admin_stats_cost_savings_response(
let mut estimated_full_cost: f64 = usage
.iter()
.map(|item| {
item.settlement_output_price_per_1m().unwrap_or(0.0)
item.settlement_input_price_per_1m().unwrap_or(0.0)
* item.cache_read_input_tokens as f64
/ 1_000_000.0
})

View File

@@ -0,0 +1,261 @@
CREATE TEMP TABLE tmp_rebuild_cost_savings_context ON COMMIT DROP AS
SELECT
NOW() AS now_utc,
(date_trunc('day', NOW() AT TIME ZONE 'UTC') AT TIME ZONE 'UTC') AS current_day_utc;
CREATE TEMP TABLE tmp_rebuild_cost_savings_source ON COMMIT DROP AS
SELECT
usage.user_id,
usage.username,
COALESCE(usage.provider_name, '') AS provider_name,
COALESCE(usage.model, '') AS model,
(date_trunc('day', usage.created_at AT TIME ZONE 'UTC') AT TIME ZONE 'UTC') AS day_utc,
GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::BIGINT AS cache_read_tokens,
COALESCE(CAST(usage.cache_read_cost_usd AS DOUBLE PRECISION), 0) AS cache_read_cost,
COALESCE(CAST(usage.cache_creation_cost_usd AS DOUBLE PRECISION), 0) AS cache_creation_cost,
COALESCE(
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION / 1000000.0
AS estimated_full_cost
FROM usage_billing_facts AS usage
LEFT JOIN usage_settlement_snapshots
ON usage_settlement_snapshots.request_id = usage.request_id
CROSS JOIN tmp_rebuild_cost_savings_context AS context
WHERE usage.created_at < context.current_day_utc;
TRUNCATE TABLE
stats_daily_cost_savings,
stats_daily_cost_savings_provider,
stats_daily_cost_savings_model,
stats_daily_cost_savings_model_provider,
stats_user_daily_cost_savings,
stats_user_daily_cost_savings_provider,
stats_user_daily_cost_savings_model,
stats_user_daily_cost_savings_model_provider;
INSERT INTO stats_daily_cost_savings (
id,
date,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-daily-cost-savings:', CAST(source.day_utc AS TEXT))),
source.day_utc,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
GROUP BY source.day_utc, context.now_utc;
INSERT INTO stats_daily_cost_savings_provider (
id,
date,
provider_name,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-daily-cost-savings-provider:', CAST(source.day_utc AS TEXT), ':', source.provider_name)),
source.day_utc,
source.provider_name,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
GROUP BY source.day_utc, source.provider_name, context.now_utc;
INSERT INTO stats_daily_cost_savings_model (
id,
date,
model,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-daily-cost-savings-model:', CAST(source.day_utc AS TEXT), ':', source.model)),
source.day_utc,
source.model,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
GROUP BY source.day_utc, source.model, context.now_utc;
INSERT INTO stats_daily_cost_savings_model_provider (
id,
date,
model,
provider_name,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-daily-cost-savings-model-provider:', CAST(source.day_utc AS TEXT), ':', source.model, ':', source.provider_name)),
source.day_utc,
source.model,
source.provider_name,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
GROUP BY source.day_utc, source.model, source.provider_name, context.now_utc;
INSERT INTO stats_user_daily_cost_savings (
id,
user_id,
username,
date,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-user-daily-cost-savings:', source.user_id, ':', CAST(source.day_utc AS TEXT))),
source.user_id,
MAX(source.username),
source.day_utc,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
WHERE source.user_id IS NOT NULL
GROUP BY source.user_id, source.day_utc, context.now_utc;
INSERT INTO stats_user_daily_cost_savings_provider (
id,
user_id,
username,
date,
provider_name,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-user-daily-cost-savings-provider:', source.user_id, ':', CAST(source.day_utc AS TEXT), ':', source.provider_name)),
source.user_id,
MAX(source.username),
source.day_utc,
source.provider_name,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
WHERE source.user_id IS NOT NULL
GROUP BY source.user_id, source.day_utc, source.provider_name, context.now_utc;
INSERT INTO stats_user_daily_cost_savings_model (
id,
user_id,
username,
date,
model,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-user-daily-cost-savings-model:', source.user_id, ':', CAST(source.day_utc AS TEXT), ':', source.model)),
source.user_id,
MAX(source.username),
source.day_utc,
source.model,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
WHERE source.user_id IS NOT NULL
GROUP BY source.user_id, source.day_utc, source.model, context.now_utc;
INSERT INTO stats_user_daily_cost_savings_model_provider (
id,
user_id,
username,
date,
model,
provider_name,
cache_read_tokens,
cache_read_cost,
cache_creation_cost,
estimated_full_cost,
created_at,
updated_at
)
SELECT
md5(CONCAT('stats-user-daily-cost-savings-model-provider:', source.user_id, ':', CAST(source.day_utc AS TEXT), ':', source.model, ':', source.provider_name)),
source.user_id,
MAX(source.username),
source.day_utc,
source.model,
source.provider_name,
COALESCE(SUM(source.cache_read_tokens), 0)::BIGINT,
CAST(COALESCE(SUM(source.cache_read_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.cache_creation_cost), 0) AS DOUBLE PRECISION),
CAST(COALESCE(SUM(source.estimated_full_cost), 0) AS DOUBLE PRECISION),
context.now_utc,
context.now_utc
FROM tmp_rebuild_cost_savings_source AS source
CROSS JOIN tmp_rebuild_cost_savings_context AS context
WHERE source.user_id IS NOT NULL
GROUP BY
source.user_id,
source.day_utc,
source.model,
source.provider_name,
context.now_utc;

View File

@@ -916,8 +916,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -982,8 +982,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -1058,8 +1058,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -1135,8 +1135,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -2611,8 +2611,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -2693,8 +2693,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -2779,8 +2779,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0
@@ -2866,8 +2866,8 @@ WITH aggregated AS (
COALESCE(
SUM(
COALESCE(
CAST(usage_settlement_snapshots.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage.output_price_per_1m AS DOUBLE PRECISION),
CAST(usage_settlement_snapshots.input_price_per_1m AS DOUBLE PRECISION),
CAST(usage.input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE(usage.cache_read_input_tokens, 0), 0)::DOUBLE PRECISION
/ 1000000.0

View File

@@ -306,7 +306,10 @@ mod tests {
.into_iter()
.map(|item| item.version)
.collect::<Vec<_>>();
assert_eq!(versions, vec![20260422110000, 20260422120000]);
assert_eq!(
versions,
vec![20260422110000, 20260422120000, 20260504120000]
);
}
#[test]
@@ -318,7 +321,7 @@ mod tests {
.into_iter()
.map(|item| item.version)
.collect::<Vec<_>>();
assert_eq!(versions, vec![20260422120000]);
assert_eq!(versions, vec![20260422120000, 20260504120000]);
}
#[tokio::test]
@@ -619,6 +622,7 @@ mod tests {
cache_read_cost_usd,
total_cost_usd,
actual_total_cost_usd,
input_price_per_1m,
output_price_per_1m,
status,
billing_status,
@@ -646,6 +650,7 @@ mod tests {
1.25,
1.10,
50.0,
50.0,
'completed',
'settled',
TIMESTAMPTZ '2024-05-06 07:18:09+00',
@@ -665,9 +670,10 @@ mod tests {
let pending_before = pending_backfills(&pool)
.await
.expect("pending backfills should load");
assert_eq!(pending_before.len(), 2);
assert_eq!(pending_before.len(), 3);
assert_eq!(pending_before[0].version, 20260422110000);
assert_eq!(pending_before[1].version, 20260422120000);
assert_eq!(pending_before[2].version, 20260504120000);
run_backfills(&pool)
.await
@@ -683,7 +689,10 @@ mod tests {
.fetch_all(&pool)
.await
.expect("applied backfill versions should load");
assert_eq!(applied_versions, vec![20260422110000, 20260422120000]);
assert_eq!(
applied_versions,
vec![20260422110000, 20260422120000, 20260504120000]
);
let api_key_total_requests: i64 = query_scalar(
"SELECT COALESCE(total_requests, 0)::BIGINT FROM public.api_keys WHERE id = 'api-key-backfill-1'",
@@ -1335,6 +1344,6 @@ mod tests {
.fetch_one(&pool)
.await
.expect("backfill count should load");
assert_eq!(applied_count, 1);
assert_eq!(applied_count, 3);
}
}

View File

@@ -1915,7 +1915,7 @@ impl UsageReadRepository for InMemoryUsageReadRepository {
.saturating_add(item.cache_read_input_tokens);
summary.cache_read_cost_usd += item.cache_read_cost_usd;
summary.cache_creation_cost_usd += item.cache_creation_cost_usd;
summary.estimated_full_cost_usd += item.settlement_output_price_per_1m().unwrap_or(0.0)
summary.estimated_full_cost_usd += item.settlement_input_price_per_1m().unwrap_or(0.0)
* item.cache_read_input_tokens as f64
/ 1_000_000.0;
}

View File

@@ -396,6 +396,99 @@ fn finalize_usage_breakdown_rows(
items
}
fn absorb_usage_audit_aggregation_rows(
target: &mut BTreeMap<String, StoredUsageAuditAggregation>,
rows: Vec<StoredUsageAuditAggregation>,
) {
for row in rows {
let group_key = row.group_key.clone();
let entry =
target
.entry(group_key.clone())
.or_insert_with(|| StoredUsageAuditAggregation {
group_key,
display_name: row.display_name.clone(),
secondary_name: row.secondary_name.clone(),
request_count: 0,
total_tokens: 0,
output_tokens: 0,
effective_input_tokens: 0,
total_input_context: 0,
cache_creation_tokens: 0,
cache_creation_ephemeral_5m_tokens: 0,
cache_creation_ephemeral_1h_tokens: 0,
cache_read_tokens: 0,
total_cost_usd: 0.0,
actual_total_cost_usd: 0.0,
avg_response_time_ms: None,
success_count: row.success_count.map(|_| 0),
});
if entry.display_name.is_none() {
entry.display_name = row.display_name;
}
if entry.secondary_name.is_none() {
entry.secondary_name = row.secondary_name;
}
let existing_request_count = entry.request_count;
let next_request_count = row.request_count;
entry.request_count = entry.request_count.saturating_add(row.request_count);
entry.total_tokens = entry.total_tokens.saturating_add(row.total_tokens);
entry.output_tokens = entry.output_tokens.saturating_add(row.output_tokens);
entry.effective_input_tokens = entry
.effective_input_tokens
.saturating_add(row.effective_input_tokens);
entry.total_input_context = entry
.total_input_context
.saturating_add(row.total_input_context);
entry.cache_creation_tokens = entry
.cache_creation_tokens
.saturating_add(row.cache_creation_tokens);
entry.cache_creation_ephemeral_5m_tokens = entry
.cache_creation_ephemeral_5m_tokens
.saturating_add(row.cache_creation_ephemeral_5m_tokens);
entry.cache_creation_ephemeral_1h_tokens = entry
.cache_creation_ephemeral_1h_tokens
.saturating_add(row.cache_creation_ephemeral_1h_tokens);
entry.cache_read_tokens = entry
.cache_read_tokens
.saturating_add(row.cache_read_tokens);
entry.total_cost_usd += row.total_cost_usd;
entry.actual_total_cost_usd += row.actual_total_cost_usd;
entry.success_count = match (entry.success_count, row.success_count) {
(Some(left), Some(right)) => Some(left.saturating_add(right)),
(Some(left), None) => Some(left),
(None, Some(right)) => Some(right),
(None, None) => None,
};
entry.avg_response_time_ms = match (entry.avg_response_time_ms, row.avg_response_time_ms) {
(Some(left), Some(right)) if entry.request_count > 0 => Some(
((left * existing_request_count as f64) + (right * next_request_count as f64))
/ entry.request_count as f64,
),
(Some(left), _) => Some(left),
(None, Some(right)) => Some(right),
(None, None) => None,
};
}
}
fn finalize_usage_audit_aggregation_rows(
grouped: BTreeMap<String, StoredUsageAuditAggregation>,
limit: usize,
) -> Vec<StoredUsageAuditAggregation> {
let mut items = grouped.into_values().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(limit);
items
}
fn decode_usage_breakdown_summary_row(
row: &PgRow,
) -> Result<StoredUsageBreakdownSummaryRow, DataLayerError> {
@@ -466,6 +559,67 @@ fn decode_usage_breakdown_summary_row(
})
}
fn decode_usage_audit_aggregation_row(
row: &PgRow,
) -> Result<StoredUsageAuditAggregation, DataLayerError> {
Ok(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),
})
}
fn absorb_usage_audit_summary(target: &mut StoredUsageAuditSummary, row: StoredUsageAuditSummary) {
target.total_requests = target.total_requests.saturating_add(row.total_requests);
target.input_tokens = target.input_tokens.saturating_add(row.input_tokens);
@@ -1336,7 +1490,7 @@ SELECT
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
COALESCE(SUM(input_tokens + output_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(cache_creation_tokens), 0)::BIGINT AS cache_creation_tokens,
COALESCE(SUM(cache_read_tokens), 0)::BIGINT AS cache_read_tokens,
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
@@ -1367,7 +1521,7 @@ SELECT
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
COALESCE(SUM(input_tokens + output_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(cache_creation_tokens), 0)::BIGINT AS cache_creation_tokens,
COALESCE(SUM(cache_read_tokens), 0)::BIGINT AS cache_read_tokens,
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
@@ -1424,7 +1578,33 @@ SELECT
END
), 0)::BIGINT AS effective_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 total_tokens,
COALESCE(SUM(
CASE
WHEN GREATEST(COALESCE("usage".input_tokens, 0), 0) <= 0 THEN 0
WHEN GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0) <= 0
THEN GREATEST(COALESCE("usage".input_tokens, 0), 0)
WHEN split_part(lower(COALESCE(COALESCE("usage".endpoint_api_format, "usage".api_format), '')), ':', 1)
IN ('openai', 'gemini', 'google')
THEN GREATEST(
GREATEST(COALESCE("usage".input_tokens, 0), 0)
- GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0),
0
)
ELSE GREATEST(COALESCE("usage".input_tokens, 0), 0)
END
+ GREATEST(COALESCE("usage".output_tokens, 0), 0)
+ 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
+ GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)
), 0)::BIGINT AS total_tokens,
COALESCE(SUM(
CASE
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
@@ -3682,7 +3862,33 @@ SELECT
"usage".model AS model,
"usage".provider_name AS provider,
COUNT(*)::BIGINT AS requests,
COALESCE(SUM(GREATEST(COALESCE("usage".total_tokens, 0), 0)), 0)::BIGINT AS total_tokens,
COALESCE(SUM(
CASE
WHEN GREATEST(COALESCE("usage".input_tokens, 0), 0) <= 0 THEN 0
WHEN GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0) <= 0
THEN GREATEST(COALESCE("usage".input_tokens, 0), 0)
WHEN split_part(lower(COALESCE(COALESCE("usage".endpoint_api_format, "usage".api_format), '')), ':', 1)
IN ('openai', 'gemini', 'google')
THEN GREATEST(
GREATEST(COALESCE("usage".input_tokens, 0), 0)
- GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0),
0
)
ELSE GREATEST(COALESCE("usage".input_tokens, 0), 0)
END
+ GREATEST(COALESCE("usage".output_tokens, 0), 0)
+ 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
+ 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,
COALESCE(SUM(
@@ -3886,7 +4092,7 @@ SELECT
{group_column} AS group_key,
COALESCE(SUM(total_requests), 0)::BIGINT AS request_count,
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(total_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_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,
@@ -4893,7 +5099,7 @@ SELECT
AS cache_creation_cost_usd,
COALESCE(SUM(
COALESCE(
CAST("usage".output_price_per_1m AS DOUBLE PRECISION),
CAST("usage".input_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)::DOUBLE PRECISION / 1000000.0
), 0) AS estimated_full_cost_usd
@@ -5945,7 +6151,85 @@ WHERE stats_daily_api_key.date >=
Ok(finalize_usage_leaderboard_rows(grouped))
}
pub async fn aggregate_usage_audits(
async fn aggregate_usage_audits_from_daily_aggregates(
&self,
start_day_utc: DateTime<Utc>,
end_day_utc: DateTime<Utc>,
group_by: UsageAuditAggregationGroupBy,
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
if start_day_utc >= end_day_utc {
return Ok(Vec::new());
}
let (table_name, group_column, display_name_expr, avg_response_time_expr, success_count_expr) =
match group_by {
UsageAuditAggregationGroupBy::Model => (
"stats_user_daily_model",
"model",
"NULL::varchar",
"NULL::DOUBLE PRECISION",
"NULL::BIGINT",
),
UsageAuditAggregationGroupBy::Provider => (
"stats_user_daily_provider",
"provider_name",
"provider_name",
"CASE WHEN COALESCE(SUM(response_time_samples), 0) > 0 THEN COALESCE(SUM(response_time_sum_ms), 0) / COALESCE(SUM(response_time_samples), 0) ELSE NULL END",
"COALESCE(SUM(success_requests), 0)::BIGINT",
),
UsageAuditAggregationGroupBy::ApiFormat => (
"stats_user_daily_api_format",
"api_format",
"NULL::varchar",
"CASE WHEN COALESCE(SUM(response_time_samples), 0) > 0 THEN COALESCE(SUM(response_time_sum_ms), 0) / COALESCE(SUM(response_time_samples), 0) ELSE NULL END",
"NULL::BIGINT",
),
UsageAuditAggregationGroupBy::User => {
return Ok(Vec::new());
}
};
let sql = format!(
r#"
SELECT
{group_column} AS group_key,
{display_name_expr} AS display_name,
NULL::varchar AS secondary_name,
COALESCE(SUM(total_requests), 0)::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), 0)::DOUBLE PRECISION AS total_cost_usd,
COALESCE(SUM(actual_total_cost), 0)::DOUBLE PRECISION AS actual_total_cost_usd,
{avg_response_time_expr} AS avg_response_time_ms,
{success_count_expr} AS success_count
FROM {table_name}
WHERE date >= $1
AND date < $2
GROUP BY {group_column}
ORDER BY request_count DESC, group_key ASC
"#,
);
let mut rows = sqlx::query(&sql)
.bind(start_day_utc)
.bind(end_day_utc)
.fetch(&self.pool);
let mut items = Vec::new();
while let Some(row) = rows.try_next().await.map_postgres_err()? {
items.push(decode_usage_audit_aggregation_row(&row)?);
}
Ok(items)
}
async fn aggregate_usage_audits_raw(
&self,
query: &UsageAuditAggregationQuery,
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
@@ -6112,66 +6396,67 @@ LIMIT $3
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),
});
items.push(decode_usage_audit_aggregation_row(&row)?);
}
Ok(items)
}
pub async fn aggregate_usage_audits(
&self,
query: &UsageAuditAggregationQuery,
) -> Result<Vec<StoredUsageAuditAggregation>, DataLayerError> {
if matches!(query.group_by, UsageAuditAggregationGroupBy::User) {
return self.aggregate_usage_audits_raw(query).await;
}
let Some(cutoff_utc) = self.read_stats_daily_cutoff_date().await? else {
return self.aggregate_usage_audits_raw(query).await;
};
let start_utc = dashboard_unix_secs_to_utc(query.created_from_unix_secs);
let end_utc = dashboard_unix_secs_to_utc(query.created_until_unix_secs);
let split = split_dashboard_daily_aggregate_range(start_utc, end_utc, cutoff_utc);
let Some(_) = split.aggregate else {
return self.aggregate_usage_audits_raw(query).await;
};
let mut grouped = BTreeMap::<String, StoredUsageAuditAggregation>::new();
let raw_merge_limit = query.limit.max(10_000);
if let Some((raw_start, raw_end)) = split.raw_leading {
let raw = self
.aggregate_usage_audits_raw(&UsageAuditAggregationQuery {
created_from_unix_secs: dashboard_utc_to_unix_secs(raw_start),
created_until_unix_secs: dashboard_utc_to_unix_secs(raw_end),
group_by: query.group_by,
limit: raw_merge_limit,
})
.await?;
absorb_usage_audit_aggregation_rows(&mut grouped, raw);
}
if let Some((aggregate_start, aggregate_end)) = split.aggregate {
let aggregate = self
.aggregate_usage_audits_from_daily_aggregates(
aggregate_start,
aggregate_end,
query.group_by,
)
.await?;
absorb_usage_audit_aggregation_rows(&mut grouped, aggregate);
}
if let Some((raw_start, raw_end)) = split.raw_trailing {
let raw = self
.aggregate_usage_audits_raw(&UsageAuditAggregationQuery {
created_from_unix_secs: dashboard_utc_to_unix_secs(raw_start),
created_until_unix_secs: dashboard_utc_to_unix_secs(raw_end),
group_by: query.group_by,
limit: raw_merge_limit,
})
.await?;
absorb_usage_audit_aggregation_rows(&mut grouped, raw);
}
Ok(finalize_usage_audit_aggregation_rows(grouped, query.limit))
}
async fn summarize_usage_daily_heatmap_raw_from_range(
&self,
start_utc: DateTime<Utc>,

View File

@@ -401,6 +401,19 @@ fn usage_sql_summarize_usage_leaderboard_supports_daily_aggregates() {
);
}
#[test]
fn usage_sql_aggregate_usage_audits_supports_daily_model_and_provider_aggregates() {
let source = include_str!("mod.rs");
assert!(source.contains("aggregate_usage_audits_from_daily_aggregates"));
assert!(source.contains("stats_user_daily_model"));
assert!(source.contains("stats_user_daily_provider"));
assert!(source.contains("stats_user_daily_api_format"));
assert!(source.contains("absorb_usage_audit_aggregation_rows"));
assert!(
source.contains("split_dashboard_daily_aggregate_range(start_utc, end_utc, cutoff_utc)")
);
}
#[test]
fn usage_sql_summarize_total_tokens_by_api_key_ids_supports_daily_aggregates() {
let source = include_str!("mod.rs");