mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-07 18:07:47 +08:00
fix(usage): unify OpenAI Live and WebSocket records
This commit is contained in:
@@ -67,6 +67,7 @@ fn parse_users_me_usage_offset(query: Option<&str>) -> Result<usize, String> {
|
||||
|
||||
#[derive(Clone, Debug, Default, PartialEq, Eq)]
|
||||
struct UsersMeUsageRecordFilter {
|
||||
api_format: Option<String>,
|
||||
statuses: Option<Vec<String>>,
|
||||
is_stream: Option<bool>,
|
||||
is_websocket: Option<bool>,
|
||||
@@ -74,14 +75,19 @@ struct UsersMeUsageRecordFilter {
|
||||
}
|
||||
|
||||
fn parse_users_me_usage_record_filter(query: Option<&str>) -> UsersMeUsageRecordFilter {
|
||||
let mut filter = UsersMeUsageRecordFilter {
|
||||
api_format: query_param_value(query, "api_format")
|
||||
.map(|value| value.trim().to_string())
|
||||
.filter(|value| !value.is_empty()),
|
||||
..UsersMeUsageRecordFilter::default()
|
||||
};
|
||||
let Some(status) = query_param_value(query, "status")
|
||||
.map(|value| value.trim().to_ascii_lowercase())
|
||||
.filter(|value| !value.is_empty())
|
||||
else {
|
||||
return UsersMeUsageRecordFilter::default();
|
||||
return filter;
|
||||
};
|
||||
|
||||
let mut filter = UsersMeUsageRecordFilter::default();
|
||||
match status.as_str() {
|
||||
"stream" => {
|
||||
filter.is_stream = Some(true);
|
||||
@@ -1134,7 +1140,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
model: None,
|
||||
api_format: None,
|
||||
api_format: record_filter.api_format.clone(),
|
||||
client_family: None,
|
||||
exclude_unknown_model_or_provider: false,
|
||||
statuses: record_filter.statuses.clone(),
|
||||
@@ -1191,7 +1197,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
model: None,
|
||||
api_format: None,
|
||||
api_format: record_filter.api_format.clone(),
|
||||
client_family: None,
|
||||
exclude_unknown_model_or_provider: false,
|
||||
statuses: record_filter.statuses.clone(),
|
||||
@@ -1221,7 +1227,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
user_id: Some(auth.user.id.clone()),
|
||||
provider_name: None,
|
||||
model: None,
|
||||
api_format: None,
|
||||
api_format: record_filter.api_format.clone(),
|
||||
client_family: None,
|
||||
exclude_unknown_model_or_provider: false,
|
||||
statuses: record_filter.statuses.clone(),
|
||||
@@ -1629,6 +1635,20 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn users_me_usage_transport_statuses_are_disjoint_server_side_filters() {
|
||||
let live_websocket = parse_users_me_usage_record_filter(Some(
|
||||
"limit=20&api_format=codex%3Alive&status=websocket",
|
||||
));
|
||||
assert_eq!(live_websocket.api_format.as_deref(), Some("codex:live"));
|
||||
assert_eq!(live_websocket.is_websocket, Some(true));
|
||||
|
||||
let live_without_status =
|
||||
parse_users_me_usage_record_filter(Some("api_format=codex%3Alive"));
|
||||
assert_eq!(
|
||||
live_without_status.api_format.as_deref(),
|
||||
Some("codex:live")
|
||||
);
|
||||
assert_eq!(live_without_status.is_websocket, None);
|
||||
|
||||
for status in ["websocket", "ws", "WS"] {
|
||||
let filter = parse_users_me_usage_record_filter(Some(
|
||||
format!("limit=20&status={status}").as_str(),
|
||||
|
||||
@@ -13,7 +13,7 @@ use aether_data_contracts::repository::candidates::{
|
||||
RequestCandidateStatus, StoredRequestCandidate,
|
||||
};
|
||||
use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UsageBodyCaptureState};
|
||||
use axum::body::{Body, Bytes};
|
||||
use axum::body::{to_bytes, Body, Bytes};
|
||||
use axum::routing::{any, get, post};
|
||||
use axum::{extract::Request, Router};
|
||||
use http::{HeaderMap, HeaderValue, StatusCode};
|
||||
@@ -2173,6 +2173,82 @@ async fn gateway_handles_admin_usage_detail_locally_with_trusted_admin_principal
|
||||
upstream_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_admin_usage_detail_preserves_live_websocket_session_metadata() {
|
||||
let live_session = json!({
|
||||
"schema_version": "1",
|
||||
"transport": "websocket",
|
||||
"mode": "direct",
|
||||
"state": "closed",
|
||||
"client_frames": 3,
|
||||
"upstream_frames": 5,
|
||||
});
|
||||
let realtime_session = json!({
|
||||
"schema_version": "1",
|
||||
"transport": "websocket",
|
||||
"usage_state": "authoritative",
|
||||
"input_audio_tokens": 11,
|
||||
"output_audio_tokens": 6,
|
||||
});
|
||||
let mut usage = sample_usage_row(
|
||||
"usage-live-detail",
|
||||
"req-live-detail",
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
"OpenAI",
|
||||
"gpt-live",
|
||||
"completed",
|
||||
0,
|
||||
0,
|
||||
0.0,
|
||||
0.0,
|
||||
DAY_1_UNIX_SECS,
|
||||
);
|
||||
usage.request_type = Some("live".to_string());
|
||||
usage.api_format = Some("codex:live".to_string());
|
||||
usage.api_family = Some("codex".to_string());
|
||||
usage.endpoint_kind = Some("live".to_string());
|
||||
usage.endpoint_api_format = Some("codex:live".to_string());
|
||||
usage.provider_api_family = Some("codex".to_string());
|
||||
usage.provider_endpoint_kind = Some("live".to_string());
|
||||
usage.is_stream = true;
|
||||
usage.request_metadata = Some(json!({
|
||||
"websocket_mode": true,
|
||||
"websocket_transport": "codex_live_direct",
|
||||
"usage_available": false,
|
||||
"usage_pricing_available": false,
|
||||
"live_session": live_session,
|
||||
"realtime_session": realtime_session,
|
||||
}));
|
||||
|
||||
let state = AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(GatewayDataState::with_usage_reader_for_tests(Arc::new(
|
||||
InMemoryUsageReadRepository::seed(vec![usage]),
|
||||
)));
|
||||
let response = local_admin_usage_response(
|
||||
&state,
|
||||
http::Method::GET,
|
||||
"/api/admin/usage/usage-live-detail?include_bodies=false",
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = to_bytes(response.into_body(), 1024 * 1024)
|
||||
.await
|
||||
.expect("usage detail body should be readable");
|
||||
let payload: serde_json::Value =
|
||||
serde_json::from_slice(&body).expect("usage detail body should be JSON");
|
||||
assert_eq!(payload["is_websocket"], true);
|
||||
assert_eq!(payload["websocket_transport"], "codex_live_direct");
|
||||
assert_eq!(payload["live_session"], live_session);
|
||||
assert_eq!(payload["realtime_session"], realtime_session);
|
||||
assert_eq!(payload["metadata"]["live_session"], live_session);
|
||||
assert_eq!(payload["metadata"]["realtime_session"], realtime_session);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_handles_admin_usage_detail_with_ref_backed_bodies() {
|
||||
let (_upstream_url, upstream_hits, upstream_handle) =
|
||||
|
||||
Reference in New Issue
Block a user