mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-10 11:19:50 +08:00
Add admin operations dashboard and usage state fixes
This commit is contained in:
@@ -7,6 +7,9 @@ use super::replay::{
|
||||
admin_usage_resolve_request_capture_body_for_item, build_admin_usage_curl_response,
|
||||
build_admin_usage_detail_payload, build_admin_usage_replay_response,
|
||||
};
|
||||
use super::summary_routes::{
|
||||
admin_usage_terminal_candidate_state_override, apply_admin_usage_state_override,
|
||||
};
|
||||
use crate::handlers::admin::request::{AdminAppState, AdminRequestContext};
|
||||
use crate::handlers::admin::shared::{attach_admin_audit_response, query_param_bool};
|
||||
use crate::GatewayError;
|
||||
@@ -233,6 +236,19 @@ pub(super) async fn maybe_build_local_admin_usage_detail_response(
|
||||
let provider_key_name = admin_usage_provider_key_name(&item, &provider_key_names);
|
||||
|
||||
let mut detail_item = item.clone();
|
||||
if matches!(detail_item.status.as_str(), "pending" | "streaming")
|
||||
&& state.has_request_candidate_data_reader()
|
||||
{
|
||||
let candidates = state
|
||||
.app()
|
||||
.read_request_candidates_by_request_id(&detail_item.request_id)
|
||||
.await?;
|
||||
if let Some(override_payload) =
|
||||
admin_usage_terminal_candidate_state_override(&candidates)
|
||||
{
|
||||
apply_admin_usage_state_override(&mut detail_item, &override_payload);
|
||||
}
|
||||
}
|
||||
let mut body_load_errors = serde_json::Map::new();
|
||||
let request_body = if include_bodies {
|
||||
let (request_body, provider_request_body, response_body, client_response_body) = tokio::join!(
|
||||
|
||||
@@ -301,7 +301,7 @@ fn admin_usage_unix_millis_to_rfc3339(unix_ms: u64) -> Option<String> {
|
||||
.map(|timestamp| timestamp.to_rfc3339())
|
||||
}
|
||||
|
||||
fn admin_usage_terminal_candidate_state_override(
|
||||
pub(super) fn admin_usage_terminal_candidate_state_override(
|
||||
candidates: &[StoredRequestCandidate],
|
||||
) -> Option<serde_json::Value> {
|
||||
let candidate = admin_usage_current_candidate(candidates)?;
|
||||
@@ -343,6 +343,49 @@ fn admin_usage_terminal_candidate_state_override(
|
||||
Some(payload)
|
||||
}
|
||||
|
||||
pub(super) fn apply_admin_usage_state_override(
|
||||
item: &mut StoredRequestUsageAudit,
|
||||
override_payload: &serde_json::Value,
|
||||
) {
|
||||
if !matches!(item.status.as_str(), "pending" | "streaming") {
|
||||
return;
|
||||
}
|
||||
|
||||
let Some(object) = override_payload.as_object() else {
|
||||
return;
|
||||
};
|
||||
|
||||
if let Some(status) = object
|
||||
.get("status")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| matches!(*value, "completed" | "failed" | "cancelled"))
|
||||
{
|
||||
item.status = status.to_string();
|
||||
}
|
||||
if let Some(status_code) = object
|
||||
.get("status_code")
|
||||
.and_then(serde_json::Value::as_u64)
|
||||
.and_then(|value| u16::try_from(value).ok())
|
||||
{
|
||||
item.status_code = Some(status_code);
|
||||
}
|
||||
if let Some(response_time_ms) = object
|
||||
.get("response_time_ms")
|
||||
.and_then(serde_json::Value::as_u64)
|
||||
{
|
||||
item.response_time_ms = Some(response_time_ms);
|
||||
}
|
||||
if let Some(error_message) = object
|
||||
.get("error_message")
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
{
|
||||
item.error_message = Some(error_message.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
fn admin_usage_matches_attempt_status(
|
||||
item: &StoredRequestUsageAudit,
|
||||
status: &str,
|
||||
@@ -901,6 +944,21 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
|
||||
}
|
||||
};
|
||||
|
||||
let active_candidate_state =
|
||||
resolve_admin_usage_active_candidate_state(state, &usage).await?;
|
||||
let usage = usage
|
||||
.into_iter()
|
||||
.map(|mut item| {
|
||||
if let Some(override_payload) = active_candidate_state
|
||||
.state_overrides_by_request_id
|
||||
.get(&item.request_id)
|
||||
{
|
||||
apply_admin_usage_state_override(&mut item, override_payload);
|
||||
}
|
||||
item
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let user_ids: Vec<String> = usage
|
||||
.iter()
|
||||
.filter_map(|item| item.user_id.clone())
|
||||
|
||||
@@ -1170,6 +1170,82 @@ async fn gateway_handles_admin_usage_active_ids_for_terminal_updates() {
|
||||
upstream_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_derives_admin_usage_records_and_detail_status_from_terminal_candidate() {
|
||||
let (_records_upstream_url, records_upstream_hits, records_upstream_handle) =
|
||||
start_usage_upstream("/api/admin/usage/records").await;
|
||||
|
||||
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![sample_usage_row(
|
||||
"usage-stale-stream",
|
||||
"req-stale-stream",
|
||||
Some("user-1"),
|
||||
Some("key-1"),
|
||||
Some("primary"),
|
||||
"OpenAI",
|
||||
"gpt-5",
|
||||
"streaming",
|
||||
20,
|
||||
5,
|
||||
0.2,
|
||||
0.24,
|
||||
DAY_1_UNIX_SECS,
|
||||
)]));
|
||||
let request_candidate_repository = Arc::new(InMemoryRequestCandidateRepository::seed(vec![
|
||||
sample_request_candidate(
|
||||
"cand-stale-stream-success",
|
||||
"req-stale-stream",
|
||||
0,
|
||||
0,
|
||||
RequestCandidateStatus::Success,
|
||||
),
|
||||
]));
|
||||
let gateway = build_router_with_state(
|
||||
AppState::new()
|
||||
.expect("gateway should build")
|
||||
.with_data_state_for_tests(
|
||||
GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
|
||||
request_candidate_repository,
|
||||
usage_repository,
|
||||
),
|
||||
),
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let records_response = admin_request(reqwest::Client::new().get(format!(
|
||||
"{gateway_url}/api/admin/usage/records?start_date=2024-03-21&end_date=2024-03-21&tz_offset_minutes=0&limit=10&offset=0"
|
||||
)))
|
||||
.send()
|
||||
.await
|
||||
.expect("records request should succeed");
|
||||
|
||||
assert_eq!(records_response.status(), StatusCode::OK);
|
||||
let records_payload: serde_json::Value = records_response
|
||||
.json()
|
||||
.await
|
||||
.expect("records json should parse");
|
||||
assert_eq!(records_payload["records"][0]["id"], "usage-stale-stream");
|
||||
assert_eq!(records_payload["records"][0]["status"], "completed");
|
||||
|
||||
let detail_response = admin_request(
|
||||
reqwest::Client::new().get(format!("{gateway_url}/api/admin/usage/usage-stale-stream")),
|
||||
)
|
||||
.send()
|
||||
.await
|
||||
.expect("detail request should succeed");
|
||||
|
||||
assert_eq!(detail_response.status(), StatusCode::OK);
|
||||
let detail_payload: serde_json::Value = detail_response
|
||||
.json()
|
||||
.await
|
||||
.expect("detail json should parse");
|
||||
assert_eq!(detail_payload["id"], "usage-stale-stream");
|
||||
assert_eq!(detail_payload["status"], "completed");
|
||||
assert_eq!(*records_upstream_hits.lock().expect("mutex should lock"), 0);
|
||||
|
||||
gateway_handle.abort();
|
||||
records_upstream_handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gateway_handles_admin_usage_records_locally_with_trusted_admin_principal() {
|
||||
let (_upstream_url, upstream_hits, upstream_handle) =
|
||||
|
||||
Reference in New Issue
Block a user