mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-02 01:10:23 +08:00
feat(admin): 将用户管理的用量统计改为全量累计口径 (#314)
- 新增按用户 ID 批量汇总 usage 的后端查询 - 在 /api/admin/users 中返回 request_count 和 total_tokens - 移除用户管理页面额外的 usage 聚合请求 - 为管理端用户列表补充累计统计回归测试
This commit is contained in:
@@ -9,15 +9,16 @@ use aether_data_contracts::repository::usage::{
|
||||
StoredUsageDashboardDailyBreakdownRow, StoredUsageDashboardProviderCount,
|
||||
StoredUsageDashboardSummary, StoredUsageErrorDistributionRow, StoredUsageLeaderboardSummary,
|
||||
StoredUsagePerformancePercentilesRow, StoredUsageSettledCostSummary,
|
||||
StoredUsageTimeSeriesBucket, UsageAuditAggregationGroupBy, UsageAuditAggregationQuery,
|
||||
UsageAuditKeywordSearchQuery, UsageAuditSummaryQuery, UsageBodyField, UsageBreakdownGroupBy,
|
||||
UsageBreakdownSummaryQuery, UsageCacheAffinityHitSummaryQuery,
|
||||
UsageCacheAffinityIntervalGroupBy, UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery,
|
||||
UsageCostSavingsSummaryQuery, UsageDashboardDailyBreakdownQuery,
|
||||
UsageDashboardProviderCountsQuery, UsageDashboardSummaryQuery, UsageErrorDistributionQuery,
|
||||
UsageLeaderboardGroupBy, UsageLeaderboardQuery, UsageMonitoringErrorCountQuery,
|
||||
UsageMonitoringErrorListQuery, UsagePerformancePercentilesQuery, UsageSettledCostSummaryQuery,
|
||||
UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
|
||||
StoredUsageTimeSeriesBucket, StoredUsageUserTotals, UsageAuditAggregationGroupBy,
|
||||
UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery, UsageAuditSummaryQuery,
|
||||
UsageBodyField, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
|
||||
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
|
||||
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCostSavingsSummaryQuery,
|
||||
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
|
||||
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
|
||||
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,
|
||||
UsagePerformancePercentilesQuery, UsageSettledCostSummaryQuery, UsageTimeSeriesGranularity,
|
||||
UsageTimeSeriesQuery,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use chrono::Utc;
|
||||
@@ -1853,6 +1854,46 @@ impl UsageReadRepository for InMemoryUsageReadRepository {
|
||||
Ok(totals)
|
||||
}
|
||||
|
||||
async fn summarize_usage_totals_by_user_ids(
|
||||
&self,
|
||||
user_ids: &[String],
|
||||
) -> Result<Vec<StoredUsageUserTotals>, DataLayerError> {
|
||||
let user_id_set = user_ids
|
||||
.iter()
|
||||
.cloned()
|
||||
.collect::<std::collections::BTreeSet<_>>();
|
||||
let mut totals = BTreeMap::<String, StoredUsageUserTotals>::new();
|
||||
for item in self
|
||||
.by_request_id
|
||||
.read()
|
||||
.expect("usage repository lock")
|
||||
.values()
|
||||
{
|
||||
if matches!(item.status.as_str(), "pending" | "streaming")
|
||||
|| matches!(item.provider_name.as_str(), "unknown" | "pending")
|
||||
{
|
||||
continue;
|
||||
}
|
||||
let Some(user_id) = item.user_id.as_deref() else {
|
||||
continue;
|
||||
};
|
||||
if !user_id_set.contains(user_id) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let entry =
|
||||
totals
|
||||
.entry(user_id.to_string())
|
||||
.or_insert_with(|| StoredUsageUserTotals {
|
||||
user_id: user_id.to_string(),
|
||||
..Default::default()
|
||||
});
|
||||
entry.request_count = entry.request_count.saturating_add(1);
|
||||
entry.total_tokens = entry.total_tokens.saturating_add(item.total_tokens);
|
||||
}
|
||||
Ok(totals.into_values().collect())
|
||||
}
|
||||
|
||||
async fn summarize_usage_by_provider_api_key_ids(
|
||||
&self,
|
||||
provider_api_key_ids: &[String],
|
||||
|
||||
@@ -11,9 +11,9 @@ pub(crate) use aether_data_contracts::repository::usage::{
|
||||
StoredUsageDashboardProviderCount, StoredUsageDashboardSummary,
|
||||
StoredUsageErrorDistributionRow, StoredUsageLeaderboardSummary,
|
||||
StoredUsagePerformancePercentilesRow, StoredUsageSettledCostSummary,
|
||||
StoredUsageTimeSeriesBucket, UpsertUsageRecord, UsageAuditAggregationGroupBy,
|
||||
UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery, UsageAuditListQuery,
|
||||
UsageAuditSummaryQuery, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
|
||||
StoredUsageTimeSeriesBucket, StoredUsageUserTotals, UpsertUsageRecord,
|
||||
UsageAuditAggregationGroupBy, UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery,
|
||||
UsageAuditListQuery, UsageAuditSummaryQuery, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
|
||||
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
|
||||
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCostSavingsSummaryQuery,
|
||||
UsageDailyHeatmapQuery, UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
|
||||
|
||||
@@ -5,15 +5,16 @@ use aether_data_contracts::repository::usage::{
|
||||
StoredUsageDashboardDailyBreakdownRow, StoredUsageDashboardProviderCount,
|
||||
StoredUsageDashboardSummary, StoredUsageErrorDistributionRow, StoredUsageLeaderboardSummary,
|
||||
StoredUsagePerformancePercentilesRow, StoredUsageSettledCostSummary,
|
||||
StoredUsageTimeSeriesBucket, UsageAuditAggregationGroupBy, UsageAuditAggregationQuery,
|
||||
UsageAuditKeywordSearchQuery, UsageAuditSummaryQuery, UsageBodyCaptureState, UsageBodyField,
|
||||
UsageBreakdownGroupBy, UsageBreakdownSummaryQuery, UsageCacheAffinityHitSummaryQuery,
|
||||
UsageCacheAffinityIntervalGroupBy, UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery,
|
||||
UsageCostSavingsSummaryQuery, UsageDashboardDailyBreakdownQuery,
|
||||
UsageDashboardProviderCountsQuery, UsageDashboardSummaryQuery, UsageErrorDistributionQuery,
|
||||
UsageLeaderboardGroupBy, UsageLeaderboardQuery, UsageMonitoringErrorCountQuery,
|
||||
UsageMonitoringErrorListQuery, UsagePerformancePercentilesQuery, UsageSettledCostSummaryQuery,
|
||||
UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
|
||||
StoredUsageTimeSeriesBucket, StoredUsageUserTotals, UsageAuditAggregationGroupBy,
|
||||
UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery, UsageAuditSummaryQuery,
|
||||
UsageBodyCaptureState, UsageBodyField, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
|
||||
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
|
||||
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCostSavingsSummaryQuery,
|
||||
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
|
||||
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
|
||||
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,
|
||||
UsagePerformancePercentilesQuery, UsageSettledCostSummaryQuery, UsageTimeSeriesGranularity,
|
||||
UsageTimeSeriesQuery,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use flate2::{read::GzDecoder, write::GzEncoder, Compression};
|
||||
@@ -546,6 +547,19 @@ GROUP BY api_key_id
|
||||
ORDER BY api_key_id ASC
|
||||
"#;
|
||||
|
||||
const SUMMARIZE_USAGE_TOTALS_BY_USER_IDS_SQL: &str = r#"
|
||||
SELECT
|
||||
"usage".user_id,
|
||||
COUNT(*)::BIGINT AS request_count,
|
||||
COALESCE(SUM(GREATEST(COALESCE("usage".total_tokens, 0), 0)), 0)::BIGINT AS total_tokens
|
||||
FROM "usage"
|
||||
WHERE "usage".user_id = ANY($1::TEXT[])
|
||||
AND "usage".status NOT IN ('pending', 'streaming')
|
||||
AND "usage".provider_name NOT IN ('unknown', 'pending')
|
||||
GROUP BY "usage".user_id
|
||||
ORDER BY "usage".user_id ASC
|
||||
"#;
|
||||
|
||||
const SUMMARIZE_USAGE_BY_PROVIDER_API_KEY_IDS_SQL: &str = r#"
|
||||
SELECT
|
||||
provider_api_key_id,
|
||||
@@ -4004,6 +4018,35 @@ WHERE "usage".created_at >= TO_TIMESTAMP($1::double precision)"#,
|
||||
Ok(totals)
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_totals_by_user_ids(
|
||||
&self,
|
||||
user_ids: &[String],
|
||||
) -> Result<Vec<StoredUsageUserTotals>, DataLayerError> {
|
||||
if user_ids.is_empty() {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let mut rows = sqlx::query(SUMMARIZE_USAGE_TOTALS_BY_USER_IDS_SQL)
|
||||
.bind(user_ids)
|
||||
.fetch(&self.pool);
|
||||
|
||||
let mut items = Vec::new();
|
||||
while let Some(row) = rows.try_next().await.map_postgres_err()? {
|
||||
items.push(StoredUsageUserTotals {
|
||||
user_id: row.try_get::<String, _>("user_id").map_postgres_err()?,
|
||||
request_count: row
|
||||
.try_get::<i64, _>("request_count")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
total_tokens: row
|
||||
.try_get::<i64, _>("total_tokens")
|
||||
.map_postgres_err()?
|
||||
.max(0) as u64,
|
||||
});
|
||||
}
|
||||
Ok(items)
|
||||
}
|
||||
|
||||
pub async fn summarize_usage_by_provider_api_key_ids(
|
||||
&self,
|
||||
provider_api_key_ids: &[String],
|
||||
@@ -4655,6 +4698,13 @@ impl UsageReadRepository for SqlxUsageReadRepository {
|
||||
Self::summarize_total_tokens_by_api_key_ids(self, api_key_ids).await
|
||||
}
|
||||
|
||||
async fn summarize_usage_totals_by_user_ids(
|
||||
&self,
|
||||
user_ids: &[String],
|
||||
) -> Result<Vec<StoredUsageUserTotals>, DataLayerError> {
|
||||
Self::summarize_usage_totals_by_user_ids(self, user_ids).await
|
||||
}
|
||||
|
||||
async fn summarize_usage_by_provider_api_key_ids(
|
||||
&self,
|
||||
provider_api_key_ids: &[String],
|
||||
|
||||
Reference in New Issue
Block a user