mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-03 01:40:21 +08:00
Merge remote-tracking branch 'origin/pr/377' into codex/pr-376-377-383-384-combined
# Conflicts: # apps/aether-gateway/src/tests/frontdoor/public_support/dashboard.rs # crates/aether-data/src/lifecycle/backfill.rs # crates/aether-data/src/repository/usage/postgres/mod.rs
This commit is contained in:
@@ -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
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -444,6 +444,41 @@ impl ApiKeyUsageDelta {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||
pub(crate) struct ModelUsageContribution {
|
||||
pub model: String,
|
||||
pub request_count: i64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||
pub(crate) struct ModelUsageDelta {
|
||||
pub request_count: i64,
|
||||
}
|
||||
|
||||
impl ModelUsageDelta {
|
||||
pub(crate) fn between(before: &ModelUsageContribution, after: &ModelUsageContribution) -> Self {
|
||||
Self {
|
||||
request_count: after.request_count - before.request_count,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn addition(after: &ModelUsageContribution) -> Self {
|
||||
Self {
|
||||
request_count: after.request_count,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn removal(before: &ModelUsageContribution) -> Self {
|
||||
Self {
|
||||
request_count: -before.request_count,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn is_noop(&self) -> bool {
|
||||
self.request_count == 0
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Default)]
|
||||
pub(crate) struct ProviderApiKeyUsageContribution {
|
||||
pub key_id: String,
|
||||
@@ -620,6 +655,23 @@ pub(crate) fn provider_api_key_usage_contribution(
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn model_usage_contribution(
|
||||
usage: &StoredRequestUsageAudit,
|
||||
) -> Option<ModelUsageContribution> {
|
||||
if matches!(usage.status.as_str(), "pending" | "streaming") {
|
||||
return None;
|
||||
}
|
||||
let model = usage.model.trim();
|
||||
if model.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
Some(ModelUsageContribution {
|
||||
model: model.to_string(),
|
||||
request_count: 1,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn api_key_usage_contribution(
|
||||
usage: &StoredRequestUsageAudit,
|
||||
) -> Option<ApiKeyUsageContribution> {
|
||||
@@ -647,9 +699,10 @@ pub(crate) fn api_key_usage_contribution(
|
||||
mod tests {
|
||||
use super::{
|
||||
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
||||
provider_api_key_usage_contribution, provider_api_key_usage_is_error,
|
||||
provider_api_key_usage_is_success, strip_deprecated_usage_display_fields,
|
||||
usage_can_recover_terminal_failure, StoredRequestUsageAudit, UpsertUsageRecord,
|
||||
model_usage_contribution, provider_api_key_usage_contribution,
|
||||
provider_api_key_usage_is_error, provider_api_key_usage_is_success,
|
||||
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure, ModelUsageDelta,
|
||||
StoredRequestUsageAudit, UpsertUsageRecord,
|
||||
};
|
||||
|
||||
#[test]
|
||||
@@ -925,4 +978,75 @@ mod tests {
|
||||
assert_eq!(contribution.total_cost_usd, 0.25);
|
||||
assert_eq!(contribution.last_used_at_unix_secs, Some(123));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn model_usage_contribution_tracks_terminal_requests_only() {
|
||||
let completed = StoredRequestUsageAudit::new(
|
||||
"usage-1".to_string(),
|
||||
"request-1".to_string(),
|
||||
Some("user-1".to_string()),
|
||||
Some("api-key-1".to_string()),
|
||||
None,
|
||||
None,
|
||||
"OpenAI".to_string(),
|
||||
" gpt-5.5 ".to_string(),
|
||||
None,
|
||||
Some("provider-1".to_string()),
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
false,
|
||||
false,
|
||||
12,
|
||||
8,
|
||||
20,
|
||||
0.25,
|
||||
0.25,
|
||||
Some(200),
|
||||
None,
|
||||
None,
|
||||
Some(120),
|
||||
None,
|
||||
"completed".to_string(),
|
||||
"settled".to_string(),
|
||||
123,
|
||||
124,
|
||||
Some(125),
|
||||
)
|
||||
.expect("usage should build");
|
||||
let contribution =
|
||||
model_usage_contribution(&completed).expect("completed usage should count");
|
||||
assert_eq!(contribution.model, "gpt-5.5");
|
||||
assert_eq!(contribution.request_count, 1);
|
||||
|
||||
let mut streaming = completed.clone();
|
||||
streaming.status = "streaming".to_string();
|
||||
assert!(model_usage_contribution(&streaming).is_none());
|
||||
|
||||
let mut pending = completed;
|
||||
pending.status = "pending".to_string();
|
||||
assert!(model_usage_contribution(&pending).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn model_usage_delta_handles_model_changes() {
|
||||
let before = super::ModelUsageContribution {
|
||||
model: "gpt-5.4".to_string(),
|
||||
request_count: 1,
|
||||
};
|
||||
let after = super::ModelUsageContribution {
|
||||
model: "gpt-5.5".to_string(),
|
||||
request_count: 1,
|
||||
};
|
||||
|
||||
assert_eq!(ModelUsageDelta::removal(&before).request_count, -1);
|
||||
assert_eq!(ModelUsageDelta::addition(&after).request_count, 1);
|
||||
assert!(ModelUsageDelta::between(&before, &before).is_noop());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,7 +36,8 @@ use uuid::Uuid;
|
||||
|
||||
use super::{
|
||||
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
|
||||
provider_api_key_usage_contribution, strip_deprecated_usage_display_fields, ApiKeyUsageDelta,
|
||||
model_usage_contribution, provider_api_key_usage_contribution,
|
||||
strip_deprecated_usage_display_fields, ApiKeyUsageDelta, ModelUsageDelta,
|
||||
PendingUsageCleanupSummary, ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary,
|
||||
StoredProviderUsageSummary, StoredRequestUsageAudit, StoredUsageDailySummary,
|
||||
UpsertUsageRecord, UsageAuditListQuery, UsageDailyHeatmapQuery, UsageReadRepository,
|
||||
@@ -885,6 +886,67 @@ fn decode_usage_audit_summary_row(row: &PgRow) -> Result<StoredUsageAuditSummary
|
||||
})
|
||||
}
|
||||
|
||||
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 decode_usage_error_distribution_row(
|
||||
row: &PgRow,
|
||||
) -> Result<StoredUsageErrorDistributionRow, DataLayerError> {
|
||||
@@ -1213,6 +1275,9 @@ const SUMMARIZE_USAGE_BY_PROVIDER_API_KEY_IDS_SQL: &str =
|
||||
const APPLY_API_KEY_USAGE_DELTA_SQL: &str =
|
||||
include_str!("queries/apply_api_key_usage_delta_sql.sql");
|
||||
|
||||
const APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL: &str =
|
||||
include_str!("queries/apply_global_model_usage_delta_sql.sql");
|
||||
|
||||
const RESET_API_KEY_USAGE_STATS_SQL: &str =
|
||||
include_str!("queries/reset_api_key_usage_stats_sql.sql");
|
||||
|
||||
@@ -4789,6 +4854,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)
|
||||
@@ -4862,6 +4938,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)
|
||||
@@ -4977,6 +5064,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)
|
||||
@@ -6216,6 +6314,11 @@ WHERE date >= $1
|
||||
GROUP BY {group_column}
|
||||
ORDER BY request_count DESC, group_key ASC
|
||||
"#,
|
||||
group_column = group_column,
|
||||
display_name_expr = display_name_expr,
|
||||
avg_response_time_expr = avg_response_time_expr,
|
||||
success_count_expr = success_count_expr,
|
||||
table_name = table_name,
|
||||
);
|
||||
|
||||
let mut rows = sqlx::query(&sql)
|
||||
@@ -7326,6 +7429,40 @@ ORDER BY "usage".user_id ASC
|
||||
}
|
||||
}
|
||||
|
||||
let before_model_contribution =
|
||||
previous_usage.as_ref().and_then(model_usage_contribution);
|
||||
let after_model_contribution = model_usage_contribution(&stored);
|
||||
match (
|
||||
before_model_contribution.as_ref(),
|
||||
after_model_contribution.as_ref(),
|
||||
) {
|
||||
(Some(before), Some(after)) if before.model == after.model => {
|
||||
let delta = ModelUsageDelta::between(before, after);
|
||||
apply_global_model_usage_delta_in_tx(tx, before.model.as_str(), &delta)
|
||||
.await?;
|
||||
}
|
||||
_ => {
|
||||
if let Some(before) = before_model_contribution.as_ref() {
|
||||
let delta = ModelUsageDelta::removal(before);
|
||||
apply_global_model_usage_delta_in_tx(
|
||||
tx,
|
||||
before.model.as_str(),
|
||||
&delta,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
if let Some(after) = after_model_contribution.as_ref() {
|
||||
let delta = ModelUsageDelta::addition(after);
|
||||
apply_global_model_usage_delta_in_tx(
|
||||
tx,
|
||||
after.model.as_str(),
|
||||
&delta,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let before_provider_contribution = previous_usage
|
||||
.as_ref()
|
||||
.and_then(provider_api_key_usage_contribution);
|
||||
@@ -7884,6 +8021,32 @@ async fn apply_api_key_usage_delta_in_tx(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn apply_global_model_usage_delta_in_tx(
|
||||
tx: &mut sqlx::Transaction<'_, Postgres>,
|
||||
model: &str,
|
||||
delta: &ModelUsageDelta,
|
||||
) -> Result<(), DataLayerError> {
|
||||
if model.trim().is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
if delta.is_noop() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
sqlx::query(APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL)
|
||||
.bind(model)
|
||||
.bind(i32::try_from(delta.request_count).map_err(|_| {
|
||||
DataLayerError::UnexpectedValue(format!(
|
||||
"global_models.usage_count delta exceeds i32: {}",
|
||||
delta.request_count
|
||||
))
|
||||
})?)
|
||||
.execute(&mut **tx)
|
||||
.await
|
||||
.map_postgres_err()?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn apply_provider_api_key_usage_delta_in_tx(
|
||||
tx: &mut sqlx::Transaction<'_, Postgres>,
|
||||
key_id: &str,
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
UPDATE global_models
|
||||
SET
|
||||
usage_count = GREATEST(COALESCE(usage_count, 0) + $2, 0),
|
||||
updated_at = NOW()
|
||||
WHERE name = $1
|
||||
@@ -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,
|
||||
|
||||
@@ -443,6 +443,33 @@ 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/postgres/20260505130000_project_upstream_stream_in_usage_billing_facts.sql"
|
||||
);
|
||||
|
||||
assert!(migration.contains("AS upstream_is_stream"));
|
||||
assert!(migration.contains("COALESCE(usage_rows.upstream_is_stream"));
|
||||
assert!(migration.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'"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_sql_reads_http_audits_for_single_record_fetches() {
|
||||
assert!(super::FIND_BY_REQUEST_ID_SQL.contains("LEFT JOIN usage_http_audits"));
|
||||
@@ -577,6 +604,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