From cb7b9c9ecd56cb08fb437d54437b9fad0ed23877 Mon Sep 17 00:00:00 2001 From: elky Date: Mon, 5 Oct 2026 00:28:31 +0800 Subject: [PATCH] feat: unify user analytics and optimize overview aggregation Merge user accounts and usage reporting into one page with a combined ranking and account table, shared precise time ranges, and simpler range labels. Parse overview metadata once through a schema-only view migration and disable JIT locally for bucket rebuilds. Preserve automatic backfills. Add redacted OAuth refresh diagnostics, bucket failure context, and regression coverage. Resolve strict Clippy warnings. --- .../planner/standard/normalize/chat.rs | 32 +- .../planner/standard/normalize/responses.rs | 36 +- apps/aether-gateway/src/codex_profile.rs | 21 +- .../observability/stats/analytics_routes.rs | 40 +- .../observability/stats/leaderboard_routes.rs | 16 +- .../handlers/admin/observability/stats/mod.rs | 5 +- .../admin/observability/stats/range.rs | 154 ++++++ .../observability/usage/summary_routes.rs | 42 +- .../runtime/oauth_token_refresh.rs | 20 +- apps/aether-gateway/src/state/oauth.rs | 98 +++- ...000000_optimize_overview_fact_metadata.sql | 78 +++ .../postgres/src/usage/overview_buckets.rs | 22 +- crates/aether-data/runtime/schema/README.md | 10 +- .../postgres/190_overview_analytics.sql | 45 +- .../runtime/src/lifecycle/migrate/tests.rs | 2 + .../migrate/tests/overview_fact_metadata.rs | 180 +++++++ .../__tests__/admin-analytics-cache.spec.ts | 9 + frontend/src/api/admin.ts | 14 +- .../__tests__/UserUsageStats.scope.spec.ts | 122 ++++- .../features/overview/__tests__/users.spec.ts | 44 +- .../overview/components/OverviewToolbar.vue | 94 ++-- .../features/overview/users/UserReports.vue | 3 +- .../overview/users/UserUsageStats.vue | 280 +++++------ .../views/admin/AdminOperationsDashboard.vue | 1 + frontend/src/views/admin/CostAnalysis.vue | 1 + frontend/src/views/admin/UserStats.vue | 455 ++++++++++-------- frontend/src/views/shared/Usage.vue | 18 +- 27 files changed, 1255 insertions(+), 587 deletions(-) create mode 100644 crates/aether-data/adapters/postgres/migrations/20261004000000_optimize_overview_fact_metadata.sql create mode 100644 crates/aether-data/runtime/src/lifecycle/migrate/tests/overview_fact_metadata.rs diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/normalize/chat.rs b/apps/aether-gateway/src/ai_serving/planner/standard/normalize/chat.rs index 5c2fbfe31..ed215957d 100644 --- a/apps/aether-gateway/src/ai_serving/planner/standard/normalize/chat.rs +++ b/apps/aether-gateway/src/ai_serving/planner/standard/normalize/chat.rs @@ -112,6 +112,22 @@ pub(crate) fn build_cross_format_openai_chat_request_body( Some(provider_request_body) } +pub(crate) fn build_cross_format_openai_chat_upstream_url( + parts: &http::request::Parts, + transport: &GatewayProviderTransportSnapshot, + mapped_model: &str, + provider_api_format: &str, + upstream_is_stream: bool, +) -> Option { + crate::ai_serving::transport::build_cross_format_openai_chat_upstream_url( + transport, + mapped_model, + provider_api_format, + upstream_is_stream, + parts.uri.query(), + ) +} + #[cfg(test)] mod antigravity_schema_tests { use super::*; @@ -147,19 +163,3 @@ mod antigravity_schema_tests { } } } - -pub(crate) fn build_cross_format_openai_chat_upstream_url( - parts: &http::request::Parts, - transport: &GatewayProviderTransportSnapshot, - mapped_model: &str, - provider_api_format: &str, - upstream_is_stream: bool, -) -> Option { - crate::ai_serving::transport::build_cross_format_openai_chat_upstream_url( - transport, - mapped_model, - provider_api_format, - upstream_is_stream, - parts.uri.query(), - ) -} diff --git a/apps/aether-gateway/src/ai_serving/planner/standard/normalize/responses.rs b/apps/aether-gateway/src/ai_serving/planner/standard/normalize/responses.rs index 5a138ac83..4ae5ed7b8 100644 --- a/apps/aether-gateway/src/ai_serving/planner/standard/normalize/responses.rs +++ b/apps/aether-gateway/src/ai_serving/planner/standard/normalize/responses.rs @@ -275,6 +275,24 @@ pub(crate) fn build_local_openai_responses_upstream_url( ) } +pub(crate) fn build_cross_format_openai_responses_upstream_url( + parts: &http::request::Parts, + transport: &GatewayProviderTransportSnapshot, + mapped_model: &str, + client_api_format: &str, + provider_api_format: &str, + upstream_is_stream: bool, +) -> Option { + crate::ai_serving::transport::build_cross_format_openai_responses_upstream_url( + transport, + mapped_model, + client_api_format, + provider_api_format, + upstream_is_stream, + parts.uri.query(), + ) +} + #[cfg(test)] mod antigravity_schema_tests { use super::*; @@ -309,21 +327,3 @@ mod antigravity_schema_tests { } } } - -pub(crate) fn build_cross_format_openai_responses_upstream_url( - parts: &http::request::Parts, - transport: &GatewayProviderTransportSnapshot, - mapped_model: &str, - client_api_format: &str, - provider_api_format: &str, - upstream_is_stream: bool, -) -> Option { - crate::ai_serving::transport::build_cross_format_openai_responses_upstream_url( - transport, - mapped_model, - client_api_format, - provider_api_format, - upstream_is_stream, - parts.uri.query(), - ) -} diff --git a/apps/aether-gateway/src/codex_profile.rs b/apps/aether-gateway/src/codex_profile.rs index aa27cc50c..d9428df72 100644 --- a/apps/aether-gateway/src/codex_profile.rs +++ b/apps/aether-gateway/src/codex_profile.rs @@ -305,13 +305,11 @@ pub(crate) fn spawn_worker(app: AppState) -> tokio::task::JoinHandle<()> { #[cfg(test)] mod tests { - use std::sync::{ - atomic::{AtomicBool, Ordering}, - Mutex, OnceLock, - }; + use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; use aether_runtime_state::{MemoryRuntimeStateConfig, RuntimeState}; + use tokio::sync::{Mutex, MutexGuard}; use super::{ cached_version_to_restore, fixed_version_from, parse_cli_release, refresh_enabled_from, @@ -322,7 +320,7 @@ mod tests { set_codex_client_profile, CodexClientProfile, }; - static PROFILE_TEST_LOCK: OnceLock> = OnceLock::new(); + static PROFILE_TEST_LOCK: Mutex<()> = Mutex::const_new(()); struct ProfileRestore(CodexClientProfile); @@ -332,9 +330,8 @@ mod tests { } } - fn profile_restore_guard() -> (std::sync::MutexGuard<'static, ()>, ProfileRestore) { - let lock = PROFILE_TEST_LOCK.get_or_init(|| Mutex::new(())); - let guard = lock.lock().expect("profile test lock"); + async fn profile_restore_guard() -> (MutexGuard<'static, ()>, ProfileRestore) { + let guard = PROFILE_TEST_LOCK.lock().await; let restore = ProfileRestore(codex_client_profile()); (guard, restore) } @@ -397,7 +394,7 @@ mod tests { #[tokio::test] async fn cache_hit_is_restored_without_network_when_refresh_is_disabled() { - let (_lock, _restore) = profile_restore_guard(); + let (_lock, _restore) = profile_restore_guard().await; let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default()); runtime .kv_set( @@ -424,7 +421,7 @@ mod tests { #[tokio::test] async fn refresh_failure_keeps_previous_profile() { - let (_lock, _restore) = profile_restore_guard(); + let (_lock, _restore) = profile_restore_guard().await; let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default()); let before = codex_client_profile(); let result = refresh_once_with_fetch(&runtime, None, true, || async { @@ -438,7 +435,7 @@ mod tests { #[tokio::test] async fn fixed_version_override_skips_network_and_publishes_profile() { - let (_lock, _restore) = profile_restore_guard(); + let (_lock, _restore) = profile_restore_guard().await; let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default()); let fetch_called = AtomicBool::new(false); let result = refresh_once_with_fetch(&runtime, Some("0.220.0"), true, || async { @@ -455,7 +452,7 @@ mod tests { #[tokio::test] async fn rollback_is_rejected_without_replacing_profile() { - let (_lock, _restore) = profile_restore_guard(); + let (_lock, _restore) = profile_restore_guard().await; set_codex_cli_version("0.220.0").unwrap(); let runtime = RuntimeState::memory(MemoryRuntimeStateConfig::default()); let result = diff --git a/apps/aether-gateway/src/handlers/admin/observability/stats/analytics_routes.rs b/apps/aether-gateway/src/handlers/admin/observability/stats/analytics_routes.rs index a039afdcd..988c3f073 100644 --- a/apps/aether-gateway/src/handlers/admin/observability/stats/analytics_routes.rs +++ b/apps/aether-gateway/src/handlers/admin/observability/stats/analytics_routes.rs @@ -1,5 +1,8 @@ use super::super::resolve_usage_user_group_scope; -use super::range::{build_comparison_range, parse_bounded_u32}; +use super::range::{ + build_comparison_range, parse_bounded_u32, precise_admin_stats_time_range, + resolve_precise_time_bounds, +}; use super::resolve_admin_usage_time_range; use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; use crate::handlers::admin::shared::{ @@ -285,12 +288,36 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response( Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), }; - let time_range = match resolve_admin_usage_time_range(request_context.query_string()) { + let legacy_time_range = match resolve_admin_usage_time_range(request_context.query_string()) + { Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), }; - if let Err(detail) = time_range.validate_for_time_series(granularity) { - return Ok(Some(admin_stats_bad_request_response(detail))); + let precise_bounds = match resolve_precise_time_bounds(request_context.query_string()) { + Ok(value) => value, + Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), + }; + let precise_time_range = match precise_bounds { + Some((from, to)) => { + match precise_admin_stats_time_range(request_context.query_string(), from, to) { + Ok(value) => Some(value), + Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), + } + } + None => None, + }; + let time_range = precise_time_range.as_ref().unwrap_or(&legacy_time_range); + if precise_bounds.is_none() { + if let Err(detail) = time_range.validate_for_time_series(granularity) { + return Ok(Some(admin_stats_bad_request_response(detail))); + } + } else if precise_bounds + .and_then(|(from, to)| to.checked_sub(from)) + .is_some_and(|seconds| seconds > 90 * 86_400) + { + return Ok(Some(admin_stats_bad_request_response( + "Query range cannot exceed 90 days".to_string(), + ))); } if !state.has_usage_data_reader() { return Ok(Some(admin_stats_time_series_empty_response())); @@ -314,7 +341,8 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response( | AdminStatsGranularity::Week | AdminStatsGranularity::Month => UsageTimeSeriesGranularity::Day, }; - let Some((created_from_unix_secs, created_until_unix_secs)) = time_range.to_unix_bounds() + let Some((created_from_unix_secs, created_until_unix_secs)) = + precise_bounds.or_else(|| time_range.to_unix_bounds()) else { return Ok(Some(admin_stats_time_series_empty_response())); }; @@ -336,7 +364,7 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response( }) .await?; return Ok(Some(build_admin_stats_time_series_response_from_summaries( - &time_range, + time_range, granularity, &buckets, ))); diff --git a/apps/aether-gateway/src/handlers/admin/observability/stats/leaderboard_routes.rs b/apps/aether-gateway/src/handlers/admin/observability/stats/leaderboard_routes.rs index 49cd8a1a8..4b61a391c 100644 --- a/apps/aether-gateway/src/handlers/admin/observability/stats/leaderboard_routes.rs +++ b/apps/aether-gateway/src/handlers/admin/observability/stats/leaderboard_routes.rs @@ -5,7 +5,7 @@ use super::leaderboard::{ build_user_leaderboard_items_from_summaries, compare_leaderboard_items, load_user_leaderboard_metadata, AdminStatsLeaderboardItem, AdminStatsLeaderboardNameMode, }; -use super::range::{parse_bounded_u32, parse_nonnegative_usize}; +use super::range::{parse_bounded_u32, parse_nonnegative_usize, resolve_precise_time_bounds}; use super::resolve_admin_usage_time_range; use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; use crate::handlers::admin::shared::{query_param_bool, query_param_value}; @@ -228,6 +228,10 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response( Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), }; + let precise_bounds = match resolve_precise_time_bounds(query) { + Ok(value) => value, + Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), + }; let metric = match AdminStatsLeaderboardMetric::parse(query) { Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), @@ -272,7 +276,8 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response( "user_id is not supported for the user group leaderboard".to_string(), ))); } - let Some((created_from_unix_secs, created_until_unix_secs)) = time_range.to_unix_bounds() + let Some((created_from_unix_secs, created_until_unix_secs)) = + precise_bounds.or_else(|| time_range.to_unix_bounds()) else { return Ok(Some(build_admin_stats_user_group_leaderboard_response( metric, @@ -402,6 +407,10 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response( Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), }; + let precise_bounds = match resolve_precise_time_bounds(query) { + Ok(value) => value, + Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), + }; let metric = match AdminStatsLeaderboardMetric::parse(query) { Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), @@ -442,7 +451,8 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response( Ok(value) => value, Err(detail) => return Ok(Some(admin_stats_bad_request_response(detail))), }; - let Some((created_from_unix_secs, created_until_unix_secs)) = time_range.to_unix_bounds() + let Some((created_from_unix_secs, created_until_unix_secs)) = + precise_bounds.or_else(|| time_range.to_unix_bounds()) else { return Ok(Some(admin_stats_leaderboard_empty_response( metric, diff --git a/apps/aether-gateway/src/handlers/admin/observability/stats/mod.rs b/apps/aether-gateway/src/handlers/admin/observability/stats/mod.rs index 112f4f457..b2dee4760 100644 --- a/apps/aether-gateway/src/handlers/admin/observability/stats/mod.rs +++ b/apps/aether-gateway/src/handlers/admin/observability/stats/mod.rs @@ -8,7 +8,10 @@ mod leaderboard; mod leaderboard_routes; mod provider_quota_routes; mod range; -pub(crate) use self::range::{parse_bounded_u32, resolve_admin_usage_time_range}; +pub(crate) use self::range::{ + parse_bounded_u32, precise_admin_stats_time_range, resolve_admin_usage_time_range, + resolve_precise_time_bounds, resolve_usage_time_bounds, +}; pub(crate) use aether_admin::observability::stats::{ admin_stats_bad_request_response, aggregate_usage_stats, round_to, AdminStatsTimeRange, AdminStatsUsageFilter, diff --git a/apps/aether-gateway/src/handlers/admin/observability/stats/range.rs b/apps/aether-gateway/src/handlers/admin/observability/stats/range.rs index d8611ea03..c803cdf2d 100644 --- a/apps/aether-gateway/src/handlers/admin/observability/stats/range.rs +++ b/apps/aether-gateway/src/handlers/admin/observability/stats/range.rs @@ -4,10 +4,14 @@ pub(super) use aether_admin::observability::stats::{ admin_usage_default_days, build_comparison_range, build_time_range_from_days, parse_naive_date, parse_nonnegative_usize, parse_tz_offset_minutes, resolve_preset_dates, user_today, }; +use chrono::{DateTime, Offset, TimeZone, Utc}; pub(crate) fn resolve_admin_usage_time_range( query: Option<&str>, ) -> Result { + if let Some((from, to)) = resolve_precise_time_bounds(query)? { + return precise_admin_stats_time_range(query, from, to); + } match AdminStatsTimeRange::resolve_optional(query)? { Some(time_range) => Ok(time_range), None => { @@ -20,3 +24,153 @@ pub(crate) fn resolve_admin_usage_time_range( } } } + +/// Resolve an exact UTC range supplied by the shared admin range picker. +/// +/// The older stats handlers use `start_date`/`end_date` and fixed offsets. Keep +/// that parser intact and only opt into this path when both RFC 3339 endpoints +/// are present, so existing callers retain their behavior. +pub(crate) fn resolve_precise_time_bounds( + query: Option<&str>, +) -> Result, String> { + let entries = + url::form_urlencoded::parse(query.unwrap_or_default().as_bytes()).collect::>(); + let from = entries + .iter() + .filter(|(key, _)| key == "from") + .collect::>(); + let to = entries + .iter() + .filter(|(key, _)| key == "to") + .collect::>(); + if from.is_empty() && to.is_empty() { + return Ok(None); + } + if from.len() != 1 || to.len() != 1 { + return Err("from and to must each be provided once".into()); + } + if entries + .iter() + .any(|(key, _)| matches!(key.as_ref(), "start_date" | "end_date" | "preset" | "days")) + { + return Err("precise from/to cannot be combined with date presets".into()); + } + if let Some(zone) = query_param_value(query, "timezone") { + zone.parse::() + .map_err(|_| "invalid timezone".to_string())?; + } + let parse = |value: &str| -> Result { + let value = DateTime::parse_from_rfc3339(value) + .map_err(|_| "from/to must be RFC 3339 timestamps".to_string())?; + if value.timestamp_subsec_nanos() != 0 { + return Err("request records support second-aligned ranges".into()); + } + u64::try_from(value.timestamp()).map_err(|_| "from/to must not precede Unix epoch".into()) + }; + let bounds = (parse(&from[0].1)?, parse(&to[0].1)?); + if bounds.0 >= bounds.1 || bounds.1 - bounds.0 > 366 * 86_400 { + return Err("from/to must define a nonempty range of at most 366 days".into()); + } + Ok(Some(bounds)) +} + +/// Return the exact range when present, otherwise preserve the legacy stats +/// date/preset behavior. +pub(crate) fn resolve_usage_time_bounds(query: Option<&str>) -> Result, String> { + if let Some(bounds) = resolve_precise_time_bounds(query)? { + return Ok(Some(bounds)); + } + Ok(resolve_admin_usage_time_range(query)?.to_unix_bounds()) +} + +/// Build the date metadata used by the existing stats response builders for an +/// exact range. The data query still uses the exact UTC bounds; this metadata +/// only supplies the local date labels and offset expected by old clients. +pub(crate) fn precise_admin_stats_time_range( + query: Option<&str>, + from: u64, + to: u64, +) -> Result { + let timezone_name = query_param_value(query, "timezone"); + let (start_date, end_date, tz_offset_minutes) = if let Some(name) = timezone_name { + let timezone = name + .parse::() + .map_err(|_| "invalid timezone".to_string())?; + let start = Utc + .timestamp_opt( + i64::try_from(from).map_err(|_| "invalid from timestamp")?, + 0, + ) + .single() + .ok_or_else(|| "invalid from timestamp".to_string())? + .with_timezone(&timezone); + let end = Utc + .timestamp_opt( + i64::try_from(to.saturating_sub(1)).map_err(|_| "invalid to timestamp")?, + 0, + ) + .single() + .ok_or_else(|| "invalid to timestamp".to_string())? + .with_timezone(&timezone); + ( + start.date_naive(), + end.date_naive(), + start.offset().fix().local_minus_utc() / 60, + ) + } else { + let offset = parse_tz_offset_minutes(query)?; + let fixed = chrono::FixedOffset::east_opt(offset * 60) + .ok_or_else(|| "invalid timezone offset".to_string())?; + let start = Utc + .timestamp_opt( + i64::try_from(from).map_err(|_| "invalid from timestamp")?, + 0, + ) + .single() + .ok_or_else(|| "invalid from timestamp".to_string())? + .with_timezone(&fixed); + let end = Utc + .timestamp_opt( + i64::try_from(to.saturating_sub(1)).map_err(|_| "invalid to timestamp")?, + 0, + ) + .single() + .ok_or_else(|| "invalid to timestamp".to_string())? + .with_timezone(&fixed); + (start.date_naive(), end.date_naive(), offset) + }; + + Ok(AdminStatsTimeRange { + start_date, + end_date, + tz_offset_minutes, + }) +} + +fn query_param_value(query: Option<&str>, key: &str) -> Option { + url::form_urlencoded::parse(query.unwrap_or_default().as_bytes()) + .find(|(name, _)| name == key) + .map(|(_, value)| value.into_owned()) +} + +#[cfg(test)] +mod tests { + use super::{precise_admin_stats_time_range, resolve_precise_time_bounds}; + + #[test] + fn precise_stats_range_preserves_subday_bounds_and_timezone_labels() { + let query = "from=2026-09-01T23:45:00Z&to=2026-09-02T00:15:00Z&timezone=Asia%2FShanghai"; + let (from, to) = resolve_precise_time_bounds(Some(query)).unwrap().unwrap(); + assert_eq!(to - from, 30 * 60); + let range = precise_admin_stats_time_range(Some(query), from, to).unwrap(); + assert_eq!(range.start_date.to_string(), "2026-09-02"); + assert_eq!(range.end_date.to_string(), "2026-09-02"); + assert_eq!(range.tz_offset_minutes, 480); + } + + #[test] + fn precise_stats_range_rejects_mixed_legacy_presets() { + let query = "from=2026-09-01T00:00:00Z&to=2026-09-02T00:00:00Z&preset=today"; + assert!(resolve_precise_time_bounds(Some(query)).is_err()); + } +} diff --git a/apps/aether-gateway/src/handlers/admin/observability/usage/summary_routes.rs b/apps/aether-gateway/src/handlers/admin/observability/usage/summary_routes.rs index 8a4c1f48a..b2c795602 100644 --- a/apps/aether-gateway/src/handlers/admin/observability/usage/summary_routes.rs +++ b/apps/aether-gateway/src/handlers/admin/observability/usage/summary_routes.rs @@ -1,5 +1,5 @@ use super::super::resolve_usage_user_group_scope; -use super::super::stats::resolve_admin_usage_time_range; +use super::super::stats::resolve_usage_time_bounds; use super::analytics::admin_usage_api_key_names; use super::analytics::admin_usage_provider_key_names; use crate::handlers::admin::request::{AdminAppState, AdminRequestContext}; @@ -37,45 +37,7 @@ const ADMIN_USAGE_ACTIVE_LIMIT: usize = 50; pub(super) fn resolve_record_time_bounds( query: Option<&str>, ) -> Result, String> { - let entries = - url::form_urlencoded::parse(query.unwrap_or_default().as_bytes()).collect::>(); - let from = entries - .iter() - .filter(|(key, _)| key == "from") - .collect::>(); - let to = entries - .iter() - .filter(|(key, _)| key == "to") - .collect::>(); - if from.is_empty() && to.is_empty() { - return resolve_admin_usage_time_range(query).map(|range| range.to_unix_bounds()); - } - if from.len() != 1 || to.len() != 1 { - return Err("from and to must each be provided once".into()); - } - if entries - .iter() - .any(|(key, _)| matches!(key.as_ref(), "start_date" | "end_date" | "preset" | "days")) - { - return Err("precise from/to cannot be combined with date presets".into()); - } - if let Some(zone) = query_param_value(query, "timezone") { - zone.parse::() - .map_err(|_| "invalid timezone".to_string())?; - } - let parse = |value: &str| -> Result { - let value = chrono::DateTime::parse_from_rfc3339(value) - .map_err(|_| "from/to must be RFC 3339 timestamps".to_string())?; - if value.timestamp_subsec_nanos() != 0 { - return Err("request records support second-aligned ranges".into()); - } - u64::try_from(value.timestamp()).map_err(|_| "from/to must not precede Unix epoch".into()) - }; - let bounds = (parse(&from[0].1)?, parse(&to[0].1)?); - if bounds.0 >= bounds.1 || bounds.1 - bounds.0 > 366 * 86_400 { - return Err("from/to must define a nonempty range of at most 366 days".into()); - } - Ok(Some(bounds)) + resolve_usage_time_bounds(query) } async fn load_admin_usage_by_ids( diff --git a/apps/aether-gateway/src/maintenance/runtime/oauth_token_refresh.rs b/apps/aether-gateway/src/maintenance/runtime/oauth_token_refresh.rs index 7eff1577e..193032475 100644 --- a/apps/aether-gateway/src/maintenance/runtime/oauth_token_refresh.rs +++ b/apps/aether-gateway/src/maintenance/runtime/oauth_token_refresh.rs @@ -155,14 +155,17 @@ pub(crate) async fn perform_oauth_token_refresh_once( Ok(None) => { summary.skipped = summary.skipped.saturating_add(1); } - Err(_) => { + Err(err) => { summary.failed = summary.failed.saturating_add(1); warn!( event_name = "oauth_token_refresh_failed", log_type = "ops", worker = "oauth_token_refresh", provider_id = %provider.id, + provider_type = %provider.provider_type, + endpoint_id = %endpoint.id, key_id = %key.id, + error = %crate::error::redact_error_debug(&err), "gateway oauth token auto refresh failed" ); } @@ -440,6 +443,7 @@ mod tests { agent_identity_needs_task_recovery, auth_config_has_refresh_token, is_nonfatal_legacy_catalog_credential_error, oauth_refresh_candidate, }; + use crate::error::redact_error_debug; use crate::GatewayError; #[test] @@ -543,4 +547,18 @@ mod tests { ) )); } + + #[test] + fn oauth_refresh_failure_detail_preserves_context_without_credentials() { + let error = GatewayError::Internal( + r#"oauth request failed: status=503 token="refresh-secret" retry=2"#.to_string(), + ); + + let detail = redact_error_debug(&error); + + assert!(detail.contains("oauth request failed")); + assert!(detail.contains("status=503")); + assert!(detail.contains("[REDACTED]")); + assert!(!detail.contains("refresh-secret")); + } } diff --git a/apps/aether-gateway/src/state/oauth.rs b/apps/aether-gateway/src/state/oauth.rs index 0accef522..1094273f1 100644 --- a/apps/aether-gateway/src/state/oauth.rs +++ b/apps/aether-gateway/src/state/oauth.rs @@ -412,6 +412,17 @@ fn normalize_local_oauth_refresh_error_message( .unwrap_or_else(|| "Token 刷新失败".to_string()) } +fn local_oauth_refresh_gateway_error( + error: &provider_transport::LocalOAuthRefreshError, +) -> GatewayError { + // Keep a bounded, credential-redacted reason for internal diagnostics. + // GatewayError::Internal still returns the generic error response to clients. + GatewayError::Internal(format!( + "local oauth refresh failed: {}", + crate::error::redact_error_detail(error) + )) +} + fn merge_local_oauth_refresh_failure_reason( current_reason: Option<&str>, refresh_reason: &str, @@ -1556,10 +1567,8 @@ impl AppState { } return Ok(None); } - Err(_) => { - return Err(GatewayError::Internal( - "local oauth refresh failed".to_string(), - )); + Err(err) => { + return Err(local_oauth_refresh_gateway_error(&err)); } }; @@ -3484,13 +3493,86 @@ mod tests { use tokio::sync::Notify; use super::{ - AgentIdentityAuthConfigFence, AppState, CodexRuntimeOAuthObservation, - ProviderTransportSnapshotCacheKey, ProviderTransportSnapshotFlight, - ProviderTransportSnapshotFlightResult, ProviderTransportSnapshotInflightRegistration, - PROVIDER_TRANSPORT_SNAPSHOT_CACHE_STALE_TTL, PROVIDER_TRANSPORT_SNAPSHOT_CACHE_TTL, + local_oauth_refresh_gateway_error, AgentIdentityAuthConfigFence, AppState, + CodexRuntimeOAuthObservation, ProviderTransportSnapshotCacheKey, + ProviderTransportSnapshotFlight, ProviderTransportSnapshotFlightResult, + ProviderTransportSnapshotInflightRegistration, PROVIDER_TRANSPORT_SNAPSHOT_CACHE_STALE_TTL, + PROVIDER_TRANSPORT_SNAPSHOT_CACHE_TTL, }; use crate::data::GatewayDataState; + #[test] + fn oauth_refresh_diagnostic_preserves_failure_reason_and_redacts_credentials() { + for (message, secret) in [ + ( + "connection refused refresh_token=refresh-secret", + "refresh-secret", + ), + ( + "connection refused accessToken=access-secret", + "access-secret", + ), + ( + "connection refused client_secret=client-secret", + "client-secret", + ), + ( + "connection refused Authorization: Bearer bearer-secret", + "bearer-secret", + ), + ( + "connection refused https://proxy-user:proxy-secret@proxy.example", + "proxy-secret", + ), + ] { + let error = crate::provider_transport::LocalOAuthRefreshError::TransportMessage { + provider_type: "codex", + message: message.to_string(), + }; + let diagnostic = local_oauth_refresh_gateway_error(&error).into_message(); + assert!(diagnostic.contains("codex oauth refresh transport failed")); + assert!(diagnostic.contains("connection refused")); + assert!(!diagnostic.contains(secret), "diagnostic: {diagnostic}"); + } + } + + #[test] + fn oauth_refresh_diagnostic_keeps_http_status_without_provider_body() { + let error = crate::provider_transport::LocalOAuthRefreshError::HttpStatus { + provider_type: "codex", + status_code: 503, + body_excerpt: "unstructured-provider-secret".to_string(), + }; + let diagnostic = local_oauth_refresh_gateway_error(&error).into_message(); + + assert!(diagnostic.contains("codex oauth refresh returned HTTP 503")); + assert!(!diagnostic.contains("unstructured-provider-secret")); + } + + #[tokio::test] + async fn oauth_refresh_diagnostic_is_hidden_from_client_response() { + use axum::body::to_bytes; + use axum::response::IntoResponse; + + let error = crate::provider_transport::LocalOAuthRefreshError::TransportMessage { + provider_type: "codex", + message: "connection refused".to_string(), + }; + let response = local_oauth_refresh_gateway_error(&error).into_response(); + + assert_eq!( + response.status(), + axum::http::StatusCode::INTERNAL_SERVER_ERROR + ); + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("error response body should read"); + assert_eq!( + serde_json::from_slice::(&body).expect("error should be JSON"), + json!({"error": {"message": "internal server error"}}), + ); + } + fn sample_provider() -> StoredProviderCatalogProvider { StoredProviderCatalogProvider::new( "provider-1".to_string(), diff --git a/crates/aether-data/adapters/postgres/migrations/20261004000000_optimize_overview_fact_metadata.sql b/crates/aether-data/adapters/postgres/migrations/20261004000000_optimize_overview_fact_metadata.sql new file mode 100644 index 000000000..51240a686 --- /dev/null +++ b/crates/aether-data/adapters/postgres/migrations/20261004000000_optimize_overview_fact_metadata.sql @@ -0,0 +1,78 @@ +-- Replace only the read model; existing facts, queues and projections are untouched. +-- Preserve billing/attribution semantics while avoiding repeated JSON parsing. +CREATE OR REPLACE VIEW public.usage_analytics_facts_v1 AS +SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at, + CASE WHEN identity.owner_id IS NOT NULL AND identity.is_standalone=false THEN identity.owner_id END AS actor_user_id, + identity.owner_id AS credential_owner_id, + CASE WHEN identity.owner_id IS NULL THEN 'unknown' WHEN identity.is_standalone THEN 'standalone' + WHEN NOT identity.is_standalone THEN 'employee' ELSE 'unknown' END AS attribution_kind, + CASE WHEN identity.owner_id IS NULL THEN 'unknown' WHEN identity.is_standalone THEN 'standalone_key' + WHEN NOT identity.is_standalone THEN 'user_account' ELSE 'unknown' END AS attribution_source, + COALESCE(a.record_kind, 'request') AS record_kind, a.parent_request_id, + u.api_key_id, u.model, u.target_model, u.provider_id, u.provider_name, + u.api_format, u.endpoint_kind, u.request_type, u.is_stream, u.has_format_conversion, + u.status, u.status_code, u.error_category, u.failure_origin, u.failure_stage, u.failure_reason, + u.failure_schema_version, u.response_time_ms, u.first_byte_time_ms, + COALESCE(s.billing_status, u.billing_status) AS settlement_status, + COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available, + COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') AS pricing_available, + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb + THEN b.input_tokens END AS input_tokens, + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb + THEN b.output_tokens END AS output_tokens, + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb + THEN b.total_tokens END AS total_tokens, + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb + THEN b.cache_read_input_tokens END AS cache_read_input_tokens, + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb + THEN b.cache_creation_input_tokens END AS cache_creation_input_tokens, + CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') + THEN round(COALESCE(s.billing_total_cost_usd::numeric, u.total_cost_usd::numeric), 8) END AS rated_amount, + CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + AND (s.billing_actual_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') + THEN round(COALESCE(s.billing_actual_total_cost_usd::numeric, u.actual_total_cost_usd::numeric), 8) END AS billable_amount, + s.quota_covered_amount_usd AS quota_covered_amount, + s.wallet_consumed_amount_usd AS wallet_consumed_amount, + s.wallet_debit_amount_usd AS wallet_debit_amount, + s.wallet_recharge_debit_usd AS wallet_recharge_debit_amount, + s.wallet_gift_debit_usd AS wallet_gift_debit_amount, + s.wallet_overdraft_usd AS wallet_overdraft_amount, + s.allocation_status, s.finalized_at AS settled_at, + CASE WHEN s.billing_total_cost_usd IS NOT NULL THEN 'settlement_snapshot' ELSE 'legacy_float' END AS amount_source, + b.upstream_is_stream, + CASE WHEN metadata.value #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed') + THEN metadata.value #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source, + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL + THEN round(s.input_price_per_1m::numeric * b.cache_read_input_tokens::numeric / 1000000,8) END AS cache_estimated_full_cost_amount, + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL + THEN round(s.billing_cache_read_cost_usd::numeric,8) END AS cache_read_cost_amount, + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_creation_cost_usd IS NOT NULL + THEN round(s.billing_cache_creation_cost_usd::numeric,8) END AS cache_creation_cost_amount +FROM public.usage u +-- OFFSET 0 keeps this projection from being flattened: large metadata is +-- detoasted and parsed once per request, rather than once per metric expression. +CROSS JOIN LATERAL (SELECT u.request_metadata::jsonb AS value OFFSET 0) metadata +LEFT JOIN public.usage_settlement_snapshots s USING (request_id) +LEFT JOIN public.usage_attribution_snapshots a USING (request_id) +JOIN public.usage_billing_facts b USING (request_id) +LEFT JOIN public.api_keys k ON k.id=u.api_key_id +CROSS JOIN LATERAL ( + SELECT CASE WHEN a.request_id IS NOT NULL THEN a.credential_owner_id + WHEN EXISTS (SELECT 1 FROM public.users WHERE id=u.user_id AND NOT is_deleted) THEN u.user_id END AS owner_id, + COALESCE(k.is_standalone, + CASE WHEN jsonb_typeof(metadata.value #> '{analytics_attribution,is_standalone}')='boolean' + THEN (metadata.value #>> '{analytics_attribution,is_standalone}')::boolean END, + CASE WHEN jsonb_typeof(metadata.value->'api_key_is_standalone')='boolean' + THEN (metadata.value->>'api_key_is_standalone')::boolean END, + CASE WHEN a.attribution_source='user_account' THEN false + WHEN a.attribution_source='standalone_key' THEN true END, + CASE WHEN u.api_key_id IS NULL THEN false END) AS is_standalone +) identity; diff --git a/crates/aether-data/adapters/postgres/src/usage/overview_buckets.rs b/crates/aether-data/adapters/postgres/src/usage/overview_buckets.rs index 648f24d96..b878143c7 100644 --- a/crates/aether-data/adapters/postgres/src/usage/overview_buckets.rs +++ b/crates/aether-data/adapters/postgres/src/usage/overview_buckets.rs @@ -88,6 +88,7 @@ impl SqlxUsageReadRepository { for row in rows { let granularity: String = row.try_get("granularity").map_postgres_err()?; let bucket: DateTime = row.try_get("bucket_start").map_postgres_err()?; + let started = std::time::Instant::now(); match self .rebuild_merged_overview_bucket(&granularity, bucket) .await @@ -95,8 +96,20 @@ impl SqlxUsageReadRepository { Ok(true) => published += 1, Ok(false) => {} Err(error) => { + let error_detail = error.to_string().chars().take(500).collect::(); + tracing::warn!( + event_name = "overview_bucket_rebuild_failed", + log_type = "ops", + projection_version = "overview-v2", + granularity = %granularity, + bucket_start = %bucket, + elapsed_ms = started.elapsed().as_millis() as u64, + retry_after_secs = 600, + error = %error_detail, + "overview bucket rebuild failed; retry deferred" + ); sqlx::query("UPDATE stats_bucket_state SET last_error=$3,last_failed_at=NOW() WHERE projection_version='overview-v2' AND granularity=$1 AND bucket_start=$2") - .bind(&granularity).bind(bucket).bind(error.to_string().chars().take(500).collect::()) + .bind(&granularity).bind(bucket).bind(error_detail) .execute(&self.pool).await.map_postgres_err()?; } } @@ -144,6 +157,13 @@ impl SqlxUsageReadRepository { .execute(&mut *tx) .await .map_postgres_err()?; + // These bounded aggregates have many expressions but run only once per + // bucket. JIT compilation consumes a significant part of their timeout. + // Keep the setting transaction-local so other pool users retain theirs. + sqlx::query("SET LOCAL jit = off") + .execute(&mut *tx) + .await + .map_postgres_err()?; let acquired: bool = sqlx::query_scalar("SELECT pg_try_advisory_xact_lock(hashtextextended($1, 19))") .bind(format!("overview-v2:{granularity}:{}", bucket.timestamp())) diff --git a/crates/aether-data/runtime/schema/README.md b/crates/aether-data/runtime/schema/README.md index f178e52c7..77fb531d5 100644 --- a/crates/aether-data/runtime/schema/README.md +++ b/crates/aether-data/runtime/schema/README.md @@ -33,12 +33,17 @@ The schema workspace has three normal source areas: | `drivers/postgres/` | Current maintenance fragments for executable SQL. | Edit only for deployment compatibility, ordering, or generator gaps. | | `bootstrap/postgres/` | Source fragments for the Postgres empty-database bootstrap snapshot. | Edit here when the bootstrap snapshot changes, then rebuild `aether-data` so `build.rs` regenerates the embedded snapshot. | -Everything else is output: +Generated schema and the composed baseline are outputs: | Path | Role | Edit policy | |---|---|---| | `generated/postgres/` | Machine-written SQL emitted from `logical/*.toml` for audit and drift detection. | Do not edit; regenerate with `compose_schema.sh generate`. | -| `../../adapters/postgres/migrations/` | Runtime SQL artifacts embedded by each database adapter. | Regenerate through `compose_schema.sh compose`; do not edit independently. | +| `../../adapters/postgres/migrations/20260403000000_baseline.sql` | Composed PostgreSQL baseline embedded by the adapter. | Regenerate through `compose_schema.sh compose`; do not edit independently. | + +Later incremental migrations are maintained directly under +`../../adapters/postgres/migrations/`; they have no compose target. Add a new +version for an upgrade and preserve the checksums of already-applied scripts. +Keep any corresponding maintained bootstrap definitions in sync. `generated/**` is deliberately checked in so reviews and CI can see exactly what the logical schema compiler emits for each driver. It is not a fourth SQL @@ -131,6 +136,7 @@ cannot run them inside a transaction. | `20260921020000` | Add the attribution-owner lookup index concurrently on existing databases. | | `20260921020100` | Create the usage metadata actor index concurrently. | | `20261001000000` | Remove deleted-user attribution from dashboard activity on future user deletion; schema-only upgrade without rewriting historical rows. | +| `20261004000000` | Parse request metadata once per overview fact; replace only the view definition without rewriting facts or statistics. | Do not remove an applied migration after folding its changes into an earlier schema definition. Existing databases retain its version in `_sqlx_migrations` diff --git a/crates/aether-data/runtime/schema/bootstrap/postgres/190_overview_analytics.sql b/crates/aether-data/runtime/schema/bootstrap/postgres/190_overview_analytics.sql index 996532093..7db83dce9 100644 --- a/crates/aether-data/runtime/schema/bootstrap/postgres/190_overview_analytics.sql +++ b/crates/aether-data/runtime/schema/bootstrap/postgres/190_overview_analytics.sql @@ -215,23 +215,23 @@ SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at, u.status, u.status_code, u.error_category, u.failure_origin, u.failure_stage, u.failure_reason, u.failure_schema_version, u.response_time_ms, u.first_byte_time_ms, COALESCE(s.billing_status, u.billing_status) AS settlement_status, - COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available, - COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb AS usage_available, + COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') AS pricing_available, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN b.input_tokens END AS input_tokens, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN b.output_tokens END AS output_tokens, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN b.total_tokens END AS total_tokens, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN b.cache_read_input_tokens END AS cache_read_input_tokens, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available', 'true'::jsonb) <> 'false'::jsonb THEN b.cache_creation_input_tokens END AS cache_creation_input_tokens, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb AND (s.billing_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') THEN round(COALESCE(s.billing_total_cost_usd::numeric, u.total_cost_usd::numeric), 8) END AS rated_amount, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_pricing_available', 'true'::jsonb) <> 'false'::jsonb AND (s.billing_actual_total_cost_usd IS NOT NULL OR COALESCE(s.billing_status, u.billing_status) = 'settled') THEN round(COALESCE(s.billing_actual_total_cost_usd::numeric, u.actual_total_cost_usd::numeric), 8) END AS billable_amount, s.quota_covered_amount_usd AS quota_covered_amount, @@ -243,21 +243,24 @@ SELECT u.request_id, COALESCE(u.id, u.request_id) AS id, u.created_at, s.allocation_status, s.finalized_at AS settled_at, CASE WHEN s.billing_total_cost_usd IS NOT NULL THEN 'settlement_snapshot' ELSE 'legacy_float' END AS amount_source, b.upstream_is_stream, - CASE WHEN u.request_metadata #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed') - THEN u.request_metadata #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb - AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + CASE WHEN metadata.value #>> '{analytics_measurement,source}' IN ('reported','estimated','mixed') + THEN metadata.value #>> '{analytics_measurement,source}' ELSE 'unknown' END AS token_source, + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL THEN round(s.input_price_per_1m::numeric * b.cache_read_input_tokens::numeric / 1000000,8) END AS cache_estimated_full_cost_amount, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb - AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_read_cost_usd IS NOT NULL THEN round(s.billing_cache_read_cost_usd::numeric,8) END AS cache_read_cost_amount, - CASE WHEN COALESCE(u.request_metadata::jsonb->'usage_available','true'::jsonb) <> 'false'::jsonb - AND COALESCE(u.request_metadata::jsonb->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb + CASE WHEN COALESCE(metadata.value->'usage_available','true'::jsonb) <> 'false'::jsonb + AND COALESCE(metadata.value->'usage_pricing_available','true'::jsonb) <> 'false'::jsonb AND s.input_price_per_1m IS NOT NULL AND s.billing_cache_creation_cost_usd IS NOT NULL THEN round(s.billing_cache_creation_cost_usd::numeric,8) END AS cache_creation_cost_amount FROM public.usage u +-- OFFSET 0 keeps this projection from being flattened: large metadata is +-- detoasted and parsed once per request, rather than once per metric expression. +CROSS JOIN LATERAL (SELECT u.request_metadata::jsonb AS value OFFSET 0) metadata LEFT JOIN public.usage_settlement_snapshots s USING (request_id) LEFT JOIN public.usage_attribution_snapshots a USING (request_id) JOIN public.usage_billing_facts b USING (request_id) @@ -266,10 +269,10 @@ CROSS JOIN LATERAL ( SELECT CASE WHEN a.request_id IS NOT NULL THEN a.credential_owner_id WHEN EXISTS (SELECT 1 FROM public.users WHERE id=u.user_id AND NOT is_deleted) THEN u.user_id END AS owner_id, COALESCE(k.is_standalone, - CASE WHEN jsonb_typeof(u.request_metadata::jsonb #> '{analytics_attribution,is_standalone}')='boolean' - THEN (u.request_metadata #>> '{analytics_attribution,is_standalone}')::boolean END, - CASE WHEN jsonb_typeof(u.request_metadata::jsonb->'api_key_is_standalone')='boolean' - THEN (u.request_metadata->>'api_key_is_standalone')::boolean END, + CASE WHEN jsonb_typeof(metadata.value #> '{analytics_attribution,is_standalone}')='boolean' + THEN (metadata.value #>> '{analytics_attribution,is_standalone}')::boolean END, + CASE WHEN jsonb_typeof(metadata.value->'api_key_is_standalone')='boolean' + THEN (metadata.value->>'api_key_is_standalone')::boolean END, CASE WHEN a.attribution_source='user_account' THEN false WHEN a.attribution_source='standalone_key' THEN true END, CASE WHEN u.api_key_id IS NULL THEN false END) AS is_standalone diff --git a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs index e89b0c215..67baa12cf 100644 --- a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs +++ b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs @@ -31,6 +31,7 @@ mod dashboard_user_anonymization; mod legacy_overview_upgrade; mod migration_deadlines; mod overview_dirty_events; +mod overview_fact_metadata; mod overview_migration_safety; mod policy_nulls; mod provider_expenses; @@ -1594,6 +1595,7 @@ fn pending_migrations_from_applied_skips_versions_already_applied() { 20260921020100, 20260923000000, 20261001000000, + 20261004000000, ] ); } diff --git a/crates/aether-data/runtime/src/lifecycle/migrate/tests/overview_fact_metadata.rs b/crates/aether-data/runtime/src/lifecycle/migrate/tests/overview_fact_metadata.rs new file mode 100644 index 000000000..ca8552f24 --- /dev/null +++ b/crates/aether-data/runtime/src/lifecycle/migrate/tests/overview_fact_metadata.rs @@ -0,0 +1,180 @@ +use super::*; +use serde_json::Value; + +const OPTIMIZATION_VERSION: i64 = 20261004000000; + +async fn read_facts(pool: &PgPool) -> Value { + query_scalar( + "SELECT jsonb_agg(to_jsonb(f) ORDER BY request_id) FROM usage_analytics_facts_v1 f", + ) + .fetch_one(pool) + .await + .unwrap() +} + +async fn stored_rows(pool: &PgPool) -> Vec { + let mut rows = Vec::new(); + for table in [ + "usage", + "usage_settlement_snapshots", + "usage_attribution_snapshots", + "stats_overview_dirty_events", + "stats_bucket_state", + "stats_overview_hourly", + "stats_overview_daily", + "dashboard_request_contributions", + "dashboard_stats_total", + "dashboard_activity_minute", + ] { + rows.push( + query_scalar(&format!( + "SELECT COALESCE(jsonb_agg(to_jsonb(t) ORDER BY to_jsonb(t)), '[]') FROM {table} t" + )) + .fetch_one(pool) + .await + .unwrap(), + ); + } + rows +} + +#[tokio::test] +async fn overview_fact_metadata_optimization_preserves_facts_without_backfill() { + let Some(server) = ManagedPostgresServer::try_start().await.unwrap() else { + return; + }; + let mut connection = PgConnection::connect(server.database_url()).await.unwrap(); + connection.ensure_migrations_table().await.unwrap(); + for migration in POSTGRES_MIGRATOR + .iter() + .filter(|migration| migration.version < OPTIMIZATION_VERSION) + { + connection.apply(migration).await.unwrap(); + } + let pool = PgPool::connect(server.database_url()).await.unwrap(); + sqlx::raw_sql( + r#" +INSERT INTO users(id,username,email_verified) VALUES('owner','owner',false); +INSERT INTO api_keys(id,user_id,key_hash,is_standalone) +VALUES('employee-key','owner',repeat('e',64),false), + ('standalone-key','owner',repeat('f',64),true); +"#, + ) + .execute(&pool) + .await + .unwrap(); + + // Preserve distinctions between absent/null flags, booleans and strings, + // malformed nested shapes, key overrides and pre-attribution legacy rows. + let metadata_cases = [ + None, + Some("null"), + Some("[]"), + Some("false"), + Some("\"metadata\""), + Some("{}"), + Some(r#"{"usage_available":false,"usage_pricing_available":false}"#), + Some(r#"{"usage_available":null,"usage_pricing_available":null}"#), + Some(r#"{"usage_available":"false","usage_pricing_available":"false"}"#), + Some( + r#"{"analytics_attribution":{"is_standalone":true},"analytics_measurement":{"source":"reported"}}"#, + ), + Some( + r#"{"analytics_attribution":{"is_standalone":false},"analytics_measurement":{"source":"estimated"}}"#, + ), + Some( + r#"{"analytics_attribution":{"is_standalone":"true"},"api_key_is_standalone":false,"analytics_measurement":{"source":"mixed"}}"#, + ), + Some( + r#"{"analytics_attribution":[],"api_key_is_standalone":true,"analytics_measurement":{"source":false}}"#, + ), + Some( + r#"{"analytics_attribution":null,"api_key_is_standalone":"true","analytics_measurement":[]}"#, + ), + Some( + r#"{"usage_available":true,"usage_available":false,"analytics_attribution":{"is_standalone":false},"analytics_attribution":{"is_standalone":true}}"#, + ), + ]; + for (index, metadata) in metadata_cases.into_iter().enumerate() { + for key in [None, Some("employee-key"), Some("standalone-key")] { + let request = format!("metadata-{index}-{}", key.unwrap_or("legacy")); + query( + r#"INSERT INTO usage(id,request_id,user_id,api_key_id,model,provider_name, + status,billing_status,input_tokens,output_tokens,total_tokens, + cache_read_input_tokens,cache_creation_input_tokens,total_cost_usd, + actual_total_cost_usd,response_time_ms,first_byte_time_ms,is_stream, + created_at,request_metadata) + VALUES($1,$1,'owner',$2,'test','test','completed','settled',100,10,110, + 20,5,0.12345678,0.11111111,1000,100,true,'2026-01-01 00:00:00+00',$3::json)"#, + ) + .bind(&request) + .bind(key) + .bind(metadata) + .execute(&pool) + .await + .unwrap(); + if index % 2 == 0 { + query( + r#"INSERT INTO usage_settlement_snapshots(request_id,billing_status, + billing_input_tokens,billing_effective_input_tokens,billing_output_tokens, + billing_cache_read_tokens,billing_cache_creation_tokens, + billing_total_cost_usd,billing_actual_total_cost_usd,input_price_per_1m, + billing_cache_read_cost_usd,billing_cache_creation_cost_usd, + quota_covered_amount_usd,wallet_consumed_amount_usd,wallet_debit_amount_usd, + wallet_recharge_debit_usd,wallet_gift_debit_usd,wallet_overdraft_usd, + allocation_status) + VALUES($1,'settled',200,175,20,15,10,0.3,0.25,2,0.00001,0.00002, + 0.05,0.2,0.2,0.1,0.1,0,'complete')"#, + ) + .bind(&request) + .execute(&pool) + .await + .unwrap(); + } + if key.is_none() { + query("DELETE FROM usage_attribution_snapshots WHERE request_id=$1") + .bind(&request) + .execute(&pool) + .await + .unwrap(); + } + } + } + // Keep a large nested payload so parity also covers toasted JSON metadata. + query("UPDATE usage SET request_metadata=json_build_object('usage_available',true,'payload',repeat('metadata payload ',4096)) WHERE request_id='metadata-5-legacy'") + .execute(&pool).await.unwrap(); + let before_facts = read_facts(&pool).await; + let before_rows = stored_rows(&pool).await; + + // The definition-only upgrade must not touch historical projections/queues. + let mut blocked_history = pool.begin().await.unwrap(); + query("LOCK TABLE stats_bucket_state,stats_overview_dirty_events,stats_overview_hourly,stats_overview_daily IN ACCESS EXCLUSIVE MODE") + .execute(&mut *blocked_history).await.unwrap(); + query("SET lock_timeout='500ms'") + .execute(&mut connection) + .await + .unwrap(); + let migration = POSTGRES_MIGRATOR + .iter() + .find(|migration| migration.version == OPTIMIZATION_VERSION) + .unwrap(); + connection.apply(migration).await.unwrap(); + blocked_history.rollback().await.unwrap(); + + assert_eq!(read_facts(&pool).await, before_facts); + assert_eq!(stored_rows(&pool).await, before_rows); + + // The maintained bootstrap fragment must install the same read model. + let bootstrap = + include_str!("../../../../schema/bootstrap/postgres/190_overview_analytics.sql"); + let view_start = bootstrap + .find("CREATE OR REPLACE VIEW public.usage_analytics_facts_v1 AS") + .unwrap(); + sqlx::raw_sql(&bootstrap[view_start..]) + .execute(&pool) + .await + .unwrap(); + assert_eq!(read_facts(&pool).await, before_facts); + assert_eq!(stored_rows(&pool).await, before_rows); + pool.close().await; +} diff --git a/frontend/src/api/__tests__/admin-analytics-cache.spec.ts b/frontend/src/api/__tests__/admin-analytics-cache.spec.ts index a4a14c82b..0e3060cd9 100644 --- a/frontend/src/api/__tests__/admin-analytics-cache.spec.ts +++ b/frontend/src/api/__tests__/admin-analytics-cache.spec.ts @@ -113,4 +113,13 @@ describe('adminApi analytics cache options', () => { }) }) + it('refreshes both leaderboards for the same precise period as user accounts', async () => { + const range = { from: '2026-09-10T15:20:42Z', to: '2026-09-10T16:20:42Z', timezone: 'Asia/Shanghai' } + await adminApi.getLeaderboardUsers(range, { skipCache: true }) + await adminApi.getLeaderboardUserGroups(range, { skipCache: true }) + expect(getMock).toHaveBeenCalledWith('/api/admin/stats/leaderboard/users', { params: range }) + expect(getMock).toHaveBeenCalledWith('/api/admin/stats/leaderboard/user-groups', { params: range }) + for (const call of cachedRequestMock.mock.calls as unknown[][]) expect(call[2]).toBe(0) + }) + }) diff --git a/frontend/src/api/admin.ts b/frontend/src/api/admin.ts index 1010f67c0..d4cec093e 100644 --- a/frontend/src/api/admin.ts +++ b/frontend/src/api/admin.ts @@ -1288,6 +1288,8 @@ export const adminApi = { // Stats / Leaderboards async getLeaderboardUsers(params?: { + from?: string + to?: string start_date?: string end_date?: string preset?: string @@ -1302,7 +1304,7 @@ export const adminApi = { include_inactive?: boolean exclude_admin?: boolean user_group_id?: string - }): Promise { + }, options?: AdminAnalyticsRequestOptions): Promise { const cacheKey = buildCacheKey('admin:stats:leaderboard:users', params) return cachedRequest( cacheKey, @@ -1312,11 +1314,13 @@ export const adminApi = { }) return response.data }, - 20 * 1000 + options?.skipCache ? 0 : 20 * 1000 ) }, async getLeaderboardUserGroups(params?: { + from?: string + to?: string start_date?: string end_date?: string preset?: string @@ -1330,7 +1334,7 @@ export const adminApi = { model?: string include_inactive?: boolean exclude_admin?: boolean - }): Promise { + }, options?: AdminAnalyticsRequestOptions): Promise { const cacheKey = buildCacheKey('admin:stats:leaderboard:user-groups', params) return cachedRequest( cacheKey, @@ -1341,7 +1345,7 @@ export const adminApi = { ) return response.data }, - 20 * 1000 + options?.skipCache ? 0 : 20 * 1000 ) }, @@ -1621,6 +1625,8 @@ export const adminApi = { async getTimeSeries( params?: { + from?: string + to?: string start_date?: string end_date?: string preset?: string diff --git a/frontend/src/features/overview/__tests__/UserUsageStats.scope.spec.ts b/frontend/src/features/overview/__tests__/UserUsageStats.scope.spec.ts index ee5b14843..a3477f46a 100644 --- a/frontend/src/features/overview/__tests__/UserUsageStats.scope.spec.ts +++ b/frontend/src/features/overview/__tests__/UserUsageStats.scope.spec.ts @@ -1,15 +1,14 @@ -import { createApp, nextTick } from 'vue' +import { createApp, h, nextTick, reactive } from 'vue' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { setI18nLocale } from '@/i18n' import UserUsageStats from '../users/UserUsageStats.vue' const api = vi.hoisted(() => ({ - users: vi.fn(), groups: vi.fn(), members: vi.fn(), - userLeaderboard: vi.fn(), groupLeaderboard: vi.fn(), summary: vi.fn(), series: vi.fn(), + users: vi.fn(), groups: vi.fn(), + userLeaderboard: vi.fn(), groupLeaderboard: vi.fn(), series: vi.fn(), })) -vi.mock('@/api/users', () => ({ usersApi: { getAllUsers: api.users, listUserGroups: api.groups, listUserGroupMembers: api.members } })) +vi.mock('@/api/users', () => ({ usersApi: { getAllUsers: api.users, listUserGroups: api.groups } })) vi.mock('@/api/admin', () => ({ adminApi: { getLeaderboardUsers: api.userLeaderboard, getLeaderboardUserGroups: api.groupLeaderboard, getTimeSeries: api.series } })) -vi.mock('@/api/usage', () => ({ usageApi: { getUsageStats: api.summary } })) vi.mock('@/components/charts/LineChart.vue', () => ({ default: { render: () => null } })) vi.mock('@/components/common', () => ({ EmptyState: { render: () => null }, LoadingState: { render: () => null }, TimeRangePicker: { render: () => null } })) vi.mock('@/components/ui', async importOriginal => { @@ -31,13 +30,22 @@ vi.mock('@/components/ui', async importOriginal => { }) let unmount = () => {} +const range = { from: '2026-09-10T15:20:42.000Z', to: '2026-09-10T16:20:42.000Z', timezone: 'Asia/Shanghai' } +let props = reactive({ range: { ...range }, revision: 0 }) async function settle() { for (let index = 0; index < 10; index += 1) await Promise.resolve() await nextTick() } async function mount() { const root = document.createElement('div') - const app = createApp(UserUsageStats) + const app = createApp({ + render: () => h(UserUsageStats, props, { + 'user-leaderboard': ({ selectUser }: { selectUser: (user: { user_id: string; username: string }) => void }) => h('button', { + 'data-select-account': '', + onClick: () => selectUser({ user_id: 'user-3', username: 'Charlie' }), + }, 'Merged user accounts'), + }), + }) app.mount(root) unmount = () => app.unmount() await settle() @@ -55,39 +63,113 @@ beforeEach(() => { vi.useFakeTimers() vi.clearAllMocks() setI18nLocale('en-US') + props = reactive({ range: { ...range }, revision: 0 }) api.users.mockResolvedValue([{ id: 'user-1', username: 'Alice', is_active: true, groups: [{ id: 'group-1' }] }, { id: 'user-2', username: 'Bob', is_active: true, groups: [] }]) api.groups.mockResolvedValue({ items: [{ id: 'group-1', name: 'Engineering' }, { id: 'group-2', name: 'Support' }] }) - api.members.mockResolvedValue([{ id: 'user-1', is_active: true, is_deleted: false }]) api.userLeaderboard.mockResolvedValue({ items: [], total: 0 }) api.groupLeaderboard.mockResolvedValue({ items: [], total: 0 }) - api.summary.mockResolvedValue({ total_requests: 5, total_tokens: 100, total_cost: 2 }) api.series.mockResolvedValue([]) }) afterEach(() => { unmount(); vi.useRealTimers(); setI18nLocale('zh-CN') }) describe('user and group usage statistics', () => { - it('applies group scope to summary, trends and member rankings, and can compare groups', async () => { + it('uses the merged account table without fetching another user leaderboard', async () => { + const root = await mount() + expect(root.textContent).toContain('Merged user accounts') + expect(api.userLeaderboard).not.toHaveBeenCalled() + expect(api.groupLeaderboard).not.toHaveBeenCalled() + expect(root.textContent).not.toContain('User summary') + expect(root.textContent).not.toContain('User-group summary') + + root.querySelector('[data-select-account]')!.click() + await nextTick() + await vi.advanceTimersByTimeAsync(120) + await settle() + expect(root.querySelectorAll('select')[1].value).toBe('user-3') + expect(root.querySelector('h3')?.parentElement?.textContent).toContain('Charlie') + expect(api.series).toHaveBeenLastCalledWith(expect.objectContaining({ ...range, user_id: 'user-3' }), { skipCache: true }) + + await select(root, 0, 'user_group') + expect(root.querySelector('[data-select-account]')).toBeNull() + await select(root, 0, 'user') + expect(root.querySelectorAll('select')[1].value).toBe('user-3') + expect(root.querySelector('h3')?.parentElement?.textContent).toContain('Charlie') + expect(api.userLeaderboard.mock.calls.every(([params]) => params.user_group_id)).toBe(true) + }) + + it('applies group scope to trends and member rankings, and can compare groups', async () => { const root = await mount() - expect(api.summary).toHaveBeenLastCalledWith(expect.objectContaining({ user_id: 'user-1' })) await select(root, 0, 'user_group') expect(api.groupLeaderboard).toHaveBeenCalled() - expect(api.summary).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: 'group-1' })) - expect(api.summary.mock.lastCall?.[0]).not.toHaveProperty('user_id') - expect(api.userLeaderboard).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: 'group-1', limit: 10 })) - expect(api.members).toHaveBeenLastCalledWith('group-1') + expect(api.userLeaderboard).toHaveBeenLastCalledWith(expect.objectContaining({ ...range, user_group_id: 'group-1', limit: 10 }), { skipCache: true }) await select(root, 2, 'group-2') - expect(api.series).toHaveBeenCalledWith(expect.objectContaining({ user_group_id: 'group-2' })) + expect(api.series).toHaveBeenCalledWith(expect.objectContaining({ ...range, user_group_id: 'group-2' }), { skipCache: true }) expect(root.textContent).toContain('Group member leaderboard') }) it('keeps the ungrouped sentinel across scoped queries without fetching a fictitious group', async () => { const root = await mount() await select(root, 0, 'user_group') - api.members.mockClear() await select(root, 1, '__ungrouped__') - expect(api.summary).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: '__ungrouped__' })) - expect(api.series).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: '__ungrouped__' })) - expect(api.userLeaderboard).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: '__ungrouped__' })) - expect(api.members).not.toHaveBeenCalled() + expect(api.series).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: '__ungrouped__' }), { skipCache: true }) + expect(api.userLeaderboard).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: '__ungrouped__' }), { skipCache: true }) + }) + + it('uses the page range and refresh for trends and comparisons without losing selection', async () => { + const root = await mount() + await select(root, 2, 'user-2') + api.userLeaderboard.mockClear() + api.series.mockClear() + props.range = { from: '2026-09-03T02:47:00.000Z', to: '2026-10-03T02:47:00.000Z', timezone: 'Asia/Shanghai' } + await nextTick() + await vi.advanceTimersByTimeAsync(120) + expect(api.userLeaderboard).not.toHaveBeenCalled() + expect(api.series).toHaveBeenCalledWith(expect.objectContaining({ ...props.range, user_id: 'user-2' }), { skipCache: true }) + props.revision += 1 + await nextTick() + await vi.advanceTimersByTimeAsync(120) + expect(api.userLeaderboard).not.toHaveBeenCalled() + expect(api.series).toHaveBeenCalledTimes(4) + expect(root.querySelector('h1')).toBeNull() + expect(root.querySelectorAll('select')[2].value).toBe('user-2') + }) + + it('discards the old trend when the page range changes during a request', async () => { + let resolveOld: (value: unknown) => void = () => {} + api.series.mockImplementationOnce(() => new Promise(resolve => { resolveOld = resolve })) + const root = await mount() + props.range = { ...range, from: '2026-09-10T16:00:00.000Z' } + await nextTick() + resolveOld([{ date: 'old', total_cost: 999 }]) + await settle() + expect(root.textContent).not.toContain('999') + await vi.advanceTimersByTimeAsync(120) + expect(api.series.mock.lastCall?.[0]).toMatchObject(props.range) + }) + + it('reports failed requests and retries them without unmounting the usage section', async () => { + api.series.mockRejectedValueOnce(new Error('Trend unavailable')) + const root = await mount() + expect(root.textContent).toContain('Trend unavailable') + const retry = [...root.querySelectorAll('button')].find(item => item.textContent?.includes('Retry')) + expect(retry).toBeDefined() + retry!.click() + await settle() + expect(api.series).toHaveBeenCalledTimes(2) + expect(api.userLeaderboard).not.toHaveBeenCalled() + expect(root.textContent).not.toContain('Trend unavailable') + }) + + it('retries failed group rankings while preserving scoped member rankings', async () => { + api.groupLeaderboard.mockRejectedValueOnce(new Error('Group rankings unavailable')) + const root = await mount() + await select(root, 0, 'user_group') + expect(root.textContent).toContain('Group rankings unavailable') + const retry = [...root.querySelectorAll('button')].find(item => item.textContent?.includes('Retry')) + retry!.click() + await settle() + expect(api.groupLeaderboard).toHaveBeenCalledTimes(2) + expect(api.userLeaderboard).toHaveBeenLastCalledWith(expect.objectContaining({ user_group_id: 'group-1' }), { skipCache: true }) + expect(root.textContent).not.toContain('Group rankings unavailable') }) }) diff --git a/frontend/src/features/overview/__tests__/users.spec.ts b/frontend/src/features/overview/__tests__/users.spec.ts index 1547b9128..f5ced0f70 100644 --- a/frontend/src/features/overview/__tests__/users.spec.ts +++ b/frontend/src/features/overview/__tests__/users.spec.ts @@ -1,4 +1,4 @@ -import { createApp, h, nextTick, type Component } from 'vue' +import { createApp, nextTick, type Component } from 'vue' import { createMemoryHistory, createRouter } from 'vue-router' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import UserStats from '@/views/admin/UserStats.vue' @@ -7,12 +7,18 @@ import { setI18nLocale } from '@/i18n' const api = vi.hoisted(() => ({ users: vi.fn(), user: vi.fn(), timeseries: vi.fn(), breakdown: vi.fn(), consumption: vi.fn(), exportCsv: vi.fn() })) const accountApi = vi.hoisted(() => ({ wallets: vi.fn(), transactions: vi.fn(), plans: vi.fn() })) +const usageStats = vi.hoisted(() => ({ selectUser: vi.fn() })) vi.mock('@/api/overview', () => ({ overviewApi: api })) vi.mock('@/api/admin-wallets', () => ({ adminWalletApi: { listWallets: accountApi.wallets, getWalletTransactions: accountApi.transactions } })) vi.mock('@/api/users', () => ({ usersApi: { listUserPlanEntitlements: accountApi.plans } })) vi.mock('@/components/charts/BarChart.vue', () => ({ default: { render: () => null } })) vi.mock('@/components/charts/LineChart.vue', () => ({ default: { render: () => null } })) -vi.mock('@/features/overview/users/UserUsageStats.vue', () => ({ __esModule: true, default: { render: () => h('div', { 'data-user-usage-stats': '' }, 'Usage statistics') } })) +vi.mock('@/features/overview/users/UserUsageStats.vue', async () => { + const { defineComponent, h } = await import('vue') + return { default: defineComponent({ + setup: (_, { slots }) => () => h('div', { 'data-user-usage-stats': '' }, slots['user-leaderboard']?.({ selectUser: usageStats.selectUser })), + }) } +}) const range = 'from=2026-09-01T00:00:00Z&to=2026-09-02T00:00:00Z&timezone=UTC' const meta = { schema_version: 1, metric_version: 'overview-v2', scope: { kind: 'installation' }, @@ -64,17 +70,15 @@ beforeEach(() => { afterEach(() => { cleanup.splice(0).forEach(fn => fn()); vi.useRealTimers(); setI18nLocale('zh-CN') }) describe('enterprise user accounts', () => { - it('loads usage and group statistics only when selected, preserving the account view', async () => { + it('shows account and usage statistics together on one page', async () => { const { root } = await mount(UserStats, `/admin/user-stats?${range}`) - expect(root.querySelector('[data-user-usage-stats]')).toBeNull() - button(root, '使用统计').click() - await settle() - await settle() expect(root.querySelector('[data-user-usage-stats]')).not.toBeNull() - root.querySelector('button[data-value="accounts"]')?.click() - await settle() - expect(root.querySelector('[data-user-usage-stats]')).toBeNull() expect(section(root, '[data-user-accounts]').textContent).toContain(employee.username) + expect(section(root, '[data-user-reports]').querySelector('[data-user-accounts]')).not.toBeNull() + expect(root.querySelectorAll('[data-user-accounts]')).toHaveLength(1) + expect(section(root, '[data-user-accounts]').textContent).toContain('用户排行与账目') + expect(root.querySelector('[role="tablist"]')).toBeNull() + expect(root.querySelector('button[data-value="accounts"]')).toBeNull() expect(api.users).toHaveBeenCalledTimes(1) }) @@ -84,25 +88,39 @@ describe('enterprise user accounts', () => { expect(section(root, '[data-user-summary="consumption"]').textContent).toContain('432.25') expect(section(root, '[data-user-summary="activity"]').textContent).toMatch(/7\s*\/ 62/) expect(root.querySelector('select[aria-label="归属"]')).toBeNull() - expect(root.querySelector('button[data-value="accounts"][data-state="active"]')).not.toBeNull() expect(api.users.mock.lastCall?.[0]).toMatchObject({ sort: 'billable_amount', order: 'desc', limit: 25, offset: 0 }) expect(api.users.mock.lastCall?.[0]).not.toHaveProperty('attribution_kind') expect(api.users.mock.lastCall?.[0]).not.toHaveProperty('model') + expect(section(root, '[data-user-rank]').textContent?.trim()).toBe('1') button(root, '第 2 页').click() await settle() expect(api.users.mock.lastCall?.[0]).toMatchObject({ offset: 25 }) + expect(section(root, '[data-user-rank]').textContent?.trim()).toBe('26') button(root, 'Tokens').click() await settle() expect(api.users.mock.lastCall?.[0]).toMatchObject({ offset: 0, sort: 'total_tokens', order: 'desc' }) + expect(section(root, '[data-user-rank]').textContent?.trim()).toBe('1') expect(root.querySelector('a[href*="employee-0"]')).toBeNull() expect(root.querySelector('a[href*="/admin/usage"]')).toBeNull() expect([...section(root, '[data-user-accounts]').querySelectorAll('a, button')].some(item => item.textContent?.trim() === '使用记录')).toBe(false) }) - it('places reports first and opens reusable account history without leaving the page', async () => { + it('opens trends from the merged table and labels ascending rows as positions', async () => { + const { root } = await mount(UserStats, `/admin/user-stats?${range}`) + const accounts = section(root, '[data-user-accounts]') + button(accounts, '使用趋势').click() + expect(usageStats.selectUser).toHaveBeenCalledWith(expect.objectContaining({ user_id: employee.user_id, username: employee.username })) + expect(accountApi.wallets).not.toHaveBeenCalled() + button(accounts, '消费').click() + await settle() + expect(api.users.mock.lastCall?.[0]).toMatchObject({ sort: 'billable_amount', order: 'asc', offset: 0 }) + expect(accounts.querySelector('thead')?.textContent).toContain('序号') + expect(accounts.querySelector('thead')?.textContent).not.toContain('排名') + }) + it('includes accounts in reports and opens reusable account history without leaving the page', async () => { const { root, router } = await mount(UserStats, `/admin/user-stats?${range}`) const reports = section(root, '[data-user-reports]') const accounts = section(root, '[data-user-accounts]') - expect(reports.compareDocumentPosition(accounts) & Node.DOCUMENT_POSITION_FOLLOWING).toBeTruthy() + expect(reports.contains(accounts)).toBe(true) expect(accountApi.wallets).not.toHaveBeenCalled() expect(accountApi.plans).not.toHaveBeenCalled() const before = router.currentRoute.value.fullPath diff --git a/frontend/src/features/overview/components/OverviewToolbar.vue b/frontend/src/features/overview/components/OverviewToolbar.vue index 2b3af0d25..1f335c308 100644 --- a/frontend/src/features/overview/components/OverviewToolbar.vue +++ b/frontend/src/features/overview/components/OverviewToolbar.vue @@ -8,52 +8,54 @@
- - - + + + + + - {{ t('全站用量 · 包含独立余额 Key · 不受用户搜索影响', 'Installation usage · Includes standalone balance keys · Independent of user search') }} + {{ t('模型用量与消费趋势为全站用量 · 包含独立余额 Key · 不受用户搜索影响', 'Model usage and consumption trends cover installation usage · Includes standalone balance keys · Independent of user search') }}