perf(usage): 将 status 过滤下推到 SQL 查询层

为 UsageAuditListQuery 新增 statuses 字段,支持在数据库端按状态
筛选用量记录。admin usage summary 路由不再全量拉取后内存过滤,
改为直接查询 pending/streaming 状态的记录。
This commit is contained in:
fawney19
2026-04-15 09:30:42 +08:00
parent adde9ff237
commit 8827c46c33
8 changed files with 34 additions and 4 deletions

View File

@@ -25,6 +25,7 @@ pub(crate) async fn list_usage_for_range(
user_id: filters.user_id.clone(), user_id: filters.user_id.clone(),
provider_name: filters.provider_name.clone(), provider_name: filters.provider_name.clone(),
model: filters.model.clone(), model: filters.model.clone(),
statuses: None,
}) })
.await .await
} }
@@ -44,6 +45,7 @@ pub(crate) async fn list_usage_for_optional_range(
user_id: filters.user_id.clone(), user_id: filters.user_id.clone(),
provider_name: filters.provider_name.clone(), provider_name: filters.provider_name.clone(),
model: filters.model.clone(), model: filters.model.clone(),
statuses: None,
}) })
.await .await
} }

View File

@@ -16,6 +16,7 @@ pub(in super::super) async fn list_recent_completed_usage_for_cache_affinity(
user_id: user_id.map(ToOwned::to_owned), user_id: user_id.map(ToOwned::to_owned),
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await?; .await?;
items.retain(|item| item.status == "completed"); items.retain(|item| item.status == "completed");

View File

@@ -67,14 +67,23 @@ pub(super) async fn maybe_build_local_admin_usage_summary_response(
let query = request_context.request_query_string.as_deref(); let query = request_context.request_query_string.as_deref();
let requested_ids = admin_usage_parse_ids(query); let requested_ids = admin_usage_parse_ids(query);
let usage = state let list_query = if requested_ids.is_some() {
.list_usage_audits(&UsageAuditListQuery::default()) UsageAuditListQuery::default()
.await?; } else {
UsageAuditListQuery {
statuses: Some(vec![
"pending".to_string(),
"streaming".to_string(),
]),
..Default::default()
}
};
let usage = state.list_usage_audits(&list_query).await?;
let mut items: Vec<_> = usage let mut items: Vec<_> = usage
.into_iter() .into_iter()
.filter(|item| match requested_ids.as_ref() { .filter(|item| match requested_ids.as_ref() {
Some(ids) => ids.contains(&item.id), Some(ids) => ids.contains(&item.id),
None => matches!(item.status.as_str(), "pending" | "streaming"), None => true,
}) })
.collect(); .collect();
items.sort_by(|left, right| { items.sort_by(|left, right| {

View File

@@ -485,6 +485,7 @@ async fn dashboard_list_usage_for_range(
user_id: user_id.map(ToOwned::to_owned), user_id: user_id.map(ToOwned::to_owned),
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await .await
{ {
@@ -1213,6 +1214,7 @@ pub(super) async fn handle_dashboard_provider_status_get(
user_id: None, user_id: None,
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await .await
{ {

View File

@@ -879,6 +879,7 @@ pub(super) async fn handle_users_me_usage_active_get(
user_id: Some(auth.user.id.clone()), user_id: Some(auth.user.id.clone()),
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await .await
{ {
@@ -950,6 +951,7 @@ pub(super) async fn handle_users_me_usage_interval_timeline_get(
user_id: Some(auth.user.id.clone()), user_id: Some(auth.user.id.clone()),
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await .await
{ {

View File

@@ -217,6 +217,7 @@ pub(super) async fn handle_wallet_today_cost(
user_id: Some(auth.user.id.clone()), user_id: Some(auth.user.id.clone()),
provider_name: None, provider_name: None,
model: None, model: None,
statuses: None,
}) })
.await .await
{ {

View File

@@ -480,6 +480,7 @@ pub struct UsageAuditListQuery {
pub user_id: Option<String>, pub user_id: Option<String>,
pub provider_name: Option<String>, pub provider_name: Option<String>,
pub model: Option<String>, pub model: Option<String>,
pub statuses: Option<Vec<String>>,
} }
#[derive(Debug, Clone, PartialEq, Default, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, PartialEq, Default, serde::Serialize, serde::Deserialize)]

View File

@@ -1195,10 +1195,22 @@ impl SqlxUsageReadRepository {
} }
if let Some(model) = query.model.as_deref() { if let Some(model) = query.model.as_deref() {
builder.push(if has_where { " AND " } else { " WHERE " }); builder.push(if has_where { " AND " } else { " WHERE " });
has_where = true;
builder builder
.push("\"usage\".model = ") .push("\"usage\".model = ")
.push_bind(model.to_string()); .push_bind(model.to_string());
} }
if let Some(statuses) = query.statuses.as_deref() {
if !statuses.is_empty() {
builder.push(if has_where { " AND " } else { " WHERE " });
builder.push("\"usage\".status IN (");
let mut separated = builder.separated(", ");
for status in statuses {
separated.push_bind(status.to_string());
}
separated.push_unseparated(")");
}
}
builder.push(" ORDER BY \"usage\".created_at ASC, \"usage\".request_id ASC"); builder.push(" ORDER BY \"usage\".created_at ASC, \"usage\".request_id ASC");
let query = builder.build(); let query = builder.build();