修复 usage 详情 body 引用解包

This commit is contained in:
MMEXA
2026-07-10 00:23:31 +08:00
parent d2ea437c1c
commit f07eb25cfc
4 changed files with 166 additions and 9 deletions
@@ -58,7 +58,9 @@ pub(super) async fn admin_usage_resolve_body_value(
| Some(UsageBodyCaptureState::Unavailable)
| Some(UsageBodyCaptureState::None) => return Ok(None),
Some(UsageBodyCaptureState::Inline) | Some(UsageBodyCaptureState::Truncated) => {
return Ok(inline_body.cloned());
if inline_body.is_some() {
return Ok(inline_body.cloned());
}
}
Some(UsageBodyCaptureState::Reference) | None => {}
}
@@ -12,7 +12,7 @@ use aether_data::repository::users::{
use aether_data_contracts::repository::candidates::{
RequestCandidateStatus, StoredRequestCandidate,
};
use aether_data_contracts::repository::usage::StoredRequestUsageAudit;
use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UsageBodyCaptureState};
use axum::body::{Body, Bytes};
use axum::routing::{any, get, post};
use axum::{extract::Request, Router};
@@ -2288,6 +2288,84 @@ async fn gateway_handles_admin_usage_detail_with_ref_backed_bodies() {
upstream_handle.abort();
}
#[tokio::test]
async fn gateway_resolves_admin_usage_detail_when_inline_state_has_body_ref() {
let (_upstream_url, upstream_hits, upstream_handle) =
start_usage_upstream("/api/admin/usage/usage-inline-state-ref-detail").await;
let mut usage = sample_usage_row(
"usage-inline-state-ref-detail",
"req-inline-state-ref-detail",
Some("user-1"),
Some("key-1"),
Some("primary"),
"Gemini",
"gemini-3.5-flash",
"completed",
120,
30,
0.3,
0.36,
DAY_1_UNIX_SECS,
);
usage.request_body = Some(json!({
"contents": [{"role": "user", "parts": [{"text": "hello"}]}],
}));
usage.provider_request_body = Some(json!({
"model": "gemini-3-flash-agent",
"request": {"contents": [{"role": "user", "parts": [{"text": "hello"}]}]},
}));
usage.response_body = Some(json!({
"chunks": [{"candidates": [{"content": {"parts": [{"text": "hello back"}]}}]}],
}));
usage.client_response_body = Some(json!({
"candidates": [{"content": {"parts": [{"text": "hello back"}]}}],
}));
usage.request_body_state = Some(UsageBodyCaptureState::Inline);
usage.provider_request_body_state = Some(UsageBodyCaptureState::Inline);
usage.response_body_state = Some(UsageBodyCaptureState::Inline);
usage.client_response_body_state = Some(UsageBodyCaptureState::Inline);
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed_with_detached_bodies(
vec![usage],
));
let data_state = GatewayDataState::with_usage_reader_for_tests(usage_repository);
let gateway = build_router_with_state(
AppState::new()
.expect("gateway should build")
.with_data_state_for_tests(data_state),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let response = admin_request(reqwest::Client::new().get(format!(
"{gateway_url}/api/admin/usage/usage-inline-state-ref-detail"
)))
.send()
.await
.expect("request should succeed");
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["request_body"]["contents"][0]["role"], "user");
assert_eq!(
payload["provider_request_body"]["model"],
"gemini-3-flash-agent"
);
assert_eq!(
payload["response_body"]["chunks"][0]["candidates"][0]["content"]["parts"][0]["text"],
"hello back"
);
assert_eq!(
payload["client_response_body"]["candidates"][0]["content"]["parts"][0]["text"],
"hello back"
);
assert!(payload["body_load_errors"].is_null());
assert_eq!(*upstream_hits.lock().expect("mutex should lock"), 0);
gateway_handle.abort();
upstream_handle.abort();
}
#[tokio::test]
async fn local_admin_usage_detail_attaches_explicit_audit() {
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![sample_usage_row(
@@ -10461,6 +10461,27 @@ fn prepare_usage_body_storage(value: Option<&Value>) -> Result<UsageBodyStorage,
})
}
fn usage_body_capture_state_for_storage(
incoming_state: Option<UsageBodyCaptureState>,
storage: &UsageBodyStorage,
body_ref: Option<&str>,
) -> Option<UsageBodyCaptureState> {
if matches!(
incoming_state,
Some(
UsageBodyCaptureState::Disabled
| UsageBodyCaptureState::Unavailable
| UsageBodyCaptureState::None
)
) {
return incoming_state;
}
if storage.has_detached_blob() || body_ref.is_some() {
return Some(UsageBodyCaptureState::Reference);
}
incoming_state
}
fn json_bind_text(value: Option<&Value>) -> Result<Option<String>, DataLayerError> {
value
.map(|value| {
@@ -10515,10 +10536,26 @@ fn prepare_usage_upsert_context(
),
};
let http_audit_states = UsageHttpAuditStates {
request_body_state: usage.request_body_state,
provider_request_body_state: usage.provider_request_body_state,
response_body_state: usage.response_body_state,
client_response_body_state: usage.client_response_body_state,
request_body_state: usage_body_capture_state_for_storage(
usage.request_body_state,
&request_body_storage,
http_audit_refs.request_body_ref.as_deref(),
),
provider_request_body_state: usage_body_capture_state_for_storage(
usage.provider_request_body_state,
&provider_request_body_storage,
http_audit_refs.provider_request_body_ref.as_deref(),
),
response_body_state: usage_body_capture_state_for_storage(
usage.response_body_state,
&response_body_storage,
http_audit_refs.response_body_ref.as_deref(),
),
client_response_body_state: usage_body_capture_state_for_storage(
usage.client_response_body_state,
&client_response_body_storage,
http_audit_refs.client_response_body_ref.as_deref(),
),
};
let request_metadata_value = prepare_request_metadata_for_body_storage(
usage.request_metadata.clone(),
@@ -6,15 +6,16 @@ use super::{
attach_usage_routing_snapshot_metadata, attach_usage_settlement_pricing_snapshot_metadata,
inflate_usage_json_value, prepare_request_metadata_for_body_storage,
prepare_usage_body_storage, resolved_read_usage_body_ref, resolved_write_usage_body_ref,
split_dashboard_daily_aggregate_range, split_dashboard_hourly_aggregate_range, usage_body_ref,
usage_http_audit_body_refs, usage_http_audit_capture_mode, usage_routing_snapshot_from_usage,
split_dashboard_daily_aggregate_range, split_dashboard_hourly_aggregate_range,
usage_body_capture_state_for_storage, usage_body_ref, usage_http_audit_body_refs,
usage_http_audit_capture_mode, usage_routing_snapshot_from_usage,
usage_settlement_pricing_snapshot_from_usage, AggregateRangeSplit, SqlxUsageReadRepository,
UsageHttpAuditRefs, UsageRoutingSnapshot, UsageSettlementPricingSnapshot,
MAX_INLINE_USAGE_BODY_BYTES,
};
use crate::driver::postgres::{PostgresPoolConfig, PostgresPoolFactory};
use crate::repository::usage::UpsertUsageRecord;
use aether_data_contracts::repository::usage::UsageBodyField;
use aether_data_contracts::repository::usage::{UsageBodyCaptureState, UsageBodyField};
#[tokio::test]
async fn repository_constructs_from_lazy_pool() {
@@ -1039,6 +1040,45 @@ fn prepare_usage_body_storage_compresses_large_payloads() {
);
}
#[test]
fn usage_body_capture_state_for_storage_marks_detached_bodies_as_reference() {
let payload = json!({"message": "hello"});
let storage = prepare_usage_body_storage(Some(&payload)).expect("storage should serialize");
assert!(storage.has_detached_blob());
assert_eq!(
usage_body_capture_state_for_storage(Some(UsageBodyCaptureState::Inline), &storage, None,),
Some(UsageBodyCaptureState::Reference)
);
assert_eq!(
usage_body_capture_state_for_storage(None, &storage, None),
Some(UsageBodyCaptureState::Reference)
);
}
#[test]
fn usage_body_capture_state_for_storage_preserves_unavailable_states() {
let payload = json!({"message": "hello"});
let storage = prepare_usage_body_storage(Some(&payload)).expect("storage should serialize");
assert_eq!(
usage_body_capture_state_for_storage(
Some(UsageBodyCaptureState::Disabled),
&storage,
Some("usage://request/req-1/request_body"),
),
Some(UsageBodyCaptureState::Disabled)
);
assert_eq!(
usage_body_capture_state_for_storage(
Some(UsageBodyCaptureState::Unavailable),
&storage,
Some("usage://request/req-1/request_body"),
),
Some(UsageBodyCaptureState::Unavailable)
);
}
#[test]
fn prepare_request_metadata_for_body_storage_strips_body_ref_compatibility_keys() {
let detached = prepare_usage_body_storage(Some(&json!({