From 07cb401fd42f9ffbbf566fbb83b57ecb9e202e2f Mon Sep 17 00:00:00 2001 From: AAEE86 Date: Tue, 22 Sep 2026 11:54:58 +0800 Subject: [PATCH] fix: extract nested provider response models --- .../contracts/src/repository/usage/types.rs | 186 ++++++++++++++++-- crates/aether-usage/runtime/src/event_wire.rs | 2 + crates/aether-usage/runtime/src/record.rs | 2 + .../runtime/src/request_metadata.rs | 14 ++ crates/aether-usage/runtime/src/runtime.rs | 2 + crates/aether-usage/runtime/src/write.rs | 6 + 6 files changed, 201 insertions(+), 11 deletions(-) diff --git a/crates/aether-data/contracts/src/repository/usage/types.rs b/crates/aether-data/contracts/src/repository/usage/types.rs index 7382b6867..a01145409 100644 --- a/crates/aether-data/contracts/src/repository/usage/types.rs +++ b/crates/aether-data/contracts/src/repository/usage/types.rs @@ -1,3 +1,4 @@ +use aether_ai_formats::normalize_api_format_alias; use async_trait::async_trait; use chrono::{DateTime, Utc}; use serde_json::Value; @@ -120,7 +121,7 @@ pub fn normalize_provider_service_tier(value: &str) -> Option { Some(value.to_ascii_lowercase()) } -/// 清洗响应体顶层 `model`,保留大小写,只去除首尾空白。 +/// 清洗模型名称,保留大小写,只去除首尾空白。 pub fn normalize_provider_response_model(value: &str) -> Option { let value = value.trim(); if value.is_empty() || value.len() > 256 { @@ -129,12 +130,117 @@ pub fn normalize_provider_response_model(value: &str) -> Option { Some(value.to_string()) } +fn extract_model_at_paths(value: &Value, paths: &[&[&str]]) -> Option { + paths.iter().find_map(|path| { + let value = path + .iter() + .try_fold(value, |current, key| current.as_object()?.get(*key))?; + value.as_str().and_then(normalize_provider_response_model) + }) +} + +fn response_model_paths(provider_api_format: Option<&str>) -> &'static [&'static [&'static str]] { + match normalize_api_format_alias(provider_api_format.unwrap_or_default()).as_str() { + "gemini:generate_content" => { + // Gemini 原生响应使用 modelVersion;部分网关会改写为 model。 + &[&["modelVersion"], &["model_version"], &["model"]] + } + "gemini:embedding" => { + // Gemini Embedding 可能返回 model、modelVersion 或 Vertex 的 deployedModelId。 + &[ + &["model"], + &["modelVersion"], + &["model_version"], + &["deployedModelId"], + &["deployed_model_id"], + ] + } + "gemini:interactions" => { + // Interactions 请求既可能叫 model,也可能叫 agent;响应优先读取 model。 + &[ + &["model"], + &["modelVersion"], + &["model_version"], + &["agent"], + ] + } + _ => &[&["model"]], + } +} + +fn extract_model_from_known_response_wrappers( + response_body: &Value, + paths: &[&[&str]], +) -> Option { + // 只展开协议中已知的 response/chunks 包装,避免在候选内容、工具参数等任意嵌套 + // 对象中搜索同名字段,误把 role="model" 一类内容当成响应模型。 + extract_model_at_paths(response_body, paths) + .or_else(|| { + response_body + .get("response") + .and_then(|response| extract_model_at_paths(response, paths)) + }) + .or_else(|| { + response_body + .get("chunks") + .and_then(Value::as_array) + .and_then(|chunks| { + chunks.iter().rev().find_map(|chunk| { + extract_model_at_paths(chunk, paths).or_else(|| { + chunk + .get("response") + .and_then(|response| extract_model_at_paths(response, paths)) + }) + }) + }) + }) + .or_else(|| { + response_body + .get("response") + .and_then(|response| response.get("chunks")) + .and_then(Value::as_array) + .and_then(|chunks| { + chunks.iter().rev().find_map(|chunk| { + extract_model_at_paths(chunk, paths).or_else(|| { + chunk + .get("response") + .and_then(|response| extract_model_at_paths(response, paths)) + }) + }) + }) + }) +} + +fn extract_provider_model_from_response_body( + response_body: &Value, + provider_api_format: Option<&str>, +) -> Option { + extract_model_from_known_response_wrappers( + response_body, + response_model_paths(provider_api_format), + ) +} + +fn extract_provider_model_from_request_body( + request_body: &Value, + request_api_format: Option<&str>, +) -> Option { + let paths: &[&[&str]] = + match normalize_api_format_alias(request_api_format.unwrap_or_default()).as_str() { + "gemini:interactions" => &[&["model"], &["agent"]], + _ => &[&["model"]], + }; + extract_model_at_paths(request_body, paths) +} + /// 只有请求体和响应体都可作为完整事实时,才计算响应模型,避免用截断内容猜测。 pub fn extract_provider_response_model_from_bodies( request_body: Option<&Value>, request_body_state: Option, + request_api_format: Option<&str>, response_body: Option<&Value>, response_body_state: Option, + provider_api_format: Option<&str>, ) -> Option { if !usage_body_capture_is_authoritative(request_body, request_body_state) || !usage_body_capture_is_authoritative(response_body, response_body_state) @@ -142,16 +248,10 @@ pub fn extract_provider_response_model_from_bodies( return None; } - let request_model = request_body - .and_then(Value::as_object) - .and_then(|body| body.get("model")) - .and_then(Value::as_str) - .and_then(normalize_provider_response_model)?; - let response_model = response_body - .and_then(Value::as_object) - .and_then(|body| body.get("model")) - .and_then(Value::as_str) - .and_then(normalize_provider_response_model)?; + let request_model = + extract_provider_model_from_request_body(request_body?, request_api_format)?; + let response_model = + extract_provider_model_from_response_body(response_body?, provider_api_format)?; (request_model != response_model).then_some(response_model) } @@ -2723,8 +2823,10 @@ mod tests { extract_provider_response_model_from_bodies( Some(&request), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), Some(&response), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), ), Some("gpt-5.1".to_string()) ); @@ -2732,8 +2834,10 @@ mod tests { extract_provider_response_model_from_bodies( Some(&request), Some(UsageBodyCaptureState::Inline), + Some("openai:responses"), Some(&json!({"model": "gpt-5"})), Some(UsageBodyCaptureState::Inline), + Some("openai:responses"), ), None ); @@ -2741,8 +2845,10 @@ mod tests { extract_provider_response_model_from_bodies( Some(&request), Some(UsageBodyCaptureState::Truncated), + Some("openai:chat"), Some(&response), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), ), None ); @@ -2760,8 +2866,66 @@ mod tests { extract_provider_response_model_from_bodies( Some(&json!({"model": "gpt-5"})), None, + Some("openai:chat"), Some(&json!({"model": 42})), None, + Some("openai:chat"), + ), + None + ); + } + + #[test] + fn response_model_uses_provider_format_specific_nested_paths() { + let request = json!({"model": "gemini-2.5-flash"}); + let response = json!({ + "response": { + "modelVersion": "gemini-2.5-flash-001", + "candidates": [{"content": {"role": "model"}}] + } + }); + assert_eq!( + extract_provider_response_model_from_bodies( + Some(&request), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), + Some(&response), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), + ), + Some("gemini-2.5-flash-001".to_string()) + ); + + let wrapped_chunks = json!({ + "chunks": [ + {"response": {"modelVersion": "gemini-old"}}, + {"response": {"modelVersion": "gemini-final"}} + ] + }); + assert_eq!( + extract_provider_response_model_from_bodies( + Some(&request), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), + Some(&wrapped_chunks), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), + ), + Some("gemini-final".to_string()) + ); + + let ambiguous = json!({ + "metadata": {"model": "do-not-use"}, + "candidates": [{"content": {"role": "model"}}] + }); + assert_eq!( + extract_provider_response_model_from_bodies( + Some(&request), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), + Some(&ambiguous), + Some(UsageBodyCaptureState::Inline), + Some("gemini:generate_content"), ), None ); diff --git a/crates/aether-usage/runtime/src/event_wire.rs b/crates/aether-usage/runtime/src/event_wire.rs index 189f1b6d1..60427df40 100644 --- a/crates/aether-usage/runtime/src/event_wire.rs +++ b/crates/aether-usage/runtime/src/event_wire.rs @@ -242,8 +242,10 @@ impl WireOverrides { metadata, request_body, data.request_body_state, + data.api_format.as_deref(), data.response_body.as_ref(), data.response_body_state, + data.endpoint_api_format.as_deref(), ); // Billing reads raw-body TTL before metadata regardless of capture state. // Preserve that precedence independently of reasoning and tier authority. diff --git a/crates/aether-usage/runtime/src/record.rs b/crates/aether-usage/runtime/src/record.rs index 8433168ab..b7436e7b0 100644 --- a/crates/aether-usage/runtime/src/record.rs +++ b/crates/aether-usage/runtime/src/record.rs @@ -90,8 +90,10 @@ pub fn build_upsert_usage_record_from_event( data.request_metadata, data.request_body.as_ref(), data.request_body_state, + data.api_format.as_deref(), data.response_body.as_ref(), data.response_body_state, + data.endpoint_api_format.as_deref(), ); let now_unix_secs = event.timestamp_ms / 1_000; diff --git a/crates/aether-usage/runtime/src/request_metadata.rs b/crates/aether-usage/runtime/src/request_metadata.rs index 708bfc2b1..3d61866e8 100644 --- a/crates/aether-usage/runtime/src/request_metadata.rs +++ b/crates/aether-usage/runtime/src/request_metadata.rs @@ -255,8 +255,10 @@ pub(crate) fn attach_provider_response_model_metadata( metadata: Option, request_body: Option<&Value>, request_body_state: Option, + request_api_format: Option<&str>, response_body: Option<&Value>, response_body_state: Option, + provider_api_format: Option<&str>, ) -> Option { let both_bodies_are_authoritative = usage_body_capture_is_authoritative(request_body, request_body_state) @@ -264,8 +266,10 @@ pub(crate) fn attach_provider_response_model_metadata( let response_model = extract_provider_response_model_from_bodies( request_body, request_body_state, + request_api_format, response_body, response_body_state, + provider_api_format, ); if !both_bodies_are_authoritative && response_model.is_none() { return metadata; @@ -293,8 +297,10 @@ pub(crate) fn refresh_provider_response_model_metadata( metadata: Option, request_body: Option<&Value>, request_body_state: Option, + request_api_format: Option<&str>, response_body: Option<&Value>, response_body_state: Option, + provider_api_format: Option<&str>, ) -> Option { let mut object = match metadata { Some(Value::Object(object)) => object, @@ -308,8 +314,10 @@ pub(crate) fn refresh_provider_response_model_metadata( if let Some(response_model) = extract_provider_response_model_from_bodies( request_body, request_body_state, + request_api_format, response_body, response_body_state, + provider_api_format, ) { object.insert( PROVIDER_RESPONSE_MODEL_METADATA_KEY.to_string(), @@ -881,8 +889,10 @@ mod tests { Some(json!({"provider_response_model": "old-model", "trace_id": "trace-1"})), Some(&json!({"model": "gpt-5"})), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), Some(&json!({"model": "gpt-5.1"})), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), ) .expect("response model should be attached"); assert_eq!(metadata["provider_response_model"], "gpt-5.1"); @@ -892,8 +902,10 @@ mod tests { Some(json!({"provider_response_model": "gpt-5.1"})), Some(&json!({"model": "gpt-5"})), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), Some(&json!({"model": "gpt-5"})), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), ); assert!(metadata.is_none()); @@ -901,8 +913,10 @@ mod tests { Some(json!({"provider_response_model": "gpt-5.1"})), None, Some(UsageBodyCaptureState::Disabled), + Some("openai:chat"), Some(&json!({"model": "gpt-5.2"})), Some(UsageBodyCaptureState::Inline), + Some("openai:chat"), ); assert!(metadata.is_none()); } diff --git a/crates/aether-usage/runtime/src/runtime.rs b/crates/aether-usage/runtime/src/runtime.rs index 9b1152f95..3a571be5e 100644 --- a/crates/aether-usage/runtime/src/runtime.rs +++ b/crates/aether-usage/runtime/src/runtime.rs @@ -5304,8 +5304,10 @@ fn preserve_provider_response_facts(event: &mut UsageEvent) { metadata, event.data.request_body.as_ref(), event.data.request_body_state, + event.data.api_format.as_deref(), event.data.response_body.as_ref(), event.data.response_body_state, + event.data.endpoint_api_format.as_deref(), ); } diff --git a/crates/aether-usage/runtime/src/write.rs b/crates/aether-usage/runtime/src/write.rs index dc193c0c3..1fb82a2be 100644 --- a/crates/aether-usage/runtime/src/write.rs +++ b/crates/aether-usage/runtime/src/write.rs @@ -729,8 +729,10 @@ fn build_terminal_usage_event_from_seed_impl( request_metadata, request_body.as_ref(), body_states.request_body_state, + Some(client_contract.as_str()), provider_response.as_ref(), body_states.response_body_state, + Some(provider_contract.as_str()), ); let mut data = UsageEventData { @@ -1059,8 +1061,10 @@ pub fn build_sync_terminal_usage_seed( request_metadata, context_seed.request_body.as_ref(), context_seed.body_states.request_body_state, + Some(context_seed.client_contract.as_str()), provider_response_full.as_ref(), provider_response_body_state, + Some(context_seed.provider_contract.as_str()), ); TerminalUsageSeed { @@ -1244,8 +1248,10 @@ pub fn build_stream_terminal_usage_seed( request_metadata, context_seed.request_body.as_ref(), context_seed.body_states.request_body_state, + Some(context_seed.client_contract.as_str()), provider_response_full.as_ref(), provider_response_body_state, + Some(context_seed.provider_contract.as_str()), ); // The parser's terminal summary is authoritative when a response body is truncated or the // body and summary disagree; attach it after the body refresh so it wins.