feat(usage): expose request timing details

This commit is contained in:
elky
2026-06-10 09:16:15 +08:00
parent 84f41dae77
commit 8edcbdcb29
41 changed files with 1884 additions and 335 deletions
+9 -3
View File
@@ -11,9 +11,15 @@ pub(crate) async fn read_request_candidate_trace(
request_id: &str,
attempted_only: bool,
) -> Result<Option<RequestCandidateTrace>, DataLayerError> {
let all_candidates = state
.list_request_candidates_by_request_id(request_id)
.await?;
let all_candidates = if attempted_only {
state
.list_attempted_request_candidates_by_request_id(request_id)
.await?
} else {
state
.list_request_candidates_by_request_id(request_id)
.await?
};
Ok(RequestCandidateTrace::from_candidates(
request_id,
all_candidates,
@@ -19,6 +19,16 @@ impl GatewayDataState {
}
}
pub(crate) async fn list_attempted_request_candidates_by_request_id(
&self,
request_id: &str,
) -> Result<Vec<StoredRequestCandidate>, DataLayerError> {
match &self.request_candidate_reader {
Some(repository) => repository.list_attempted_by_request_id(request_id).await,
None => Ok(Vec::new()),
}
}
pub(crate) async fn list_request_candidates_by_provider_id(
&self,
provider_id: &str,
@@ -1124,6 +1124,16 @@ impl GatewayDataState {
}
}
pub(crate) async fn find_request_usage_by_request_id_shallow(
&self,
request_id: &str,
) -> Result<Option<StoredRequestUsageAudit>, DataLayerError> {
match &self.usage_reader {
Some(repository) => repository.find_by_request_id_shallow(request_id).await,
None => Ok(None),
}
}
pub(crate) async fn find_request_usage_by_id(
&self,
usage_id: &str,
@@ -1958,6 +1968,14 @@ impl GatewayDataState {
self.find_request_usage_by_request_id(request_id).await
}
pub(crate) async fn read_request_usage_audit_shallow(
&self,
request_id: &str,
) -> Result<Option<StoredRequestUsageAudit>, DataLayerError> {
self.find_request_usage_by_request_id_shallow(request_id)
.await
}
pub(crate) async fn read_request_audit_bundle(
&self,
request_id: &str,
@@ -89,7 +89,7 @@ async fn resolve_admin_monitoring_trace(
{
let usage = app
.data
.read_request_usage_audit(request_id)
.read_request_usage_audit_shallow(request_id)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))?;
return Ok(Some(ResolvedAdminMonitoringTrace { trace, usage }));
@@ -98,7 +98,7 @@ async fn resolve_admin_monitoring_trace(
let mut usage_candidates = Vec::new();
if let Some(usage) = app
.data
.read_request_usage_audit(request_id)
.read_request_usage_audit_shallow(request_id)
.await
.map_err(|err| GatewayError::Internal(err.to_string()))?
{
@@ -14,17 +14,76 @@ use aether_admin::observability::usage::{
admin_usage_bad_request_response, admin_usage_data_unavailable_response,
admin_usage_provider_key_name, ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL,
};
use aether_data_contracts::repository::usage::UsageBodyField;
use aether_data_contracts::repository::usage::{StoredRequestUsageAudit, UsageBodyField};
use axum::{
body::Body,
http,
response::{IntoResponse, Response},
Json,
};
use serde_json::json;
use serde_json::{json, Value};
use std::collections::BTreeMap;
use tokio::try_join;
struct AdminUsageDetailBodyValue {
value: Option<Value>,
load_failed: bool,
}
async fn resolve_admin_usage_detail_request_body(
state: &AdminAppState<'_>,
item: &StoredRequestUsageAudit,
) -> AdminUsageDetailBodyValue {
match admin_usage_resolve_request_capture_body_for_item(state, item, None).await {
Ok(body) => AdminUsageDetailBodyValue {
value: body,
load_failed: false,
},
Err(err) => {
tracing::warn!(
error = ?err,
usage_id = %item.id,
request_id = %item.request_id,
field = UsageBodyField::RequestBody.as_storage_field(),
"failed to resolve admin usage detail body"
);
let value = admin_usage_resolve_request_capture_body(item, None);
AdminUsageDetailBodyValue {
load_failed: value.is_none(),
value,
}
}
}
}
async fn resolve_admin_usage_detail_body_value(
state: &AdminAppState<'_>,
item: &StoredRequestUsageAudit,
field: UsageBodyField,
) -> AdminUsageDetailBodyValue {
let inline_body = item.body_value(field);
match admin_usage_resolve_body_value(state, item, inline_body, field).await {
Ok(body) => AdminUsageDetailBodyValue {
value: body,
load_failed: false,
},
Err(err) => {
tracing::warn!(
error = ?err,
usage_id = %item.id,
request_id = %item.request_id,
field = field.as_storage_field(),
"failed to resolve admin usage detail body"
);
let value = inline_body.cloned();
AdminUsageDetailBodyValue {
load_failed: value.is_none(),
value,
}
}
}
}
pub(super) async fn maybe_build_local_admin_usage_detail_response(
state: &AdminAppState<'_>,
request_context: &AdminRequestContext<'_>,
@@ -174,32 +233,40 @@ 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();
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) = try_join!(
admin_usage_resolve_request_capture_body_for_item(state, &item, None),
admin_usage_resolve_body_value(
let (request_body, provider_request_body, response_body, client_response_body) = tokio::join!(
resolve_admin_usage_detail_request_body(state, &item),
resolve_admin_usage_detail_body_value(
state,
&item,
item.provider_request_body.as_ref(),
UsageBodyField::ProviderRequestBody,
),
admin_usage_resolve_body_value(
resolve_admin_usage_detail_body_value(
state,
&item,
item.response_body.as_ref(),
UsageBodyField::ResponseBody,
),
admin_usage_resolve_body_value(
resolve_admin_usage_detail_body_value(
state,
&item,
item.client_response_body.as_ref(),
UsageBodyField::ClientResponseBody,
),
)?;
detail_item.provider_request_body = provider_request_body;
detail_item.response_body = response_body;
detail_item.client_response_body = client_response_body;
request_body
);
for (field, resolved) in [
(UsageBodyField::RequestBody, &request_body),
(UsageBodyField::ProviderRequestBody, &provider_request_body),
(UsageBodyField::ResponseBody, &response_body),
(UsageBodyField::ClientResponseBody, &client_response_body),
] {
if resolved.load_failed {
body_load_errors.insert(field.as_storage_field().to_string(), json!(true));
}
}
detail_item.provider_request_body = provider_request_body.value;
detail_item.response_body = response_body.value;
detail_item.client_response_body = client_response_body.value;
request_body.value
} else {
None
};
@@ -207,7 +274,7 @@ pub(super) async fn maybe_build_local_admin_usage_detail_response(
// request_body 已通过 request capture 解析;其余 detached body 在上方并行加载。
}
let default_headers = admin_usage_curl_headers();
let payload = build_admin_usage_detail_payload(
let mut payload = build_admin_usage_detail_payload(
&detail_item,
&users_by_id,
&api_key_names,
@@ -218,6 +285,11 @@ pub(super) async fn maybe_build_local_admin_usage_detail_response(
request_body,
&default_headers,
);
payload["body_load_errors"] = if include_bodies && !body_load_errors.is_empty() {
Value::Object(body_load_errors)
} else {
Value::Null
};
return Ok(Some(attach_admin_audit_response(
Json(payload).into_response(),
@@ -5,13 +5,12 @@ 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_client_family,
admin_usage_data_unavailable_response, admin_usage_has_fallback, admin_usage_is_failed,
admin_usage_matches_search, admin_usage_matches_username, admin_usage_parse_ids,
admin_usage_parse_limit, admin_usage_parse_offset, admin_usage_provider_key_name,
admin_usage_record_json, 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_is_failed, admin_usage_matches_search,
admin_usage_matches_username, admin_usage_parse_ids, admin_usage_parse_limit,
admin_usage_parse_offset, admin_usage_provider_key_name, admin_usage_record_json,
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::repository::users::StoredUserSummary;
use aether_data_contracts::repository::{
@@ -195,29 +194,51 @@ async fn resolve_admin_usage_attempt_flags_by_usage_id(
.collect())
}
async fn resolve_admin_usage_image_progress_by_request_id(
#[derive(Default)]
struct AdminUsageActiveCandidateState {
image_progress_by_request_id: BTreeMap<String, serde_json::Value>,
state_overrides_by_request_id: BTreeMap<String, serde_json::Value>,
}
async fn resolve_admin_usage_active_candidate_state(
state: &AdminAppState<'_>,
items: &[StoredRequestUsageAudit],
) -> Result<BTreeMap<String, serde_json::Value>, GatewayError> {
) -> Result<AdminUsageActiveCandidateState, GatewayError> {
if !state.has_request_candidate_data_reader() || items.is_empty() {
return Ok(BTreeMap::new());
return Ok(AdminUsageActiveCandidateState::default());
}
let request_ids = items
.iter()
.map(|item| item.request_id.clone())
.collect::<BTreeSet<_>>();
let mut progress_by_request_id = BTreeMap::new();
let active_usage_by_request_id = items
.iter()
.filter(|item| matches!(item.status.as_str(), "pending" | "streaming"))
.map(|item| (item.request_id.clone(), item))
.collect::<BTreeMap<_, _>>();
let mut candidate_state = AdminUsageActiveCandidateState::default();
for request_id in request_ids {
let candidates = state
.app()
.read_request_candidates_by_request_id(&request_id)
.await?;
if let Some(progress) = latest_admin_usage_image_progress(&candidates) {
progress_by_request_id.insert(request_id, progress);
candidate_state
.image_progress_by_request_id
.insert(request_id.clone(), progress);
}
if active_usage_by_request_id.contains_key(&request_id) {
if let Some(override_payload) =
admin_usage_terminal_candidate_state_override(&candidates)
{
candidate_state
.state_overrides_by_request_id
.insert(request_id, override_payload);
}
}
}
Ok(progress_by_request_id)
Ok(candidate_state)
}
fn latest_admin_usage_image_progress(
@@ -246,6 +267,82 @@ fn latest_admin_usage_image_progress(
.map(|(_, _, _, progress)| progress)
}
fn admin_usage_current_candidate(
candidates: &[StoredRequestCandidate],
) -> Option<&StoredRequestCandidate> {
candidates
.iter()
.filter(|candidate| {
!matches!(
candidate.status,
RequestCandidateStatus::Available
| RequestCandidateStatus::Unused
| RequestCandidateStatus::Skipped
)
})
.max_by_key(|candidate| {
(
candidate.candidate_index,
candidate.retry_index,
candidate
.started_at_unix_ms
.or(candidate.finished_at_unix_ms)
.unwrap_or(candidate.created_at_unix_ms),
)
})
}
fn admin_usage_unix_millis_to_rfc3339(unix_ms: u64) -> Option<String> {
let secs = i64::try_from(unix_ms / 1_000).ok()?;
let nanos = u32::try_from(unix_ms % 1_000)
.ok()?
.saturating_mul(1_000_000);
chrono::DateTime::<chrono::Utc>::from_timestamp(secs, nanos)
.map(|timestamp| timestamp.to_rfc3339())
}
fn admin_usage_terminal_candidate_state_override(
candidates: &[StoredRequestCandidate],
) -> Option<serde_json::Value> {
let candidate = admin_usage_current_candidate(candidates)?;
let status = match candidate.status {
RequestCandidateStatus::Success => "completed",
RequestCandidateStatus::Failed => "failed",
RequestCandidateStatus::Cancelled => "cancelled",
_ => return None,
};
let latency_ms = candidate.latency_ms.or_else(|| {
Some(
candidate
.finished_at_unix_ms?
.saturating_sub(candidate.started_at_unix_ms?),
)
});
let mut payload = json!({ "status": status });
if let Some(latency_ms) = latency_ms {
payload["response_time_ms"] = json!(latency_ms);
if let Some(response_time_updated_at) = candidate
.finished_at_unix_ms
.or_else(|| {
candidate
.started_at_unix_ms
.map(|started_at| started_at.saturating_add(latency_ms))
})
.and_then(admin_usage_unix_millis_to_rfc3339)
{
payload["response_time_updated_at"] = json!(response_time_updated_at);
}
}
if let Some(status_code) = candidate.status_code {
payload["status_code"] = json!(status_code);
}
if let Some(error_message) = candidate.error_message.as_ref() {
payload["error_message"] = json!(error_message);
}
Some(payload)
}
fn admin_usage_matches_attempt_status(
item: &StoredRequestUsageAudit,
status: &str,
@@ -264,19 +361,6 @@ fn admin_usage_matches_attempt_status(
}
}
fn admin_usage_matches_client_family(
item: &StoredRequestUsageAudit,
client_family: Option<&str>,
) -> bool {
let Some(client_family) = client_family
.map(str::trim)
.filter(|value| !value.is_empty())
else {
return true;
};
admin_usage_client_family(item).is_some_and(|value| value.eq_ignore_ascii_case(client_family))
}
fn admin_usage_bool_query_param(query: Option<&str>, name: &str) -> bool {
query_param_value(query, name)
.as_deref()
@@ -290,15 +374,24 @@ fn admin_usage_bool_query_param(query: Option<&str>, name: &str) -> bool {
.unwrap_or(false)
}
fn admin_usage_is_unknown_label(value: &str) -> bool {
matches!(
value.trim().to_ascii_lowercase().as_str(),
"unknown" | "unknow"
)
fn admin_usage_include_total_query_param(query: Option<&str>) -> bool {
query_param_value(query, "include_total")
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(|value| {
!(value == "0"
|| value.eq_ignore_ascii_case("false")
|| value.eq_ignore_ascii_case("no")
|| value.eq_ignore_ascii_case("off"))
})
.unwrap_or(true)
}
fn admin_usage_has_unknown_model_or_provider(item: &StoredRequestUsageAudit) -> bool {
admin_usage_is_unknown_label(&item.model) || admin_usage_is_unknown_label(&item.provider_name)
fn admin_usage_fast_page_total(offset: usize, limit: usize, record_count: usize) -> usize {
offset
.saturating_add(record_count)
.saturating_add(usize::from(limit > 0 && record_count == limit))
}
#[allow(clippy::too_many_arguments)]
@@ -314,6 +407,7 @@ fn build_admin_usage_records_response_with_attempt_flags(
total: usize,
limit: usize,
offset: usize,
total_is_estimated: bool,
) -> Response<Body> {
let records: Vec<_> = items
.iter()
@@ -343,6 +437,7 @@ fn build_admin_usage_records_response_with_attempt_flags(
"total": total,
"limit": limit,
"offset": offset,
"total_is_estimated": total_is_estimated,
}))
.into_response()
}
@@ -485,6 +580,8 @@ fn build_admin_usage_keyword_search_query(
provider_name: base_query.provider_name.clone(),
model: base_query.model.clone(),
api_format: base_query.api_format.clone(),
client_family: base_query.client_family.clone(),
exclude_unknown_model_or_provider: base_query.exclude_unknown_model_or_provider,
statuses: base_query.statuses.clone(),
is_stream: base_query.is_stream,
error_only: base_query.error_only,
@@ -582,6 +679,7 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
state.has_auth_api_key_data_reader(),
&BTreeMap::new(),
&BTreeMap::new(),
&BTreeMap::new(),
)));
};
state
@@ -605,15 +703,16 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
};
let api_key_names = admin_usage_api_key_names(state, &items).await?;
let provider_key_names = admin_usage_provider_key_names(state, &items).await?;
let image_progress_by_request_id =
resolve_admin_usage_image_progress_by_request_id(state, &items).await?;
let active_candidate_state =
resolve_admin_usage_active_candidate_state(state, &items).await?;
return Ok(Some(build_admin_usage_active_requests_response(
&items,
&api_key_names,
state.has_auth_api_key_data_reader(),
&provider_key_names,
&image_progress_by_request_id,
&active_candidate_state.image_progress_by_request_id,
&active_candidate_state.state_overrides_by_request_id,
)));
}
Some("records")
@@ -641,6 +740,8 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
let client_family_filter = query_param_value(query, "client_family");
let hide_unknown_records = admin_usage_bool_query_param(query, "hide_unknown")
|| admin_usage_bool_query_param(query, "hide_unknown_records");
let include_total = admin_usage_include_total_query_param(query);
let total_only = admin_usage_bool_query_param(query, "total_only");
let limit = match admin_usage_parse_limit(query) {
Ok(value) => value,
Err(detail) => return Ok(Some(admin_usage_bad_request_response(detail))),
@@ -664,13 +765,6 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
offset,
)));
};
let base_query = build_admin_usage_records_query(
created_from_unix_secs,
created_until_unix_secs,
query,
None,
None,
);
let active_search = search.as_deref().filter(|value| !value.trim().is_empty());
let active_username_filter = username_filter
.as_deref()
@@ -678,10 +772,16 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
let active_client_family_filter = client_family_filter
.as_deref()
.filter(|value| !value.trim().is_empty());
let (usage, total) = if hide_unknown_records
|| attempt_status_filter.is_some()
|| active_client_family_filter.is_some()
{
let mut base_query = build_admin_usage_records_query(
created_from_unix_secs,
created_until_unix_secs,
query,
None,
None,
);
base_query.client_family = active_client_family_filter.map(str::to_owned);
base_query.exclude_unknown_model_or_provider = hide_unknown_records;
let (usage, total, total_is_estimated) = if attempt_status_filter.is_some() {
let mut usage = state.list_usage_audits(&base_query).await?;
let user_ids: Vec<String> = usage
.iter()
@@ -718,18 +818,20 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
&attempt_flags_by_usage_id,
request_candidate_reader_available,
)
}) && admin_usage_matches_client_family(item, active_client_family_filter)
&& (!hide_unknown_records
|| !admin_usage_has_unknown_model_or_provider(item))
})
});
sort_usage_newest_first(&mut usage);
let total = usage.len();
let records = usage
.into_iter()
.skip(offset)
.take(limit)
.collect::<Vec<_>>();
(records, total)
let records = if total_only {
Vec::new()
} else {
usage
.into_iter()
.skip(offset)
.take(limit)
.collect::<Vec<_>>()
};
(records, total, false)
} else if active_search.is_some() || active_username_filter.is_some() {
let keywords = active_search
.map(parse_admin_usage_search_keywords)
@@ -749,30 +851,53 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
None,
None,
);
let total = usize::try_from(
state
.count_usage_audits_by_keyword_search(&keyword_query)
.await?,
)
.unwrap_or(usize::MAX);
let paged_query = UsageAuditKeywordSearchQuery {
limit: Some(limit),
offset: Some(offset),
..keyword_query
};
(
state
.list_usage_audits_by_keyword_search(&paged_query)
.await?,
total,
)
} else {
let total = usize::try_from(state.count_usage_audits(&base_query).await?)
if total_only {
let total = usize::try_from(
state
.count_usage_audits_by_keyword_search(&keyword_query)
.await?,
)
.unwrap_or(usize::MAX);
let mut paged_query = base_query.clone();
paged_query.limit = Some(limit);
paged_query.offset = Some(offset);
(state.list_usage_audits(&paged_query).await?, total)
(Vec::new(), total, false)
} else {
let paged_query = UsageAuditKeywordSearchQuery {
limit: Some(limit),
offset: Some(offset),
..keyword_query.clone()
};
let records = state
.list_usage_audits_by_keyword_search(&paged_query)
.await?;
let total = if include_total {
usize::try_from(
state
.count_usage_audits_by_keyword_search(&keyword_query)
.await?,
)
.unwrap_or(usize::MAX)
} else {
admin_usage_fast_page_total(offset, limit, records.len())
};
(records, total, !include_total)
}
} else {
if total_only {
let total = usize::try_from(state.count_usage_audits(&base_query).await?)
.unwrap_or(usize::MAX);
(Vec::new(), total, false)
} else {
let mut paged_query = base_query.clone();
paged_query.limit = Some(limit);
paged_query.offset = Some(offset);
let records = state.list_usage_audits(&paged_query).await?;
let total = if include_total {
usize::try_from(state.count_usage_audits(&base_query).await?)
.unwrap_or(usize::MAX)
} else {
admin_usage_fast_page_total(offset, limit, records.len())
};
(records, total, !include_total)
}
};
let user_ids: Vec<String> = usage
@@ -800,6 +925,7 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
total,
limit,
offset,
total_is_estimated,
)));
}
_ => {}
@@ -807,3 +933,88 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
Ok(None)
}
#[cfg(test)]
mod tests {
use aether_data_contracts::repository::candidates::{
RequestCandidateStatus, StoredRequestCandidate,
};
use super::admin_usage_terminal_candidate_state_override;
fn sample_candidate(
candidate_index: i32,
status: RequestCandidateStatus,
status_code: Option<i32>,
latency_ms: Option<i32>,
error_message: Option<&str>,
) -> StoredRequestCandidate {
StoredRequestCandidate::new(
format!("candidate-{candidate_index}"),
"req-1".to_string(),
Some("user-1".to_string()),
Some("api-key-1".to_string()),
Some("alice".to_string()),
Some("default".to_string()),
candidate_index,
0,
Some("provider-1".to_string()),
Some("endpoint-1".to_string()),
Some("provider-key-1".to_string()),
status,
None,
false,
status_code,
None,
error_message.map(str::to_string),
latency_ms,
None,
None,
None,
1_000,
Some(1_000),
Some(10_210),
)
.expect("candidate should build")
}
#[test]
fn admin_usage_active_override_uses_current_terminal_candidate_latency() {
let candidate = sample_candidate(
0,
RequestCandidateStatus::Success,
Some(200),
Some(9_210),
None,
);
let payload =
admin_usage_terminal_candidate_state_override(&[candidate]).expect("override");
assert_eq!(payload["status"], "completed");
assert_eq!(payload["response_time_ms"], 9_210);
assert_eq!(
payload["response_time_updated_at"],
"1970-01-01T00:00:10.210+00:00"
);
}
#[test]
fn admin_usage_active_override_ignores_terminal_candidate_when_newer_attempt_is_live() {
let failed = sample_candidate(
0,
RequestCandidateStatus::Failed,
Some(503),
Some(1_000),
Some("first attempt failed"),
);
let mut streaming =
sample_candidate(1, RequestCandidateStatus::Streaming, None, None, None);
streaming.started_at_unix_ms = Some(10_500);
streaming.finished_at_unix_ms = None;
let payload = admin_usage_terminal_candidate_state_override(&[failed, streaming]);
assert!(payload.is_none());
}
}
@@ -4,11 +4,14 @@ use aether_ai_serving::UPSTREAM_IS_STREAM_KEY;
use aether_billing::{
normalize_input_tokens_for_billing, normalize_total_input_context_for_cache_hit_rate,
};
use aether_data_contracts::repository::usage::{
StoredRequestUsageAudit, StoredUsageBreakdownSummaryRow, StoredUsageDailySummary,
UsageAuditKeywordSearchQuery, UsageAuditListQuery, UsageBreakdownGroupBy,
UsageBreakdownSummaryQuery, UsageCacheAffinityIntervalGroupBy, UsageCacheAffinityIntervalQuery,
UsageDashboardSummaryQuery,
use aether_data_contracts::repository::{
candidates::{RequestCandidateStatus, StoredRequestCandidate},
usage::{
StoredRequestUsageAudit, StoredUsageBreakdownSummaryRow, StoredUsageDailySummary,
UsageAuditKeywordSearchQuery, UsageAuditListQuery, UsageBreakdownGroupBy,
UsageBreakdownSummaryQuery, UsageCacheAffinityIntervalGroupBy,
UsageCacheAffinityIntervalQuery, UsageDashboardSummaryQuery,
},
};
use axum::{
body::Body,
@@ -17,7 +20,7 @@ use axum::{
Json,
};
use chrono::Utc;
use serde_json::json;
use serde_json::{json, Value};
use crate::GatewayError;
@@ -531,6 +534,8 @@ fn build_users_me_usage_active_payload(item: &StoredRequestUsageAudit) -> serde_
"rate_multiplier": item.settlement_rate_multiplier(),
"response_time_ms": item.response_time_ms,
"first_byte_time_ms": item.first_byte_time_ms,
"updated_at": unix_secs_to_rfc3339(item.updated_at_unix_secs),
"response_time_updated_at": users_me_usage_response_time_updated_at(item),
"status_code": item.status_code,
"error_message": item.error_message,
"api_format": item.api_format,
@@ -573,6 +578,118 @@ fn build_users_me_usage_active_payload(item: &StoredRequestUsageAudit) -> serde_
payload
}
fn users_me_usage_response_time_updated_at(item: &StoredRequestUsageAudit) -> Option<String> {
item.response_time_ms?;
if matches!(item.status.as_str(), "pending" | "streaming")
&& item.updated_at_unix_secs <= item.created_at_unix_ms
{
return None;
}
unix_secs_to_rfc3339(item.updated_at_unix_secs)
}
fn unix_millis_to_rfc3339(unix_ms: u64) -> Option<String> {
let secs = i64::try_from(unix_ms / 1_000).ok()?;
let nanos = u32::try_from(unix_ms % 1_000)
.ok()?
.saturating_mul(1_000_000);
chrono::DateTime::<Utc>::from_timestamp(secs, nanos).map(|timestamp| timestamp.to_rfc3339())
}
fn users_me_usage_current_candidate(
candidates: &[StoredRequestCandidate],
) -> Option<&StoredRequestCandidate> {
candidates
.iter()
.filter(|candidate| {
!matches!(
candidate.status,
RequestCandidateStatus::Available
| RequestCandidateStatus::Unused
| RequestCandidateStatus::Skipped
)
})
.max_by_key(|candidate| {
(
candidate.candidate_index,
candidate.retry_index,
candidate
.started_at_unix_ms
.or(candidate.finished_at_unix_ms)
.unwrap_or(candidate.created_at_unix_ms),
)
})
}
fn users_me_usage_terminal_candidate_state_override(
candidates: &[StoredRequestCandidate],
) -> Option<Value> {
let candidate = users_me_usage_current_candidate(candidates)?;
let status = match candidate.status {
RequestCandidateStatus::Success => "completed",
RequestCandidateStatus::Failed => "failed",
RequestCandidateStatus::Cancelled => "cancelled",
_ => return None,
};
let latency_ms = candidate.latency_ms.or_else(|| {
Some(
candidate
.finished_at_unix_ms?
.saturating_sub(candidate.started_at_unix_ms?),
)
});
let mut payload = json!({ "status": status });
if let Some(latency_ms) = latency_ms {
payload["response_time_ms"] = json!(latency_ms);
if let Some(response_time_updated_at) = candidate
.finished_at_unix_ms
.or_else(|| {
candidate
.started_at_unix_ms
.map(|started_at| started_at.saturating_add(latency_ms))
})
.and_then(unix_millis_to_rfc3339)
{
payload["response_time_updated_at"] = json!(response_time_updated_at);
}
}
if let Some(status_code) = candidate.status_code {
payload["status_code"] = json!(status_code);
}
if let Some(error_message) = candidate.error_message.as_ref() {
payload["error_message"] = json!(error_message);
}
Some(payload)
}
async fn resolve_users_me_usage_active_state_overrides_by_request_id(
state: &AppState,
items: &[StoredRequestUsageAudit],
) -> Result<BTreeMap<String, Value>, GatewayError> {
if !state.has_request_candidate_data_reader() || items.is_empty() {
return Ok(BTreeMap::new());
}
let active_request_ids = items
.iter()
.filter(|item| matches!(item.status.as_str(), "pending" | "streaming"))
.map(|item| item.request_id.clone())
.collect::<BTreeSet<_>>();
let mut overrides = BTreeMap::new();
for request_id in active_request_ids {
let candidates = state
.read_request_candidates_by_request_id(&request_id)
.await?;
if let Some(override_payload) =
users_me_usage_terminal_candidate_state_override(&candidates)
{
overrides.insert(request_id, override_payload);
}
}
Ok(overrides)
}
fn users_me_usage_is_failed(item: &StoredRequestUsageAudit) -> bool {
let has_failure_signal = item.status_code.is_some_and(|value| value >= 400)
|| item
@@ -925,6 +1042,8 @@ pub(super) async fn handle_users_me_usage_get(
provider_name: None,
model: None,
api_format: None,
client_family: None,
exclude_unknown_model_or_provider: false,
statuses: None,
is_stream: None,
error_only: false,
@@ -978,6 +1097,8 @@ pub(super) async fn handle_users_me_usage_get(
provider_name: None,
model: None,
api_format: None,
client_family: None,
exclude_unknown_model_or_provider: false,
statuses: None,
is_stream: None,
error_only: false,
@@ -1004,6 +1125,8 @@ pub(super) async fn handle_users_me_usage_get(
provider_name: None,
model: None,
api_format: None,
client_family: None,
exclude_unknown_model_or_provider: false,
statuses: None,
is_stream: None,
error_only: false,
@@ -1140,6 +1263,8 @@ pub(super) async fn handle_users_me_usage_active_get(
provider_name: None,
model: None,
api_format: None,
client_family: None,
exclude_unknown_model_or_provider: false,
statuses: Some(vec!["pending".to_string(), "streaming".to_string()]),
is_stream: None,
error_only: false,
@@ -1168,11 +1293,35 @@ pub(super) async fn handle_users_me_usage_active_get(
.filter(|item| !users_me_usage_is_failed(item))
.collect::<Vec<_>>()
};
let active_state_overrides =
match resolve_users_me_usage_active_state_overrides_by_request_id(state, &items).await {
Ok(value) => value,
Err(err) => {
return build_auth_error_response(
http::StatusCode::INTERNAL_SERVER_ERROR,
format!("user active usage candidate lookup failed: {err:?}"),
false,
);
}
};
Json(json!({
"requests": items
.iter()
.map(build_users_me_usage_active_payload)
.map(|item| {
let mut payload = build_users_me_usage_active_payload(item);
if let (Some(payload), Some(overrides)) = (
payload.as_object_mut(),
active_state_overrides
.get(&item.request_id)
.and_then(Value::as_object),
) {
for (key, value) in overrides {
payload.insert(key.clone(), value.clone());
}
}
payload
})
.collect::<Vec<_>>(),
}))
.into_response()
@@ -1364,13 +1513,16 @@ async fn build_usage_heatmap_summaries(
mod tests {
use std::collections::BTreeMap;
use aether_data_contracts::repository::usage::StoredRequestUsageAudit;
use aether_data_contracts::repository::{
candidates::{RequestCandidateStatus, StoredRequestCandidate},
usage::StoredRequestUsageAudit,
};
use serde_json::json;
use super::{
build_users_me_usage_active_payload, build_users_me_usage_record_payload,
users_me_usage_client_is_stream, users_me_usage_is_failed,
users_me_usage_upstream_is_stream,
users_me_usage_terminal_candidate_state_override, users_me_usage_upstream_is_stream,
};
fn sample_usage(status: &str) -> StoredRequestUsageAudit {
@@ -1415,6 +1567,41 @@ mod tests {
.expect("usage should build")
}
fn sample_candidate(
status: RequestCandidateStatus,
status_code: Option<i32>,
latency_ms: Option<i32>,
error_message: Option<&str>,
) -> StoredRequestCandidate {
StoredRequestCandidate::new(
"candidate-1".to_string(),
"req-1".to_string(),
Some("user-1".to_string()),
Some("api-key-1".to_string()),
Some("alice".to_string()),
Some("default".to_string()),
0,
0,
Some("provider-1".to_string()),
Some("endpoint-1".to_string()),
Some("provider-key-1".to_string()),
status,
None,
false,
status_code,
None,
error_message.map(str::to_string),
latency_ms,
None,
None,
None,
1_000,
Some(1_000),
Some(10_210),
)
.expect("candidate should build")
}
#[test]
fn user_usage_record_payload_rehydrates_cache_creation_total_from_classified_fields() {
let item = StoredRequestUsageAudit {
@@ -1447,6 +1634,45 @@ mod tests {
assert_eq!(payload["cache_creation_ephemeral_1h_input_tokens"], 6);
}
#[test]
fn user_usage_active_override_uses_terminal_candidate_latency() {
let candidate = sample_candidate(
RequestCandidateStatus::Success,
Some(200),
Some(9_210),
None,
);
let payload =
users_me_usage_terminal_candidate_state_override(&[candidate]).expect("override");
assert_eq!(payload["status"], "completed");
assert_eq!(payload["response_time_ms"], 9_210);
assert_eq!(payload["status_code"], 200);
assert_eq!(
payload["response_time_updated_at"],
"1970-01-01T00:00:10.210+00:00"
);
}
#[test]
fn user_usage_active_override_ignores_terminal_candidate_when_newer_attempt_is_live() {
let failed = sample_candidate(
RequestCandidateStatus::Failed,
Some(503),
Some(1_000),
Some("first attempt failed"),
);
let mut streaming = sample_candidate(RequestCandidateStatus::Streaming, None, None, None);
streaming.candidate_index = 1;
streaming.started_at_unix_ms = Some(10_500);
streaming.finished_at_unix_ms = None;
let payload = users_me_usage_terminal_candidate_state_override(&[failed, streaming]);
assert!(payload.is_none());
}
#[test]
fn user_usage_payload_keeps_claude_effective_input_when_cache_read_is_large() {
let item = StoredRequestUsageAudit {
@@ -1172,7 +1172,7 @@ async fn gateway_handles_admin_usage_active_ids_for_terminal_updates() {
#[tokio::test]
async fn gateway_handles_admin_usage_records_locally_with_trusted_admin_principal() {
let (upstream_url, upstream_hits, upstream_handle) =
let (_upstream_url, upstream_hits, upstream_handle) =
start_usage_upstream("/api/admin/usage/records").await;
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![
@@ -1336,6 +1336,100 @@ async fn gateway_filters_admin_usage_records_with_unknown_model_or_provider() {
upstream_handle.abort();
}
#[tokio::test]
async fn gateway_supports_fast_admin_usage_record_totals() {
let (upstream_url, upstream_hits, upstream_handle) =
start_usage_upstream("/api/admin/usage/records").await;
let usage_repository = Arc::new(InMemoryUsageReadRepository::seed(vec![
sample_usage_row(
"usage-a",
"req-a",
Some("user-1"),
Some("key-1"),
Some("primary"),
"OpenAI",
"gpt-5",
"completed",
120,
30,
0.3,
0.36,
DAY_2_UNIX_SECS,
),
sample_usage_row(
"usage-b",
"req-b",
Some("user-1"),
Some("key-1"),
Some("primary"),
"OpenAI",
"gpt-5-mini",
"completed",
80,
20,
0.2,
0.24,
DAY_2_UNIX_SECS - 1,
),
sample_usage_row(
"usage-c",
"req-c",
Some("user-1"),
Some("key-1"),
Some("primary"),
"Anthropic",
"claude-sonnet",
"completed",
60,
10,
0.1,
0.12,
DAY_1_UNIX_SECS,
),
]));
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,
)),
);
let (gateway_url, gateway_handle) = start_server(gateway).await;
let fast_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&include_total=false&limit=2&offset=0"
)))
.send()
.await
.expect("request should succeed");
assert_eq!(fast_response.status(), StatusCode::OK);
let fast_payload: serde_json::Value =
fast_response.json().await.expect("json body should parse");
assert_eq!(fast_payload["records"].as_array().unwrap().len(), 2);
assert_eq!(fast_payload["total"], 3);
assert_eq!(fast_payload["total_is_estimated"], true);
let total_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&total_only=true&limit=2&offset=0"
)))
.send()
.await
.expect("request should succeed");
assert_eq!(total_response.status(), StatusCode::OK);
let total_payload: serde_json::Value =
total_response.json().await.expect("json body should parse");
assert_eq!(total_payload["records"].as_array().unwrap().len(), 0);
assert_eq!(total_payload["total"], 3);
assert_eq!(total_payload["total_is_estimated"], false);
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_provider_key_name_fallback_from_request_metadata()
{