mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
Fix/usage transfer filter (#307)
* fix(usage): 恢复 usage 列表中的 fallback 路由信号 * fix(admin): 支持 usage 记录展示和筛选 fallback 转移 * fix(usage): 在用户 usage 记录中暴露 fallback 标记 * fix(usage): 共享 usage 页面支持 fallback 筛选 * fix(usage): 对齐 fallback 筛选相关前端类型 * fix(usage): propagate has_fallback through active polling --------- Co-authored-by: fawney19 <elky0401@gmail.com>
This commit is contained in:
@@ -5,10 +5,11 @@ use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
||||
use crate::handlers::admin::shared::query_param_value;
|
||||
use crate::GatewayError;
|
||||
use aether_admin::observability::usage::{
|
||||
admin_usage_bad_request_response, admin_usage_data_unavailable_response, admin_usage_parse_ids,
|
||||
admin_usage_parse_limit, admin_usage_parse_offset, build_admin_usage_active_requests_response,
|
||||
build_admin_usage_records_response, build_admin_usage_summary_stats_response_from_summary,
|
||||
ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL,
|
||||
admin_usage_bad_request_response, admin_usage_data_unavailable_response,
|
||||
admin_usage_has_fallback, admin_usage_matches_search, admin_usage_matches_username,
|
||||
admin_usage_parse_ids, admin_usage_parse_limit, admin_usage_parse_offset,
|
||||
build_admin_usage_active_requests_response, build_admin_usage_records_response,
|
||||
build_admin_usage_summary_stats_response_from_summary, ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL,
|
||||
};
|
||||
use aether_data_contracts::repository::usage::{
|
||||
StoredRequestUsageAudit, UsageAuditKeywordSearchQuery, UsageAuditListQuery,
|
||||
@@ -54,6 +55,7 @@ fn apply_admin_usage_status_filter(query: &mut UsageAuditListQuery, status: Opti
|
||||
"pending" | "streaming" | "completed" | "cancelled" => {
|
||||
query.statuses = Some(vec![status.to_string()]);
|
||||
}
|
||||
"has_fallback" => {}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
@@ -331,6 +333,9 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
|
||||
Ok(value) => value,
|
||||
Err(detail) => return Ok(Some(admin_usage_bad_request_response(detail))),
|
||||
};
|
||||
let has_fallback_only = query_param_value(query, "status")
|
||||
.as_deref()
|
||||
.is_some_and(|value| value.trim().eq_ignore_ascii_case("has_fallback"));
|
||||
let search = query_param_value(query, "search");
|
||||
let username_filter = query_param_value(query, "username");
|
||||
let limit = match admin_usage_parse_limit(query) {
|
||||
@@ -367,7 +372,44 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
|
||||
let active_username_filter = username_filter
|
||||
.as_deref()
|
||||
.filter(|value| !value.trim().is_empty());
|
||||
let (usage, total) = if active_search.is_some() || active_username_filter.is_some() {
|
||||
let (usage, total) = if has_fallback_only {
|
||||
let mut usage = state.list_usage_audits(&base_query).await?;
|
||||
let user_ids: Vec<String> = usage
|
||||
.iter()
|
||||
.filter_map(|item| item.user_id.clone())
|
||||
.collect::<BTreeSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
let users_by_id: BTreeMap<
|
||||
String,
|
||||
aether_data::repository::users::StoredUserSummary,
|
||||
> = state.resolve_auth_user_summaries_by_ids(&user_ids).await?;
|
||||
let api_key_names = admin_usage_api_key_names(state, &usage).await?;
|
||||
|
||||
usage.retain(|item| {
|
||||
admin_usage_matches_search(
|
||||
item,
|
||||
active_search,
|
||||
&users_by_id,
|
||||
&api_key_names,
|
||||
state.has_auth_user_data_reader(),
|
||||
state.has_auth_api_key_data_reader(),
|
||||
) && admin_usage_matches_username(
|
||||
item,
|
||||
active_username_filter,
|
||||
&users_by_id,
|
||||
state.has_auth_user_data_reader(),
|
||||
) && admin_usage_has_fallback(item)
|
||||
});
|
||||
sort_usage_newest_first(&mut usage);
|
||||
let total = usage.len();
|
||||
let records = usage
|
||||
.into_iter()
|
||||
.skip(offset)
|
||||
.take(limit)
|
||||
.collect::<Vec<_>>();
|
||||
(records, total)
|
||||
} else if active_search.is_some() || active_username_filter.is_some() {
|
||||
let keywords = active_search
|
||||
.map(parse_admin_usage_search_keywords)
|
||||
.unwrap_or_default();
|
||||
|
||||
@@ -217,6 +217,7 @@ fn build_users_me_usage_record_payload(
|
||||
"first_byte_time_ms": item.first_byte_time_ms,
|
||||
"is_stream": item.is_stream,
|
||||
"status": item.status,
|
||||
"has_fallback": item.has_fallback(),
|
||||
"created_at": unix_secs_to_rfc3339(item.created_at_unix_ms),
|
||||
"cache_creation_input_tokens": item.cache_creation_input_tokens,
|
||||
"cache_creation_ephemeral_5m_input_tokens": item.cache_creation_ephemeral_5m_input_tokens,
|
||||
@@ -265,6 +266,7 @@ fn build_users_me_usage_active_payload(item: &StoredRequestUsageAudit) -> serde_
|
||||
"endpoint_api_format": item.endpoint_api_format,
|
||||
"has_format_conversion": item.has_format_conversion,
|
||||
"target_model": item.target_model,
|
||||
"has_fallback": item.has_fallback(),
|
||||
});
|
||||
if item.api_format.is_none() {
|
||||
payload
|
||||
|
||||
@@ -745,21 +745,25 @@ async fn gateway_handles_admin_usage_active_locally_with_trusted_admin_principal
|
||||
start_usage_upstream("/api/admin/usage/active").await;
|
||||
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![
|
||||
sample_usage_row(
|
||||
"usage-pending",
|
||||
"req-pending",
|
||||
Some("user-1"),
|
||||
Some("key-1"),
|
||||
Some("primary"),
|
||||
"OpenAI",
|
||||
"gpt-5",
|
||||
"pending",
|
||||
10,
|
||||
0,
|
||||
0.0,
|
||||
0.0,
|
||||
DAY_2_UNIX_SECS,
|
||||
),
|
||||
{
|
||||
let mut row = sample_usage_row(
|
||||
"usage-pending",
|
||||
"req-pending",
|
||||
Some("user-1"),
|
||||
Some("key-1"),
|
||||
Some("primary"),
|
||||
"OpenAI",
|
||||
"gpt-5",
|
||||
"pending",
|
||||
10,
|
||||
0,
|
||||
0.0,
|
||||
0.0,
|
||||
DAY_2_UNIX_SECS,
|
||||
);
|
||||
row.candidate_index = Some(1);
|
||||
row
|
||||
},
|
||||
sample_usage_row(
|
||||
"usage-done",
|
||||
"req-done",
|
||||
@@ -819,6 +823,7 @@ async fn gateway_handles_admin_usage_active_locally_with_trusted_admin_principal
|
||||
assert_eq!(payload["requests"][0]["effective_input_tokens"], 5);
|
||||
assert_eq!(payload["requests"][0]["provider"], "OpenAI");
|
||||
assert_eq!(payload["requests"][0]["api_key_name"], "fresh-primary");
|
||||
assert_eq!(payload["requests"][0]["has_fallback"], true);
|
||||
assert_eq!(
|
||||
payload["requests"][0]["provider_key_name"],
|
||||
"upstream-primary"
|
||||
@@ -1107,6 +1112,77 @@ async fn gateway_handles_admin_usage_records_with_provider_key_name_fallback_fro
|
||||
upstream_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_filters_admin_usage_records_by_has_fallback_status() {
|
||||
let (_upstream_url, upstream_hits, upstream_handle) =
|
||||
start_usage_upstream("/api/admin/usage/records").await;
|
||||
|
||||
let mut fallback_usage = sample_usage_row(
|
||||
"usage-has-fallback",
|
||||
"req-has-fallback",
|
||||
Some("user-1"),
|
||||
Some("key-1"),
|
||||
Some("primary"),
|
||||
"OpenAI",
|
||||
"gpt-5",
|
||||
"completed",
|
||||
12,
|
||||
8,
|
||||
0.02,
|
||||
0.02,
|
||||
DAY_1_UNIX_SECS,
|
||||
);
|
||||
fallback_usage.candidate_index = Some(1);
|
||||
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![
|
||||
fallback_usage,
|
||||
sample_usage_row(
|
||||
"usage-no-fallback",
|
||||
"req-no-fallback",
|
||||
Some("user-1"),
|
||||
Some("key-1"),
|
||||
Some("primary"),
|
||||
"OpenAI",
|
||||
"gpt-5",
|
||||
"completed",
|
||||
12,
|
||||
8,
|
||||
0.02,
|
||||
0.02,
|
||||
DAY_1_UNIX_SECS,
|
||||
),
|
||||
]));
|
||||
let user_repository = Arc::new(InMemoryUserReadRepository::seed(vec![sample_user_summary(
|
||||
"user-1", "alice",
|
||||
)]));
|
||||
let gateway = build_router_with_state(
|
||||
AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
GatewayDataState::with_usage_reader_for_tests(usage_repository)
|
||||
.with_user_reader(user_repository),
|
||||
),
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let response = admin_request(reqwest::Client::new().get(format!(
|
||||
"{gateway_url}/api/admin/usage/records?start_date=2024-03-21&end_date=2024-03-22&tz_offset_minutes=0&status=has_fallback&limit=10&offset=0"
|
||||
)))
|
||||
.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["total"], 1);
|
||||
assert_eq!(payload["records"][0]["id"], "usage-has-fallback");
|
||||
assert_eq!(payload["records"][0]["has_fallback"], true);
|
||||
assert_eq!(*upstream_hits.lock().expect("mutex should lock"), 0);
|
||||
|
||||
gateway_handle.abort();
|
||||
upstream_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_handles_admin_usage_records_with_snapshot_first_user_and_api_key_names() {
|
||||
let (_upstream_url, upstream_hits, upstream_handle) =
|
||||
|
||||
@@ -4795,6 +4795,7 @@ async fn gateway_handles_users_me_usage_locally_without_proxying_upstream() {
|
||||
"streaming",
|
||||
now - chrono::Duration::minutes(5),
|
||||
);
|
||||
streaming_usage.candidate_index = Some(2);
|
||||
streaming_usage.request_metadata = Some(json!({
|
||||
"rate_multiplier": 0.5,
|
||||
"input_price_per_1m": 3.0,
|
||||
@@ -4886,6 +4887,7 @@ async fn gateway_handles_users_me_usage_locally_without_proxying_upstream() {
|
||||
assert_eq!(payload["records"][0]["output_price_per_1m"], 9.0);
|
||||
assert_eq!(payload["records"][0]["cache_creation_price_per_1m"], 3.75);
|
||||
assert_eq!(payload["records"][0]["cache_read_price_per_1m"], 0.3);
|
||||
assert_eq!(payload["records"][0]["has_fallback"], true);
|
||||
assert_eq!(payload["records"][0]["api_key"]["name"], "renamed-key");
|
||||
assert_eq!(payload["records"][0]["api_key"]["display"], "renamed-key");
|
||||
assert_eq!(
|
||||
|
||||
Reference in New Issue
Block a user