mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-10 05:00:19 +08:00
fix(usage): calibrate active elapsed clock efficiently
This commit is contained in:
@@ -9,9 +9,9 @@ use aether_admin::observability::usage::{
|
||||
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_record_json, attach_usage_server_now_header,
|
||||
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::{
|
||||
@@ -314,13 +314,15 @@ fn build_admin_usage_records_response_with_attempt_flags(
|
||||
})
|
||||
.collect();
|
||||
|
||||
Json(json!({
|
||||
"records": records,
|
||||
"total": total,
|
||||
"limit": limit,
|
||||
"offset": offset,
|
||||
}))
|
||||
.into_response()
|
||||
attach_usage_server_now_header(
|
||||
Json(json!({
|
||||
"records": records,
|
||||
"total": total,
|
||||
"limit": limit,
|
||||
"offset": offset,
|
||||
}))
|
||||
.into_response(),
|
||||
)
|
||||
}
|
||||
|
||||
fn build_admin_usage_records_query(
|
||||
|
||||
@@ -4,6 +4,7 @@ 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_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
|
||||
use aether_data_contracts::repository::usage::{
|
||||
StoredRequestUsageAudit, StoredUsageBreakdownSummaryRow, StoredUsageDailySummary,
|
||||
UsageAuditKeywordSearchQuery, UsageAuditListQuery, UsageBreakdownGroupBy,
|
||||
@@ -29,6 +30,21 @@ use super::{
|
||||
|
||||
const USERS_ME_USAGE_DATA_UNAVAILABLE_DETAIL: &str = "用户用量数据暂不可用";
|
||||
|
||||
fn users_me_usage_server_now_unix_ms() -> u64 {
|
||||
u64::try_from(Utc::now().timestamp_millis()).unwrap_or_default()
|
||||
}
|
||||
|
||||
fn attach_users_me_usage_server_now_header(mut response: Response<Body>) -> Response<Body> {
|
||||
if let Ok(value) = http::HeaderValue::from_str(&users_me_usage_server_now_unix_ms().to_string())
|
||||
{
|
||||
response.headers_mut().insert(
|
||||
http::HeaderName::from_static(USAGE_SERVER_NOW_UNIX_MS_HEADER),
|
||||
value,
|
||||
);
|
||||
}
|
||||
response
|
||||
}
|
||||
|
||||
fn build_users_me_usage_reader_unavailable_response() -> Response<Body> {
|
||||
build_auth_error_response(
|
||||
http::StatusCode::SERVICE_UNAVAILABLE,
|
||||
@@ -1094,7 +1110,7 @@ pub(super) async fn handle_users_me_usage_get(
|
||||
&summary_by_provider
|
||||
));
|
||||
}
|
||||
Json(payload).into_response()
|
||||
attach_users_me_usage_server_now_header(Json(payload).into_response())
|
||||
}
|
||||
|
||||
pub(super) async fn handle_users_me_usage_active_get(
|
||||
@@ -1169,13 +1185,15 @@ pub(super) async fn handle_users_me_usage_active_get(
|
||||
.collect::<Vec<_>>()
|
||||
};
|
||||
|
||||
Json(json!({
|
||||
"requests": items
|
||||
.iter()
|
||||
.map(build_users_me_usage_active_payload)
|
||||
.collect::<Vec<_>>(),
|
||||
}))
|
||||
.into_response()
|
||||
attach_users_me_usage_server_now_header(
|
||||
Json(json!({
|
||||
"requests": items
|
||||
.iter()
|
||||
.map(build_users_me_usage_active_payload)
|
||||
.collect::<Vec<_>>(),
|
||||
}))
|
||||
.into_response(),
|
||||
)
|
||||
}
|
||||
|
||||
pub(super) async fn handle_users_me_usage_interval_timeline_get(
|
||||
@@ -1364,12 +1382,20 @@ async fn build_usage_heatmap_summaries(
|
||||
mod tests {
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
|
||||
use aether_data_contracts::repository::usage::StoredRequestUsageAudit;
|
||||
use axum::{
|
||||
body::Body,
|
||||
response::{IntoResponse, Response},
|
||||
Json,
|
||||
};
|
||||
use chrono::Utc;
|
||||
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,
|
||||
attach_users_me_usage_server_now_header, 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_server_now_unix_ms,
|
||||
users_me_usage_upstream_is_stream,
|
||||
};
|
||||
|
||||
@@ -1415,6 +1441,39 @@ mod tests {
|
||||
.expect("usage should build")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn users_me_usage_server_now_unix_ms_uses_epoch_millis() {
|
||||
let before = u64::try_from(Utc::now().timestamp_millis()).unwrap_or_default();
|
||||
let value = users_me_usage_server_now_unix_ms();
|
||||
let after = u64::try_from(Utc::now().timestamp_millis()).unwrap_or_default();
|
||||
|
||||
assert!(value >= before);
|
||||
assert!(value <= after);
|
||||
assert!(value > 1_000_000_000_000);
|
||||
}
|
||||
|
||||
fn assert_users_me_usage_server_now_header(response: &Response<Body>) {
|
||||
let value = response
|
||||
.headers()
|
||||
.get(USAGE_SERVER_NOW_UNIX_MS_HEADER)
|
||||
.expect("user usage response should include server now header")
|
||||
.to_str()
|
||||
.expect("server now header should be valid ASCII")
|
||||
.parse::<u64>()
|
||||
.expect("server now header should be epoch millis");
|
||||
|
||||
assert!(value > 1_000_000_000_000);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn users_me_usage_server_now_header_is_added_to_response() {
|
||||
let response = attach_users_me_usage_server_now_header(
|
||||
Json(json!({ "requests": [] })).into_response(),
|
||||
);
|
||||
|
||||
assert_users_me_usage_server_now_header(&response);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn user_usage_record_payload_rehydrates_cache_creation_total_from_classified_fields() {
|
||||
let item = StoredRequestUsageAudit {
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
|
||||
use axum::body::Body;
|
||||
use axum::extract::{Request, State};
|
||||
use axum::http::{self, HeaderValue, Response};
|
||||
@@ -6,6 +7,8 @@ use axum::middleware::Next;
|
||||
use crate::headers::header_value_str;
|
||||
use crate::state::{AppState, FrontdoorCorsConfig};
|
||||
|
||||
const FRONTDOOR_CREDENTIALS_EXPOSE_HEADERS: &str = "*, x-aether-server-now-unix-ms";
|
||||
|
||||
fn append_vary(headers: &mut http::HeaderMap, value: &'static str) {
|
||||
headers.append(http::header::VARY, HeaderValue::from_static(value));
|
||||
}
|
||||
@@ -31,7 +34,11 @@ fn apply_frontdoor_cors_headers(
|
||||
);
|
||||
headers.insert(
|
||||
http::header::ACCESS_CONTROL_EXPOSE_HEADERS,
|
||||
HeaderValue::from_static("*"),
|
||||
HeaderValue::from_static(if cors.allow_credentials() {
|
||||
FRONTDOOR_CREDENTIALS_EXPOSE_HEADERS
|
||||
} else {
|
||||
"*"
|
||||
}),
|
||||
);
|
||||
if let Some(value) = requested_headers {
|
||||
if let Ok(value) = HeaderValue::from_str(value) {
|
||||
@@ -109,3 +116,52 @@ pub(crate) async fn frontdoor_cors_middleware(
|
||||
);
|
||||
response
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn assert_exposes_header(value: &HeaderValue, expected: &str) {
|
||||
let exposed_headers = value
|
||||
.to_str()
|
||||
.expect("expose headers should be valid ASCII");
|
||||
assert!(
|
||||
exposed_headers
|
||||
.split(',')
|
||||
.map(str::trim)
|
||||
.any(|header| header.eq_ignore_ascii_case(expected)),
|
||||
"{exposed_headers} should include {expected}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn frontdoor_cors_explicitly_exposes_usage_server_time_for_credentials() {
|
||||
let cors = FrontdoorCorsConfig::new(vec!["http://localhost:5173".to_string()], true)
|
||||
.expect("cors config should build");
|
||||
let mut headers = http::HeaderMap::new();
|
||||
|
||||
apply_frontdoor_cors_headers(&mut headers, &cors, "http://localhost:5173", None);
|
||||
|
||||
let expose_headers = headers
|
||||
.get(http::header::ACCESS_CONTROL_EXPOSE_HEADERS)
|
||||
.expect("expose headers should be set");
|
||||
assert_exposes_header(expose_headers, "*");
|
||||
assert_exposes_header(expose_headers, USAGE_SERVER_NOW_UNIX_MS_HEADER);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn frontdoor_cors_keeps_wildcard_expose_headers_without_credentials() {
|
||||
let cors = FrontdoorCorsConfig::new(vec!["http://localhost:5173".to_string()], false)
|
||||
.expect("cors config should build");
|
||||
let mut headers = http::HeaderMap::new();
|
||||
|
||||
apply_frontdoor_cors_headers(&mut headers, &cors, "http://localhost:5173", None);
|
||||
|
||||
assert_eq!(
|
||||
headers
|
||||
.get(http::header::ACCESS_CONTROL_EXPOSE_HEADERS)
|
||||
.expect("expose headers should be set"),
|
||||
"*"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,6 +33,7 @@ use crate::constants::{
|
||||
};
|
||||
use crate::control::resolve_public_request_context;
|
||||
use crate::data::GatewayDataState;
|
||||
use crate::tests::{assert_usage_server_now_header_between, unix_epoch_millis_for_tests};
|
||||
|
||||
const ADMIN_USAGE_DATA_UNAVAILABLE_DETAIL: &str = "Admin usage data unavailable";
|
||||
const DAY_1_UNIX_SECS: i64 = 1_711_000_000;
|
||||
@@ -1014,6 +1015,7 @@ async fn gateway_handles_admin_usage_active_locally_with_trusted_admin_principal
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let client_send_unix_ms = unix_epoch_millis_for_tests();
|
||||
let response =
|
||||
admin_request(reqwest::Client::new().get(format!(
|
||||
"{gateway_url}/api/admin/usage/active?start_date=2024-03-21&end_date=2024-03-22&tz_offset_minutes=0"
|
||||
@@ -1021,9 +1023,17 @@ async fn gateway_handles_admin_usage_active_locally_with_trusted_admin_principal
|
||||
.send()
|
||||
.await
|
||||
.expect("request should succeed");
|
||||
let client_receive_unix_ms = unix_epoch_millis_for_tests();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let response_headers = response.headers().clone();
|
||||
assert_usage_server_now_header_between(
|
||||
&response_headers,
|
||||
client_send_unix_ms,
|
||||
client_receive_unix_ms,
|
||||
);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert!(payload.get("server_now_unix_ms").is_none());
|
||||
assert_eq!(payload["requests"].as_array().expect("array").len(), 1);
|
||||
assert_eq!(payload["requests"][0]["id"], "usage-pending");
|
||||
assert_eq!(payload["requests"][0]["effective_input_tokens"], 5);
|
||||
@@ -1233,15 +1243,24 @@ async fn gateway_handles_admin_usage_records_locally_with_trusted_admin_principa
|
||||
);
|
||||
let (gateway_url, gateway_handle) = start_server(gateway).await;
|
||||
|
||||
let client_send_unix_ms = unix_epoch_millis_for_tests();
|
||||
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=failed&provider=Anthropic&limit=10&offset=0"
|
||||
)))
|
||||
.send()
|
||||
.await
|
||||
.expect("request should succeed");
|
||||
let client_receive_unix_ms = unix_epoch_millis_for_tests();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let response_headers = response.headers().clone();
|
||||
assert_usage_server_now_header_between(
|
||||
&response_headers,
|
||||
client_send_unix_ms,
|
||||
client_receive_unix_ms,
|
||||
);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert!(payload.get("server_now_unix_ms").is_none());
|
||||
assert_eq!(payload["total"], 1);
|
||||
assert_eq!(payload["records"][0]["id"], "usage-b");
|
||||
assert_eq!(payload["records"][0]["username"], "bob");
|
||||
|
||||
@@ -7,6 +7,7 @@ use crate::tests::{
|
||||
build_state_with_execution_runtime_override, json, start_server, AppState, Arc, Body,
|
||||
FrontdoorCorsConfig, Mutex, Request, Router, StatusCode, FRONTDOOR_MANIFEST_PATH, READYZ_PATH,
|
||||
};
|
||||
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
|
||||
use aether_crypto::DEVELOPMENT_ENCRYPTION_KEY;
|
||||
use aether_data::repository::auth::InMemoryAuthApiKeySnapshotRepository;
|
||||
use aether_data::repository::candidate_selection::InMemoryMinimalCandidateSelectionReadRepository;
|
||||
@@ -440,12 +441,19 @@ async fn gateway_adds_cors_headers_to_proxied_responses() {
|
||||
.expect("allow origin header"),
|
||||
"http://localhost:3000"
|
||||
);
|
||||
assert_eq!(
|
||||
response_headers
|
||||
.get("access-control-expose-headers")
|
||||
.expect("expose headers header"),
|
||||
"*"
|
||||
);
|
||||
let expose_headers = response_headers
|
||||
.get("access-control-expose-headers")
|
||||
.expect("expose headers header")
|
||||
.to_str()
|
||||
.expect("expose headers should be valid ASCII");
|
||||
assert!(expose_headers
|
||||
.split(',')
|
||||
.map(str::trim)
|
||||
.any(|header| header == "*"));
|
||||
assert!(expose_headers
|
||||
.split(',')
|
||||
.map(str::trim)
|
||||
.any(|header| header.eq_ignore_ascii_case(USAGE_SERVER_NOW_UNIX_MS_HEADER)));
|
||||
assert_eq!(
|
||||
*execution_runtime_hits.lock().expect("mutex should lock"),
|
||||
1
|
||||
|
||||
@@ -11,10 +11,10 @@ use super::{
|
||||
};
|
||||
use crate::data::GatewayDataState;
|
||||
use crate::tests::{
|
||||
any, build_router, build_router_with_state, json, start_server, to_bytes, AppState, Arc, Body,
|
||||
Json, Mutex, Request, Router, StatusCode, CONTROL_ROUTE_FAMILY_HEADER,
|
||||
CONTROL_ROUTE_KIND_HEADER, TRUSTED_ADMIN_SESSION_ID_HEADER, TRUSTED_ADMIN_USER_ID_HEADER,
|
||||
TRUSTED_ADMIN_USER_ROLE_HEADER,
|
||||
any, assert_usage_server_now_header_between, build_router, build_router_with_state, json,
|
||||
start_server, to_bytes, unix_epoch_millis_for_tests, AppState, Arc, Body, Json, Mutex, Request,
|
||||
Router, StatusCode, CONTROL_ROUTE_FAMILY_HEADER, CONTROL_ROUTE_KIND_HEADER,
|
||||
TRUSTED_ADMIN_SESSION_ID_HEADER, TRUSTED_ADMIN_USER_ID_HEADER, TRUSTED_ADMIN_USER_ROLE_HEADER,
|
||||
};
|
||||
use aether_crypto::{encrypt_python_fernet_plaintext, DEVELOPMENT_ENCRYPTION_KEY};
|
||||
use aether_data::repository::announcements::{AnnouncementListQuery, AnnouncementReadRepository};
|
||||
@@ -5345,6 +5345,7 @@ async fn gateway_handles_users_me_usage_locally_without_proxying_upstream() {
|
||||
})
|
||||
.await;
|
||||
|
||||
let client_send_unix_ms = unix_epoch_millis_for_tests();
|
||||
let response = reqwest::Client::new()
|
||||
.get(format!(
|
||||
"{gateway_url}/api/users/me/usage?limit=10&offset=0&search=renamed-key"
|
||||
@@ -5355,9 +5356,17 @@ async fn gateway_handles_users_me_usage_locally_without_proxying_upstream() {
|
||||
.send()
|
||||
.await
|
||||
.expect("request should succeed");
|
||||
let client_receive_unix_ms = unix_epoch_millis_for_tests();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let response_headers = response.headers().clone();
|
||||
assert_usage_server_now_header_between(
|
||||
&response_headers,
|
||||
client_send_unix_ms,
|
||||
client_receive_unix_ms,
|
||||
);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert!(payload.get("server_now_unix_ms").is_none());
|
||||
assert_eq!(payload["total_requests"], 2);
|
||||
assert_eq!(payload["total_input_tokens"], 240);
|
||||
assert_eq!(payload["pagination"]["total"], 3);
|
||||
@@ -5667,6 +5676,7 @@ async fn gateway_handles_users_me_usage_active_locally_without_proxying_upstream
|
||||
)
|
||||
.await;
|
||||
|
||||
let client_send_unix_ms = unix_epoch_millis_for_tests();
|
||||
let response = reqwest::Client::new()
|
||||
.get(format!("{gateway_url}/api/users/me/usage/active"))
|
||||
.header("authorization", format!("Bearer {access_token}"))
|
||||
@@ -5675,9 +5685,17 @@ async fn gateway_handles_users_me_usage_active_locally_without_proxying_upstream
|
||||
.send()
|
||||
.await
|
||||
.expect("request should succeed");
|
||||
let client_receive_unix_ms = unix_epoch_millis_for_tests();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let response_headers = response.headers().clone();
|
||||
assert_usage_server_now_header_between(
|
||||
&response_headers,
|
||||
client_send_unix_ms,
|
||||
client_receive_unix_ms,
|
||||
);
|
||||
let payload: serde_json::Value = response.json().await.expect("json body should parse");
|
||||
assert!(payload.get("server_now_unix_ms").is_none());
|
||||
let requests = payload["requests"].as_array().expect("requests array");
|
||||
assert_eq!(requests.len(), 2);
|
||||
assert_eq!(requests[0]["status"], "streaming");
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
pub(super) use std::convert::Infallible;
|
||||
pub(super) use std::sync::{Arc, Mutex};
|
||||
|
||||
use aether_contracts::USAGE_SERVER_NOW_UNIX_MS_HEADER;
|
||||
pub(super) use axum::body::{to_bytes, Body, Bytes};
|
||||
pub(super) use axum::response::Response;
|
||||
pub(super) use axum::routing::any;
|
||||
@@ -29,6 +32,38 @@ pub(super) use super::router::{attach_static_frontend, build_router, build_route
|
||||
pub(super) use super::state::{AppState, FrontdoorCorsConfig};
|
||||
pub(super) use super::usage::UsageRuntimeConfig;
|
||||
|
||||
const SERVER_NOW_HEADER_TEST_TOLERANCE_MS: u64 = 1_000;
|
||||
|
||||
pub(super) fn unix_epoch_millis_for_tests() -> u64 {
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.expect("current time should be after epoch")
|
||||
.as_millis()
|
||||
.try_into()
|
||||
.expect("current epoch millis should fit in u64")
|
||||
}
|
||||
|
||||
pub(super) fn assert_usage_server_now_header_between(
|
||||
headers: &reqwest::header::HeaderMap,
|
||||
lower_bound_unix_ms: u64,
|
||||
upper_bound_unix_ms: u64,
|
||||
) {
|
||||
let server_now_unix_ms = headers
|
||||
.get(USAGE_SERVER_NOW_UNIX_MS_HEADER)
|
||||
.expect("usage response should include server timing header")
|
||||
.to_str()
|
||||
.expect("server timing header should be valid ASCII")
|
||||
.parse::<u64>()
|
||||
.expect("server timing header should be epoch millis");
|
||||
|
||||
assert!(
|
||||
server_now_unix_ms >= lower_bound_unix_ms.saturating_sub(SERVER_NOW_HEADER_TEST_TOLERANCE_MS)
|
||||
&& server_now_unix_ms
|
||||
<= upper_bound_unix_ms.saturating_add(SERVER_NOW_HEADER_TEST_TOLERANCE_MS),
|
||||
"server timing header {server_now_unix_ms} should be near request window {lower_bound_unix_ms}..={upper_bound_unix_ms}"
|
||||
);
|
||||
}
|
||||
|
||||
pub(super) async fn start_server(app: Router) -> (String, tokio::task::JoinHandle<()>) {
|
||||
let listener = crate::test_support::bind_loopback_listener()
|
||||
.await
|
||||
|
||||
Reference in New Issue
Block a user