chore: resolve pr 377 checks

This commit is contained in:
fawney19
2026-05-06 02:22:32 +08:00
502 changed files with 89553 additions and 21893 deletions

View File

@@ -1,11 +1,360 @@
macro_rules! impl_materialized_usage_read_repository {
($repository:ty) => {
#[async_trait::async_trait]
impl $crate::repository::usage::UsageReadRepository for $repository {
async fn find_by_id(
&self,
id: &str,
) -> Result<
Option<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::find_by_id(&repository, id).await
}
async fn list_by_ids(
&self,
ids: &[String],
) -> Result<
Vec<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_by_ids(&repository, ids).await
}
async fn find_by_request_id(
&self,
request_id: &str,
) -> Result<
Option<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::find_by_request_id(&repository, request_id).await
}
async fn resolve_body_ref(
&self,
body_ref: &str,
) -> Result<Option<serde_json::Value>, $crate::DataLayerError> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::resolve_body_ref(&repository, body_ref).await
}
async fn list_usage_audits(
&self,
query: &$crate::repository::usage::UsageAuditListQuery,
) -> Result<
Vec<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_usage_audits(&repository, query).await
}
async fn count_usage_audits(
&self,
query: &$crate::repository::usage::UsageAuditListQuery,
) -> Result<u64, $crate::DataLayerError> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::count_usage_audits(&repository, query).await
}
async fn list_usage_audits_by_keyword_search(
&self,
query: &$crate::repository::usage::UsageAuditKeywordSearchQuery,
) -> Result<
Vec<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_usage_audits_by_keyword_search(&repository, query).await
}
async fn count_usage_audits_by_keyword_search(
&self,
query: &$crate::repository::usage::UsageAuditKeywordSearchQuery,
) -> Result<u64, $crate::DataLayerError> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::count_usage_audits_by_keyword_search(&repository, query).await
}
async fn aggregate_usage_audits(
&self,
query: &$crate::repository::usage::UsageAuditAggregationQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageAuditAggregation>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::aggregate_usage_audits(&repository, query).await
}
async fn summarize_usage_audits(
&self,
query: &$crate::repository::usage::UsageAuditSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageAuditSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_audits(&repository, query).await
}
async fn summarize_usage_totals_by_user_ids(
&self,
user_ids: &[String],
) -> Result<Vec<$crate::repository::usage::StoredUsageUserTotals>, $crate::DataLayerError>
{
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_totals_by_user_ids(&repository, user_ids).await
}
async fn summarize_usage_cache_hit_summary(
&self,
query: &$crate::repository::usage::UsageCacheHitSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageCacheHitSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_cache_hit_summary(&repository, query).await
}
async fn summarize_usage_settled_cost(
&self,
query: &$crate::repository::usage::UsageSettledCostSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageSettledCostSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_settled_cost(&repository, query).await
}
async fn summarize_usage_cache_affinity_hit_summary(
&self,
query: &$crate::repository::usage::UsageCacheAffinityHitSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageCacheAffinityHitSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_cache_affinity_hit_summary(&repository, query).await
}
async fn list_usage_cache_affinity_intervals(
&self,
query: &$crate::repository::usage::UsageCacheAffinityIntervalQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageCacheAffinityIntervalRow>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_usage_cache_affinity_intervals(&repository, query).await
}
async fn summarize_dashboard_usage(
&self,
query: &$crate::repository::usage::UsageDashboardSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageDashboardSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_dashboard_usage(&repository, query).await
}
async fn list_dashboard_daily_breakdown(
&self,
query: &$crate::repository::usage::UsageDashboardDailyBreakdownQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageDashboardDailyBreakdownRow>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_dashboard_daily_breakdown(&repository, query).await
}
async fn summarize_dashboard_provider_counts(
&self,
query: &$crate::repository::usage::UsageDashboardProviderCountsQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageDashboardProviderCount>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_dashboard_provider_counts(&repository, query).await
}
async fn summarize_usage_breakdown(
&self,
query: &$crate::repository::usage::UsageBreakdownSummaryQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageBreakdownSummaryRow>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_breakdown(&repository, query).await
}
async fn count_monitoring_usage_errors(
&self,
query: &$crate::repository::usage::UsageMonitoringErrorCountQuery,
) -> Result<u64, $crate::DataLayerError> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::count_monitoring_usage_errors(&repository, query).await
}
async fn list_monitoring_usage_errors(
&self,
query: &$crate::repository::usage::UsageMonitoringErrorListQuery,
) -> Result<
Vec<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_monitoring_usage_errors(&repository, query).await
}
async fn summarize_usage_error_distribution(
&self,
query: &$crate::repository::usage::UsageErrorDistributionQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageErrorDistributionRow>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_error_distribution(&repository, query).await
}
async fn summarize_usage_performance_percentiles(
&self,
query: &$crate::repository::usage::UsagePerformancePercentilesQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsagePerformancePercentilesRow>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_performance_percentiles(&repository, query).await
}
async fn summarize_usage_provider_performance(
&self,
query: &$crate::repository::usage::UsageProviderPerformanceQuery,
) -> Result<
$crate::repository::usage::StoredUsageProviderPerformance,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_provider_performance(&repository, query).await
}
async fn summarize_usage_cost_savings(
&self,
query: &$crate::repository::usage::UsageCostSavingsSummaryQuery,
) -> Result<
$crate::repository::usage::StoredUsageCostSavingsSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_cost_savings(&repository, query).await
}
async fn summarize_usage_time_series(
&self,
query: &$crate::repository::usage::UsageTimeSeriesQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageTimeSeriesBucket>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_time_series(&repository, query).await
}
async fn summarize_usage_leaderboard(
&self,
query: &$crate::repository::usage::UsageLeaderboardQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageLeaderboardSummary>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_leaderboard(&repository, query).await
}
async fn list_recent_usage_audits(
&self,
user_id: Option<&str>,
limit: usize,
) -> Result<
Vec<$crate::repository::usage::StoredRequestUsageAudit>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::list_recent_usage_audits(&repository, user_id, limit).await
}
async fn summarize_total_tokens_by_api_key_ids(
&self,
api_key_ids: &[String],
) -> Result<std::collections::BTreeMap<String, u64>, $crate::DataLayerError> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_total_tokens_by_api_key_ids(&repository, api_key_ids).await
}
async fn summarize_usage_by_provider_api_key_ids(
&self,
provider_api_key_ids: &[String],
) -> Result<
std::collections::BTreeMap<
String,
$crate::repository::usage::StoredProviderApiKeyUsageSummary,
>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$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_provider_usage_since(
&self,
provider_id: &str,
since_unix_secs: u64,
) -> Result<
$crate::repository::usage::StoredProviderUsageSummary,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_provider_usage_since(&repository, provider_id, since_unix_secs).await
}
async fn summarize_usage_daily_heatmap(
&self,
query: &$crate::repository::usage::UsageDailyHeatmapQuery,
) -> Result<
Vec<$crate::repository::usage::StoredUsageDailySummary>,
$crate::DataLayerError,
> {
let repository = self.materialize_read_model().await?;
<$crate::repository::usage::InMemoryUsageReadRepository as $crate::repository::usage::UsageReadRepository>::summarize_usage_daily_heatmap(&repository, query).await
}
}
};
}
mod memory;
mod sql;
mod mysql;
mod postgres;
mod sqlite;
#[allow(unused_imports)]
pub(crate) use aether_data_contracts::repository::usage::{
StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary, StoredProviderUsageWindow,
StoredRequestUsageAudit, StoredUsageAuditAggregation, StoredUsageAuditSummary,
StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
PendingUsageCleanupSummary, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
StoredProviderUsageWindow, StoredRequestUsageAudit, StoredUsageAuditAggregation,
StoredUsageAuditSummary, StoredUsageBreakdownSummaryRow, StoredUsageCacheAffinityHitSummary,
StoredUsageCacheAffinityIntervalRow, StoredUsageCacheHitSummary, StoredUsageCostSavingsSummary,
StoredUsageDailySummary, StoredUsageDashboardDailyBreakdownRow,
StoredUsageDashboardProviderCount, StoredUsageDashboardSummary,
@@ -17,16 +366,22 @@ pub(crate) use aether_data_contracts::repository::usage::{
UsageAuditAggregationGroupBy, UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery,
UsageAuditListQuery, UsageAuditSummaryQuery, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCostSavingsSummaryQuery,
UsageDailyHeatmapQuery, UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupSummary,
UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDailyHeatmapQuery,
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,
UsagePerformancePercentilesQuery, UsageProviderPerformanceQuery, UsageReadRepository,
UsageRepository, UsageSettledCostSummaryQuery, UsageTimeSeriesGranularity,
UsageTimeSeriesQuery, UsageWriteRepository,
};
pub mod cleanup {
pub use super::postgres::cleanup::*;
}
pub use memory::InMemoryUsageReadRepository;
pub use sql::SqlxUsageReadRepository;
pub use mysql::{MysqlUsageReadRepository, MysqlUsageWriteRepository};
pub use postgres::SqlxUsageReadRepository;
pub use sqlite::{SqliteUsageReadRepository, SqliteUsageWriteRepository};
#[derive(Debug, Clone, PartialEq, Default)]
pub(crate) struct ApiKeyUsageContribution {

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -11,12 +11,12 @@ use aether_data_contracts::repository::usage::{
UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery, UsageAuditSummaryQuery,
UsageBodyCaptureState, UsageBodyField, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCostSavingsSummaryQuery,
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,
UsagePerformancePercentilesQuery, UsageProviderPerformanceQuery, UsageSettledCostSummaryQuery,
UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupSummary,
UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDashboardDailyBreakdownQuery,
UsageDashboardProviderCountsQuery, UsageDashboardSummaryQuery, UsageErrorDistributionQuery,
UsageLeaderboardGroupBy, UsageLeaderboardQuery, UsageMonitoringErrorCountQuery,
UsageMonitoringErrorListQuery, UsagePerformancePercentilesQuery, UsageProviderPerformanceQuery,
UsageSettledCostSummaryQuery, UsageTimeSeriesGranularity, UsageTimeSeriesQuery,
};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
@@ -38,16 +38,19 @@ 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,
ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary, StoredProviderUsageSummary,
StoredRequestUsageAudit, StoredUsageDailySummary, UpsertUsageRecord, UsageAuditListQuery,
UsageDailyHeatmapQuery, UsageReadRepository, UsageWriteRepository,
PendingUsageCleanupSummary, ProviderApiKeyUsageDelta, StoredProviderApiKeyUsageSummary,
StoredProviderUsageSummary, StoredRequestUsageAudit, StoredUsageDailySummary,
UpsertUsageRecord, UsageAuditListQuery, UsageDailyHeatmapQuery, UsageReadRepository,
UsageWriteRepository,
};
use crate::postgres::PostgresTransactionRunner;
use crate::driver::postgres::PostgresTransactionRunner;
use crate::{
error::{postgres_error, SqlxResultExt},
DataLayerError,
};
pub mod cleanup;
// Legacy inline body columns on public.usage are deprecated. Keep the threshold at zero so
// newly captured bodies always spill to usage_body_blobs and resolve through usage_http_audits.
const MAX_INLINE_USAGE_BODY_BYTES: usize = 0;
@@ -557,67 +560,6 @@ fn decode_usage_breakdown_summary_row(
})
}
fn decode_usage_audit_aggregation_row(
row: &PgRow,
) -> Result<StoredUsageAuditAggregation, DataLayerError> {
Ok(StoredUsageAuditAggregation {
group_key: row.try_get::<String, _>("group_key").map_postgres_err()?,
display_name: row
.try_get::<Option<String>, _>("display_name")
.map_postgres_err()?,
secondary_name: row
.try_get::<Option<String>, _>("secondary_name")
.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,
output_tokens: row
.try_get::<i64, _>("output_tokens")
.map_postgres_err()?
.max(0) as u64,
effective_input_tokens: row
.try_get::<i64, _>("effective_input_tokens")
.map_postgres_err()?
.max(0) as u64,
total_input_context: row
.try_get::<i64, _>("total_input_context")
.map_postgres_err()?
.max(0) as u64,
cache_creation_tokens: row
.try_get::<i64, _>("cache_creation_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_5m_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_1h_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_read_tokens: row
.try_get::<i64, _>("cache_read_tokens")
.map_postgres_err()?
.max(0) as u64,
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
actual_total_cost_usd: row
.try_get::<f64, _>("actual_total_cost_usd")
.map_postgres_err()?,
avg_response_time_ms: row
.try_get::<Option<f64>, _>("avg_response_time_ms")
.map_postgres_err()?,
success_count: row
.try_get::<Option<i64>, _>("success_count")
.map_postgres_err()?
.map(|value| value.max(0) as u64),
})
}
fn absorb_usage_audit_summary(target: &mut StoredUsageAuditSummary, row: StoredUsageAuditSummary) {
target.total_requests = target.total_requests.saturating_add(row.total_requests);
target.input_tokens = target.input_tokens.saturating_add(row.input_tokens);
@@ -883,6 +825,67 @@ fn decode_usage_audit_summary_row(row: &PgRow) -> Result<StoredUsageAuditSummary
})
}
fn decode_usage_audit_aggregation_row(
row: &PgRow,
) -> Result<StoredUsageAuditAggregation, DataLayerError> {
Ok(StoredUsageAuditAggregation {
group_key: row.try_get::<String, _>("group_key").map_postgres_err()?,
display_name: row
.try_get::<Option<String>, _>("display_name")
.map_postgres_err()?,
secondary_name: row
.try_get::<Option<String>, _>("secondary_name")
.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,
output_tokens: row
.try_get::<i64, _>("output_tokens")
.map_postgres_err()?
.max(0) as u64,
effective_input_tokens: row
.try_get::<i64, _>("effective_input_tokens")
.map_postgres_err()?
.max(0) as u64,
total_input_context: row
.try_get::<i64, _>("total_input_context")
.map_postgres_err()?
.max(0) as u64,
cache_creation_tokens: row
.try_get::<i64, _>("cache_creation_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_5m_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_1h_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_read_tokens: row
.try_get::<i64, _>("cache_read_tokens")
.map_postgres_err()?
.max(0) as u64,
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
actual_total_cost_usd: row
.try_get::<f64, _>("actual_total_cost_usd")
.map_postgres_err()?,
avg_response_time_ms: row
.try_get::<Option<f64>, _>("avg_response_time_ms")
.map_postgres_err()?,
success_count: row
.try_get::<Option<i64>, _>("success_count")
.map_postgres_err()?
.map(|value| value.max(0) as u64),
})
}
fn decode_usage_error_distribution_row(
row: &PgRow,
) -> Result<StoredUsageErrorDistributionRow, DataLayerError> {
@@ -1211,6 +1214,9 @@ const SUMMARIZE_USAGE_BY_PROVIDER_API_KEY_IDS_SQL: &str =
const APPLY_API_KEY_USAGE_DELTA_SQL: &str =
include_str!("queries/apply_api_key_usage_delta_sql.sql");
const APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL: &str =
include_str!("queries/apply_global_model_usage_delta_sql.sql");
const RESET_API_KEY_USAGE_STATS_SQL: &str =
include_str!("queries/reset_api_key_usage_stats_sql.sql");
@@ -1220,9 +1226,6 @@ const REBUILD_API_KEY_USAGE_STATS_SQL: &str =
const APPLY_PROVIDER_API_KEY_USAGE_DELTA_SQL: &str =
include_str!("queries/apply_provider_api_key_usage_delta_sql.sql");
const APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL: &str =
include_str!("queries/apply_global_model_usage_delta_sql.sql");
const RESET_PROVIDER_API_KEY_USAGE_STATS_SQL: &str =
include_str!("queries/reset_provider_api_key_usage_stats_sql.sql");
@@ -1323,6 +1326,99 @@ const LIST_RECENT_USAGE_AUDITS_PREFIX: &str =
const UPSERT_SQL: &str = include_str!("queries/upsert_sql.sql");
const SELECT_STALE_PENDING_USAGE_BATCH_SQL: &str = r#"
SELECT
usage.request_id,
usage.status,
COALESCE(usage_settlement_snapshots.billing_status, usage.billing_status) AS billing_status
FROM usage
LEFT JOIN usage_settlement_snapshots
ON usage_settlement_snapshots.request_id = usage.request_id
WHERE usage.status IN ('pending', 'streaming')
AND usage.created_at < $1
ORDER BY usage.created_at ASC, usage.request_id ASC
LIMIT $2
FOR UPDATE OF usage SKIP LOCKED
"#;
const SELECT_COMPLETED_PENDING_REQUEST_IDS_SQL: &str = r#"
SELECT DISTINCT request_id
FROM request_candidates
WHERE request_id = ANY($1)
AND (
status = 'streaming'
OR (
status = 'success'
AND COALESCE(extra_data->>'stream_completed', 'false') = 'true'
)
)
"#;
const UPDATE_RECOVERED_STALE_USAGE_SQL: &str = r#"
UPDATE usage
SET status = 'completed',
status_code = 200,
error_message = NULL
WHERE request_id = $1
"#;
const UPDATE_FAILED_STALE_USAGE_SQL: &str = r#"
UPDATE usage
SET status = 'failed',
status_code = 504,
error_message = $2
WHERE request_id = $1
"#;
const UPDATE_FAILED_VOID_STALE_USAGE_SQL: &str = r#"
WITH updated_usage AS (
UPDATE usage
SET status = 'failed',
status_code = 504,
error_message = $2,
billing_status = 'void',
finalized_at = $3,
total_cost_usd = 0,
request_cost_usd = 0,
actual_total_cost_usd = 0,
actual_request_cost_usd = 0
WHERE request_id = $1
RETURNING request_id
)
INSERT INTO usage_settlement_snapshots (
request_id,
billing_status,
finalized_at
)
SELECT request_id, 'void', $3
FROM updated_usage
ON CONFLICT (request_id)
DO UPDATE SET
billing_status = EXCLUDED.billing_status,
finalized_at = COALESCE(
usage_settlement_snapshots.finalized_at,
EXCLUDED.finalized_at
),
updated_at = NOW()
"#;
const UPDATE_RECOVERED_STREAMING_CANDIDATES_SQL: &str = r#"
UPDATE request_candidates
SET status = 'success',
finished_at = $2
WHERE request_id = $1
AND status = 'streaming'
"#;
const UPDATE_FAILED_PENDING_CANDIDATES_SQL: &str = r#"
UPDATE request_candidates
SET status = 'failed',
finished_at = $2,
error_message = ''
WHERE request_id = $1
AND status IN ('pending', 'streaming')
"#;
#[derive(Debug, Clone)]
pub struct SqlxUsageReadRepository {
pool: PgPool,
@@ -1398,7 +1494,7 @@ SELECT
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(input_tokens + output_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(cache_creation_tokens), 0)::BIGINT AS cache_creation_tokens,
COALESCE(SUM(cache_read_tokens), 0)::BIGINT AS cache_read_tokens,
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
@@ -1407,7 +1503,7 @@ SELECT
COALESCE(SUM(total_cost), 0)::DOUBLE PRECISION AS total_cost_usd,
COALESCE(SUM(actual_total_cost), 0)::DOUBLE PRECISION AS actual_total_cost_usd,
COALESCE(SUM(error_requests), 0)::BIGINT AS error_requests,
COALESCE(SUM(response_time_sum_ms), 0) AS response_time_sum_ms,
COALESCE(SUM(response_time_sum_ms), 0)::DOUBLE PRECISION AS response_time_sum_ms,
COALESCE(SUM(response_time_samples), 0)::BIGINT AS response_time_samples
FROM stats_user_daily
WHERE user_id = $1
@@ -1429,7 +1525,7 @@ SELECT
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(input_tokens + output_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(cache_creation_tokens), 0)::BIGINT AS cache_creation_tokens,
COALESCE(SUM(cache_read_tokens), 0)::BIGINT AS cache_read_tokens,
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
@@ -1438,7 +1534,7 @@ SELECT
COALESCE(SUM(total_cost), 0)::DOUBLE PRECISION AS total_cost_usd,
COALESCE(SUM(actual_total_cost), 0)::DOUBLE PRECISION AS actual_total_cost_usd,
COALESCE(SUM(error_requests), 0)::BIGINT AS error_requests,
COALESCE(SUM(response_time_sum_ms), 0) AS response_time_sum_ms,
COALESCE(SUM(response_time_sum_ms), 0)::DOUBLE PRECISION AS response_time_sum_ms,
COALESCE(SUM(response_time_samples), 0)::BIGINT AS response_time_samples
FROM stats_daily
WHERE date >= $1
@@ -1486,33 +1582,7 @@ SELECT
END
), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(GREATEST(COALESCE("usage".output_tokens, 0), 0)), 0)::BIGINT AS output_tokens,
COALESCE(SUM(
CASE
WHEN GREATEST(COALESCE("usage".input_tokens, 0), 0) <= 0 THEN 0
WHEN GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0) <= 0
THEN GREATEST(COALESCE("usage".input_tokens, 0), 0)
WHEN split_part(lower(COALESCE(COALESCE("usage".endpoint_api_format, "usage".api_format), '')), ':', 1)
IN ('openai', 'gemini', 'google')
THEN GREATEST(
GREATEST(COALESCE("usage".input_tokens, 0), 0)
- GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0),
0
)
ELSE GREATEST(COALESCE("usage".input_tokens, 0), 0)
END
+ GREATEST(COALESCE("usage".output_tokens, 0), 0)
+ CASE
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
AND (
COALESCE("usage".cache_creation_input_tokens_5m, 0)
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
) > 0
THEN COALESCE("usage".cache_creation_input_tokens_5m, 0)
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
ELSE COALESCE("usage".cache_creation_input_tokens, 0)
END
+ GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)
), 0)::BIGINT AS total_tokens,
COALESCE(SUM(GREATEST(COALESCE("usage".total_tokens, 0), 0)), 0)::BIGINT AS total_tokens,
COALESCE(SUM(
CASE
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
@@ -2462,7 +2532,7 @@ SELECT
COALESCE(SUM(actual_total_cost), 0)::DOUBLE PRECISION AS actual_total_cost_usd,
COALESCE(SUM(cache_creation_cost), 0)::DOUBLE PRECISION AS cache_creation_cost_usd,
COALESCE(SUM(cache_read_cost), 0)::DOUBLE PRECISION AS cache_read_cost_usd,
COALESCE(SUM(response_time_sum_ms), 0) AS total_response_time_ms,
COALESCE(SUM(response_time_sum_ms), 0)::DOUBLE PRECISION AS total_response_time_ms,
COALESCE(SUM(error_requests), 0)::BIGINT AS error_requests
FROM stats_user_daily
WHERE user_id = $1
@@ -2494,7 +2564,7 @@ SELECT
COALESCE(SUM(actual_total_cost), 0)::DOUBLE PRECISION AS actual_total_cost_usd,
COALESCE(SUM(cache_creation_cost), 0)::DOUBLE PRECISION AS cache_creation_cost_usd,
COALESCE(SUM(cache_read_cost), 0)::DOUBLE PRECISION AS cache_read_cost_usd,
COALESCE(SUM(response_time_sum_ms), 0) AS total_response_time_ms,
COALESCE(SUM(response_time_sum_ms), 0)::DOUBLE PRECISION AS total_response_time_ms,
COALESCE(SUM(error_requests), 0)::BIGINT AS error_requests
FROM stats_daily
WHERE date >= $1
@@ -3770,33 +3840,7 @@ SELECT
"usage".model AS model,
"usage".provider_name AS provider,
COUNT(*)::BIGINT AS requests,
COALESCE(SUM(
CASE
WHEN GREATEST(COALESCE("usage".input_tokens, 0), 0) <= 0 THEN 0
WHEN GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0) <= 0
THEN GREATEST(COALESCE("usage".input_tokens, 0), 0)
WHEN split_part(lower(COALESCE(COALESCE("usage".endpoint_api_format, "usage".api_format), '')), ':', 1)
IN ('openai', 'gemini', 'google')
THEN GREATEST(
GREATEST(COALESCE("usage".input_tokens, 0), 0)
- GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0),
0
)
ELSE GREATEST(COALESCE("usage".input_tokens, 0), 0)
END
+ GREATEST(COALESCE("usage".output_tokens, 0), 0)
+ CASE
WHEN COALESCE("usage".cache_creation_input_tokens, 0) = 0
AND (
COALESCE("usage".cache_creation_input_tokens_5m, 0)
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
) > 0
THEN COALESCE("usage".cache_creation_input_tokens_5m, 0)
+ COALESCE("usage".cache_creation_input_tokens_1h, 0)
ELSE COALESCE("usage".cache_creation_input_tokens, 0)
END
+ GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)
), 0)::BIGINT AS total_tokens,
COALESCE(SUM(GREATEST(COALESCE("usage".total_tokens, 0), 0)), 0)::BIGINT AS total_tokens,
COALESCE(SUM(COALESCE(CAST("usage".total_cost_usd AS DOUBLE PRECISION), 0)), 0)
AS total_cost_usd,
COALESCE(SUM(
@@ -4000,7 +4044,7 @@ SELECT
{group_column} AS group_key,
COALESCE(SUM(total_requests), 0)::BIGINT AS request_count,
COALESCE(SUM(input_tokens), 0)::BIGINT AS input_tokens,
COALESCE(SUM(effective_input_tokens + output_tokens + cache_creation_tokens + cache_read_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(total_tokens), 0)::BIGINT AS total_tokens,
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
COALESCE(SUM(effective_input_tokens), 0)::BIGINT AS effective_input_tokens,
COALESCE(SUM(total_input_context), 0)::BIGINT AS total_input_context,
@@ -4735,17 +4779,17 @@ SELECT
COALESCE(SUM(success_flag), 0)::BIGINT AS success_count,
CASE
WHEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0) > 0
THEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN output_tokens
ELSE 0
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0)::DOUBLE PRECISION
ELSE NULL
@@ -4822,17 +4866,17 @@ SELECT
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
CASE
WHEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0) > 0
THEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN output_tokens
ELSE 0
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0)::DOUBLE PRECISION
ELSE NULL
@@ -4854,7 +4898,7 @@ SELECT
ELSE NULL
END AS p90_first_byte_time_ms,
COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN 1
ELSE 0
END), 0)::BIGINT AS tps_sample_count,
@@ -4954,17 +4998,17 @@ SELECT
COALESCE(SUM(output_tokens), 0)::BIGINT AS output_tokens,
CASE
WHEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0) > 0
THEN COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN output_tokens
ELSE 0
END), 0)::DOUBLE PRECISION * 1000.0 / COALESCE(SUM(CASE
WHEN success_flag = 1 AND output_tps_duration_ms > 0 AND output_tokens > 0
THEN output_tps_duration_ms
WHEN success_flag = 1 AND response_time_ms > 0 AND output_tokens > 0
THEN response_time_ms
ELSE 0
END), 0)::DOUBLE PRECISION
ELSE NULL
@@ -5040,7 +5084,7 @@ SELECT
AS cache_creation_cost_usd,
COALESCE(SUM(
COALESCE(
CAST("usage".input_price_per_1m AS DOUBLE PRECISION),
CAST("usage".output_price_per_1m AS DOUBLE PRECISION),
0
) * GREATEST(COALESCE("usage".cache_read_input_tokens, 0), 0)::DOUBLE PRECISION / 1000000.0
), 0) AS estimated_full_cost_usd
@@ -6157,6 +6201,11 @@ WHERE date >= $1
GROUP BY {group_column}
ORDER BY request_count DESC, group_key ASC
"#,
group_column = group_column,
display_name_expr = display_name_expr,
avg_response_time_expr = avg_response_time_expr,
success_count_expr = success_count_expr,
table_name = table_name,
);
let mut rows = sqlx::query(&sql)
@@ -6337,7 +6386,62 @@ LIMIT $3
let mut items = Vec::new();
while let Some(row) = rows.try_next().await.map_postgres_err()? {
items.push(decode_usage_audit_aggregation_row(&row)?);
items.push(StoredUsageAuditAggregation {
group_key: row.try_get::<String, _>("group_key").map_postgres_err()?,
display_name: row
.try_get::<Option<String>, _>("display_name")
.map_postgres_err()?,
secondary_name: row
.try_get::<Option<String>, _>("secondary_name")
.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,
output_tokens: row
.try_get::<i64, _>("output_tokens")
.map_postgres_err()?
.max(0) as u64,
effective_input_tokens: row
.try_get::<i64, _>("effective_input_tokens")
.map_postgres_err()?
.max(0) as u64,
total_input_context: row
.try_get::<i64, _>("total_input_context")
.map_postgres_err()?
.max(0) as u64,
cache_creation_tokens: row
.try_get::<i64, _>("cache_creation_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_5m_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_creation_ephemeral_1h_tokens: row
.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")
.map_postgres_err()?
.max(0) as u64,
cache_read_tokens: row
.try_get::<i64, _>("cache_read_tokens")
.map_postgres_err()?
.max(0) as u64,
total_cost_usd: row.try_get::<f64, _>("total_cost_usd").map_postgres_err()?,
actual_total_cost_usd: row
.try_get::<f64, _>("actual_total_cost_usd")
.map_postgres_err()?,
avg_response_time_ms: row
.try_get::<Option<f64>, _>("avg_response_time_ms")
.map_postgres_err()?,
success_count: row
.try_get::<Option<i64>, _>("success_count")
.map_postgres_err()?
.map(|value| value.max(0) as u64),
});
}
Ok(items)
}
@@ -6474,13 +6578,13 @@ WHERE "usage".created_at >= $1
r#"
SELECT
date,
total_requests,
input_tokens,
output_tokens,
cache_creation_tokens,
cache_read_tokens,
total_cost,
COALESCE(actual_total_cost, 0) AS actual_total_cost
total_requests::BIGINT AS total_requests,
input_tokens::BIGINT AS input_tokens,
output_tokens::BIGINT AS output_tokens,
cache_creation_tokens::BIGINT AS cache_creation_tokens,
cache_read_tokens::BIGINT AS cache_read_tokens,
total_cost::DOUBLE PRECISION AS total_cost,
COALESCE(actual_total_cost, 0)::DOUBLE PRECISION AS actual_total_cost
FROM stats_user_daily
WHERE user_id = $1
AND date >= $2
@@ -6497,13 +6601,13 @@ ORDER BY date ASC
r#"
SELECT
date,
total_requests,
input_tokens,
output_tokens,
cache_creation_tokens,
cache_read_tokens,
total_cost,
actual_total_cost
total_requests::BIGINT AS total_requests,
input_tokens::BIGINT AS input_tokens,
output_tokens::BIGINT AS output_tokens,
cache_creation_tokens::BIGINT AS cache_creation_tokens,
cache_read_tokens::BIGINT AS cache_read_tokens,
total_cost::DOUBLE PRECISION AS total_cost,
actual_total_cost::DOUBLE PRECISION AS actual_total_cost
FROM stats_daily
WHERE date >= $1
AND date < $2
@@ -7267,6 +7371,40 @@ ORDER BY "usage".user_id ASC
}
}
let before_model_contribution =
previous_usage.as_ref().and_then(model_usage_contribution);
let after_model_contribution = model_usage_contribution(&stored);
match (
before_model_contribution.as_ref(),
after_model_contribution.as_ref(),
) {
(Some(before), Some(after)) if before.model == after.model => {
let delta = ModelUsageDelta::between(before, after);
apply_global_model_usage_delta_in_tx(tx, before.model.as_str(), &delta)
.await?;
}
_ => {
if let Some(before) = before_model_contribution.as_ref() {
let delta = ModelUsageDelta::removal(before);
apply_global_model_usage_delta_in_tx(
tx,
before.model.as_str(),
&delta,
)
.await?;
}
if let Some(after) = after_model_contribution.as_ref() {
let delta = ModelUsageDelta::addition(after);
apply_global_model_usage_delta_in_tx(
tx,
after.model.as_str(),
&delta,
)
.await?;
}
}
}
let before_provider_contribution = previous_usage
.as_ref()
.and_then(provider_api_key_usage_contribution);
@@ -7305,46 +7443,140 @@ ORDER BY "usage".user_id ASC
}
}
}
let before_model_contribution =
previous_usage.as_ref().and_then(model_usage_contribution);
let after_model_contribution = model_usage_contribution(&stored);
match (
before_model_contribution.as_ref(),
after_model_contribution.as_ref(),
) {
(Some(before), Some(after)) if before.model == after.model => {
let delta = ModelUsageDelta::between(before, after);
apply_global_model_usage_delta_in_tx(tx, before.model.as_str(), &delta)
.await?;
}
_ => {
if let Some(before) = before_model_contribution.as_ref() {
let delta = ModelUsageDelta::removal(before);
apply_global_model_usage_delta_in_tx(
tx,
before.model.as_str(),
&delta,
)
.await?;
}
if let Some(after) = after_model_contribution.as_ref() {
let delta = ModelUsageDelta::addition(after);
apply_global_model_usage_delta_in_tx(
tx,
after.model.as_str(),
&delta,
)
.await?;
}
}
}
Ok(stored)
}) as BoxFuture<'_, Result<StoredRequestUsageAudit, DataLayerError>>
})
.await
}
pub async fn cleanup_stale_pending_requests(
&self,
cutoff_unix_secs: u64,
now_unix_secs: u64,
timeout_minutes: u64,
batch_size: usize,
) -> Result<PendingUsageCleanupSummary, DataLayerError> {
if batch_size == 0 {
return Ok(PendingUsageCleanupSummary::default());
}
let cutoff_timestamp = i64::try_from(cutoff_unix_secs).map_err(|_| {
DataLayerError::InvalidInput(format!(
"invalid stale pending usage cutoff: {cutoff_unix_secs}"
))
})?;
let now_timestamp = i64::try_from(now_unix_secs).map_err(|_| {
DataLayerError::InvalidInput(format!(
"invalid stale pending usage timestamp: {now_unix_secs}"
))
})?;
let cutoff = DateTime::<Utc>::from_timestamp(cutoff_timestamp, 0).ok_or_else(|| {
DataLayerError::InvalidInput(format!(
"invalid stale pending usage cutoff: {cutoff_unix_secs}"
))
})?;
let now = DateTime::<Utc>::from_timestamp(now_timestamp, 0).ok_or_else(|| {
DataLayerError::InvalidInput(format!(
"invalid stale pending usage timestamp: {now_unix_secs}"
))
})?;
let mut summary = PendingUsageCleanupSummary::default();
let batch_size_i64 = i64::try_from(batch_size).map_err(|_| {
DataLayerError::InvalidInput(format!(
"invalid stale pending usage batch size: {batch_size}"
))
})?;
loop {
let mut tx = self.pool.begin().await.map_postgres_err()?;
let stale_rows = sqlx::query(SELECT_STALE_PENDING_USAGE_BATCH_SQL)
.bind(cutoff)
.bind(batch_size_i64)
.fetch_all(&mut *tx)
.await
.map_postgres_err()?;
if stale_rows.is_empty() {
tx.rollback().await.map_postgres_err()?;
break;
}
let stale_rows = stale_rows
.iter()
.map(|row| {
Ok(StalePendingUsageRow {
request_id: row.try_get("request_id").map_postgres_err()?,
status: row.try_get("status").map_postgres_err()?,
billing_status: row.try_get("billing_status").map_postgres_err()?,
})
})
.collect::<Result<Vec<_>, DataLayerError>>()?;
let request_ids = stale_rows
.iter()
.map(|row| row.request_id.clone())
.collect::<Vec<_>>();
let completed_request_ids = if request_ids.is_empty() {
Vec::new()
} else {
sqlx::query(SELECT_COMPLETED_PENDING_REQUEST_IDS_SQL)
.bind(request_ids)
.fetch_all(&mut *tx)
.await
.map_postgres_err()?
.iter()
.map(|row| row.try_get("request_id").map_postgres_err())
.collect::<Result<Vec<String>, DataLayerError>>()?
};
for row in stale_rows {
if completed_request_ids.contains(&row.request_id) {
sqlx::query(UPDATE_RECOVERED_STALE_USAGE_SQL)
.bind(&row.request_id)
.execute(&mut *tx)
.await
.map_postgres_err()?;
sqlx::query(UPDATE_RECOVERED_STREAMING_CANDIDATES_SQL)
.bind(&row.request_id)
.bind(now)
.execute(&mut *tx)
.await
.map_postgres_err()?;
summary.recovered += 1;
continue;
}
let error_message = stale_pending_error_message(&row.status, timeout_minutes);
if row.billing_status == "pending" {
sqlx::query(UPDATE_FAILED_VOID_STALE_USAGE_SQL)
.bind(&row.request_id)
.bind(&error_message)
.bind(now)
.execute(&mut *tx)
.await
.map_postgres_err()?;
} else {
sqlx::query(UPDATE_FAILED_STALE_USAGE_SQL)
.bind(&row.request_id)
.bind(&error_message)
.execute(&mut *tx)
.await
.map_postgres_err()?;
}
sqlx::query(UPDATE_FAILED_PENDING_CANDIDATES_SQL)
.bind(&row.request_id)
.bind(now)
.execute(&mut *tx)
.await
.map_postgres_err()?;
summary.failed += 1;
}
tx.commit().await.map_postgres_err()?;
}
Ok(summary)
}
pub async fn rebuild_api_key_usage_stats(&self) -> Result<u64, DataLayerError> {
self.tx_runner
.run_read_write(|tx| {
@@ -7624,6 +7856,42 @@ impl UsageWriteRepository for SqlxUsageReadRepository {
async fn rebuild_provider_api_key_usage_stats(&self) -> Result<u64, DataLayerError> {
Self::rebuild_provider_api_key_usage_stats(self).await
}
async fn cleanup_stale_pending_requests(
&self,
cutoff_unix_secs: u64,
now_unix_secs: u64,
timeout_minutes: u64,
batch_size: usize,
) -> Result<PendingUsageCleanupSummary, DataLayerError> {
Self::cleanup_stale_pending_requests(
self,
cutoff_unix_secs,
now_unix_secs,
timeout_minutes,
batch_size,
)
.await
}
async fn cleanup_usage(
&self,
window: &UsageCleanupWindow,
batch_size: usize,
auto_delete_expired_keys: bool,
) -> Result<UsageCleanupSummary, DataLayerError> {
Self::cleanup_usage(self, window, batch_size, auto_delete_expired_keys).await
}
}
struct StalePendingUsageRow {
request_id: String,
status: String,
billing_status: String,
}
fn stale_pending_error_message(status: &str, timeout_minutes: u64) -> String {
format!("请求超时: 状态 '{status}' 超过 {timeout_minutes} 分钟未完成")
}
async fn find_usage_by_request_id_in_tx(
@@ -7695,6 +7963,32 @@ async fn apply_api_key_usage_delta_in_tx(
Ok(())
}
async fn apply_global_model_usage_delta_in_tx(
tx: &mut sqlx::Transaction<'_, Postgres>,
model: &str,
delta: &ModelUsageDelta,
) -> Result<(), DataLayerError> {
if model.trim().is_empty() {
return Ok(());
}
if delta.is_noop() {
return Ok(());
}
sqlx::query(APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL)
.bind(model)
.bind(i32::try_from(delta.request_count).map_err(|_| {
DataLayerError::UnexpectedValue(format!(
"global_models.usage_count delta exceeds i32: {}",
delta.request_count
))
})?)
.execute(&mut **tx)
.await
.map_postgres_err()?;
Ok(())
}
async fn apply_provider_api_key_usage_delta_in_tx(
tx: &mut sqlx::Transaction<'_, Postgres>,
key_id: &str,
@@ -7757,33 +8051,6 @@ async fn apply_provider_api_key_usage_delta_in_tx(
Ok(())
}
async fn apply_global_model_usage_delta_in_tx(
tx: &mut sqlx::Transaction<'_, Postgres>,
model: &str,
delta: &ModelUsageDelta,
) -> Result<(), DataLayerError> {
let model = model.trim();
if model.is_empty() {
return Ok(());
}
if delta.is_noop() {
return Ok(());
}
sqlx::query(APPLY_GLOBAL_MODEL_USAGE_DELTA_SQL)
.bind(model)
.bind(i32::try_from(delta.request_count).map_err(|_| {
DataLayerError::UnexpectedValue(format!(
"global_models.usage_count delta exceeds i32: {}",
delta.request_count
))
})?)
.execute(&mut **tx)
.await
.map_postgres_err()?;
Ok(())
}
// Build the usage read model from the split storage layout.
//
// Query projections already prefer the newer audit/snapshot owners and only fall back to

View File

@@ -12,7 +12,7 @@ use super::{
UsageHttpAuditRefs, UsageRoutingSnapshot, UsageSettlementPricingSnapshot,
MAX_INLINE_USAGE_BODY_BYTES,
};
use crate::postgres::{PostgresPoolConfig, PostgresPoolFactory};
use crate::driver::postgres::{PostgresPoolConfig, PostgresPoolFactory};
use crate::repository::usage::UpsertUsageRecord;
use aether_data_contracts::repository::usage::UsageBodyField;
@@ -382,6 +382,8 @@ fn usage_sql_summarize_usage_daily_heatmap_supports_daily_aggregates() {
assert!(source.contains("summarize_usage_daily_heatmap_from_daily_aggregates"));
assert!(source.contains("FROM stats_daily"));
assert!(source.contains("FROM stats_user_daily"));
assert!(source.contains("total_requests::BIGINT AS total_requests"));
assert!(source.contains("total_cost::DOUBLE PRECISION AS total_cost"));
assert!(
source.contains("split_dashboard_daily_aggregate_range(start_utc, end_utc, cutoff_utc)")
);
@@ -458,18 +460,14 @@ fn usage_sql_provider_performance_reads_upstream_stream_from_billing_facts() {
#[test]
fn usage_billing_facts_projects_upstream_stream_mode() {
let migration = include_str!(
"../../../../migrations/20260505130000_project_upstream_stream_in_usage_billing_facts.sql"
"../../../../migrations/postgres/20260505130000_project_upstream_stream_in_usage_billing_facts.sql"
);
let baseline = include_str!("../../../../bootstrap/20260413020000_baseline_v2.sql");
for source in [migration, baseline] {
assert!(source.contains("AS upstream_is_stream"));
assert!(source.contains("COALESCE(usage_rows.upstream_is_stream"));
assert!(source.contains("COALESCE(usage_rows.is_stream, FALSE)"));
}
assert!(migration.contains("AS upstream_is_stream"));
assert!(migration.contains("COALESCE(usage_rows.upstream_is_stream"));
assert!(migration.contains("COALESCE(usage_rows.is_stream, FALSE)"));
assert!(migration.contains("ADD COLUMN IF NOT EXISTS upstream_is_stream boolean"));
assert!(migration.contains("request_metadata->>'upstream_is_stream'"));
assert!(baseline.contains("upstream_is_stream boolean"));
}
#[test]

File diff suppressed because it is too large Load Diff