Revert "feat: combine usage quota and pool stats updates"

This reverts commit 1f5b294bd0.
This commit is contained in:
fawney19
2026-05-07 01:33:40 +08:00
parent 1f5b294bd0
commit df9d30340d
35 changed files with 123 additions and 3473 deletions

View File

@@ -30,10 +30,9 @@ use super::{
api_key_usage_contribution, provider_api_key_usage_contribution,
strip_deprecated_usage_display_fields, usage_can_recover_terminal_failure,
ApiKeyUsageContribution, ApiKeyUsageDelta, ProviderApiKeyUsageContribution,
ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
StoredProviderUsageWindow, StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord,
UsageAuditListQuery, UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
};
use crate::repository::auth::InMemoryAuthApiKeySnapshotRepository;
use crate::repository::provider_catalog::InMemoryProviderCatalogReadRepository;
@@ -2328,59 +2327,6 @@ impl UsageReadRepository for InMemoryUsageReadRepository {
Ok(summaries)
}
async fn summarize_usage_by_provider_api_key_windows(
&self,
requests: &[ProviderApiKeyWindowUsageRequest],
) -> Result<Vec<StoredProviderApiKeyWindowUsageSummary>, DataLayerError> {
let usage = self.by_request_id.read().expect("usage repository lock");
let mut summaries = Vec::with_capacity(requests.len());
for request in requests {
let provider_api_key_id = request.provider_api_key_id.trim();
if provider_api_key_id.is_empty() {
return Err(DataLayerError::InvalidInput(
"provider api key window usage provider_api_key_id cannot be empty".to_string(),
));
}
let window_code = request.window_code.trim();
if window_code.is_empty() {
return Err(DataLayerError::InvalidInput(
"provider api key window usage window_code cannot be empty".to_string(),
));
}
if request.start_unix_secs >= request.end_unix_secs {
return Err(DataLayerError::InvalidInput(
"provider api key window usage range must be non-empty".to_string(),
));
}
let mut summary = StoredProviderApiKeyWindowUsageSummary {
provider_api_key_id: provider_api_key_id.to_string(),
window_code: window_code.to_string(),
..StoredProviderApiKeyWindowUsageSummary::default()
};
for item in usage.values() {
if item.provider_api_key_id.as_deref() != Some(provider_api_key_id) {
continue;
}
if item.created_at_unix_ms < request.start_unix_secs
|| item.created_at_unix_ms >= request.end_unix_secs
{
continue;
}
summary.request_count = summary.request_count.saturating_add(1);
summary.total_tokens = summary.total_tokens.saturating_add(item.total_tokens);
summary.total_cost_usd += item.total_cost_usd;
}
summaries.push(summary);
}
Ok(summaries)
}
async fn summarize_provider_usage_since(
&self,
provider_id: &str,
@@ -2988,9 +2934,8 @@ mod tests {
UsageWriteRepository,
};
use aether_data_contracts::repository::usage::{
usage_body_ref, ProviderApiKeyWindowUsageRequest, UsageAuditAggregationGroupBy,
UsageAuditAggregationQuery, UsageBodyField, UsageProviderPerformanceQuery,
UsageTimeSeriesGranularity,
usage_body_ref, UsageAuditAggregationGroupBy, UsageAuditAggregationQuery, UsageBodyField,
UsageProviderPerformanceQuery, UsageTimeSeriesGranularity,
};
use serde_json::json;
@@ -4588,44 +4533,6 @@ mod tests {
assert_eq!(item.last_used_at_unix_secs, Some(1_711_000_250));
}
#[tokio::test]
async fn summarizes_provider_api_key_window_usage_with_zero_rows() {
let repository = InMemoryUsageReadRepository::seed(vec![
sample_usage("req-1", 1_711_000_000),
sample_usage("req-2", 1_711_000_250),
]);
let usage = repository
.summarize_usage_by_provider_api_key_windows(&[
ProviderApiKeyWindowUsageRequest {
provider_api_key_id: "provider-key-1".to_string(),
window_code: "5h".to_string(),
start_unix_secs: 1_711_000_000,
end_unix_secs: 1_711_000_300,
},
ProviderApiKeyWindowUsageRequest {
provider_api_key_id: "provider-key-empty".to_string(),
window_code: "weekly".to_string(),
start_unix_secs: 1_711_000_000,
end_unix_secs: 1_711_000_300,
},
])
.await
.expect("window summary should succeed");
assert_eq!(usage.len(), 2);
assert_eq!(usage[0].provider_api_key_id, "provider-key-1");
assert_eq!(usage[0].window_code, "5h");
assert_eq!(usage[0].request_count, 2);
assert_eq!(usage[0].total_tokens, 300);
assert_eq!(usage[0].total_cost_usd, 0.24);
assert_eq!(usage[1].provider_api_key_id, "provider-key-empty");
assert_eq!(usage[1].window_code, "weekly");
assert_eq!(usage[1].request_count, 0);
assert_eq!(usage[1].total_tokens, 0);
assert_eq!(usage[1].total_cost_usd, 0.0);
}
#[tokio::test]
async fn list_usage_audits_applies_second_based_time_filters() {
let repository = InMemoryUsageReadRepository::seed(vec![

View File

@@ -319,17 +319,6 @@ macro_rules! impl_materialized_usage_read_repository {
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_by_provider_api_key_ids(&repository, provider_api_key_ids).await
}
async fn summarize_usage_by_provider_api_key_windows(
&self,
requests: &[$crate::repository::usage::ProviderApiKeyWindowUsageRequest],
) -> Result<
Vec<$crate::repository::usage::StoredProviderApiKeyWindowUsageSummary>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_by_provider_api_key_windows(&repository, requests).await
}
async fn summarize_provider_usage_since(
&self,
provider_id: &str,
@@ -363,10 +352,9 @@ mod sqlite;
#[allow(unused_imports)]
pub(crate) use aether_data_contracts::repository::usage::{
PendingUsageCleanupSummary, ProviderApiKeyWindowUsageRequest, StoredProviderApiKeyUsageSummary,
StoredProviderApiKeyWindowUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
PendingUsageCleanupSummary, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
StoredProviderUsageWindow, StoredRequestUsageAudit, StoredUsageAuditAggregation,
StoredUsageAuditSummary, StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
StoredUsageCacheAffinityIntervalRow, StoredUsageCacheHitSummary, StoredUsageCostSavingsSummary,
StoredUsageDailySummary, StoredUsageDashboardDailyBreakdownRow,
StoredUsageDashboardProviderCount, StoredUsageDashboardSummary,

View File

@@ -38,8 +38,7 @@ use super::{
api_key_usage_contribution, incoming_usage_can_recover_terminal_failure,
model_usage_contribution, provider_api_key_usage_contribution,
strip_deprecated_usage_display_fields, ApiKeyUsageDelta, ModelUsageDelta,
PendingUsageCleanupSummary, ProviderApiKeyUsageDelta, ProviderApiKeyWindowUsageRequest,
StoredProviderApiKeyUsageSummary, StoredProviderApiKeyWindowUsageSummary,
PendingUsageCleanupSummary, ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary,
StoredProviderUsageSummary, StoredRequestUsageAudit, StoredUsageDailySummary,
UpsertUsageRecord, UsageAuditListQuery, UsageDailyHeatmapQuery, UsageReadRepository,
UsageWriteRepository,
@@ -1315,9 +1314,6 @@ const SUMMARIZE_USAGE_TOTALS_BY_USER_IDS_SQL: &str =
const SUMMARIZE_USAGE_BY_PROVIDER_API_KEY_IDS_SQL: &str =
include_str!("queries/summarize_usage_by_provider_api_key_ids_sql.sql");
const SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL: &str =
include_str!("queries/summarize_provider_api_key_window_usage_sql.sql");
const APPLY_API_KEY_USAGE_DELTA_SQL: &str =
include_str!("queries/apply_api_key_usage_delta_sql.sql");
@@ -7229,98 +7225,6 @@ ORDER BY "usage".user_id ASC
Ok(summaries)
}
pub async fn summarize_usage_by_provider_api_key_windows(
&self,
requests: &[ProviderApiKeyWindowUsageRequest],
) -> Result<Vec<StoredProviderApiKeyWindowUsageSummary>, DataLayerError> {
if requests.is_empty() {
return Ok(Vec::new());
}
let mut provider_api_key_ids = Vec::with_capacity(requests.len());
let mut window_codes = Vec::with_capacity(requests.len());
let mut start_unix_secs = Vec::with_capacity(requests.len());
let mut end_unix_secs = Vec::with_capacity(requests.len());
for request in requests {
let provider_api_key_id = request.provider_api_key_id.trim();
if provider_api_key_id.is_empty() {
return Err(DataLayerError::InvalidInput(
"provider api key window usage provider_api_key_id cannot be empty".to_string(),
));
}
let window_code = request.window_code.trim();
if window_code.is_empty() {
return Err(DataLayerError::InvalidInput(
"provider api key window usage window_code cannot be empty".to_string(),
));
}
if request.start_unix_secs >= request.end_unix_secs {
return Err(DataLayerError::InvalidInput(
"provider api key window usage range must be non-empty".to_string(),
));
}
provider_api_key_ids.push(provider_api_key_id.to_string());
window_codes.push(window_code.to_string());
start_unix_secs.push(i64::try_from(request.start_unix_secs).map_err(|_| {
DataLayerError::InvalidInput(
"provider api key window usage start_unix_secs is out of range".to_string(),
)
})?);
end_unix_secs.push(i64::try_from(request.end_unix_secs).map_err(|_| {
DataLayerError::InvalidInput(
"provider api key window usage end_unix_secs is out of range".to_string(),
)
})?);
}
let mut rows = sqlx::query(SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL)
.bind(&provider_api_key_ids)
.bind(&window_codes)
.bind(&start_unix_secs)
.bind(&end_unix_secs)
.fetch(&self.pool);
let mut summaries = Vec::new();
while let Some(row) = rows.try_next().await.map_postgres_err()? {
let total_cost_usd = row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?;
if !total_cost_usd.is_finite() {
return Err(DataLayerError::UnexpectedValue(
"usage.total_cost_usd window aggregate is not finite".to_string(),
));
}
summaries.push(StoredProviderApiKeyWindowUsageSummary {
provider_api_key_id: row
.try_get::<String, _>("provider_api_key_id")
.map_postgres_err()?,
window_code: row.try_get::<String, _>("window_code").map_postgres_err()?,
request_count: row
.try_get::<i64, _>("request_count")
.map_postgres_err()?
.try_into()
.map_err(|_| {
DataLayerError::UnexpectedValue(
"usage.request_count window aggregate is negative".to_string(),
)
})?,
total_tokens: row
.try_get::<i64, _>("total_tokens")
.map_postgres_err()?
.try_into()
.map_err(|_| {
DataLayerError::UnexpectedValue(
"usage.total_tokens window aggregate is negative".to_string(),
)
})?,
total_cost_usd,
});
}
Ok(summaries)
}
pub async fn upsert(
&self,
usage: UpsertUsageRecord,
@@ -8140,13 +8044,6 @@ impl UsageReadRepository for SqlxUsageReadRepository {
Self::summarize_usage_by_provider_api_key_ids(self, provider_api_key_ids).await
}
async fn summarize_usage_by_provider_api_key_windows(
&self,
requests: &[ProviderApiKeyWindowUsageRequest],
) -> Result<Vec<StoredProviderApiKeyWindowUsageSummary>, DataLayerError> {
Self::summarize_usage_by_provider_api_key_windows(self, requests).await
}
async fn summarize_provider_usage_since(
&self,
provider_id: &str,

View File

@@ -1,36 +0,0 @@
WITH requested AS (
SELECT
request_row.provider_api_key_id,
request_row.window_code,
request_row.start_unix_secs,
request_row.end_unix_secs,
request_row.ordinality
FROM UNNEST(
$1::TEXT[],
$2::TEXT[],
$3::BIGINT[],
$4::BIGINT[]
) WITH ORDINALITY AS request_row(
provider_api_key_id,
window_code,
start_unix_secs,
end_unix_secs,
ordinality
)
)
SELECT
requested.provider_api_key_id,
requested.window_code,
COUNT("usage".id)::BIGINT AS request_count,
COALESCE(SUM("usage".total_tokens), 0)::BIGINT AS total_tokens,
CAST(COALESCE(SUM("usage".total_cost_usd), 0) AS DOUBLE PRECISION) AS total_cost_usd
FROM requested
LEFT JOIN usage_billing_facts AS "usage"
ON "usage".provider_api_key_id = requested.provider_api_key_id
AND "usage".created_at >= to_timestamp(requested.start_unix_secs::DOUBLE PRECISION)
AND "usage".created_at < to_timestamp(requested.end_unix_secs::DOUBLE PRECISION)
GROUP BY
requested.provider_api_key_id,
requested.window_code,
requested.ordinality
ORDER BY requested.ordinality ASC

View File

@@ -234,21 +234,6 @@ fn usage_sql_summarizes_usage_by_provider_api_key_ids_in_database() {
assert!(super::SUMMARIZE_USAGE_BY_PROVIDER_API_KEY_IDS_SQL.contains("ANY($1::TEXT[])"));
}
#[test]
fn usage_sql_summarizes_provider_key_window_usage_from_billing_facts() {
assert!(super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL.contains("UNNEST"));
assert!(super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL
.contains("LEFT JOIN usage_billing_facts AS \"usage\""));
assert!(
super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL.contains("created_at >= to_timestamp")
);
assert!(
super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL.contains("created_at < to_timestamp")
);
assert!(super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL
.contains("COUNT(\"usage\".id)::BIGINT AS request_count"));
}
#[test]
fn usage_sql_serializes_request_id_upserts_before_reading_previous_usage() {
assert!(super::LOCK_USAGE_REQUEST_ID_SQL.contains("pg_advisory_xact_lock"));
@@ -470,8 +455,6 @@ fn usage_sql_raw_aggregates_use_canonical_billing_facts() {
.contains("FROM usage_billing_facts AS \"usage\""));
assert!(super::SUMMARIZE_USAGE_TOTALS_BY_USER_IDS_SQL
.contains("FROM usage_billing_facts AS \"usage\""));
assert!(super::SUMMARIZE_PROVIDER_API_KEY_WINDOW_USAGE_SQL
.contains("usage_billing_facts AS \"usage\""));
}
#[test]