mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
feat(usage): 新增输出速度统计与请求详情性能分析
- 把 upstream_is_stream 物化到 usage 与 billing facts,避免 Provider 聚合回连 public.usage - 统一前端标准/流式 TPS 计算与显示,流式按首字后生成耗时计算 - 新增请求详情抽屉展示单请求输出速度与 Provider 聚合 TPS
This commit is contained in:
@@ -1048,6 +1048,7 @@ CREATE TABLE IF NOT EXISTS public.usage (
|
||||
request_type character varying(50),
|
||||
api_format character varying(50),
|
||||
is_stream boolean DEFAULT false,
|
||||
upstream_is_stream boolean,
|
||||
status_code integer,
|
||||
error_message text,
|
||||
response_time_ms integer,
|
||||
@@ -4939,13 +4940,18 @@ SELECT
|
||||
) AS price_per_request,
|
||||
settlement.billing_pricing_source,
|
||||
settlement.billing_rule_id,
|
||||
settlement.billing_rule_version
|
||||
settlement.billing_rule_version,
|
||||
COALESCE(usage_rows.upstream_is_stream, COALESCE(usage_rows.is_stream, FALSE)) AS upstream_is_stream
|
||||
FROM public."usage" AS usage_rows
|
||||
LEFT JOIN public.usage_settlement_snapshots AS settlement
|
||||
ON settlement.request_id = usage_rows.request_id;
|
||||
|
||||
COMMENT ON VIEW public.usage_billing_facts IS
|
||||
'Canonical billing read model. Token/cost fields prefer usage_settlement_snapshots.billing_* and fall back to deprecated usage mirrors for legacy rows.';
|
||||
COMMENT ON COLUMN public.usage.upstream_is_stream IS
|
||||
'Resolved upstream stream mode from request_metadata.upstream_is_stream, falling back to is_stream for legacy rows.';
|
||||
COMMENT ON COLUMN public.usage_billing_facts.upstream_is_stream IS
|
||||
'Resolved upstream stream mode from public.usage.upstream_is_stream, falling back to usage.is_stream for legacy rows.';
|
||||
|
||||
COMMENT ON COLUMN public.usage.input_tokens IS
|
||||
'DEPRECATED: billing dimension mirror. Use public.usage_settlement_snapshots.billing_input_tokens or public.usage_billing_facts.input_tokens.';
|
||||
|
||||
@@ -0,0 +1,233 @@
|
||||
ALTER TABLE IF EXISTS public.usage
|
||||
ADD COLUMN IF NOT EXISTS upstream_is_stream boolean;
|
||||
|
||||
UPDATE public.usage
|
||||
SET upstream_is_stream = COALESCE(
|
||||
CASE
|
||||
WHEN (request_metadata->>'upstream_is_stream') IN ('true', 'false')
|
||||
THEN (request_metadata->>'upstream_is_stream')::boolean
|
||||
ELSE NULL
|
||||
END,
|
||||
COALESCE(is_stream, FALSE)
|
||||
)
|
||||
WHERE upstream_is_stream IS NULL;
|
||||
|
||||
ANALYZE public.usage;
|
||||
|
||||
COMMENT ON COLUMN public.usage.upstream_is_stream IS
|
||||
'Resolved upstream stream mode from request_metadata.upstream_is_stream, falling back to is_stream for legacy rows.';
|
||||
|
||||
CREATE OR REPLACE VIEW public.usage_billing_facts AS
|
||||
SELECT
|
||||
usage_rows.id,
|
||||
usage_rows.request_id,
|
||||
usage_rows.user_id,
|
||||
usage_rows.api_key_id,
|
||||
usage_rows.username,
|
||||
usage_rows.api_key_name,
|
||||
usage_rows.provider_name,
|
||||
usage_rows.model,
|
||||
usage_rows.target_model,
|
||||
usage_rows.provider_id,
|
||||
usage_rows.provider_endpoint_id,
|
||||
usage_rows.provider_api_key_id,
|
||||
usage_rows.request_type,
|
||||
usage_rows.api_format,
|
||||
usage_rows.api_family,
|
||||
usage_rows.endpoint_kind,
|
||||
usage_rows.endpoint_api_format,
|
||||
usage_rows.provider_api_family,
|
||||
usage_rows.provider_endpoint_kind,
|
||||
COALESCE(usage_rows.has_format_conversion, FALSE) AS has_format_conversion,
|
||||
COALESCE(usage_rows.is_stream, FALSE) AS is_stream,
|
||||
usage_rows.status_code,
|
||||
usage_rows.error_message,
|
||||
usage_rows.error_category,
|
||||
usage_rows.response_time_ms,
|
||||
usage_rows.first_byte_time_ms,
|
||||
usage_rows.status,
|
||||
COALESCE(settlement.billing_status, usage_rows.billing_status) AS billing_status,
|
||||
usage_rows.created_at,
|
||||
COALESCE(settlement.finalized_at, usage_rows.finalized_at) AS finalized_at,
|
||||
GREATEST(COALESCE(settlement.billing_input_tokens, usage_rows.input_tokens, 0), 0)::bigint
|
||||
AS input_tokens,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
settlement.billing_effective_input_tokens,
|
||||
CASE
|
||||
WHEN GREATEST(COALESCE(usage_rows.input_tokens, 0), 0) <= 0 THEN 0
|
||||
WHEN GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0) <= 0
|
||||
THEN GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
WHEN split_part(lower(COALESCE(COALESCE(usage_rows.endpoint_api_format, usage_rows.api_format), '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN GREATEST(
|
||||
GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
- GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0),
|
||||
0
|
||||
)
|
||||
ELSE GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
END
|
||||
),
|
||||
0
|
||||
)::bigint AS effective_input_tokens,
|
||||
GREATEST(COALESCE(settlement.billing_output_tokens, usage_rows.output_tokens, 0), 0)::bigint
|
||||
AS output_tokens,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
settlement.billing_cache_creation_tokens,
|
||||
CASE
|
||||
WHEN COALESCE(usage_rows.cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE(usage_rows.cache_creation_input_tokens, 0)
|
||||
END,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint AS cache_creation_input_tokens,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
settlement.billing_cache_creation_5m_tokens,
|
||||
usage_rows.cache_creation_input_tokens_5m,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint AS cache_creation_input_tokens_5m,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
settlement.billing_cache_creation_1h_tokens,
|
||||
usage_rows.cache_creation_input_tokens_1h,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint AS cache_creation_input_tokens_1h,
|
||||
GREATEST(COALESCE(settlement.billing_cache_read_tokens, usage_rows.cache_read_input_tokens, 0), 0)::bigint
|
||||
AS cache_read_input_tokens,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
CASE
|
||||
WHEN settlement.billing_input_tokens IS NOT NULL
|
||||
OR settlement.billing_output_tokens IS NOT NULL
|
||||
OR settlement.billing_cache_creation_tokens IS NOT NULL
|
||||
OR settlement.billing_cache_creation_5m_tokens IS NOT NULL
|
||||
OR settlement.billing_cache_creation_1h_tokens IS NOT NULL
|
||||
OR settlement.billing_cache_read_tokens IS NOT NULL
|
||||
THEN COALESCE(settlement.billing_input_tokens, 0)
|
||||
+ COALESCE(settlement.billing_output_tokens, 0)
|
||||
+ COALESCE(
|
||||
settlement.billing_cache_creation_tokens,
|
||||
COALESCE(settlement.billing_cache_creation_5m_tokens, 0)
|
||||
+ COALESCE(settlement.billing_cache_creation_1h_tokens, 0),
|
||||
0
|
||||
)
|
||||
+ COALESCE(settlement.billing_cache_read_tokens, 0)
|
||||
END,
|
||||
usage_rows.total_tokens,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint AS total_tokens,
|
||||
GREATEST(
|
||||
COALESCE(
|
||||
settlement.billing_total_input_context,
|
||||
CASE
|
||||
WHEN split_part(lower(COALESCE(COALESCE(usage_rows.endpoint_api_format, usage_rows.api_format), '')), ':', 1)
|
||||
IN ('claude', 'anthropic')
|
||||
THEN GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
+ CASE
|
||||
WHEN COALESCE(usage_rows.cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE(usage_rows.cache_creation_input_tokens, 0)
|
||||
END
|
||||
+ GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0)
|
||||
WHEN split_part(lower(COALESCE(COALESCE(usage_rows.endpoint_api_format, usage_rows.api_format), '')), ':', 1)
|
||||
IN ('openai', 'gemini', 'google')
|
||||
THEN CASE
|
||||
WHEN GREATEST(COALESCE(usage_rows.input_tokens, 0), 0) <= 0 THEN 0
|
||||
WHEN GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0) <= 0
|
||||
THEN GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
ELSE GREATEST(
|
||||
GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
- GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0),
|
||||
0
|
||||
)
|
||||
END
|
||||
+ GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0)
|
||||
ELSE GREATEST(COALESCE(usage_rows.input_tokens, 0), 0)
|
||||
+ CASE
|
||||
WHEN COALESCE(usage_rows.cache_creation_input_tokens, 0) = 0
|
||||
AND (
|
||||
COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
) > 0
|
||||
THEN COALESCE(usage_rows.cache_creation_input_tokens_5m, 0)
|
||||
+ COALESCE(usage_rows.cache_creation_input_tokens_1h, 0)
|
||||
ELSE COALESCE(usage_rows.cache_creation_input_tokens, 0)
|
||||
END
|
||||
+ GREATEST(COALESCE(usage_rows.cache_read_input_tokens, 0), 0)
|
||||
END,
|
||||
0
|
||||
),
|
||||
0
|
||||
)::bigint AS total_input_context,
|
||||
COALESCE(CAST(usage_rows.input_cost_usd AS DOUBLE PRECISION), 0) AS input_cost_usd,
|
||||
COALESCE(CAST(usage_rows.output_cost_usd AS DOUBLE PRECISION), 0) AS output_cost_usd,
|
||||
COALESCE(
|
||||
CAST(settlement.billing_cache_creation_cost_usd AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.cache_creation_cost_usd AS DOUBLE PRECISION),
|
||||
0
|
||||
) AS cache_creation_cost_usd,
|
||||
COALESCE(
|
||||
CAST(settlement.billing_cache_read_cost_usd AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.cache_read_cost_usd AS DOUBLE PRECISION),
|
||||
0
|
||||
) AS cache_read_cost_usd,
|
||||
COALESCE(
|
||||
CAST(settlement.billing_total_cost_usd AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.total_cost_usd AS DOUBLE PRECISION),
|
||||
0
|
||||
) AS total_cost_usd,
|
||||
COALESCE(
|
||||
CAST(settlement.billing_actual_total_cost_usd AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.actual_total_cost_usd AS DOUBLE PRECISION),
|
||||
0
|
||||
) AS actual_total_cost_usd,
|
||||
COALESCE(
|
||||
CAST(settlement.output_price_per_1m AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.output_price_per_1m AS DOUBLE PRECISION)
|
||||
) AS output_price_per_1m,
|
||||
COALESCE(
|
||||
CAST(settlement.input_price_per_1m AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.input_price_per_1m AS DOUBLE PRECISION)
|
||||
) AS input_price_per_1m,
|
||||
COALESCE(
|
||||
CAST(settlement.cache_creation_price_per_1m AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.cache_creation_price_per_1m AS DOUBLE PRECISION)
|
||||
) AS cache_creation_price_per_1m,
|
||||
COALESCE(
|
||||
CAST(settlement.cache_read_price_per_1m AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.cache_read_price_per_1m AS DOUBLE PRECISION)
|
||||
) AS cache_read_price_per_1m,
|
||||
COALESCE(
|
||||
CAST(settlement.price_per_request AS DOUBLE PRECISION),
|
||||
CAST(usage_rows.price_per_request AS DOUBLE PRECISION)
|
||||
) AS price_per_request,
|
||||
settlement.billing_pricing_source,
|
||||
settlement.billing_rule_id,
|
||||
settlement.billing_rule_version,
|
||||
COALESCE(usage_rows.upstream_is_stream, COALESCE(usage_rows.is_stream, FALSE)) AS upstream_is_stream
|
||||
FROM public."usage" AS usage_rows
|
||||
LEFT JOIN public.usage_settlement_snapshots AS settlement
|
||||
ON settlement.request_id = usage_rows.request_id;
|
||||
|
||||
COMMENT ON COLUMN public.usage_billing_facts.upstream_is_stream IS
|
||||
'Resolved upstream stream mode from public.usage.upstream_is_stream, falling back to usage.is_stream for legacy rows.';
|
||||
@@ -8,7 +8,7 @@ use tracing::{error, info, warn};
|
||||
|
||||
static MIGRATOR: Migrator = sqlx::migrate!("./migrations");
|
||||
static BASELINE_V2_SQL: &str = include_str!("../bootstrap/20260413020000_baseline_v2.sql");
|
||||
const BASELINE_V2_CUTOFF_VERSION: i64 = 20260502000000;
|
||||
const BASELINE_V2_CUTOFF_VERSION: i64 = 20260505130000;
|
||||
const MIGRATIONS_TABLE_EXISTS_SQL: &str =
|
||||
"SELECT to_regclass('public._sqlx_migrations') IS NOT NULL";
|
||||
const PUBLIC_BASE_TABLE_COUNT_SQL: &str = r#"
|
||||
@@ -666,6 +666,7 @@ SELECT EXISTS (
|
||||
20260424000000,
|
||||
20260428000000,
|
||||
20260502000000,
|
||||
20260505130000,
|
||||
]
|
||||
);
|
||||
}
|
||||
@@ -689,6 +690,8 @@ SELECT EXISTS (
|
||||
assert!(BASELINE_V2_SQL.contains("settlement_snapshot_schema_version"));
|
||||
assert!(BASELINE_V2_SQL.contains("billing_effective_input_tokens"));
|
||||
assert!(BASELINE_V2_SQL.contains("CREATE OR REPLACE VIEW public.usage_billing_facts"));
|
||||
assert!(BASELINE_V2_SQL.contains("upstream_is_stream boolean"));
|
||||
assert!(BASELINE_V2_SQL.contains("COALESCE(usage_rows.upstream_is_stream"));
|
||||
assert!(BASELINE_V2_SQL.contains("usage_settlement_snapshots.billing_total_cost_usd"));
|
||||
assert!(BASELINE_V2_SQL.contains("candidate_index integer"));
|
||||
assert!(BASELINE_V2_SQL.contains("CREATE TABLE IF NOT EXISTS public.stats_user_summary"));
|
||||
@@ -1271,6 +1274,7 @@ ORDER BY id
|
||||
20260424000000,
|
||||
20260428000000,
|
||||
20260502000000,
|
||||
20260505130000,
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
@@ -902,6 +902,32 @@ fn usage_is_success(item: &StoredRequestUsageAudit) -> bool {
|
||||
) && item.status_code.is_none_or(|code| code < 400)
|
||||
}
|
||||
|
||||
fn usage_output_tps_uses_generation_time(item: &StoredRequestUsageAudit) -> bool {
|
||||
item.request_metadata
|
||||
.as_ref()
|
||||
.and_then(Value::as_object)
|
||||
.and_then(|metadata| metadata.get("upstream_is_stream"))
|
||||
.and_then(Value::as_bool)
|
||||
.unwrap_or(item.is_stream)
|
||||
}
|
||||
|
||||
fn usage_output_tps_duration_ms(item: &StoredRequestUsageAudit) -> Option<u64> {
|
||||
let response_time_ms = item.response_time_ms?;
|
||||
if response_time_ms == 0 {
|
||||
return None;
|
||||
}
|
||||
|
||||
if !usage_output_tps_uses_generation_time(item) {
|
||||
return Some(response_time_ms);
|
||||
}
|
||||
|
||||
let first_byte_time_ms = item.first_byte_time_ms?;
|
||||
if first_byte_time_ms >= response_time_ms {
|
||||
return None;
|
||||
}
|
||||
Some(response_time_ms - first_byte_time_ms)
|
||||
}
|
||||
|
||||
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") {
|
||||
@@ -1720,13 +1746,15 @@ impl UsageReadRepository for InMemoryUsageReadRepository {
|
||||
self.response_time_sample_count =
|
||||
self.response_time_sample_count.saturating_add(1);
|
||||
self.response_times.push(response_time_ms);
|
||||
if response_time_ms > 0 && item.output_tokens > 0 {
|
||||
self.tps_output_tokens =
|
||||
self.tps_output_tokens.saturating_add(item.output_tokens);
|
||||
self.tps_response_time_ms_sum = self
|
||||
.tps_response_time_ms_sum
|
||||
.saturating_add(response_time_ms);
|
||||
self.tps_sample_count = self.tps_sample_count.saturating_add(1);
|
||||
if let Some(output_tps_duration_ms) = usage_output_tps_duration_ms(item) {
|
||||
if item.output_tokens > 0 {
|
||||
self.tps_output_tokens =
|
||||
self.tps_output_tokens.saturating_add(item.output_tokens);
|
||||
self.tps_response_time_ms_sum = self
|
||||
.tps_response_time_ms_sum
|
||||
.saturating_add(output_tps_duration_ms);
|
||||
self.tps_sample_count = self.tps_sample_count.saturating_add(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(first_byte_time_ms) = item.first_byte_time_ms {
|
||||
@@ -4822,6 +4850,7 @@ mod tests {
|
||||
second.output_tokens = 40;
|
||||
second.response_time_ms = Some(1000);
|
||||
second.first_byte_time_ms = Some(200);
|
||||
second.request_metadata = Some(json!({ "upstream_is_stream": true }));
|
||||
|
||||
let mut failed = sample_usage("req-provider-perf-failed", 1_711_000_400);
|
||||
failed.output_tokens = 999;
|
||||
@@ -4852,7 +4881,7 @@ mod tests {
|
||||
|
||||
assert_eq!(summary.summary.request_count, 4);
|
||||
assert_eq!(summary.summary.success_count, 3);
|
||||
assert!((summary.summary.avg_output_tps.expect("summary tps") - 18.571_428).abs() < 0.001);
|
||||
assert!((summary.summary.avg_output_tps.expect("summary tps") - 19.117_647).abs() < 0.001);
|
||||
assert_eq!(summary.summary.avg_first_byte_time_ms, Some(150.0));
|
||||
assert!(
|
||||
(summary
|
||||
@@ -4870,7 +4899,7 @@ mod tests {
|
||||
assert_eq!(provider.request_count, 3);
|
||||
assert_eq!(provider.success_count, 2);
|
||||
assert_eq!(provider.output_tokens, 1099);
|
||||
assert_eq!(provider.avg_output_tps, Some(25.0));
|
||||
assert!((provider.avg_output_tps.expect("provider tps") - 26.315_789).abs() < 0.001);
|
||||
assert_eq!(provider.avg_first_byte_time_ms, Some(150.0));
|
||||
assert_eq!(provider.avg_response_time_ms, Some(2000.0));
|
||||
assert_eq!(provider.p90_response_time_ms, None);
|
||||
@@ -4880,6 +4909,8 @@ mod tests {
|
||||
assert_eq!(summary.timeline.len(), 1);
|
||||
assert_eq!(summary.timeline[0].date, "2024-03-21T05:00:00+00:00");
|
||||
assert_eq!(summary.timeline[0].provider_id, "provider-1");
|
||||
assert_eq!(summary.timeline[0].avg_output_tps, Some(25.0));
|
||||
assert!(
|
||||
(summary.timeline[0].avg_output_tps.expect("timeline tps") - 26.315_789).abs() < 0.001
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4697,6 +4697,17 @@ WITH filtered_usage AS (
|
||||
GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0) AS first_byte_time_ms,
|
||||
"usage".response_time_ms IS NOT NULL AS has_response_time,
|
||||
"usage".first_byte_time_ms IS NOT NULL AS has_first_byte_time,
|
||||
CASE
|
||||
WHEN COALESCE("usage".upstream_is_stream, "usage".is_stream, false)
|
||||
THEN CASE
|
||||
WHEN "usage".response_time_ms IS NOT NULL
|
||||
AND "usage".first_byte_time_ms IS NOT NULL
|
||||
AND GREATEST(COALESCE("usage".response_time_ms, 0), 0) > GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
THEN GREATEST(COALESCE("usage".response_time_ms, 0), 0) - GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
ELSE 0
|
||||
END
|
||||
ELSE GREATEST(COALESCE("usage".response_time_ms, 0), 0)
|
||||
END AS output_tps_duration_ms,
|
||||
CASE
|
||||
WHEN lower(COALESCE("usage".status, '')) IN ('completed', 'success', 'ok', 'billed', 'settled')
|
||||
AND ("usage".status_code IS NULL OR "usage".status_code < 400)
|
||||
@@ -4724,17 +4735,17 @@ SELECT
|
||||
COALESCE(SUM(success_flag), 0)::BIGINT AS success_count,
|
||||
CASE
|
||||
WHEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0) > 0
|
||||
THEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tokens
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION
|
||||
ELSE NULL
|
||||
@@ -4770,6 +4781,17 @@ WITH filtered_usage AS (
|
||||
GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0) AS first_byte_time_ms,
|
||||
"usage".response_time_ms IS NOT NULL AS has_response_time,
|
||||
"usage".first_byte_time_ms IS NOT NULL AS has_first_byte_time,
|
||||
CASE
|
||||
WHEN COALESCE("usage".upstream_is_stream, "usage".is_stream, false)
|
||||
THEN CASE
|
||||
WHEN "usage".response_time_ms IS NOT NULL
|
||||
AND "usage".first_byte_time_ms IS NOT NULL
|
||||
AND GREATEST(COALESCE("usage".response_time_ms, 0), 0) > GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
THEN GREATEST(COALESCE("usage".response_time_ms, 0), 0) - GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
ELSE 0
|
||||
END
|
||||
ELSE GREATEST(COALESCE("usage".response_time_ms, 0), 0)
|
||||
END AS output_tps_duration_ms,
|
||||
CASE
|
||||
WHEN lower(COALESCE("usage".status, '')) IN ('completed', 'success', 'ok', 'billed', 'settled')
|
||||
AND ("usage".status_code IS NULL OR "usage".status_code < 400)
|
||||
@@ -4800,17 +4822,17 @@ SELECT
|
||||
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
|
||||
CASE
|
||||
WHEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0) > 0
|
||||
THEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tokens
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION
|
||||
ELSE NULL
|
||||
@@ -4832,7 +4854,7 @@ SELECT
|
||||
ELSE NULL
|
||||
END AS p90_first_byte_time_ms,
|
||||
COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN 1
|
||||
ELSE 0
|
||||
END), 0)::BIGINT AS tps_sample_count,
|
||||
@@ -4885,6 +4907,17 @@ ORDER BY request_count DESC, provider_id ASC
|
||||
GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0) AS first_byte_time_ms,
|
||||
"usage".response_time_ms IS NOT NULL AS has_response_time,
|
||||
"usage".first_byte_time_ms IS NOT NULL AS has_first_byte_time,
|
||||
CASE
|
||||
WHEN COALESCE("usage".upstream_is_stream, "usage".is_stream, false)
|
||||
THEN CASE
|
||||
WHEN "usage".response_time_ms IS NOT NULL
|
||||
AND "usage".first_byte_time_ms IS NOT NULL
|
||||
AND GREATEST(COALESCE("usage".response_time_ms, 0), 0) > GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
THEN GREATEST(COALESCE("usage".response_time_ms, 0), 0) - GREATEST(COALESCE("usage".first_byte_time_ms, 0), 0)
|
||||
ELSE 0
|
||||
END
|
||||
ELSE GREATEST(COALESCE("usage".response_time_ms, 0), 0)
|
||||
END AS output_tps_duration_ms,
|
||||
CASE
|
||||
WHEN lower(COALESCE("usage".status, '')) IN ('completed', 'success', 'ok', 'billed', 'settled')
|
||||
AND ("usage".status_code IS NULL OR "usage".status_code < 400)
|
||||
@@ -4921,17 +4954,17 @@ SELECT
|
||||
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
|
||||
CASE
|
||||
WHEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0) > 0
|
||||
THEN COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tokens
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
|
||||
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
|
||||
THEN response_time_ms
|
||||
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
|
||||
THEN output_tps_duration_ms
|
||||
ELSE 0
|
||||
END), 0)::DOUBLE PRECISION
|
||||
ELSE NULL
|
||||
|
||||
@@ -20,6 +20,7 @@ INSERT INTO "usage" (
|
||||
provider_endpoint_kind,
|
||||
has_format_conversion,
|
||||
is_stream,
|
||||
upstream_is_stream,
|
||||
input_tokens,
|
||||
output_tokens,
|
||||
total_tokens,
|
||||
@@ -76,6 +77,14 @@ INSERT INTO "usage" (
|
||||
$19,
|
||||
COALESCE($20, FALSE),
|
||||
COALESCE($21, FALSE),
|
||||
COALESCE(
|
||||
CASE
|
||||
WHEN ($53::json->>'upstream_is_stream') IN ('true', 'false')
|
||||
THEN ($53::json->>'upstream_is_stream')::boolean
|
||||
ELSE NULL
|
||||
END,
|
||||
COALESCE($21, FALSE)
|
||||
),
|
||||
COALESCE($22, 0),
|
||||
COALESCE($23, 0),
|
||||
COALESCE($24, COALESCE($22, 0) + COALESCE($23, 0)),
|
||||
@@ -135,6 +144,7 @@ DO UPDATE SET
|
||||
provider_endpoint_kind = CASE WHEN "usage".billing_status = 'pending' THEN COALESCE(EXCLUDED.provider_endpoint_kind, "usage".provider_endpoint_kind) ELSE "usage".provider_endpoint_kind END,
|
||||
has_format_conversion = CASE WHEN "usage".billing_status = 'pending' THEN COALESCE(EXCLUDED.has_format_conversion, "usage".has_format_conversion) ELSE "usage".has_format_conversion END,
|
||||
is_stream = CASE WHEN "usage".billing_status = 'pending' THEN COALESCE(EXCLUDED.is_stream, "usage".is_stream) ELSE "usage".is_stream END,
|
||||
upstream_is_stream = CASE WHEN "usage".billing_status = 'pending' THEN COALESCE(EXCLUDED.upstream_is_stream, "usage".upstream_is_stream, "usage".is_stream, false) ELSE "usage".upstream_is_stream END,
|
||||
input_tokens = CASE WHEN "usage".billing_status = 'pending' AND EXCLUDED.status IN ('completed', 'failed', 'cancelled') THEN GREATEST("usage".input_tokens, EXCLUDED.input_tokens) ELSE "usage".input_tokens END,
|
||||
output_tokens = CASE WHEN "usage".billing_status = 'pending' AND EXCLUDED.status IN ('completed', 'failed', 'cancelled') THEN GREATEST("usage".output_tokens, EXCLUDED.output_tokens) ELSE "usage".output_tokens END,
|
||||
total_tokens = CASE WHEN "usage".billing_status = 'pending' AND EXCLUDED.status IN ('completed', 'failed', 'cancelled') THEN GREATEST("usage".total_tokens, EXCLUDED.total_tokens) ELSE "usage".total_tokens END,
|
||||
|
||||
@@ -441,6 +441,37 @@ fn usage_sql_raw_aggregates_use_canonical_billing_facts() {
|
||||
.contains("FROM usage_billing_facts AS \"usage\""));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_sql_provider_performance_reads_upstream_stream_from_billing_facts() {
|
||||
let source = include_str!("mod.rs");
|
||||
assert!(source.contains("\"usage\".upstream_is_stream"));
|
||||
assert!(
|
||||
!source.contains("usage_base.request_metadata->>'upstream_is_stream'"),
|
||||
"provider performance queries should not rejoin public.usage to resolve upstream stream mode"
|
||||
);
|
||||
assert!(
|
||||
!source.contains("LEFT JOIN public.usage AS usage_base"),
|
||||
"provider performance queries should stay on usage_billing_facts to avoid an extra usage scan"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_billing_facts_projects_upstream_stream_mode() {
|
||||
let migration = include_str!(
|
||||
"../../../../migrations/20260505130000_project_upstream_stream_in_usage_billing_facts.sql"
|
||||
);
|
||||
let baseline = include_str!("../../../../bootstrap/20260413020000_baseline_v2.sql");
|
||||
|
||||
for source in [migration, baseline] {
|
||||
assert!(source.contains("AS upstream_is_stream"));
|
||||
assert!(source.contains("COALESCE(usage_rows.upstream_is_stream"));
|
||||
assert!(source.contains("COALESCE(usage_rows.is_stream, FALSE)"));
|
||||
}
|
||||
assert!(migration.contains("ADD COLUMN IF NOT EXISTS upstream_is_stream boolean"));
|
||||
assert!(migration.contains("request_metadata->>'upstream_is_stream'"));
|
||||
assert!(baseline.contains("upstream_is_stream boolean"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_sql_reads_http_audits_for_single_record_fetches() {
|
||||
assert!(super::FIND_BY_REQUEST_ID_SQL.contains("LEFT JOIN usage_http_audits"));
|
||||
@@ -575,6 +606,14 @@ fn usage_sql_insert_values_aligns_request_metadata_and_timestamps() {
|
||||
assert!(super::UPSERT_SQL.contains("TO_TIMESTAMP($55::double precision)"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_sql_upsert_materializes_upstream_stream_mode() {
|
||||
assert!(super::UPSERT_SQL.contains("upstream_is_stream,"));
|
||||
assert!(super::UPSERT_SQL.contains("$53::json->>'upstream_is_stream'"));
|
||||
assert!(super::UPSERT_SQL.contains("COALESCE($21, FALSE)"));
|
||||
assert!(super::UPSERT_SQL.contains("upstream_is_stream = CASE"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_sql_upsert_returning_includes_routing_placeholders() {
|
||||
assert!(super::UPSERT_SQL.contains("NULL::varchar AS http_request_body_state"));
|
||||
|
||||
Reference in New Issue
Block a user