mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-11 03:39:49 +08:00
fix: align dashboard charts with customer billing and history coverage
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
use super::{envelope, metrics_value, parse_overview_query, OverviewRequest};
|
||||
use aether_data_contracts::repository::usage::{
|
||||
StoredUsageAnalytics, StoredUsageDashboardAnalytics, UsageAnalyticsQuery, UsageAnalyticsRow,
|
||||
UsageAnalyticsView, UsageDashboardAnalyticsQuery,
|
||||
StoredUsageAnalytics, StoredUsageDashboardAnalytics, UsageAnalyticsMetrics,
|
||||
UsageAnalyticsQuery, UsageAnalyticsRow, UsageAnalyticsView, UsageDashboardAnalyticsQuery,
|
||||
};
|
||||
use chrono::DateTime;
|
||||
use serde_json::{json, Value};
|
||||
@@ -39,22 +39,42 @@ pub fn parse_dashboard_charts_query(raw: Option<&str>) -> Result<OverviewRequest
|
||||
const OVERVIEW_CHART_LIMIT: u32 = 10_000;
|
||||
|
||||
pub fn dashboard_charts_value(snapshot: &StoredUsageAnalytics) -> Value {
|
||||
// The chart read deliberately computes only its displayed metrics. Avoid
|
||||
// presenting unmeasured diagnostics as zero through the shared serializer.
|
||||
let chart_metrics = |metrics: &UsageAnalyticsMetrics| {
|
||||
let mut value = metrics_value(metrics);
|
||||
for field in [
|
||||
"input_tokens",
|
||||
"output_tokens",
|
||||
"usage_active_users",
|
||||
"slow_request_count",
|
||||
"unclassified_failure_count",
|
||||
] {
|
||||
value[field] = Value::Null;
|
||||
}
|
||||
value["usage_source"] = json!("unknown");
|
||||
value["usage_source_counts"] = json!({
|
||||
"reported": 0, "estimated": 0, "mixed": 0, "unknown": metrics.request_count,
|
||||
});
|
||||
value
|
||||
};
|
||||
let rows = |items: &[UsageAnalyticsRow]| {
|
||||
items
|
||||
.iter()
|
||||
.map(|row| {
|
||||
let mut value = metrics_value(&row.metrics);
|
||||
let mut value = chart_metrics(&row.metrics);
|
||||
if snapshot.unrecoverable_bucket_count > 0 {
|
||||
mark_incomplete_amounts(&mut value);
|
||||
}
|
||||
value["id"] = json!(row.id);
|
||||
value["label"] = json!(row.label);
|
||||
value["bucket_start"] = json!(row.bucket_start);
|
||||
value["unique_providers"] = json!(row.metrics.unique_providers);
|
||||
value
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
};
|
||||
let mut summary = metrics_value(&snapshot.summary);
|
||||
let mut summary = chart_metrics(&snapshot.summary);
|
||||
if snapshot.unrecoverable_bucket_count > 0 {
|
||||
mark_incomplete_amounts(&mut summary);
|
||||
}
|
||||
@@ -293,6 +313,9 @@ mod tests {
|
||||
..Default::default()
|
||||
};
|
||||
snapshot.summary.billable_amount = Some("2.00000000".into());
|
||||
snapshot.summary.request_count = 1;
|
||||
snapshot.summary.total_tokens = 42;
|
||||
snapshot.summary.usage_available_count = 1;
|
||||
let row = UsageAnalyticsRow {
|
||||
id: Some("model-1".into()),
|
||||
label: None,
|
||||
@@ -303,6 +326,25 @@ mod tests {
|
||||
snapshot.model_rows.push(row.clone());
|
||||
snapshot.provider_rows.push(row);
|
||||
let data = dashboard_charts_value(&snapshot);
|
||||
for metrics in [
|
||||
&data["summary"],
|
||||
&data["series"][0],
|
||||
&data["models"][0],
|
||||
&data["providers"][0],
|
||||
] {
|
||||
assert_eq!(metrics["total_tokens"], 42);
|
||||
assert_eq!(metrics["usage_source"], "unknown");
|
||||
assert_eq!(metrics["usage_source_counts"]["unknown"], 1);
|
||||
for field in [
|
||||
"input_tokens",
|
||||
"output_tokens",
|
||||
"usage_active_users",
|
||||
"slow_request_count",
|
||||
"unclassified_failure_count",
|
||||
] {
|
||||
assert!(metrics[field].is_null(), "{field} was not measured");
|
||||
}
|
||||
}
|
||||
assert_eq!(
|
||||
data["summary"]["billable_amount"]["status"],
|
||||
"known_subtotal"
|
||||
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
-- PostgreSQL 15/16 cannot enter the function's EXCEPTION subtransactions in a
|
||||
-- parallel operation. Keep its validation and immutable amount semantics intact.
|
||||
ALTER FUNCTION public.usage_customer_billable_amount(jsonb, numeric, numeric) PARALLEL UNSAFE;
|
||||
@@ -312,6 +312,9 @@ impl SqlxUsageReadRepository {
|
||||
query: &UsageAnalyticsQuery,
|
||||
) -> Result<StoredUsageAnalytics, DataLayerError> {
|
||||
query.validate()?;
|
||||
if query.view == UsageAnalyticsView::DashboardCharts {
|
||||
return self.query_dashboard_charts(query).await;
|
||||
}
|
||||
let mut tx = self.pool.begin().await.map_postgres_err()?;
|
||||
sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
|
||||
.execute(&mut *tx)
|
||||
@@ -370,7 +373,6 @@ impl SqlxUsageReadRepository {
|
||||
match query.view {
|
||||
UsageAnalyticsView::Timeseries
|
||||
| UsageAnalyticsView::Performance
|
||||
| UsageAnalyticsView::DashboardCharts
|
||||
| UsageAnalyticsView::Breakdown => {
|
||||
let timeseries = query.view != UsageAnalyticsView::Breakdown;
|
||||
// 提供商分组的行需要额外带出名称快照,并在最外层关联提供商目录解析展示名。
|
||||
@@ -531,29 +533,21 @@ impl SqlxUsageReadRepository {
|
||||
.map_postgres_err()?;
|
||||
result.consumption = decode(row.try_get("items").map_postgres_err()?)?;
|
||||
}
|
||||
UsageAnalyticsView::Summary => unreachable!(),
|
||||
UsageAnalyticsView::Summary | UsageAnalyticsView::DashboardCharts => unreachable!(),
|
||||
}
|
||||
}
|
||||
if matches!(
|
||||
query.view,
|
||||
UsageAnalyticsView::Timeseries
|
||||
| UsageAnalyticsView::Performance
|
||||
| UsageAnalyticsView::DashboardCharts
|
||||
UsageAnalyticsView::Timeseries | UsageAnalyticsView::Performance
|
||||
) {
|
||||
fill_usage_analytics_timeseries(query, &mut result.rows);
|
||||
result.total = result.rows.len() as u64;
|
||||
}
|
||||
if matches!(
|
||||
query.view,
|
||||
UsageAnalyticsView::Performance | UsageAnalyticsView::DashboardCharts
|
||||
) {
|
||||
if query.view == UsageAnalyticsView::Performance {
|
||||
let mut providers = QueryBuilder::<Postgres>::new("SELECT COALESCE(jsonb_agg(jsonb_build_object('id',provider_id,'label',provider_label,'bucket_start',NULL,'metrics',to_jsonb(m)-'provider_id'-'provider_label')), '[]'::jsonb) AS items FROM (SELECT provider_id, max(provider_name) AS provider_label, ");
|
||||
providers.push(&metrics_sql);
|
||||
push_analytics_filter(&mut providers, query);
|
||||
providers.push(" GROUP BY provider_id ORDER BY count(*) DESC,provider_id");
|
||||
if query.view == UsageAnalyticsView::DashboardCharts {
|
||||
providers.push(" LIMIT 10001");
|
||||
}
|
||||
providers.push(") m");
|
||||
let row = providers
|
||||
.build()
|
||||
@@ -561,38 +555,6 @@ impl SqlxUsageReadRepository {
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
result.provider_rows = decode(row.try_get("items").map_postgres_err()?)?;
|
||||
if result.provider_rows.len() > USAGE_DASHBOARD_CHART_ROW_LIMIT
|
||||
&& query.view == UsageAnalyticsView::DashboardCharts
|
||||
{
|
||||
return Err(DataLayerError::InvalidInput(
|
||||
"dashboard provider chart exceeds 10000 groups".into(),
|
||||
));
|
||||
}
|
||||
}
|
||||
if query.view == UsageAnalyticsView::DashboardCharts {
|
||||
let timezone = if query.granularity == UsageAnalyticsGranularity::Hour {
|
||||
"UTC".into()
|
||||
} else {
|
||||
query.timezone.clone()
|
||||
};
|
||||
let mut models = QueryBuilder::<Postgres>::new("WITH filtered AS (SELECT *");
|
||||
push_analytics_filter(&mut models, query);
|
||||
models.push("), dated AS (SELECT *,date_trunc(")
|
||||
.push_bind(if query.granularity==UsageAnalyticsGranularity::Hour {"hour"}else{"day"})
|
||||
.push(",created_at AT TIME ZONE ").push_bind(timezone.clone())
|
||||
.push(") AT TIME ZONE ").push_bind(timezone).push(" AS bucket FROM filtered) SELECT COALESCE(jsonb_agg(jsonb_build_object('id',model,'label',model,'bucket_start',to_char(bucket AT TIME ZONE 'UTC','YYYY-MM-DD\"T\"HH24:MI:SS\"Z\"'),'metrics',to_jsonb(m)-'model'-'bucket')),'[]'::jsonb) AS items FROM (SELECT model,bucket,")
|
||||
.push(&metrics_sql).push(" FROM dated GROUP BY model,bucket ORDER BY bucket,model LIMIT 10001) m");
|
||||
let row = models
|
||||
.build()
|
||||
.fetch_one(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
result.model_rows = decode(row.try_get("items").map_postgres_err()?)?;
|
||||
if result.model_rows.len() > USAGE_DASHBOARD_CHART_ROW_LIMIT {
|
||||
return Err(DataLayerError::InvalidInput(
|
||||
"dashboard model chart exceeds 10000 groups; narrow the range".into(),
|
||||
));
|
||||
}
|
||||
}
|
||||
if query.view == UsageAnalyticsView::Performance {
|
||||
// Aggregate requested models across providers in the same read snapshot.
|
||||
|
||||
@@ -0,0 +1,369 @@
|
||||
use super::{
|
||||
analytics::{analytics_metrics_sql, push_analytics_source_filter},
|
||||
projection_reader::read_projection_coverage,
|
||||
SqlxUsageReadRepository,
|
||||
};
|
||||
use crate::error::SqlxResultExt;
|
||||
use aether_data_contracts::{repository::usage::*, DataLayerError};
|
||||
use chrono::{DateTime, Duration, Utc};
|
||||
use serde::de::DeserializeOwned;
|
||||
use serde_json::Value;
|
||||
use sqlx::{Postgres, QueryBuilder, Row};
|
||||
|
||||
// Committed request contributions are updated in the same transaction as the
|
||||
// top dashboard counters. Reuse them without detoasting large request metadata.
|
||||
// Rows predating that ledger retain DASHBOARD_TOTAL_FACTS_SQL token precedence.
|
||||
const DASHBOARD_CHART_FACTS_SQL: &str = r#"(
|
||||
SELECT u.created_at,u.model,u.provider_id,u.provider_name,u.status,u.response_time_ms,
|
||||
COALESCE(a.record_kind,'request') AS record_kind,
|
||||
COALESCE(s.billing_status,u.billing_status) AS settlement_status,s.allocation_status,
|
||||
CASE WHEN c.request_id IS NOT NULL THEN (c.metrics->>'usage_available_count')::bigint>0
|
||||
ELSE COALESCE(metadata.value->'usage_available','true'::jsonb)<>'false'::jsonb END AS usage_available,
|
||||
CASE WHEN c.request_id IS NOT NULL THEN
|
||||
CASE WHEN (c.metrics->>'usage_available_count')::bigint>0 THEN (c.metrics->>'total_tokens')::bigint END
|
||||
WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
CASE
|
||||
WHEN s.billing_effective_input_tokens IS NOT NULL
|
||||
THEN GREATEST(s.billing_effective_input_tokens, 0)
|
||||
+ GREATEST(COALESCE(s.billing_output_tokens, u.output_tokens, 0), 0)
|
||||
+ GREATEST(
|
||||
COALESCE(
|
||||
s.billing_cache_creation_tokens,
|
||||
CASE
|
||||
WHEN s.billing_cache_creation_5m_tokens IS NOT NULL
|
||||
OR s.billing_cache_creation_1h_tokens IS NOT NULL
|
||||
THEN COALESCE(s.billing_cache_creation_5m_tokens, 0)
|
||||
+ COALESCE(s.billing_cache_creation_1h_tokens, 0)
|
||||
END,
|
||||
CASE
|
||||
WHEN COALESCE(u.cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE(u.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(u.cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE(u.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(u.cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE(u.cache_creation_input_tokens, 0)
|
||||
END,
|
||||
0
|
||||
),
|
||||
0
|
||||
)
|
||||
+ GREATEST(
|
||||
COALESCE(
|
||||
s.billing_cache_read_tokens,
|
||||
u.cache_read_input_tokens,
|
||||
0
|
||||
),
|
||||
0
|
||||
)
|
||||
WHEN s.billing_total_input_context IS NOT NULL
|
||||
THEN GREATEST(s.billing_total_input_context, 0)
|
||||
+ GREATEST(COALESCE(s.billing_output_tokens, u.output_tokens, 0), 0)
|
||||
END,
|
||||
NULLIF(GREATEST(COALESCE(u.total_tokens, 0), 0), 0),
|
||||
CASE
|
||||
WHEN split_part(lower(COALESCE(COALESCE(u.endpoint_api_format, u.api_format), '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN GREATEST(COALESCE(u.input_tokens, 0), 0)
|
||||
+ GREATEST(COALESCE(u.output_tokens, 0), 0)
|
||||
ELSE GREATEST(COALESCE(u.input_tokens, 0), 0)
|
||||
+ GREATEST(COALESCE(u.output_tokens, 0), 0)
|
||||
+ GREATEST(
|
||||
CASE
|
||||
WHEN COALESCE(u.cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE(u.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(u.cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE(u.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(u.cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE(u.cache_creation_input_tokens, 0)
|
||||
END,
|
||||
0
|
||||
)
|
||||
+ GREATEST(COALESCE(u.cache_read_input_tokens, 0), 0)
|
||||
END,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint END AS total_tokens,
|
||||
CASE WHEN CASE WHEN c.request_id IS NOT NULL THEN (c.metrics->>'pricing_available_count')::bigint>0
|
||||
ELSE COALESCE(metadata.value->'usage_pricing_available','true'::jsonb)<>'false'::jsonb END
|
||||
AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status,u.billing_status)='settled')
|
||||
THEN round(COALESCE(s.billing_total_cost_usd::numeric,u.total_cost_usd::numeric),8) END AS rated_amount,
|
||||
CASE WHEN c.request_id IS NOT NULL THEN
|
||||
CASE WHEN (c.metrics->>'pricing_available_count')::bigint>0
|
||||
THEN round((c.metrics->>'billable_amount')::numeric,8) END
|
||||
WHEN COALESCE(metadata.value->'usage_pricing_available','true'::jsonb)<>'false'::jsonb
|
||||
AND (s.billing_actual_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status,u.billing_status)='settled')
|
||||
THEN public.usage_customer_billable_amount(metadata.value,
|
||||
COALESCE(s.billing_total_cost_usd::numeric,u.total_cost_usd::numeric),
|
||||
COALESCE(s.billing_actual_total_cost_usd::numeric,u.actual_total_cost_usd::numeric)) END AS billable_amount
|
||||
FROM public.usage u
|
||||
LEFT JOIN public.dashboard_request_contributions c USING(request_id)
|
||||
LEFT JOIN public.usage_settlement_snapshots s USING(request_id)
|
||||
LEFT JOIN public.usage_attribution_snapshots a USING(request_id)
|
||||
LEFT JOIN LATERAL (SELECT u.request_metadata::jsonb AS value
|
||||
WHERE c.request_id IS NULL OFFSET 0) metadata ON true
|
||||
) AS chart_facts"#;
|
||||
|
||||
const DASHBOARD_CHART_METRICS_SQL: &str = r#"
|
||||
count(*)::bigint AS request_count,
|
||||
count(*) FILTER (WHERE status='completed')::bigint AS successful_request_count,
|
||||
count(*) FILTER (WHERE status='failed')::bigint AS failed_request_count,
|
||||
count(*) FILTER (WHERE status='cancelled')::bigint AS cancelled_request_count,
|
||||
count(*) FILTER (WHERE status NOT IN ('completed','failed','cancelled'))::bigint AS in_flight_request_count,
|
||||
COALESCE(sum(total_tokens),0)::bigint AS total_tokens,
|
||||
count(*) FILTER (WHERE usage_available)::bigint AS usage_available_count,
|
||||
count(billable_amount)::bigint AS pricing_available_count,
|
||||
count(*) FILTER (WHERE settlement_status='settled')::bigint AS settled_count,
|
||||
count(*) FILTER (WHERE allocation_status='complete')::bigint AS allocation_available_count,
|
||||
sum(rated_amount)::text AS rated_amount,
|
||||
sum(billable_amount)::text AS billable_amount,
|
||||
count(response_time_ms)::bigint AS latency_sample_count,
|
||||
COALESCE(sum(response_time_ms),0)::double precision AS latency_sum_ms
|
||||
"#;
|
||||
|
||||
fn decode<T: DeserializeOwned>(value: Value) -> Result<T, DataLayerError> {
|
||||
serde_json::from_value(value)
|
||||
.map_err(|error| DataLayerError::UnexpectedValue(error.to_string()))
|
||||
}
|
||||
|
||||
fn can_compare_daily_history(query: &UsageAnalyticsQuery) -> bool {
|
||||
query.actor_user_id.is_none()
|
||||
&& query.credential_owner_id.is_none()
|
||||
&& query.attribution_kind.is_none()
|
||||
&& query.api_key_id.is_none()
|
||||
&& query.model.is_none()
|
||||
&& query.provider_id.is_none()
|
||||
&& query.api_format.is_none()
|
||||
&& query.endpoint_kind.is_none()
|
||||
&& query.request_type.is_none()
|
||||
&& query.status.is_none()
|
||||
&& query.is_stream.is_none()
|
||||
&& query.has_format_conversion.is_none()
|
||||
&& query.search.is_none()
|
||||
&& query.user_is_active.is_none()
|
||||
&& query.has_usage.is_none()
|
||||
}
|
||||
|
||||
impl SqlxUsageReadRepository {
|
||||
/// Keep every chart and its summary on the same bounded canonical fact scan.
|
||||
/// Historical projections can outlive raw facts and must not be mixed into
|
||||
/// this summary while its model/provider breakdown still reads raw rows.
|
||||
pub(super) async fn query_dashboard_charts(
|
||||
&self,
|
||||
query: &UsageAnalyticsQuery,
|
||||
) -> Result<StoredUsageAnalytics, DataLayerError> {
|
||||
let mut tx = self.pool.begin().await.map_postgres_err()?;
|
||||
sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
sqlx::query("SET LOCAL statement_timeout = '15s'")
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
sqlx::query("SET LOCAL jit = off")
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
let state = sqlx::query("SELECT pg_current_snapshot()::text AS revision, NOW() AS generated_at, (SELECT count(*) FROM users WHERE is_active AND NOT is_deleted) AS enabled_users")
|
||||
.fetch_one(&mut *tx).await.map_postgres_err()?;
|
||||
let granularity = match query.granularity {
|
||||
UsageAnalyticsGranularity::Hour => "hour",
|
||||
UsageAnalyticsGranularity::Day => "day",
|
||||
};
|
||||
let timezone = match query.granularity {
|
||||
UsageAnalyticsGranularity::Hour => "UTC",
|
||||
UsageAnalyticsGranularity::Day => query.timezone.as_str(),
|
||||
};
|
||||
let compare_history = can_compare_daily_history(query);
|
||||
// Like dashboard_request_fact, coverage describes customer charges;
|
||||
// a rated request with an invalid factor snapshot is still unpriced.
|
||||
let metrics_sql = if compare_history {
|
||||
DASHBOARD_CHART_METRICS_SQL.to_string()
|
||||
} else {
|
||||
analytics_metrics_sql(query.slow_threshold_ms.unwrap_or(5000)).replace(
|
||||
"count(*) FILTER (WHERE pricing_available)::bigint AS pricing_available_count",
|
||||
"count(billable_amount)::bigint AS pricing_available_count",
|
||||
)
|
||||
};
|
||||
let from = DateTime::<Utc>::from_timestamp_millis(query.from_unix_ms as i64)
|
||||
.expect("validated timestamp");
|
||||
let to = DateTime::<Utc>::from_timestamp_millis(query.to_unix_ms as i64)
|
||||
.expect("validated timestamp");
|
||||
let source_from = from
|
||||
.date_naive()
|
||||
.and_hms_opt(0, 0, 0)
|
||||
.expect("UTC midnight")
|
||||
.and_utc();
|
||||
let source_to_floor = to
|
||||
.date_naive()
|
||||
.and_hms_opt(0, 0, 0)
|
||||
.expect("UTC midnight")
|
||||
.and_utc();
|
||||
let source_to = if source_to_floor == to {
|
||||
to
|
||||
} else {
|
||||
source_to_floor + Duration::days(1)
|
||||
};
|
||||
let mut builder = QueryBuilder::<Postgres>::new("WITH ");
|
||||
if compare_history {
|
||||
// Legacy daily totals include sessions. Keep them in this one raw
|
||||
// scan for the coverage check, then exclude them from every chart.
|
||||
builder
|
||||
.push("source AS MATERIALIZED (SELECT * FROM ")
|
||||
.push(DASHBOARD_CHART_FACTS_SQL)
|
||||
.push(" WHERE created_at >= ")
|
||||
.push_bind(source_from)
|
||||
.push(" AND created_at < ")
|
||||
.push_bind(source_to)
|
||||
.push("), ");
|
||||
}
|
||||
builder.push("filtered AS MATERIALIZED (SELECT *, date_trunc(");
|
||||
builder
|
||||
.push_bind(granularity)
|
||||
.push(", created_at AT TIME ZONE ")
|
||||
.push_bind(timezone)
|
||||
.push(") AT TIME ZONE ")
|
||||
.push_bind(timezone)
|
||||
.push(" AS bucket, COALESCE(NULLIF(provider_id,''),NULLIF(provider_name,'')) AS chart_provider_id");
|
||||
push_analytics_source_filter(
|
||||
&mut builder,
|
||||
query,
|
||||
if compare_history {
|
||||
"source"
|
||||
} else {
|
||||
"public.usage_analytics_facts_v1"
|
||||
},
|
||||
);
|
||||
builder
|
||||
.push("), summary AS (SELECT ")
|
||||
.push(&metrics_sql)
|
||||
.push(" FROM filtered), series AS (SELECT bucket, ")
|
||||
.push(&metrics_sql)
|
||||
.push(", count(DISTINCT COALESCE(NULLIF(provider_id,''),NULLIF(provider_name,''))) FILTER (WHERE NULLIF(provider_id,'') IS NOT NULL OR provider_name NOT IN ('unknown','pending'))::bigint AS unique_providers FROM filtered GROUP BY bucket ORDER BY bucket LIMIT 10001), models AS (SELECT model, bucket, ")
|
||||
.push(&metrics_sql)
|
||||
.push(" FROM filtered GROUP BY model,bucket ORDER BY bucket,model LIMIT 10001), providers AS (SELECT chart_provider_id, max(provider_name) AS provider_label, ")
|
||||
.push(&metrics_sql)
|
||||
.push(" FROM filtered GROUP BY chart_provider_id ORDER BY count(*) DESC,chart_provider_id LIMIT 10001)");
|
||||
if compare_history {
|
||||
builder.push(r#", historical_cutoff AS (
|
||||
SELECT CASE WHEN pg_typeof(cutoff_date) IN ('bigint'::regtype,'integer'::regtype)
|
||||
THEN to_timestamp(cutoff_date::text::double precision)
|
||||
ELSE cutoff_date::text::timestamptz END AS cutoff
|
||||
FROM stats_summary ORDER BY updated_at DESC,created_at DESC LIMIT 1
|
||||
), historical_days AS (
|
||||
SELECT CASE WHEN pg_typeof(date) IN ('bigint'::regtype,'integer'::regtype)
|
||||
THEN to_timestamp(date::text::double precision)
|
||||
ELSE date::text::timestamptz END AS day,total_requests
|
||||
FROM stats_daily WHERE is_complete AND total_requests>0
|
||||
), retained_days AS (
|
||||
SELECT date_trunc('day',created_at AT TIME ZONE 'UTC') AT TIME ZONE 'UTC' AS day,
|
||||
count(*) AS requests
|
||||
FROM source
|
||||
WHERE status NOT IN ('pending','streaming') AND provider_name NOT IN ('unknown','pending')
|
||||
GROUP BY 1
|
||||
), lost_days AS (
|
||||
SELECT d.day FROM historical_days d CROSS JOIN historical_cutoff c
|
||||
LEFT JOIN retained_days r ON r.day=d.day
|
||||
WHERE d.day=date_trunc('day',d.day AT TIME ZONE 'UTC') AT TIME ZONE 'UTC'
|
||||
AND d.day < "#)
|
||||
.push_bind(to)
|
||||
.push(" AND d.day+INTERVAL '24 hours' > ")
|
||||
.push_bind(from)
|
||||
.push(" AND d.day+INTERVAL '24 hours' <= c.cutoff AND d.total_requests>COALESCE(r.requests,0))");
|
||||
}
|
||||
// An archived day with missing details has unknown coverage in each
|
||||
// constituent hour. Union with existing markers to avoid double-counting.
|
||||
builder.push(
|
||||
r#", lost_hours AS (
|
||||
SELECT bucket_start FROM stats_bucket_state
|
||||
WHERE projection_version IN ('overview-v1','overview-v2')
|
||||
AND granularity='hour' AND coverage_status='unrecoverable'
|
||||
UNION SELECT bucket_start FROM stats_overview_dirty_events
|
||||
WHERE projection_version='overview-v2' AND granularity='hour' AND unrecoverable
|
||||
"#,
|
||||
);
|
||||
if compare_history {
|
||||
builder.push(" UNION SELECT hour FROM lost_days CROSS JOIN LATERAL generate_series(day,day+INTERVAL '23 hours',INTERVAL '1 hour') AS hours(hour)");
|
||||
}
|
||||
builder.push("), lost AS (SELECT count(*)::bigint AS hours FROM lost_hours WHERE bucket_start < ")
|
||||
.push_bind(to)
|
||||
.push(" AND bucket_start+INTERVAL '1 hour' > ")
|
||||
.push_bind(from)
|
||||
.push(")")
|
||||
.push(r#"
|
||||
SELECT
|
||||
(SELECT hours FROM lost) AS unrecoverable_bucket_count,
|
||||
(SELECT to_jsonb(s) FROM summary s) AS summary,
|
||||
(SELECT COALESCE(jsonb_agg(jsonb_build_object(
|
||||
'id', bucket::text, 'label', bucket::text,
|
||||
'bucket_start', to_char(bucket AT TIME ZONE 'UTC','YYYY-MM-DD"T"HH24:MI:SS"Z"'),
|
||||
'metrics', to_jsonb(s)-'bucket') ORDER BY bucket),'[]'::jsonb) FROM series s) AS series,
|
||||
(SELECT COALESCE(jsonb_agg(jsonb_build_object(
|
||||
'id',model,'label',model,
|
||||
'bucket_start',to_char(bucket AT TIME ZONE 'UTC','YYYY-MM-DD"T"HH24:MI:SS"Z"'),
|
||||
'metrics',to_jsonb(m)-'model'-'bucket') ORDER BY bucket,model),'[]'::jsonb) FROM models m) AS models,
|
||||
(SELECT COALESCE(jsonb_agg(jsonb_build_object(
|
||||
'id',chart_provider_id,'label',provider_label,'bucket_start',NULL,
|
||||
'metrics',to_jsonb(p)-'chart_provider_id'-'provider_label') ORDER BY request_count DESC,chart_provider_id),'[]'::jsonb) FROM providers p) AS providers
|
||||
"#);
|
||||
let row = builder
|
||||
.build()
|
||||
.fetch_one(&mut *tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
let mut result = StoredUsageAnalytics {
|
||||
summary: decode(row.try_get("summary").map_postgres_err()?)?,
|
||||
rows: decode(row.try_get("series").map_postgres_err()?)?,
|
||||
model_rows: decode(row.try_get("models").map_postgres_err()?)?,
|
||||
provider_rows: decode(row.try_get("providers").map_postgres_err()?)?,
|
||||
unrecoverable_bucket_count: row
|
||||
.try_get::<i64, _>("unrecoverable_bucket_count")
|
||||
.map_postgres_err()? as u64,
|
||||
read_revision: state.try_get("revision").map_postgres_err()?,
|
||||
generated_at: state
|
||||
.try_get::<DateTime<Utc>, _>("generated_at")
|
||||
.map_postgres_err()?
|
||||
.to_rfc3339(),
|
||||
..Default::default()
|
||||
};
|
||||
result.summary.enabled_users = state
|
||||
.try_get::<i64, _>("enabled_users")
|
||||
.map_postgres_err()? as u64;
|
||||
for (dimension, count) in [
|
||||
("daily", result.rows.len()),
|
||||
("model", result.model_rows.len()),
|
||||
("provider", result.provider_rows.len()),
|
||||
] {
|
||||
if count > USAGE_DASHBOARD_CHART_ROW_LIMIT {
|
||||
return Err(DataLayerError::InvalidInput(format!(
|
||||
"dashboard {dimension} chart exceeds 10000 groups; narrow the range"
|
||||
)));
|
||||
}
|
||||
}
|
||||
result.coverage = read_projection_coverage(&mut tx, query, false).await?;
|
||||
fill_usage_analytics_timeseries(query, &mut result.rows);
|
||||
// Missing raw history is not evidence of zero spend. Keep populated
|
||||
// subtotals, but do not present filled empty buckets as measured zeros.
|
||||
if result.unrecoverable_bucket_count > 0 {
|
||||
for row in &mut result.rows {
|
||||
if row.metrics.request_count == 0 {
|
||||
row.metrics.rated_amount = None;
|
||||
row.metrics.billable_amount = None;
|
||||
}
|
||||
}
|
||||
} else if result.summary.request_count == 0 {
|
||||
result.summary.rated_amount = Some("0.00000000".into());
|
||||
result.summary.billable_amount = Some("0.00000000".into());
|
||||
}
|
||||
result.total = result.rows.len() as u64;
|
||||
tx.commit().await.map_postgres_err()?;
|
||||
Ok(result)
|
||||
}
|
||||
}
|
||||
@@ -66,6 +66,7 @@ mod analytics_tests;
|
||||
mod attribution;
|
||||
pub mod cleanup;
|
||||
mod dashboard;
|
||||
mod dashboard_charts;
|
||||
mod dashboard_history;
|
||||
#[cfg(test)]
|
||||
mod dashboard_history_tests;
|
||||
|
||||
@@ -196,6 +196,8 @@ impl UsageAnalyticsQuery {
|
||||
#[serde(default)]
|
||||
pub struct UsageAnalyticsMetrics {
|
||||
pub request_count: u64,
|
||||
// Populated for dashboard chart buckets; distinct counts are not additive.
|
||||
pub unique_providers: Option<u64>,
|
||||
pub successful_request_count: u64,
|
||||
pub failed_request_count: u64,
|
||||
pub cancelled_request_count: u64,
|
||||
@@ -429,6 +431,7 @@ pub fn fill_usage_analytics_timeseries(
|
||||
label: Some(start.clone()),
|
||||
bucket_start: Some(start),
|
||||
metrics: UsageAnalyticsMetrics {
|
||||
unique_providers: Some(0),
|
||||
rated_amount: Some("0.00000000".into()),
|
||||
billable_amount: Some("0.00000000".into()),
|
||||
..Default::default()
|
||||
|
||||
@@ -205,7 +205,7 @@ CREATE TRIGGER overview_usage_delete_attribution BEFORE DELETE ON public.usage
|
||||
-- procurement cost remains in actual_total_cost_usd for legacy reporting.
|
||||
CREATE OR REPLACE FUNCTION public.usage_customer_billable_amount(
|
||||
metadata jsonb, base_cost numeric, legacy_cost numeric
|
||||
) RETURNS numeric LANGUAGE plpgsql IMMUTABLE PARALLEL SAFE AS $$
|
||||
) RETURNS numeric LANGUAGE plpgsql IMMUTABLE PARALLEL UNSAFE AS $$
|
||||
DECLARE factor jsonb; multiplier numeric; amount numeric;
|
||||
factor_name text; factor_value jsonb; factor_number double precision;
|
||||
expected_multiplier double precision := 1.0; factor_count integer := 0;
|
||||
|
||||
@@ -28,6 +28,8 @@ use crate::lifecycle::bootstrap::postgres::{
|
||||
};
|
||||
|
||||
mod customer_billing_upgrade;
|
||||
mod dashboard_chart_billing;
|
||||
mod dashboard_chart_history;
|
||||
mod dashboard_user_anonymization;
|
||||
mod legacy_overview_upgrade;
|
||||
mod migration_deadlines;
|
||||
@@ -1600,6 +1602,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() {
|
||||
20261001000000,
|
||||
20261004000000,
|
||||
20261007000000,
|
||||
20261008000000,
|
||||
20261009000000,
|
||||
]
|
||||
);
|
||||
|
||||
@@ -0,0 +1,278 @@
|
||||
use super::*;
|
||||
use aether_data_contracts::repository::usage::{
|
||||
UsageAnalyticsQuery, UsageAnalyticsView, UsageDashboardAnalyticsQuery,
|
||||
};
|
||||
use chrono::{Duration, Utc};
|
||||
use serde_json::json;
|
||||
|
||||
#[tokio::test]
|
||||
async fn dashboard_charts_match_summary_customer_charges_and_request_scope() {
|
||||
let Some(server) = ManagedPostgresServer::try_start().await.unwrap() else {
|
||||
return;
|
||||
};
|
||||
let pool = PgPool::connect(server.database_url()).await.unwrap();
|
||||
prepare_and_apply_clean_postgres_database(&pool).await;
|
||||
let repo = aether_data_postgres::SqlxUsageReadRepository::new(pool.clone());
|
||||
let summary_query = UsageDashboardAnalyticsQuery {
|
||||
timezone: "Asia/Shanghai".into(),
|
||||
};
|
||||
let start = summary_query.today_start(Utc::now()).unwrap();
|
||||
query("UPDATE dashboard_stats_state SET stats_since=$1")
|
||||
.bind(start)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
let composite = json!({"billing_multiplier_snapshot":{"version":1,"factors":{"routing_group":2,"user_group":0.25},"multiplier":0.5}});
|
||||
for (id, provider, status, billing, base, actual, metadata) in [
|
||||
(
|
||||
"composite",
|
||||
"alpha",
|
||||
"completed",
|
||||
"settled",
|
||||
10.0,
|
||||
Some(3.0),
|
||||
composite,
|
||||
),
|
||||
(
|
||||
"legacy",
|
||||
"alpha",
|
||||
"completed",
|
||||
"settled",
|
||||
2.0,
|
||||
Some(1.25),
|
||||
json!({}),
|
||||
),
|
||||
(
|
||||
"free",
|
||||
"alpha",
|
||||
"completed",
|
||||
"settled",
|
||||
50.0,
|
||||
Some(40.0),
|
||||
json!({"routing_group_billing_multiplier":0}),
|
||||
),
|
||||
(
|
||||
"invalid",
|
||||
"alpha",
|
||||
"completed",
|
||||
"settled",
|
||||
99.0,
|
||||
Some(90.0),
|
||||
json!({"billing_multiplier_snapshot":null}),
|
||||
),
|
||||
(
|
||||
"unpriced",
|
||||
"beta",
|
||||
"completed",
|
||||
"settled",
|
||||
80.0,
|
||||
Some(70.0),
|
||||
json!({"usage_pricing_available":false}),
|
||||
),
|
||||
(
|
||||
"failed",
|
||||
"unknown",
|
||||
"failed",
|
||||
"settled",
|
||||
0.0,
|
||||
Some(0.0),
|
||||
json!({}),
|
||||
),
|
||||
(
|
||||
"pending",
|
||||
"pending",
|
||||
"pending",
|
||||
"pending",
|
||||
0.0,
|
||||
None,
|
||||
json!({}),
|
||||
),
|
||||
(
|
||||
"streaming",
|
||||
"beta",
|
||||
"streaming",
|
||||
"pending",
|
||||
0.0,
|
||||
None,
|
||||
json!({}),
|
||||
),
|
||||
(
|
||||
"session",
|
||||
"alpha",
|
||||
"completed",
|
||||
"settled",
|
||||
500.0,
|
||||
Some(500.0),
|
||||
json!({}),
|
||||
),
|
||||
] {
|
||||
query("INSERT INTO usage(id,request_id,model,provider_name,status,billing_status,total_cost_usd,actual_total_cost_usd,total_tokens,created_at,request_metadata,response_time_ms) VALUES($1,$1,'model',$2,$3,$4,$5,$6,10,$7,$8,1000)")
|
||||
.bind(id).bind(provider).bind(status).bind(billing).bind(base).bind(actual)
|
||||
.bind(start + Duration::seconds(1)).bind(metadata).execute(&pool).await.unwrap();
|
||||
}
|
||||
// The captured settlement cost, rather than the mutable audit float, is rated.
|
||||
query("INSERT INTO usage_settlement_snapshots(request_id,billing_status,billing_total_cost_usd,billing_actual_total_cost_usd) VALUES('composite','settled',12,4)")
|
||||
.execute(&pool).await.unwrap();
|
||||
query(
|
||||
"UPDATE usage_attribution_snapshots SET record_kind='session' WHERE request_id='session'",
|
||||
)
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
// Exclude both the previous local day and the next day's boundary.
|
||||
for (id, at) in [
|
||||
("before", start - Duration::seconds(1)),
|
||||
("after", start + Duration::days(1)),
|
||||
] {
|
||||
query("INSERT INTO usage(id,request_id,model,provider_name,status,billing_status,total_cost_usd,actual_total_cost_usd,created_at) VALUES($1,$1,'outside','outside','completed','settled',999,999,$2)")
|
||||
.bind(id).bind(at).execute(&pool).await.unwrap();
|
||||
}
|
||||
let summary = repo.query_dashboard_summary(&summary_query).await.unwrap();
|
||||
let charts = repo
|
||||
.query_usage_analytics(&UsageAnalyticsQuery {
|
||||
from_unix_ms: start.timestamp_millis() as u64,
|
||||
to_unix_ms: (start + Duration::days(1)).timestamp_millis() as u64,
|
||||
timezone: summary_query.timezone.clone(),
|
||||
view: UsageAnalyticsView::DashboardCharts,
|
||||
limit: 10_000,
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(summary.today.request_count, 8);
|
||||
assert_eq!(summary.today.billable_amount.as_deref(), Some("7.25000000"));
|
||||
assert_eq!(charts.summary.request_count, summary.today.request_count);
|
||||
assert_eq!(
|
||||
charts.summary.billable_amount,
|
||||
summary.today.billable_amount
|
||||
);
|
||||
assert_eq!(
|
||||
charts.summary.pricing_available_count,
|
||||
summary.today.pricing_available_count
|
||||
);
|
||||
assert_eq!(charts.summary.in_flight_request_count, 2);
|
||||
assert_eq!(charts.rows.len(), 1);
|
||||
assert_eq!(charts.rows[0].metrics.request_count, 8);
|
||||
assert_eq!(
|
||||
charts.rows[0].metrics.billable_amount,
|
||||
summary.today.billable_amount
|
||||
);
|
||||
assert_eq!(charts.rows[0].metrics.unique_providers, Some(2));
|
||||
assert_eq!(
|
||||
charts.rows[0].bucket_start.as_deref(),
|
||||
Some(
|
||||
start
|
||||
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
|
||||
.as_str()
|
||||
)
|
||||
);
|
||||
for rows in [&charts.model_rows, &charts.provider_rows] {
|
||||
assert_eq!(rows.iter().map(|r| r.metrics.request_count).sum::<u64>(), 8);
|
||||
let charges: f64 = rows
|
||||
.iter()
|
||||
.filter_map(|r| r.metrics.billable_amount.as_ref())
|
||||
.map(|v| v.parse::<f64>().unwrap())
|
||||
.sum();
|
||||
assert_eq!(charges, 7.25);
|
||||
}
|
||||
assert_eq!(
|
||||
charts.provider_rows.len(),
|
||||
4,
|
||||
"legacy provider names must not collapse into one null-ID group"
|
||||
);
|
||||
// Retained older requests can have no contribution ledger entry. Mixing
|
||||
// fallback rows with current cached rows must keep unknown pricing distinct
|
||||
// from free usage and preserve the immutable composite settlement amount.
|
||||
query("DELETE FROM dashboard_request_contributions WHERE request_id IN ('composite','free','invalid')")
|
||||
.execute(&pool).await.unwrap();
|
||||
let mixed = repo
|
||||
.query_usage_analytics(&UsageAnalyticsQuery {
|
||||
from_unix_ms: start.timestamp_millis() as u64,
|
||||
to_unix_ms: (start + Duration::days(1)).timestamp_millis() as u64,
|
||||
timezone: summary_query.timezone.clone(),
|
||||
view: UsageAnalyticsView::DashboardCharts,
|
||||
limit: 10_000,
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
mixed.summary.billable_amount,
|
||||
charts.summary.billable_amount
|
||||
);
|
||||
assert_eq!(mixed.summary.total_tokens, charts.summary.total_tokens);
|
||||
assert_eq!(mixed.summary.request_count, charts.summary.request_count);
|
||||
assert_eq!(
|
||||
mixed.summary.pricing_available_count,
|
||||
charts.summary.pricing_available_count
|
||||
);
|
||||
assert_eq!(mixed.rows.len(), charts.rows.len());
|
||||
for (mixed, cached) in mixed.rows.iter().zip(&charts.rows) {
|
||||
assert_eq!(mixed.bucket_start, cached.bucket_start);
|
||||
assert_eq!(
|
||||
mixed.metrics.billable_amount,
|
||||
cached.metrics.billable_amount
|
||||
);
|
||||
assert_eq!(mixed.metrics.total_tokens, cached.metrics.total_tokens);
|
||||
assert_eq!(
|
||||
mixed.metrics.pricing_available_count,
|
||||
cached.metrics.pricing_available_count
|
||||
);
|
||||
}
|
||||
let empty = repo
|
||||
.query_usage_analytics(&UsageAnalyticsQuery {
|
||||
from_unix_ms: (start - Duration::days(2)).timestamp_millis() as u64,
|
||||
to_unix_ms: (start - Duration::days(1)).timestamp_millis() as u64,
|
||||
timezone: summary_query.timezone,
|
||||
view: UsageAnalyticsView::DashboardCharts,
|
||||
limit: 10_000,
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(empty.summary.request_count, 0);
|
||||
assert_eq!(empty.summary.billable_amount.as_deref(), Some("0.00000000"));
|
||||
assert_eq!(empty.rows[0].metrics.unique_providers, Some(0));
|
||||
pool.close().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn customer_billing_parallel_fix_preserves_history_without_backfill() {
|
||||
let Some(server) = ManagedPostgresServer::try_start().await.unwrap() else {
|
||||
return;
|
||||
};
|
||||
let mut connection = PgConnection::connect(server.database_url()).await.unwrap();
|
||||
connection.ensure_migrations_table().await.unwrap();
|
||||
for migration in POSTGRES_MIGRATOR
|
||||
.iter()
|
||||
.filter(|m| m.version < 20261008000000)
|
||||
{
|
||||
connection.apply(migration).await.unwrap();
|
||||
}
|
||||
let pool = PgPool::connect(server.database_url()).await.unwrap();
|
||||
query("INSERT INTO stats_daily(id,date,total_requests,total_cost,actual_total_cost,is_complete) VALUES('untouched','2020-01-01',1,10,3,true)")
|
||||
.execute(&pool).await.unwrap();
|
||||
query("INSERT INTO usage(id,request_id,model,provider_name,status,billing_status,total_cost_usd,actual_total_cost_usd,created_at,request_metadata) VALUES('untouched','untouched','m','p','completed','settled',10,3,'2020-01-01','{\"routing_group_billing_multiplier\":2}')")
|
||||
.execute(&pool).await.unwrap();
|
||||
let before: (serde_json::Value,serde_json::Value) = sqlx::query_as("SELECT (SELECT to_jsonb(d) FROM stats_daily d WHERE id='untouched'), (SELECT to_jsonb(u) FROM usage u WHERE id='untouched')").fetch_one(&pool).await.unwrap();
|
||||
connection
|
||||
.apply(
|
||||
POSTGRES_MIGRATOR
|
||||
.iter()
|
||||
.find(|m| m.version == 20261008000000)
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let after = sqlx::query_as("SELECT (SELECT to_jsonb(d) FROM stats_daily d WHERE id='untouched'), (SELECT to_jsonb(u) FROM usage u WHERE id='untouched')").fetch_one(&pool).await.unwrap();
|
||||
assert_eq!(before, after);
|
||||
let parallel: String=query_scalar("SELECT proparallel::text FROM pg_proc WHERE oid='public.usage_customer_billable_amount(jsonb,numeric,numeric)'::regprocedure").fetch_one(&pool).await.unwrap();
|
||||
assert_eq!(parallel, "u");
|
||||
// Encourage a parallel scan; the exception-handling function must keep it serial.
|
||||
let mut tx = pool.begin().await.unwrap();
|
||||
sqlx::raw_sql("CREATE TABLE billing_parallel_probe AS SELECT i::numeric AS cost FROM generate_series(1,10000) i; ALTER TABLE billing_parallel_probe SET (parallel_workers=2); ANALYZE billing_parallel_probe; SET LOCAL min_parallel_table_scan_size=0; SET LOCAL parallel_setup_cost=0; SET LOCAL parallel_tuple_cost=0;").execute(&mut *tx).await.unwrap();
|
||||
let charge:String=query_scalar("SELECT sum(public.usage_customer_billable_amount('{\"routing_group_billing_multiplier\":2}'::jsonb,cost,1))::text FROM billing_parallel_probe").fetch_one(&mut *tx).await.unwrap();
|
||||
assert_eq!(charge, "100010000.00000000");
|
||||
tx.rollback().await.unwrap();
|
||||
pool.close().await;
|
||||
}
|
||||
@@ -0,0 +1,144 @@
|
||||
use super::*;
|
||||
use aether_data_contracts::repository::usage::{UsageAnalyticsQuery, UsageAnalyticsView};
|
||||
use chrono::{DateTime, Duration, Utc};
|
||||
use serde_json::json;
|
||||
|
||||
#[tokio::test]
|
||||
async fn dashboard_chart_history_coverage_matches_legacy_request_scope_without_raw_rescans() {
|
||||
let Some(server) = ManagedPostgresServer::try_start().await.unwrap() else {
|
||||
return;
|
||||
};
|
||||
let pool = PgPool::connect(server.database_url()).await.unwrap();
|
||||
prepare_and_apply_clean_postgres_database(&pool).await;
|
||||
let repo = aether_data_postgres::SqlxUsageReadRepository::new(pool.clone());
|
||||
let start = "2020-01-01T00:00:00Z".parse::<DateTime<Utc>>().unwrap();
|
||||
query("INSERT INTO stats_summary(id,cutoff_date) VALUES('history',$1)")
|
||||
.bind(start + Duration::days(4))
|
||||
.execute(&pool)
|
||||
.await
|
||||
.unwrap();
|
||||
for (day, count, complete) in [
|
||||
(0, 2, true),
|
||||
(1, 2, true),
|
||||
(2, 1, true),
|
||||
(3, 1, false),
|
||||
(4, 1, true),
|
||||
] {
|
||||
query("INSERT INTO stats_daily(id,date,total_requests,total_cost,actual_total_cost,is_complete) VALUES($1,$2,$3,9999,7777,$4)")
|
||||
.bind(format!("day-{day}"))
|
||||
.bind(start + Duration::days(day))
|
||||
.bind(count)
|
||||
.bind(complete)
|
||||
.execute(&pool).await.unwrap();
|
||||
}
|
||||
for (id, day, status, provider, session) in [
|
||||
("retained-request", 0, "completed", "provider", false),
|
||||
("retained-session", 0, "completed", "provider", true),
|
||||
("partial-request", 1, "completed", "provider", false),
|
||||
("partial-pending", 1, "pending", "provider", false),
|
||||
("partial-unknown", 1, "failed", "unknown", false),
|
||||
] {
|
||||
query("INSERT INTO usage(id,request_id,model,provider_name,status,billing_status,total_cost_usd,actual_total_cost_usd,created_at,request_metadata) VALUES($1,$1,'model',$2,$3,'settled',2,1,$4,$5)")
|
||||
.bind(id).bind(provider).bind(status)
|
||||
.bind(start + Duration::days(day) + Duration::seconds(1))
|
||||
.bind(json!({"routing_group_billing_multiplier":2,"analytics_attribution":{"record_kind":if session { "session" } else { "request" }}}))
|
||||
.execute(&pool).await.unwrap();
|
||||
}
|
||||
query("UPDATE usage SET total_tokens=999,input_tokens=90,output_tokens=20,cache_creation_input_tokens=11,cache_read_input_tokens=13 WHERE id='retained-request'")
|
||||
.execute(&pool).await.unwrap();
|
||||
query("INSERT INTO usage_settlement_snapshots(request_id,billing_status,billing_effective_input_tokens,billing_output_tokens,billing_cache_creation_tokens,billing_cache_read_tokens) VALUES('retained-request','settled',7,3,5,2)")
|
||||
.execute(&pool).await.unwrap();
|
||||
let day_query = |day| UsageAnalyticsQuery {
|
||||
from_unix_ms: (start + Duration::days(day)).timestamp_millis() as u64,
|
||||
to_unix_ms: (start + Duration::days(day + 1)).timestamp_millis() as u64,
|
||||
timezone: "UTC".into(),
|
||||
view: UsageAnalyticsView::DashboardCharts,
|
||||
limit: 10_000,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
// A retained session belongs to the old rollup, but never to chart totals.
|
||||
let complete = repo.query_usage_analytics(&day_query(0)).await.unwrap();
|
||||
assert_eq!(complete.unrecoverable_bucket_count, 0);
|
||||
assert_eq!(complete.summary.request_count, 1);
|
||||
assert_eq!(complete.summary.total_tokens, 17);
|
||||
assert_eq!(
|
||||
complete.summary.billable_amount.as_deref(),
|
||||
Some("4.00000000")
|
||||
);
|
||||
let mut canonical_query = day_query(0);
|
||||
canonical_query.model = Some("model".into());
|
||||
let canonical = repo.query_usage_analytics(&canonical_query).await.unwrap();
|
||||
assert_eq!(
|
||||
complete.summary.total_tokens,
|
||||
canonical.summary.total_tokens
|
||||
);
|
||||
assert_eq!(
|
||||
complete.summary.billable_amount,
|
||||
canonical.summary.billable_amount
|
||||
);
|
||||
|
||||
// Pending and unknown-provider rows cannot disguise a missing legacy request.
|
||||
let partial = repo.query_usage_analytics(&day_query(1)).await.unwrap();
|
||||
assert_eq!(partial.unrecoverable_bucket_count, 24);
|
||||
assert_eq!(partial.summary.request_count, 3);
|
||||
assert_eq!(
|
||||
partial.summary.billable_amount.as_deref(),
|
||||
Some("12.00000000")
|
||||
);
|
||||
|
||||
let missing = repo.query_usage_analytics(&day_query(2)).await.unwrap();
|
||||
assert_eq!(missing.unrecoverable_bucket_count, 24);
|
||||
assert_eq!(missing.summary.request_count, 0);
|
||||
assert_eq!(missing.summary.billable_amount, None);
|
||||
assert_eq!(missing.rows[0].metrics.billable_amount, None);
|
||||
|
||||
// Incomplete rollups, unpublished days and filtered views cannot establish
|
||||
// that request details are missing from global legacy daily totals.
|
||||
for day in [3, 4] {
|
||||
let result = repo.query_usage_analytics(&day_query(day)).await.unwrap();
|
||||
assert_eq!(result.unrecoverable_bucket_count, 0);
|
||||
assert_eq!(
|
||||
result.summary.billable_amount.as_deref(),
|
||||
Some("0.00000000")
|
||||
);
|
||||
}
|
||||
let mut filtered_query = day_query(2);
|
||||
filtered_query.model = Some("model".into());
|
||||
let filtered = repo.query_usage_analytics(&filtered_query).await.unwrap();
|
||||
assert_eq!(filtered.unrecoverable_bucket_count, 0);
|
||||
let mut partial_day_query = day_query(2);
|
||||
partial_day_query.from_unix_ms += 60 * 60 * 1000;
|
||||
let partial_day = repo
|
||||
.query_usage_analytics(&partial_day_query)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(partial_day.unrecoverable_bucket_count, 23);
|
||||
assert_eq!(partial_day.summary.billable_amount, None);
|
||||
// A local day crosses two UTC archive days. Read both complete UTC days for
|
||||
// coverage, while chart totals remain bounded to the original local range.
|
||||
let local_day_query = UsageAnalyticsQuery {
|
||||
from_unix_ms: (start + Duration::days(1) + Duration::hours(16)).timestamp_millis() as u64,
|
||||
to_unix_ms: (start + Duration::days(2) + Duration::hours(16)).timestamp_millis() as u64,
|
||||
timezone: "Asia/Shanghai".into(),
|
||||
..day_query(2)
|
||||
};
|
||||
let local_day = repo.query_usage_analytics(&local_day_query).await.unwrap();
|
||||
assert_eq!(local_day.unrecoverable_bucket_count, 24);
|
||||
assert_eq!(local_day.summary.request_count, 0);
|
||||
assert_eq!(local_day.rows.len(), 1);
|
||||
assert_eq!(local_day.rows[0].metrics.billable_amount, None);
|
||||
let mut retained_partial = day_query(0);
|
||||
retained_partial.from_unix_ms += 60 * 60 * 1000;
|
||||
let retained_partial = repo.query_usage_analytics(&retained_partial).await.unwrap();
|
||||
assert_eq!(retained_partial.unrecoverable_bucket_count, 0);
|
||||
assert_eq!(retained_partial.summary.request_count, 0);
|
||||
|
||||
// Existing lost-hour evidence and the inferred day's coverage are one set.
|
||||
query("INSERT INTO stats_overview_dirty_events(transaction_id,projection_version,granularity,bucket_start,unrecoverable) VALUES(txid_current(),'overview-v2','hour',$1,true)")
|
||||
.bind(start + Duration::days(2) + Duration::hours(3))
|
||||
.execute(&pool).await.unwrap();
|
||||
let deduplicated = repo.query_usage_analytics(&day_query(2)).await.unwrap();
|
||||
assert_eq!(deduplicated.unrecoverable_bucket_count, 24);
|
||||
pool.close().await;
|
||||
}
|
||||
@@ -116,6 +116,7 @@ WHERE version=20260919000000;
|
||||
20261001000000,
|
||||
20261004000000,
|
||||
20261007000000,
|
||||
20261008000000,
|
||||
20261009000000,
|
||||
]
|
||||
);
|
||||
|
||||
@@ -435,6 +435,20 @@ impl InMemoryUsageReadRepository {
|
||||
let metrics = |rows: &[&StoredRequestUsageAudit], slow| {
|
||||
let mut result = metrics(rows, slow, &keys);
|
||||
apply_allocations(&mut result, rows, &allocations);
|
||||
if query.view == UsageAnalyticsView::DashboardCharts {
|
||||
result.pricing_available_count = rows
|
||||
.iter()
|
||||
.filter(|row| {
|
||||
row.billing_status == "settled"
|
||||
&& available(row, USAGE_PRICING_AVAILABLE_METADATA_KEY)
|
||||
&& row.billing_cost().is_some()
|
||||
})
|
||||
.count() as u64;
|
||||
if rows.is_empty() {
|
||||
result.billable_amount = Some("0.00000000".into());
|
||||
result.rated_amount = Some("0.00000000".into());
|
||||
}
|
||||
}
|
||||
result
|
||||
};
|
||||
let mut summary = metrics(&filtered, query.slow_threshold_ms.unwrap_or(5000));
|
||||
@@ -682,17 +696,39 @@ impl InMemoryUsageReadRepository {
|
||||
&& query.group_by == UsageAnalyticsGroupBy::Provider;
|
||||
let mut grouped = groups
|
||||
.into_iter()
|
||||
.map(|(id, rows)| UsageAnalyticsRow {
|
||||
label: if provider_breakdown {
|
||||
provider_display_label(&rows, id.as_deref())
|
||||
} else {
|
||||
id.clone()
|
||||
},
|
||||
bucket_start: (query.view != UsageAnalyticsView::Breakdown)
|
||||
.then(|| id.clone())
|
||||
.flatten(),
|
||||
id,
|
||||
metrics: metrics(&rows, query.slow_threshold_ms.unwrap_or(5000)),
|
||||
.map(|(id, rows)| {
|
||||
let mut metrics = metrics(&rows, query.slow_threshold_ms.unwrap_or(5000));
|
||||
if query.view == UsageAnalyticsView::DashboardCharts {
|
||||
metrics.unique_providers = Some(
|
||||
rows.iter()
|
||||
.filter_map(|row| {
|
||||
row.provider_id
|
||||
.as_deref()
|
||||
.filter(|id| !id.is_empty())
|
||||
.or_else(|| {
|
||||
(!matches!(
|
||||
row.provider_name.as_str(),
|
||||
"" | "unknown" | "pending"
|
||||
))
|
||||
.then_some(row.provider_name.as_str())
|
||||
})
|
||||
})
|
||||
.collect::<BTreeSet<_>>()
|
||||
.len() as u64,
|
||||
);
|
||||
}
|
||||
UsageAnalyticsRow {
|
||||
label: if provider_breakdown {
|
||||
provider_display_label(&rows, id.as_deref())
|
||||
} else {
|
||||
id.clone()
|
||||
},
|
||||
bucket_start: (query.view != UsageAnalyticsView::Breakdown)
|
||||
.then(|| id.clone())
|
||||
.flatten(),
|
||||
id,
|
||||
metrics,
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
if query.view == UsageAnalyticsView::Breakdown {
|
||||
@@ -754,7 +790,14 @@ impl InMemoryUsageReadRepository {
|
||||
.today_start(at)?
|
||||
};
|
||||
providers
|
||||
.entry(row.provider_id.clone())
|
||||
.entry(
|
||||
row.provider_id
|
||||
.clone()
|
||||
.filter(|id| !id.is_empty())
|
||||
.or_else(|| {
|
||||
(!row.provider_name.is_empty()).then(|| row.provider_name.clone())
|
||||
}),
|
||||
)
|
||||
.or_default()
|
||||
.push(row);
|
||||
models
|
||||
|
||||
@@ -3,9 +3,11 @@ use std::collections::BTreeMap;
|
||||
use aether_ai_formats::UPSTREAM_IS_STREAM_KEY;
|
||||
use aether_contracts::{ExecutionPlan, ExecutionTelemetry};
|
||||
use aether_data_contracts::repository::usage::{
|
||||
UpsertUsageRecord, UsageBodyCaptureState, LIVE_SESSION_METADATA_KEY,
|
||||
USAGE_AVAILABLE_METADATA_KEY, USAGE_PRICING_AVAILABLE_METADATA_KEY,
|
||||
WEBSOCKET_MODE_METADATA_KEY, WEBSOCKET_TRANSPORT_METADATA_KEY,
|
||||
UpsertUsageRecord, UsageBodyCaptureState, BILLING_MULTIPLIER_SNAPSHOT_METADATA_KEY,
|
||||
LIVE_SESSION_METADATA_KEY, ROUTING_GROUP_BILLING_MULTIPLIER_METADATA_KEY,
|
||||
ROUTING_GROUP_ID_METADATA_KEY, ROUTING_GROUP_NAME_METADATA_KEY, USAGE_AVAILABLE_METADATA_KEY,
|
||||
USAGE_PRICING_AVAILABLE_METADATA_KEY, WEBSOCKET_MODE_METADATA_KEY,
|
||||
WEBSOCKET_TRANSPORT_METADATA_KEY,
|
||||
};
|
||||
use aether_data_contracts::DataLayerError;
|
||||
use serde_json::{json, Map, Value};
|
||||
@@ -2214,6 +2216,23 @@ fn build_runtime_request_metadata_seed_from_parts(
|
||||
provider_request_body_base64: Option<&str>,
|
||||
) -> Option<Value> {
|
||||
let mut metadata = Map::new();
|
||||
// Lifecycle writes need the same immutable routing and billing identity as terminal
|
||||
// writes, without retaining the larger report-context payloads.
|
||||
let routing_snapshot = Map::from_iter(
|
||||
[
|
||||
ROUTING_GROUP_ID_METADATA_KEY,
|
||||
ROUTING_GROUP_NAME_METADATA_KEY,
|
||||
ROUTING_GROUP_BILLING_MULTIPLIER_METADATA_KEY,
|
||||
BILLING_MULTIPLIER_SNAPSHOT_METADATA_KEY,
|
||||
]
|
||||
.into_iter()
|
||||
.filter_map(|key| context_value_ref(context, key).map(|value| (key.into(), value.clone()))),
|
||||
);
|
||||
if let Some(Value::Object(snapshot)) =
|
||||
sanitize_usage_request_metadata(Some(Value::Object(routing_snapshot)))
|
||||
{
|
||||
metadata.extend(snapshot);
|
||||
}
|
||||
for key in ["analytics_attribution", "analytics_failure"] {
|
||||
if let Some(value) = context_value_ref(context, key) {
|
||||
metadata.insert(key.into(), value.clone());
|
||||
@@ -3927,6 +3946,14 @@ mod tests {
|
||||
Some(&json!({
|
||||
"candidate_id": "cand-pending-event-1",
|
||||
"candidate_index": 3,
|
||||
"routing_group_id": "group-free",
|
||||
"routing_group_name": "免费分组",
|
||||
"routing_group_billing_multiplier": 0.0,
|
||||
"billing_multiplier_snapshot": {
|
||||
"version": 1,
|
||||
"factors": {"routing_group": 0.0},
|
||||
"multiplier": 0.0
|
||||
},
|
||||
"websocket_mode": true,
|
||||
"websocket_transport": "responses",
|
||||
"original_request_body": {"messages": [{"content": "omit me"}]},
|
||||
@@ -3951,6 +3978,17 @@ mod tests {
|
||||
assert!(record.provider_request_body.is_none());
|
||||
assert_eq!(record.candidate_id.as_deref(), Some("cand-pending-event-1"));
|
||||
assert_eq!(record.candidate_index, Some(3));
|
||||
let metadata = record.request_metadata.as_ref().expect("routing snapshot");
|
||||
assert_eq!(metadata["routing_group_id"], "group-free");
|
||||
assert_eq!(metadata["routing_group_name"], "免费分组");
|
||||
assert_eq!(metadata["routing_group_billing_multiplier"], 0.0);
|
||||
assert_eq!(
|
||||
aether_data_contracts::repository::usage::billing_multiplier_snapshot(Some(metadata))
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.multiplier(),
|
||||
0.0
|
||||
);
|
||||
assert_eq!(
|
||||
record
|
||||
.request_metadata
|
||||
@@ -3999,6 +4037,14 @@ mod tests {
|
||||
Some(&json!({
|
||||
"candidate_id": "cand-streaming-event-1",
|
||||
"candidate_index": 4,
|
||||
"routing_group_id": "group-discount",
|
||||
"routing_group_name": "折扣分组",
|
||||
"routing_group_billing_multiplier": 0.5,
|
||||
"billing_multiplier_snapshot": {
|
||||
"version": 1,
|
||||
"factors": {"routing_group": 0.5},
|
||||
"multiplier": 0.5
|
||||
},
|
||||
"provider_request_body": {"input": "omit me"}
|
||||
})),
|
||||
),
|
||||
@@ -4023,6 +4069,115 @@ mod tests {
|
||||
Some("cand-streaming-event-1")
|
||||
);
|
||||
assert_eq!(record.candidate_index, Some(4));
|
||||
let metadata = record.request_metadata.as_ref().expect("routing snapshot");
|
||||
assert_eq!(metadata["routing_group_id"], "group-discount");
|
||||
assert_eq!(metadata["routing_group_name"], "折扣分组");
|
||||
assert_eq!(metadata["routing_group_billing_multiplier"], 0.5);
|
||||
assert_eq!(
|
||||
aether_data_contracts::repository::usage::billing_multiplier_snapshot(Some(metadata))
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.multiplier(),
|
||||
0.5
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lifecycle_routing_snapshot_preserves_legacy_factors_and_invalid_markers() {
|
||||
use aether_data_contracts::repository::usage::billing_multiplier_snapshot;
|
||||
|
||||
let plan = ExecutionPlan {
|
||||
request_id: "req-lifecycle-snapshot".to_string(),
|
||||
candidate_id: None,
|
||||
provider_name: Some("OpenAI".to_string()),
|
||||
provider_id: "provider-1".to_string(),
|
||||
endpoint_id: "endpoint-1".to_string(),
|
||||
key_id: "key-1".to_string(),
|
||||
method: "POST".to_string(),
|
||||
url: "https://example.com/v1/responses".to_string(),
|
||||
headers: BTreeMap::new(),
|
||||
content_type: Some("application/json".to_string()),
|
||||
content_encoding: None,
|
||||
body: RequestBody::from_json(json!({"model": "gpt-5.4"})),
|
||||
stream: true,
|
||||
client_api_format: "openai:responses".to_string(),
|
||||
provider_api_format: "openai:responses".to_string(),
|
||||
model_name: Some("gpt-5.4".to_string()),
|
||||
proxy: None,
|
||||
transport_profile: None,
|
||||
timeouts: None,
|
||||
};
|
||||
|
||||
for (snapshot, expected_multiplier) in [
|
||||
(json!({"routing_group_billing_multiplier": 2.0}), Some(2.0)),
|
||||
(json!({"routing_group_billing_multiplier": -1}), None),
|
||||
(
|
||||
json!({
|
||||
"routing_group_billing_multiplier": 0.5,
|
||||
"billing_multiplier_snapshot": null
|
||||
}),
|
||||
None,
|
||||
),
|
||||
(
|
||||
json!({
|
||||
"routing_group_billing_multiplier": 0.5,
|
||||
"billing_multiplier_snapshot": {
|
||||
"version": 1,
|
||||
"factors": {"routing_group": 2.0},
|
||||
"multiplier": 1.0
|
||||
}
|
||||
}),
|
||||
None,
|
||||
),
|
||||
] {
|
||||
let mut context = snapshot;
|
||||
context["routing_group_id"] = json!("group-1");
|
||||
context["routing_group_name"] = json!("请求时分组");
|
||||
context["original_request_body"] = json!({"secret": "do not capture"});
|
||||
context["billing_snapshot"] = json!({"payload": "x".repeat(32 * 1024)});
|
||||
let seed = super::build_lifecycle_usage_seed(&plan, Some(&context));
|
||||
let pending_event =
|
||||
build_pending_usage_event_from_owned_seed(seed.clone(), 1_700_000_000).unwrap();
|
||||
let streaming_event =
|
||||
build_streaming_usage_event_from_owned_seed(seed.clone(), 200, None, 1_700_000_001)
|
||||
.unwrap();
|
||||
let records = [
|
||||
build_pending_usage_record(&plan, Some(&context), 1_700_000_000).unwrap(),
|
||||
build_streaming_usage_record(&plan, Some(&context), 200, None, 1_700_000_001)
|
||||
.unwrap(),
|
||||
build_upsert_usage_record_from_event(&pending_event).unwrap(),
|
||||
build_upsert_usage_record_from_event(&streaming_event).unwrap(),
|
||||
];
|
||||
|
||||
for metadata in std::iter::once(seed.request_metadata.as_ref()).chain(
|
||||
records.iter().map(|record| {
|
||||
assert!(record.request_body.is_none());
|
||||
assert!(record.provider_request_body.is_none());
|
||||
record.request_metadata.as_ref()
|
||||
}),
|
||||
) {
|
||||
let metadata = metadata.expect("lifecycle routing snapshot");
|
||||
assert_eq!(metadata["routing_group_id"], "group-1");
|
||||
assert_eq!(metadata["routing_group_name"], "请求时分组");
|
||||
assert!(metadata.get("original_request_body").is_none());
|
||||
assert!(metadata.get("billing_snapshot").is_none());
|
||||
if let Some(expected) = expected_multiplier {
|
||||
assert_eq!(
|
||||
billing_multiplier_snapshot(Some(metadata))
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.multiplier(),
|
||||
expected
|
||||
);
|
||||
} else {
|
||||
assert_eq!(
|
||||
metadata.get("billing_multiplier_snapshot"),
|
||||
Some(&Value::Null)
|
||||
);
|
||||
assert!(billing_multiplier_snapshot(Some(metadata)).is_err());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user