Merge pull request #838 from dalamudx/feat/user-group-provider-stats

fix(stats): scope group usage by providers and add ungrouped view
This commit is contained in:
ZheFox
2026-09-29 12:23:27 +08:00
committed by GitHub
15 changed files with 637 additions and 100 deletions
@@ -31,6 +31,16 @@ pub(crate) async fn resolve_usage_user_group_scope(
return Ok(Err("user group data is unavailable".to_string()));
}
if group_id == UNGROUPED_USAGE_ID {
let ids = ungrouped_usage_users(state)
.await?
.into_iter()
.filter(|user| include_inactive || user.is_active)
.filter(|user| !exclude_admin || !user.role.eq_ignore_ascii_case("admin"))
.map(|user| user.id)
.collect();
return Ok(Ok(Some(ids)));
}
match state
.resolve_usage_user_group_member_ids(&group_id, include_inactive, exclude_admin)
.await?
@@ -39,3 +49,100 @@ pub(crate) async fn resolve_usage_user_group_scope(
None => Ok(Err("user_group_id does not exist".to_string())),
}
}
/// Reserved statistics-only scope; never a permission group.
pub(crate) const UNGROUPED_USAGE_ID: &str = "__ungrouped__";
pub(crate) async fn ungrouped_usage_users(
state: &crate::handlers::admin::request::AdminAppState<'_>,
) -> Result<Vec<aether_data::repository::users::StoredUserSummary>, crate::GatewayError> {
use aether_data::repository::users::UserExportListQuery;
let mut users = Vec::new();
let mut skip = 0;
loop {
let page = state
.list_export_users_page(&UserExportListQuery {
skip,
limit: 500,
..Default::default()
})
.await?;
let count = page.len();
if count == 0 {
break;
}
let ids = page.into_iter().map(|user| user.id).collect::<Vec<_>>();
let grouped = state
.list_user_group_memberships_by_user_ids(&ids)
.await?
.into_iter()
.map(|membership| membership.user_id)
.collect::<std::collections::BTreeSet<_>>();
let ids = ids
.into_iter()
.filter(|id| !grouped.contains(id))
.collect::<Vec<_>>();
users.extend(
state
.list_users_by_ids(&ids)
.await?
.into_iter()
.filter(|user| !user.is_deleted),
);
skip += count;
if count < 500 {
break;
}
}
Ok(users)
}
/// Current group provider policy, resolved to the provider-name dimension used by usage rollups.
/// None is unrestricted; Some(empty) deliberately matches no usage.
pub(crate) async fn usage_group_provider_names(
state: &crate::handlers::admin::request::AdminAppState<'_>,
group: &aether_data::repository::users::StoredUserGroup,
) -> Result<Option<Vec<String>>, crate::GatewayError> {
if matches!(
group.allowed_providers_mode.as_str(),
"unrestricted" | "inherit"
) {
return Ok(None);
}
if group.allowed_providers_mode != "specific" {
return Ok(Some(Vec::new()));
}
let allowed = group.allowed_providers.as_deref().unwrap_or_default();
let providers = state.list_provider_catalog_providers(false).await?;
let mut names = providers
.into_iter()
.filter(|provider| {
allowed.iter().any(|value| {
let value = value.trim();
value.eq_ignore_ascii_case(&provider.id)
|| value.eq_ignore_ascii_case(&provider.name)
|| value.eq_ignore_ascii_case(&provider.provider_type)
})
})
.map(|provider| provider.name)
.collect::<Vec<_>>();
names.sort();
names.dedup();
Ok(Some(names))
}
pub(crate) async fn resolve_usage_group_provider_names(
state: &crate::handlers::admin::request::AdminAppState<'_>,
query: Option<&str>,
) -> Result<Option<Vec<String>>, crate::GatewayError> {
let Some(id) = crate::handlers::admin::shared::query_param_value(query, "user_group_id") else {
return Ok(None);
};
if id == UNGROUPED_USAGE_ID {
return Ok(None);
}
let Some(group) = state.find_user_group_by_id(&id).await? else {
return Ok(Some(Vec::new()));
};
usage_group_provider_names(state, &group).await
}
@@ -169,6 +169,7 @@ pub(super) async fn build_admin_monitoring_system_status_response(
let today_usage = state
.summarize_usage_audits(&UsageAuditSummaryQuery {
provider_names: None,
created_from_unix_secs: today_start.timestamp().max(0) as u64,
created_until_unix_secs: now_unix_secs.saturating_add(1),
user_id: None,
@@ -110,6 +110,7 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response(
};
let current_summary = state
.summarize_usage_audits(&UsageAuditSummaryQuery {
provider_names: None,
created_from_unix_secs: current_from_unix_secs,
created_until_unix_secs: current_until_unix_secs,
..Default::default()
@@ -117,6 +118,7 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response(
.await?;
let comparison_summary = state
.summarize_usage_audits(&UsageAuditSummaryQuery {
provider_names: None,
created_from_unix_secs: comparison_from_unix_secs,
created_until_unix_secs: comparison_until_unix_secs,
..Default::default()
@@ -318,6 +320,11 @@ pub(super) async fn maybe_build_local_admin_stats_analytics_response(
};
let buckets = state
.summarize_usage_time_series(&UsageTimeSeriesQuery {
provider_names: super::super::resolve_usage_group_provider_names(
state,
request_context.query_string(),
)
.await?,
created_from_unix_secs,
created_until_unix_secs,
granularity: query_granularity,
@@ -72,6 +72,7 @@ pub(super) async fn maybe_build_local_admin_stats_cost_response(
};
let buckets = state
.summarize_usage_time_series(&UsageTimeSeriesQuery {
provider_names: None,
created_from_unix_secs,
created_until_unix_secs,
granularity: UsageTimeSeriesGranularity::Day,
@@ -78,6 +78,7 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
};
let summaries = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
provider_names: None,
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::Model,
@@ -156,6 +157,7 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
};
let summaries = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
provider_names: None,
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::ApiKey,
@@ -283,34 +285,6 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
)));
};
let summaries = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::User,
user_id: None,
user_ids: None,
provider_name: filters.provider_name,
model: filters.model,
})
.await?;
let user_ids = summaries
.iter()
.map(|item| item.group_key.clone())
.collect::<Vec<_>>();
let user_metadata = load_user_leaderboard_metadata(state, &user_ids).await?;
let user_usage = build_user_leaderboard_items_from_summaries(
&summaries,
&user_metadata,
state.has_auth_user_data_reader(),
state.has_user_data_reader(),
include_inactive,
exclude_admin,
)
.into_iter()
.map(|item| (item.id.clone(), item))
.collect::<BTreeMap<_, _>>();
let mut leaderboard = Vec::new();
let mut member_counts = BTreeMap::new();
let mut active_member_counts = BTreeMap::new();
@@ -321,13 +295,38 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
.iter()
.filter(|member| !member.is_deleted && member.is_active)
.count();
let scoped_user_ids = members
let user_ids = members
.iter()
.filter(|member| !member.is_deleted)
.filter(|member| include_inactive || member.is_active)
.filter(|member| !exclude_admin || !member.role.eq_ignore_ascii_case("admin"))
.map(|member| member.user_id.as_str())
.collect::<BTreeSet<_>>();
.map(|member| member.user_id.clone())
.collect::<Vec<_>>();
let summaries = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::User,
user_id: None,
user_ids: Some(user_ids),
provider_names: super::super::usage_group_provider_names(state, &group).await?,
provider_name: filters.provider_name.clone(),
model: filters.model.clone(),
})
.await?;
let user_ids = summaries
.iter()
.map(|row| row.group_key.clone())
.collect::<Vec<_>>();
let metadata = load_user_leaderboard_metadata(state, &user_ids).await?;
let users = build_user_leaderboard_items_from_summaries(
&summaries,
&metadata,
state.has_auth_user_data_reader(),
state.has_user_data_reader(),
include_inactive,
exclude_admin,
);
let mut item = AdminStatsLeaderboardItem {
id: group.id.clone(),
name: group.name,
@@ -335,17 +334,47 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
tokens: 0,
cost: 0.0,
};
for user_id in scoped_user_ids {
if let Some(user) = user_usage.get(user_id) {
item.requests = item.requests.saturating_add(user.requests);
item.tokens = item.tokens.saturating_add(user.tokens);
item.cost += user.cost;
}
for user in users {
item.requests = item.requests.saturating_add(user.requests);
item.tokens = item.tokens.saturating_add(user.tokens);
item.cost += user.cost;
}
member_counts.insert(group.id.clone(), member_count);
active_member_counts.insert(group.id, active_member_count);
leaderboard.push(item);
}
let ungrouped = super::super::ungrouped_usage_users(state).await?;
let id = super::super::UNGROUPED_USAGE_ID.to_string();
member_counts.insert(id.clone(), ungrouped.len());
active_member_counts.insert(
id.clone(),
ungrouped.iter().filter(|user| user.is_active).count(),
);
let user_ids = ungrouped
.into_iter()
.filter(|user| include_inactive || user.is_active)
.filter(|user| !exclude_admin || !user.role.eq_ignore_ascii_case("admin"))
.map(|user| user.id)
.collect();
let rows = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::User,
user_id: None,
user_ids: Some(user_ids),
provider_names: None,
provider_name: filters.provider_name.clone(),
model: filters.model.clone(),
})
.await?;
leaderboard.push(AdminStatsLeaderboardItem {
id,
name: "Ungrouped".to_string(),
requests: rows.iter().map(|row| row.request_count).sum(),
tokens: rows.iter().map(|row| row.total_tokens).sum(),
cost: rows.iter().map(|row| row.total_cost_usd).sum(),
});
leaderboard.sort_by(|left, right| compare_leaderboard_items(metric, order, left, right));
return Ok(Some(build_admin_stats_user_group_leaderboard_response(
@@ -422,6 +451,8 @@ pub(super) async fn maybe_build_local_admin_stats_leaderboard_response(
};
let summaries = state
.summarize_usage_leaderboard(&UsageLeaderboardQuery {
provider_names: super::super::resolve_usage_group_provider_names(state, query)
.await?,
created_from_unix_secs,
created_until_unix_secs,
group_by: UsageLeaderboardGroupBy::User,
@@ -758,6 +758,8 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
};
let summary = state
.summarize_usage_audits(&UsageAuditSummaryQuery {
provider_names: super::super::resolve_usage_group_provider_names(state, query)
.await?,
created_from_unix_secs,
created_until_unix_secs,
user_id: query_param_value(query, "user_id"),
@@ -2031,7 +2031,12 @@ async fn gateway_aggregates_admin_stats_by_current_user_group_membership() {
assert_eq!(response.status(), StatusCode::OK);
let payload: serde_json::Value = response.json().await.expect("json body should parse");
assert_eq!(payload["attribution"], "current_membership");
assert_eq!(payload["total"], 1);
assert_eq!(payload["total"], 2);
let mut payload = payload;
payload["items"]
.as_array_mut()
.unwrap()
.retain(|item| item["id"] == group.id);
assert_eq!(payload["items"][0]["id"], group.id);
assert_eq!(payload["items"][0]["name"], "Engineering");
assert_eq!(payload["items"][0]["requests"], 2);
@@ -2067,6 +2072,222 @@ async fn gateway_aggregates_admin_stats_by_current_user_group_membership() {
upstream_handle.abort();
}
#[tokio::test]
async fn gateway_group_usage_intersects_members_and_providers_with_shared_overlap() {
let rows = vec![
sample_usage_row(
"g",
"g",
Some("user-1"),
None,
None,
"Gemini",
"model",
10,
2,
0.4,
0.4,
DAY_1_UNIX_SECS,
),
sample_usage_row(
"s",
"s",
Some("user-1"),
None,
None,
"Shared",
"model",
10,
2,
0.35,
0.35,
DAY_1_UNIX_SECS,
),
sample_usage_row(
"o",
"o",
Some("user-1"),
None,
None,
"Other",
"model",
10,
2,
1.5,
1.5,
DAY_1_UNIX_SECS,
),
sample_usage_row(
"x",
"x",
Some("user-2"),
None,
None,
"Gemini",
"model",
10,
2,
10.0,
10.0,
DAY_1_UNIX_SECS,
),
];
let users = InMemoryUserReadRepository::seed_auth_users([
sample_auth_user("user-1", "alice", "user", true),
sample_auth_user("user-2", "bob", "user", true),
]);
let export_users = users.list_export_users().await.unwrap();
let users = users.with_export_users(export_users);
let mut group_ids = Vec::new();
for (name, allowed, mode) in [
(
"Gemini group",
vec!["provider-gemini", "Shared", "provider-shared"],
"specific",
),
("Other group", vec!["Other", "shared-type"], "specific"),
("Denied", vec!["Gemini"], "deny_all"),
("Empty", vec![], "specific"),
("Inherited", vec![], "inherit"),
("Unrestricted", vec![], "unrestricted"),
] {
let group = users
.create_user_group(UpsertUserGroupRecord {
name: name.to_string(),
description: None,
priority: 0,
allowed_providers: Some(allowed.into_iter().map(str::to_string).collect()),
allowed_providers_mode: mode.to_string(),
allowed_api_formats: None,
allowed_api_formats_mode: "inherit".to_string(),
allowed_models: None,
allowed_models_mode: "inherit".to_string(),
rate_limit: None,
rate_limit_mode: "inherit".to_string(),
})
.await
.unwrap()
.unwrap();
users
.replace_user_group_members(&group.id, &["user-1".to_string()])
.await
.unwrap();
group_ids.push(group.id);
}
let mut shared = sample_provider("provider-shared", "Shared", 0);
shared.provider_type = "shared-type".to_string();
let providers = InMemoryProviderCatalogReadRepository::seed(
vec![
sample_provider("provider-gemini", "Gemini", 0),
shared,
sample_provider("provider-other", "Other", 0),
],
vec![],
vec![],
);
let data = GatewayDataState::with_usage_reader_for_tests(Arc::new(
InMemoryUsageReadRepository::seed(rows),
))
.with_user_reader(Arc::new(users))
.with_provider_catalog_reader(Arc::new(providers));
let gateway = build_router_with_state(AppState::new().unwrap().with_data_state_for_tests(data));
let (url, handle) = start_server(gateway).await;
let client = reqwest::Client::new();
let range = "start_date=2024-03-21&end_date=2024-03-21&tz_offset_minutes=0";
let paths = [
(
format!("usage/stats?{range}&user_group_id=__ungrouped__"),
Some(10.0),
),
(
format!("stats/time-series?{range}&granularity=day&user_group_id=__ungrouped__"),
Some(10.0),
),
(
format!("stats/leaderboard/users?{range}&metric=cost&user_group_id=__ungrouped__"),
Some(10.0),
),
(
format!("stats/leaderboard/user-groups?{range}&metric=cost"),
None,
),
(
format!("usage/stats?{range}&user_group_id={}", group_ids[0]),
Some(0.75),
),
(
format!(
"stats/time-series?{range}&granularity=day&user_group_id={}",
group_ids[1]
),
Some(1.85),
),
(
format!(
"stats/leaderboard/users?{range}&metric=cost&user_group_id={}",
group_ids[0]
),
Some(0.75),
),
(format!("usage/stats?{range}&user_id=user-1"), Some(2.25)),
(
format!(
"usage/stats?{range}&user_group_id={}&provider=Other",
group_ids[0]
),
Some(0.0),
),
];
for (path, expected) in paths {
let response = admin_request(client.get(format!("{url}/api/admin/{path}")))
.send()
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK, "{path}");
let body: serde_json::Value = response.json().await.unwrap();
if let Some(expected) = expected {
let value = if path.starts_with("stats/time-series") {
&body[0]["total_cost"]
} else if path.starts_with("stats/leaderboard") {
&body["items"][0]["cost"]
} else {
&body["total_cost"]
};
assert!(
(value.as_f64().unwrap() - expected).abs() < 1e-9,
"{path}: {body}"
);
} else {
let items = body["items"].as_array().unwrap();
assert_eq!(items.len(), 7);
let ungrouped = items
.iter()
.find(|item| item["id"] == "__ungrouped__")
.unwrap();
assert_eq!(ungrouped["cost"], 10.0);
assert_eq!(ungrouped["member_count"], 1);
assert_eq!(ungrouped["active_member_count"], 1);
for (id, cost, requests) in [
(&group_ids[0], 0.75, 2),
(&group_ids[1], 1.85, 2),
(&group_ids[2], 0.0, 0),
(&group_ids[3], 0.0, 0),
(&group_ids[4], 2.25, 3),
(&group_ids[5], 2.25, 3),
] {
let row = items
.iter()
.find(|item| item["id"].as_str() == Some(id.as_str()))
.unwrap();
assert!((row["cost"].as_f64().unwrap() - cost).abs() < 1e-9);
assert_eq!(row["requests"], requests);
}
}
}
handle.abort();
}
#[tokio::test]
async fn gateway_handles_admin_stats_leaderboard_users_locally_without_usage_reader() {
let (upstream_url, upstream_hits, upstream_handle) =