mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
Merge branch 'fawney19:main' into main
This commit is contained in:
@@ -29,11 +29,12 @@ use serde_json::Value;
|
||||
use super::{
|
||||
api_key_usage_contribution, provider_api_key_usage_contribution,
|
||||
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
|
||||
ApiKeyUsageContribution, ApiKeyUsageDelta, ProviderApiKeyUsageContribution,
|
||||
ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
|
||||
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
|
||||
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
|
||||
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
|
||||
usage_request_metadata_client_family, ApiKeyUsageContribution, ApiKeyUsageDelta,
|
||||
ProviderApiKeyUsageContribution, ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest,
|
||||
StoredProviderApiKeyUsageSummary, StoredProviderApiKeyWindowUsageSummary,
|
||||
StoredProviderUsageSummary, StoredProviderUsageWindow, StoredRequestUsageAudit,
|
||||
StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery, UsageDailyHeatmapQuery,
|
||||
UsageReadRepository, UsageWriteRepository,
|
||||
};
|
||||
use crate::repository::auth::InMemoryAuthApiKeySnapshotRepository;
|
||||
use crate::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
|
||||
@@ -56,6 +57,7 @@ impl InMemoryUsageReadRepository {
|
||||
let mut by_request_id = BTreeMap::new();
|
||||
for mut item in items {
|
||||
hydrate_legacy_body_refs(&mut item);
|
||||
hydrate_client_family(&mut item);
|
||||
by_request_id.insert(item.request_id.clone(), item);
|
||||
}
|
||||
Self {
|
||||
@@ -75,6 +77,7 @@ impl InMemoryUsageReadRepository {
|
||||
let mut detached_bodies = BTreeMap::new();
|
||||
for mut item in items {
|
||||
hydrate_legacy_body_refs(&mut item);
|
||||
hydrate_client_family(&mut item);
|
||||
let request_id = item.request_id.clone();
|
||||
if let Some(body_ref) = detach_usage_body(
|
||||
&request_id,
|
||||
@@ -2546,6 +2549,13 @@ fn hydrate_legacy_body_refs(item: &mut StoredRequestUsageAudit) {
|
||||
}
|
||||
}
|
||||
|
||||
fn hydrate_client_family(item: &mut StoredRequestUsageAudit) {
|
||||
if item.client_family.is_none() {
|
||||
item.client_family = usage_request_metadata_client_family(item.request_metadata.as_ref())
|
||||
.map(ToOwned::to_owned);
|
||||
}
|
||||
}
|
||||
|
||||
fn persisted_usage_body_ref(
|
||||
incoming_ref: Option<&str>,
|
||||
incoming_body: Option<&Value>,
|
||||
@@ -2851,6 +2861,13 @@ impl UsageWriteRepository for InMemoryUsageReadRepository {
|
||||
})
|
||||
},
|
||||
),
|
||||
client_family: usage_request_metadata_client_family(request_metadata.as_ref())
|
||||
.map(ToOwned::to_owned)
|
||||
.or_else(|| {
|
||||
existing
|
||||
.as_ref()
|
||||
.and_then(|existing| existing.client_family.clone())
|
||||
}),
|
||||
request_metadata,
|
||||
created_at_unix_ms,
|
||||
updated_at_unix_secs: usage.updated_at_unix_secs,
|
||||
@@ -4399,6 +4416,7 @@ mod tests {
|
||||
provider_endpoint_kind: Some("chat".to_string()),
|
||||
has_format_conversion: false,
|
||||
is_stream: false,
|
||||
client_family: None,
|
||||
input_tokens: 10,
|
||||
output_tokens: 20,
|
||||
total_tokens: 30,
|
||||
|
||||
@@ -363,7 +363,8 @@ mod sqlite;
|
||||
|
||||
#[allow(unused_imports)]
|
||||
pub(crate) use aether_data_contracts::repository::usage::{
|
||||
PendingUsageCleanupSummary, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
|
||||
usage_request_metadata_client_family, PendingUsageCleanupSummary,
|
||||
ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
|
||||
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
|
||||
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
|
||||
StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
|
||||
|
||||
@@ -6,8 +6,8 @@ use sqlx::{mysql::MySqlRow, Row};
|
||||
use super::{
|
||||
provider_api_key_usage_is_error, provider_api_key_usage_is_success,
|
||||
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
|
||||
InMemoryUsageReadRepository, PendingUsageCleanupSummary, StoredRequestUsageAudit,
|
||||
UpsertUsageRecord, UsageWriteRepository,
|
||||
usage_request_metadata_client_family, InMemoryUsageReadRepository, PendingUsageCleanupSummary,
|
||||
StoredRequestUsageAudit, UpsertUsageRecord, UsageWriteRepository,
|
||||
};
|
||||
use crate::driver::mysql::MysqlPool;
|
||||
use crate::error::SqlResultExt;
|
||||
@@ -850,6 +850,8 @@ fn map_usage_row(row: &MySqlRow) -> Result<StoredRequestUsageAudit, DataLayerErr
|
||||
.map(|raw| serde_json::from_str(&raw))
|
||||
.transpose()
|
||||
.map_err(|err| DataLayerError::UnexpectedValue(err.to_string()))?;
|
||||
audit.client_family = usage_request_metadata_client_family(audit.request_metadata.as_ref())
|
||||
.map(ToOwned::to_owned);
|
||||
let upstream_is_stream = row
|
||||
.try_get::<Option<bool>, _>("upstream_is_stream")
|
||||
.map_sql_err()?;
|
||||
|
||||
@@ -8591,6 +8591,9 @@ fn map_usage_row(
|
||||
.try_get::<f64, _>("cache_read_cost_usd")
|
||||
.map_postgres_err()?;
|
||||
usage.output_price_per_1m = row.try_get("output_price_per_1m").map_postgres_err()?;
|
||||
usage.client_family = row
|
||||
.try_get::<Option<String>, _>("client_family")
|
||||
.map_postgres_err()?;
|
||||
usage.request_headers = row.try_get("request_headers").map_postgres_err()?;
|
||||
let request_body = usage_json_column(
|
||||
row,
|
||||
|
||||
@@ -93,6 +93,10 @@ SELECT
|
||||
"usage".first_byte_time_ms,
|
||||
"usage".status,
|
||||
COALESCE(usage_settlement_snapshots.billing_status, "usage".billing_status) AS billing_status,
|
||||
COALESCE(
|
||||
NULLIF(BTRIM("usage".request_metadata->'client_session_affinity'->>'client_family'), ''),
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_family'), '')
|
||||
) AS client_family,
|
||||
COALESCE(usage_http_audits.request_headers, "usage".request_headers) AS request_headers,
|
||||
"usage".request_body,
|
||||
"usage".request_body_compressed,
|
||||
|
||||
@@ -93,6 +93,10 @@ SELECT
|
||||
"usage".first_byte_time_ms,
|
||||
"usage".status,
|
||||
COALESCE(usage_settlement_snapshots.billing_status, "usage".billing_status) AS billing_status,
|
||||
COALESCE(
|
||||
NULLIF(BTRIM("usage".request_metadata->'client_session_affinity'->>'client_family'), ''),
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_family'), '')
|
||||
) AS client_family,
|
||||
COALESCE(usage_http_audits.request_headers, "usage".request_headers) AS request_headers,
|
||||
"usage".request_body,
|
||||
"usage".request_body_compressed,
|
||||
|
||||
@@ -93,6 +93,10 @@ SELECT
|
||||
"usage".first_byte_time_ms,
|
||||
"usage".status,
|
||||
COALESCE(usage_settlement_snapshots.billing_status, "usage".billing_status) AS billing_status,
|
||||
COALESCE(
|
||||
NULLIF(BTRIM("usage".request_metadata->'client_session_affinity'->>'client_family'), ''),
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_family'), '')
|
||||
) AS client_family,
|
||||
NULL::json AS request_headers,
|
||||
NULL::json AS request_body,
|
||||
NULL::bytea AS request_body_compressed,
|
||||
@@ -106,9 +110,21 @@ SELECT
|
||||
NULL::json AS client_response_body,
|
||||
NULL::bytea AS client_response_body_compressed,
|
||||
CASE
|
||||
WHEN ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
WHEN NULLIF(BTRIM("usage".request_metadata->>'client_ip'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'user_agent'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'request_path'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'request_path_and_query'), '') IS NOT NULL
|
||||
OR ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
OR ("usage".request_metadata->>'upstream_is_stream') IN ('true', 'false')
|
||||
THEN jsonb_build_object(
|
||||
THEN jsonb_strip_nulls(jsonb_build_object(
|
||||
'client_ip',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_ip'), ''),
|
||||
'user_agent',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'user_agent'), ''),
|
||||
'request_path',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'request_path'), ''),
|
||||
'request_path_and_query',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'request_path_and_query'), ''),
|
||||
'client_requested_stream',
|
||||
CASE
|
||||
WHEN ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
@@ -121,7 +137,7 @@ SELECT
|
||||
THEN ("usage".request_metadata->>'upstream_is_stream')::boolean
|
||||
ELSE NULL
|
||||
END
|
||||
)::json
|
||||
))::json
|
||||
ELSE NULL::json
|
||||
END AS request_metadata,
|
||||
NULL::varchar AS http_request_body_ref,
|
||||
|
||||
@@ -93,6 +93,10 @@ SELECT
|
||||
"usage".first_byte_time_ms,
|
||||
"usage".status,
|
||||
COALESCE(usage_settlement_snapshots.billing_status, "usage".billing_status) AS billing_status,
|
||||
COALESCE(
|
||||
NULLIF(BTRIM("usage".request_metadata->'client_session_affinity'->>'client_family'), ''),
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_family'), '')
|
||||
) AS client_family,
|
||||
NULL::json AS request_headers,
|
||||
NULL::json AS request_body,
|
||||
NULL::bytea AS request_body_compressed,
|
||||
@@ -106,9 +110,21 @@ SELECT
|
||||
NULL::json AS client_response_body,
|
||||
NULL::bytea AS client_response_body_compressed,
|
||||
CASE
|
||||
WHEN ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
WHEN NULLIF(BTRIM("usage".request_metadata->>'client_ip'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'user_agent'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'request_path'), '') IS NOT NULL
|
||||
OR NULLIF(BTRIM("usage".request_metadata->>'request_path_and_query'), '') IS NOT NULL
|
||||
OR ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
OR ("usage".request_metadata->>'upstream_is_stream') IN ('true', 'false')
|
||||
THEN jsonb_build_object(
|
||||
THEN jsonb_strip_nulls(jsonb_build_object(
|
||||
'client_ip',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'client_ip'), ''),
|
||||
'user_agent',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'user_agent'), ''),
|
||||
'request_path',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'request_path'), ''),
|
||||
'request_path_and_query',
|
||||
NULLIF(BTRIM("usage".request_metadata->>'request_path_and_query'), ''),
|
||||
'client_requested_stream',
|
||||
CASE
|
||||
WHEN ("usage".request_metadata->>'client_requested_stream') IN ('true', 'false')
|
||||
@@ -121,7 +137,7 @@ SELECT
|
||||
THEN ("usage".request_metadata->>'upstream_is_stream')::boolean
|
||||
ELSE NULL
|
||||
END
|
||||
)::json
|
||||
))::json
|
||||
ELSE NULL::json
|
||||
END AS request_metadata,
|
||||
NULL::varchar AS http_request_body_ref,
|
||||
|
||||
@@ -633,6 +633,14 @@ fn usage_sql_uses_json_null_placeholders_for_usage_payload_columns() {
|
||||
super::LIST_USAGE_AUDITS_PREFIX,
|
||||
super::LIST_RECENT_USAGE_AUDITS_PREFIX,
|
||||
] {
|
||||
assert!(sql.contains("jsonb_strip_nulls(jsonb_build_object("));
|
||||
assert!(sql.contains("'client_ip'"));
|
||||
assert!(sql.contains("request_metadata->>'client_ip'"));
|
||||
assert!(sql.contains("'user_agent'"));
|
||||
assert!(sql.contains("request_metadata->>'user_agent'"));
|
||||
assert!(sql.contains("AS client_family"));
|
||||
assert!(sql.contains("request_metadata->'client_session_affinity'->>'client_family'"));
|
||||
assert!(sql.contains("request_metadata->>'client_family'"));
|
||||
assert!(sql.contains("CAST(\"usage\".input_tokens AS INTEGER) AS input_tokens"));
|
||||
assert!(sql.contains(
|
||||
"usage_settlement_snapshots.billing_input_tokens AS settlement_billing_input_tokens"
|
||||
|
||||
@@ -8,7 +8,8 @@ use sqlx::{sqlite::SqliteRow, QueryBuilder, Row, Sqlite};
|
||||
|
||||
use super::{
|
||||
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
|
||||
PendingUsageCleanupSummary, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
|
||||
usage_request_metadata_client_family, PendingUsageCleanupSummary,
|
||||
ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
|
||||
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredRequestUsageAudit,
|
||||
StoredUsageAuditAggregation, StoredUsageAuditSummary, StoredUsageBreakdownSummaryRow,
|
||||
StoredUsageCacheAffinityHitSummary, StoredUsageCacheAffinityIntervalRow,
|
||||
@@ -3812,6 +3813,8 @@ fn map_usage_row(row: &SqliteRow) -> Result<StoredRequestUsageAudit, DataLayerEr
|
||||
.map(|raw| serde_json::from_str(&raw))
|
||||
.transpose()
|
||||
.map_err(|err| DataLayerError::UnexpectedValue(err.to_string()))?;
|
||||
audit.client_family = usage_request_metadata_client_family(audit.request_metadata.as_ref())
|
||||
.map(ToOwned::to_owned);
|
||||
let upstream_is_stream = row
|
||||
.try_get::<Option<i64>, _>("upstream_is_stream")
|
||||
.map_sql_err()?
|
||||
|
||||
Reference in New Issue
Block a user