Add usage client family read model field

This commit is contained in:
RWDai
2026-05-18 19:28:36 +08:00
parent 0a26accca4
commit 6d994917a0
6 changed files with 54 additions and 12 deletions
@@ -1,8 +1,8 @@
mod types; mod types;
pub use types::{ pub use types::{
parse_usage_body_ref, usage_body_ref, PendingUsageCleanupSummary, parse_usage_body_ref, usage_body_ref, usage_request_metadata_client_family,
ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary, PendingUsageCleanupSummary, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow, StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary, StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary, StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
@@ -34,6 +34,8 @@ pub struct StoredRequestUsageAudit {
pub provider_endpoint_kind: Option<String>, pub provider_endpoint_kind: Option<String>,
pub has_format_conversion: bool, pub has_format_conversion: bool,
pub is_stream: bool, pub is_stream: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub client_family: Option<String>,
pub input_tokens: u64, pub input_tokens: u64,
pub output_tokens: u64, pub output_tokens: u64,
pub total_tokens: u64, pub total_tokens: u64,
@@ -211,6 +213,7 @@ impl StoredRequestUsageAudit {
provider_endpoint_kind, provider_endpoint_kind,
has_format_conversion, has_format_conversion,
is_stream, is_stream,
client_family: None,
input_tokens: parse_u64(input_tokens, "usage.input_tokens")?, input_tokens: parse_u64(input_tokens, "usage.input_tokens")?,
output_tokens: parse_u64(output_tokens, "usage.output_tokens")?, output_tokens: parse_u64(output_tokens, "usage.output_tokens")?,
total_tokens: parse_u64(total_tokens, "usage.total_tokens")?, total_tokens: parse_u64(total_tokens, "usage.total_tokens")?,
@@ -311,6 +314,10 @@ impl StoredRequestUsageAudit {
.filter(|value| !value.is_empty()) .filter(|value| !value.is_empty())
} }
pub fn request_metadata_client_family(&self) -> Option<&str> {
usage_request_metadata_client_family(self.request_metadata.as_ref())
}
fn billing_snapshot_resolved_number(&self, key: &str) -> Option<f64> { fn billing_snapshot_resolved_number(&self, key: &str) -> Option<f64> {
self.request_metadata_object() self.request_metadata_object()
.and_then(|metadata| metadata.get("billing_snapshot")) .and_then(|metadata| metadata.get("billing_snapshot"))
@@ -1794,6 +1801,18 @@ pub struct UsageCleanupPreviewCounts {
pub log: u64, pub log: u64,
} }
pub fn usage_request_metadata_client_family(value: Option<&Value>) -> Option<&str> {
let metadata = value.and_then(Value::as_object)?;
metadata
.get("client_session_affinity")
.and_then(Value::as_object)
.and_then(|affinity| affinity.get("client_family"))
.and_then(Value::as_str)
.or_else(|| metadata.get("client_family").and_then(Value::as_str))
.map(str::trim)
.filter(|value| !value.is_empty())
}
fn parse_u64(value: i32, field_name: &str) -> Result<u64, crate::DataLayerError> { fn parse_u64(value: i32, field_name: &str) -> Result<u64, crate::DataLayerError> {
u64::try_from(value).map_err(|_| { u64::try_from(value).map_err(|_| {
crate::DataLayerError::UnexpectedValue(format!("invalid {field_name}: {value}")) crate::DataLayerError::UnexpectedValue(format!("invalid {field_name}: {value}"))
@@ -29,11 +29,12 @@ use serde_json::Value;
use super::{ use super::{
api_key_usage_contribution, provider_api_key_usage_contribution, api_key_usage_contribution, provider_api_key_usage_contribution,
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure, strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
ApiKeyUsageContribution, ApiKeyUsageDelta, ProviderApiKeyUsageContribution, usage_request_metadata_client_family, ApiKeyUsageContribution, ApiKeyUsageDelta,
ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary, ProviderApiKeyUsageContribution, ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest,
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow, StoredProviderApiKeyUsageSummary, StoredProviderApiKeyWindowUsageSummary,
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery, StoredProviderUsageSummary, StoredProviderUsageWindow, StoredRequestUsageAudit,
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery, UsageDailyHeatmapQuery,
UsageReadRepository, UsageWriteRepository,
}; };
use crate::repository::auth::InMemoryAuthApiKeySnapshotRepository; use crate::repository::auth::InMemoryAuthApiKeySnapshotRepository;
use crate::repository::provider_catalog::InMemoryProviderCatalogReadRepository; use crate::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
@@ -56,6 +57,7 @@ impl InMemoryUsageReadRepository {
let mut by_request_id = BTreeMap::new(); let mut by_request_id = BTreeMap::new();
for mut item in items { for mut item in items {
hydrate_legacy_body_refs(&mut item); hydrate_legacy_body_refs(&mut item);
hydrate_client_family(&mut item);
by_request_id.insert(item.request_id.clone(), item); by_request_id.insert(item.request_id.clone(), item);
} }
Self { Self {
@@ -75,6 +77,7 @@ impl InMemoryUsageReadRepository {
let mut detached_bodies = BTreeMap::new(); let mut detached_bodies = BTreeMap::new();
for mut item in items { for mut item in items {
hydrate_legacy_body_refs(&mut item); hydrate_legacy_body_refs(&mut item);
hydrate_client_family(&mut item);
let request_id = item.request_id.clone(); let request_id = item.request_id.clone();
if let Some(body_ref) = detach_usage_body( if let Some(body_ref) = detach_usage_body(
&request_id, &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( fn persisted_usage_body_ref(
incoming_ref: Option<&str>, incoming_ref: Option<&str>,
incoming_body: Option<&Value>, 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, request_metadata,
created_at_unix_ms, created_at_unix_ms,
updated_at_unix_secs: usage.updated_at_unix_secs, updated_at_unix_secs: usage.updated_at_unix_secs,
@@ -4399,6 +4416,7 @@ mod tests {
provider_endpoint_kind: Some("chat".to_string()), provider_endpoint_kind: Some("chat".to_string()),
has_format_conversion: false, has_format_conversion: false,
is_stream: false, is_stream: false,
client_family: None,
input_tokens: 10, input_tokens: 10,
output_tokens: 20, output_tokens: 20,
total_tokens: 30, total_tokens: 30,
@@ -363,7 +363,8 @@ mod sqlite;
#[allow(unused_imports)] #[allow(unused_imports)]
pub(crate) use aether_data_contracts::repository::usage::{ pub(crate) use aether_data_contracts::repository::usage::{
PendingUsageCleanupSummary, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary, usage_request_metadata_client_family, PendingUsageCleanupSummary,
ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow, StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary, StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary, StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
@@ -6,8 +6,8 @@ use sqlx::{mysql::MySqlRow, Row};
use super::{ use super::{
provider_api_key_usage_is_error, provider_api_key_usage_is_success, provider_api_key_usage_is_error, provider_api_key_usage_is_success,
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure, strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
InMemoryUsageReadRepository, PendingUsageCleanupSummary, StoredRequestUsageAudit, usage_request_metadata_client_family, InMemoryUsageReadRepository, PendingUsageCleanupSummary,
UpsertUsageRecord, UsageWriteRepository, StoredRequestUsageAudit, UpsertUsageRecord, UsageWriteRepository,
}; };
use crate::driver::mysql::MysqlPool; use crate::driver::mysql::MysqlPool;
use crate::error::SqlResultExt; use crate::error::SqlResultExt;
@@ -850,6 +850,8 @@ fn map_usage_row(row: &MySqlRow) -> Result<StoredRequestUsageAudit, DataLayerErr
.map(|raw| serde_json::from_str(&raw)) .map(|raw| serde_json::from_str(&raw))
.transpose() .transpose()
.map_err(|err| DataLayerError::UnexpectedValue(err.to_string()))?; .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 let upstream_is_stream = row
.try_get::<Option<bool>, _>("upstream_is_stream") .try_get::<Option<bool>, _>("upstream_is_stream")
.map_sql_err()?; .map_sql_err()?;
@@ -6,8 +6,8 @@ use sqlx::{sqlite::SqliteRow, Row};
use super::{ use super::{
provider_api_key_usage_is_error, provider_api_key_usage_is_success, provider_api_key_usage_is_error, provider_api_key_usage_is_success,
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure, strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
InMemoryUsageReadRepository, PendingUsageCleanupSummary, StoredRequestUsageAudit, usage_request_metadata_client_family, InMemoryUsageReadRepository, PendingUsageCleanupSummary,
UpsertUsageRecord, UsageWriteRepository, StoredRequestUsageAudit, UpsertUsageRecord, UsageWriteRepository,
}; };
use crate::driver::sqlite::{sqlite_optional_real, sqlite_real, SqlitePool}; use crate::driver::sqlite::{sqlite_optional_real, sqlite_real, SqlitePool};
use crate::error::SqlResultExt; use crate::error::SqlResultExt;
@@ -847,6 +847,8 @@ fn map_usage_row(row: &SqliteRow) -> Result<StoredRequestUsageAudit, DataLayerEr
.map(|raw| serde_json::from_str(&raw)) .map(|raw| serde_json::from_str(&raw))
.transpose() .transpose()
.map_err(|err| DataLayerError::UnexpectedValue(err.to_string()))?; .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 let upstream_is_stream = row
.try_get::<Option<i64>, _>("upstream_is_stream") .try_get::<Option<i64>, _>("upstream_is_stream")
.map_sql_err()? .map_sql_err()?