feat(stats): 新增聚合读路径与回填机制并重构 dashboard/usage 读取链路

- 新增 stats_user_summary 及 user_daily_provider/api_format/cost_savings 等聚合表
- 扩展 stats_daily/hourly 有效 token 与响应时间等字段,maintenance runtime 同步写入
- 新增 backfill 模块与 --apply-backfills 命令补齐历史聚合数据
- 重写 dashboard_filters、usage_heatmap、user_rollups 查询改走聚合表
- 同步更新 baseline_v2.sql 与 migration 集,README/dev.sh 补充回填用法
This commit is contained in:
fawney19
2026-04-22 17:29:39 +08:00
parent 063ef02306
commit a5e6bd3b62
45 changed files with 15323 additions and 1039 deletions

File diff suppressed because it is too large Load Diff

View File

@@ -165,9 +165,9 @@ pub(super) async fn run_wallet_daily_usage_aggregation_once(
pub(super) async fn run_stats_aggregation_once(
data: &GatewayDataState,
) -> Result<(), DataLayerError> {
) -> Result<bool, DataLayerError> {
let Some(summary) = perform_stats_aggregation_once(data).await? else {
return Ok(());
return Ok(false);
};
info!(
@@ -183,7 +183,7 @@ pub(super) async fn run_stats_aggregation_once(
user_rows = summary.user_rows,
"gateway aggregated daily stats tables"
);
Ok(())
Ok(true)
}
pub(super) async fn run_usage_cleanup_once(data: &GatewayDataState) -> Result<(), DataLayerError> {
@@ -248,9 +248,9 @@ pub(super) async fn run_pending_cleanup_once(
pub(super) async fn run_stats_hourly_aggregation_once(
data: &GatewayDataState,
) -> Result<(), DataLayerError> {
) -> Result<bool, DataLayerError> {
let Some(summary) = perform_stats_hourly_aggregation_once(data).await? else {
return Ok(());
return Ok(false);
};
info!(
@@ -260,11 +260,12 @@ pub(super) async fn run_stats_hourly_aggregation_once(
hour_utc = %summary.hour_utc,
total_requests = summary.total_requests,
user_rows = summary.user_rows,
user_model_rows = summary.user_model_rows,
model_rows = summary.model_rows,
provider_rows = summary.provider_rows,
"gateway aggregated stats hourly tables"
);
Ok(())
Ok(true)
}
pub(super) async fn run_provider_checkin_once(state: &AppState) -> Result<(), GatewayError> {

View File

@@ -9,12 +9,22 @@ use super::{
postgres_error, stats_aggregation_target_day, system_config_bool, PercentileSummary,
StatsAggregationSummary, DELETE_STATS_DAILY_ERRORS_FOR_DATE_SQL, INSERT_STATS_DAILY_ERROR_SQL,
INSERT_STATS_SUMMARY_SQL, SELECT_EXISTING_STATS_SUMMARY_ID_SQL,
SELECT_LATEST_STATS_DAILY_DATE_SQL, SELECT_NEXT_STATS_DAILY_BUCKET_SQL,
SELECT_STATS_DAILY_AGGREGATE_SQL, SELECT_STATS_DAILY_FALLBACK_COUNT_SQL,
SELECT_STATS_DAILY_FIRST_BYTE_PERCENTILES_SQL,
SELECT_STATS_DAILY_RESPONSE_TIME_PERCENTILES_SQL, SELECT_STATS_SUMMARY_ENTITY_COUNTS_SQL,
SELECT_STATS_SUMMARY_TOTALS_SQL, UPDATE_STATS_SUMMARY_SQL, UPSERT_STATS_DAILY_API_KEY_SQL,
UPSERT_STATS_DAILY_MODEL_SQL, UPSERT_STATS_DAILY_PROVIDER_SQL, UPSERT_STATS_DAILY_SQL,
UPSERT_STATS_USER_DAILY_SQL,
UPSERT_STATS_DAILY_COST_SAVINGS_MODEL_PROVIDER_SQL, UPSERT_STATS_DAILY_COST_SAVINGS_MODEL_SQL,
UPSERT_STATS_DAILY_COST_SAVINGS_PROVIDER_SQL, UPSERT_STATS_DAILY_COST_SAVINGS_SQL,
UPSERT_STATS_DAILY_MODEL_PROVIDER_SQL, UPSERT_STATS_DAILY_MODEL_SQL,
UPSERT_STATS_DAILY_PROVIDER_SQL, UPSERT_STATS_DAILY_SQL,
UPSERT_STATS_USER_DAILY_API_FORMAT_SQL,
UPSERT_STATS_USER_DAILY_COST_SAVINGS_MODEL_PROVIDER_SQL,
UPSERT_STATS_USER_DAILY_COST_SAVINGS_MODEL_SQL,
UPSERT_STATS_USER_DAILY_COST_SAVINGS_PROVIDER_SQL, UPSERT_STATS_USER_DAILY_COST_SAVINGS_SQL,
UPSERT_STATS_USER_DAILY_MODEL_PROVIDER_SQL, UPSERT_STATS_USER_DAILY_MODEL_SQL,
UPSERT_STATS_USER_DAILY_PROVIDER_SQL, UPSERT_STATS_USER_DAILY_SQL,
UPSERT_STATS_USER_SUMMARY_SQL,
};
pub(super) async fn perform_stats_aggregation_once(
@@ -28,109 +38,125 @@ pub(super) async fn perform_stats_aggregation_once(
}
let now_utc = Utc::now();
let day_start_utc = stats_aggregation_target_day(now_utc);
let target_day_utc = stats_aggregation_target_day(now_utc);
let Some(day_start_utc) = next_stats_aggregation_day(&pool, target_day_utc)
.await
.map_err(postgres_error)?
else {
return Ok(None);
};
perform_stats_aggregation_for_day(&pool, day_start_utc, now_utc)
.await
.map(Some)
.map_err(postgres_error)
}
async fn next_stats_aggregation_day(
pool: &aether_data::postgres::PostgresPool,
target_day_utc: DateTime<Utc>,
) -> Result<Option<DateTime<Utc>>, sqlx::Error> {
let latest_row = sqlx::query(SELECT_LATEST_STATS_DAILY_DATE_SQL)
.fetch_one(pool)
.await?;
let latest_day = latest_row.try_get::<Option<DateTime<Utc>>, _>("latest_date")?;
let search_from = latest_day
.map(|value| value + chrono::Duration::days(1))
.unwrap_or_else(|| {
DateTime::<Utc>::from_timestamp(0, 0).expect("unix epoch should be valid")
});
let search_until = target_day_utc + chrono::Duration::days(1);
if search_from >= search_until {
return Ok(None);
}
let next_row = sqlx::query(SELECT_NEXT_STATS_DAILY_BUCKET_SQL)
.bind(search_from)
.bind(search_until)
.fetch_one(pool)
.await?;
let next_bucket = next_row.try_get::<Option<DateTime<Utc>>, _>("next_bucket")?;
Ok(next_bucket.filter(|value| *value <= target_day_utc))
}
async fn perform_stats_aggregation_for_day(
pool: &aether_data::postgres::PostgresPool,
day_start_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<StatsAggregationSummary, sqlx::Error> {
let day_end_utc = day_start_utc + chrono::Duration::days(1);
let mut tx = pool.begin().await.map_err(postgres_error)?;
let mut tx = pool.begin().await?;
let aggregate_row = sqlx::query(SELECT_STATS_DAILY_AGGREGATE_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.fetch_one(&mut *tx)
.await
.map_err(postgres_error)?;
let total_requests = aggregate_row
.try_get::<i64, _>("total_requests")
.map_err(postgres_error)?;
let error_requests = aggregate_row
.try_get::<i64, _>("error_requests")
.map_err(postgres_error)?;
.await?;
let total_requests = aggregate_row.try_get::<i64, _>("total_requests")?;
let error_requests = aggregate_row.try_get::<i64, _>("error_requests")?;
let success_requests = total_requests.saturating_sub(error_requests);
let fallback_count = sqlx::query(SELECT_STATS_DAILY_FALLBACK_COUNT_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(vec!["success", "failed"])
.fetch_one(&mut *tx)
.await
.map_err(postgres_error)?
.try_get::<i64, _>("fallback_count")
.map_err(postgres_error)?;
.await?
.try_get::<i64, _>("fallback_count")?;
let response_percentiles = fetch_stats_daily_percentiles(
&mut tx,
SELECT_STATS_DAILY_RESPONSE_TIME_PERCENTILES_SQL,
day_start_utc,
day_end_utc,
)
.await
.map_err(postgres_error)?;
.await?;
let first_byte_percentiles = fetch_stats_daily_percentiles(
&mut tx,
SELECT_STATS_DAILY_FIRST_BYTE_PERCENTILES_SQL,
day_start_utc,
day_end_utc,
)
.await
.map_err(postgres_error)?;
.await?;
sqlx::query(UPSERT_STATS_DAILY_SQL)
.bind(Uuid::new_v4().to_string())
.bind(day_start_utc)
.bind(total_requests)
.bind(aggregate_row.try_get::<i64, _>("cache_hit_total_requests")?)
.bind(aggregate_row.try_get::<i64, _>("cache_hit_requests")?)
.bind(aggregate_row.try_get::<i64, _>("completed_total_requests")?)
.bind(aggregate_row.try_get::<i64, _>("completed_cache_hit_requests")?)
.bind(aggregate_row.try_get::<i64, _>("completed_input_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("completed_cache_creation_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("completed_cache_read_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("completed_total_input_context")?)
.bind(aggregate_row.try_get::<f64, _>("completed_cache_creation_cost")?)
.bind(aggregate_row.try_get::<f64, _>("completed_cache_read_cost")?)
.bind(aggregate_row.try_get::<f64, _>("settled_total_cost")?)
.bind(aggregate_row.try_get::<i64, _>("settled_total_requests")?)
.bind(aggregate_row.try_get::<i64, _>("settled_input_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("settled_output_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("settled_cache_creation_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("settled_cache_read_tokens")?)
.bind(aggregate_row.try_get::<Option<i64>, _>("settled_first_finalized_at_unix_secs")?)
.bind(aggregate_row.try_get::<Option<i64>, _>("settled_last_finalized_at_unix_secs")?)
.bind(success_requests)
.bind(error_requests)
.bind(
aggregate_row
.try_get::<i64, _>("input_tokens")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<i64, _>("output_tokens")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<i64, _>("cache_creation_tokens")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<i64, _>("cache_read_tokens")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("total_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("actual_total_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("input_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("output_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("cache_creation_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("cache_read_cost")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<f64, _>("avg_response_time_ms")
.map_err(postgres_error)?,
)
.bind(aggregate_row.try_get::<i64, _>("input_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("effective_input_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("output_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("cache_creation_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("cache_creation_ephemeral_5m_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("cache_creation_ephemeral_1h_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("cache_read_tokens")?)
.bind(aggregate_row.try_get::<i64, _>("total_input_context")?)
.bind(aggregate_row.try_get::<f64, _>("total_cost")?)
.bind(aggregate_row.try_get::<f64, _>("actual_total_cost")?)
.bind(aggregate_row.try_get::<f64, _>("input_cost")?)
.bind(aggregate_row.try_get::<f64, _>("output_cost")?)
.bind(aggregate_row.try_get::<f64, _>("cache_creation_cost")?)
.bind(aggregate_row.try_get::<f64, _>("cache_read_cost")?)
.bind(aggregate_row.try_get::<f64, _>("response_time_sum_ms")?)
.bind(aggregate_row.try_get::<i64, _>("response_time_samples")?)
.bind(aggregate_row.try_get::<f64, _>("avg_response_time_ms")?)
.bind(response_percentiles.p50)
.bind(response_percentiles.p90)
.bind(response_percentiles.p99)
@@ -138,47 +164,65 @@ pub(super) async fn perform_stats_aggregation_once(
.bind(first_byte_percentiles.p90)
.bind(first_byte_percentiles.p99)
.bind(fallback_count)
.bind(
aggregate_row
.try_get::<i64, _>("unique_models")
.map_err(postgres_error)?,
)
.bind(
aggregate_row
.try_get::<i64, _>("unique_providers")
.map_err(postgres_error)?,
)
.bind(aggregate_row.try_get::<i64, _>("unique_models")?)
.bind(aggregate_row.try_get::<i64, _>("unique_providers")?)
.bind(true)
.bind(now_utc)
.bind(now_utc)
.bind(now_utc)
.execute(&mut *tx)
.await
.map_err(postgres_error)?;
.await?;
let model_rows = upsert_stats_daily_model_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
let model_rows =
upsert_stats_daily_model_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
let provider_rows =
upsert_stats_daily_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
upsert_stats_daily_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_daily_model_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_daily_cost_savings_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_daily_cost_savings_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await?;
upsert_stats_daily_cost_savings_model_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await?;
upsert_stats_daily_cost_savings_model_provider_rows(
&mut tx,
day_start_utc,
day_end_utc,
now_utc,
)
.await?;
let api_key_rows =
upsert_stats_daily_api_key_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
let error_rows = refresh_stats_daily_error_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
let user_rows = upsert_stats_user_daily_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
refresh_stats_summary_row(&mut tx, day_end_utc, now_utc)
.await
.map_err(postgres_error)?;
tx.commit().await.map_err(postgres_error)?;
upsert_stats_daily_api_key_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
let error_rows =
refresh_stats_daily_error_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
let user_rows =
upsert_stats_user_daily_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_user_daily_model_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_user_daily_model_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await?;
upsert_stats_user_daily_provider_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_user_daily_cost_savings_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
upsert_stats_user_daily_cost_savings_provider_rows(
&mut tx,
day_start_utc,
day_end_utc,
now_utc,
)
.await?;
upsert_stats_user_daily_cost_savings_model_rows(&mut tx, day_start_utc, day_end_utc, now_utc)
.await?;
upsert_stats_user_daily_cost_savings_model_provider_rows(
&mut tx,
day_start_utc,
day_end_utc,
now_utc,
)
.await?;
upsert_stats_user_daily_api_format_rows(&mut tx, day_start_utc, day_end_utc, now_utc).await?;
refresh_stats_summary_row(&mut tx, day_end_utc, now_utc).await?;
refresh_stats_user_summary_rows(&mut tx, day_end_utc, now_utc).await?;
tx.commit().await?;
Ok(Some(StatsAggregationSummary {
Ok(StatsAggregationSummary {
day_start_utc,
total_requests,
model_rows,
@@ -186,7 +230,7 @@ pub(super) async fn perform_stats_aggregation_once(
api_key_rows,
error_rows,
user_rows,
}))
})
}
async fn fetch_stats_daily_percentiles(
@@ -250,6 +294,91 @@ async fn upsert_stats_daily_provider_rows(
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_model_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_DAILY_MODEL_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_cost_savings_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_DAILY_COST_SAVINGS_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_cost_savings_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_DAILY_COST_SAVINGS_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_cost_savings_model_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_DAILY_COST_SAVINGS_MODEL_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_cost_savings_model_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_DAILY_COST_SAVINGS_MODEL_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_daily_api_key_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
@@ -305,6 +434,142 @@ async fn upsert_stats_user_daily_rows(
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_model_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_MODEL_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_model_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_MODEL_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_cost_savings_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_COST_SAVINGS_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_cost_savings_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_COST_SAVINGS_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_cost_savings_model_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_COST_SAVINGS_MODEL_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_cost_savings_model_provider_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_COST_SAVINGS_MODEL_PROVIDER_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_user_daily_api_format_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
day_start_utc: DateTime<Utc>,
day_end_utc: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_USER_DAILY_API_FORMAT_SQL)
.bind(day_start_utc)
.bind(day_end_utc)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn refresh_stats_summary_row(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
cutoff_date: DateTime<Utc>,
@@ -381,3 +646,17 @@ async fn refresh_stats_summary_row(
Ok(())
}
async fn refresh_stats_user_summary_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
cutoff_date: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<(), sqlx::Error> {
sqlx::query(UPSERT_STATS_USER_SUMMARY_SQL)
.bind(cutoff_date)
.bind(now_utc)
.execute(&mut **tx)
.await?;
Ok(())
}

View File

@@ -6,9 +6,10 @@ use crate::data::GatewayDataState;
use aether_data_contracts::DataLayerError;
use super::{
stats_hourly_aggregation_target_hour, system_config_bool, SELECT_STATS_HOURLY_AGGREGATE_SQL,
stats_hourly_aggregation_target_hour, system_config_bool, SELECT_LATEST_STATS_HOURLY_HOUR_SQL,
SELECT_NEXT_STATS_HOURLY_BUCKET_SQL, SELECT_STATS_HOURLY_AGGREGATE_SQL,
UPSERT_STATS_HOURLY_MODEL_SQL, UPSERT_STATS_HOURLY_PROVIDER_SQL, UPSERT_STATS_HOURLY_SQL,
UPSERT_STATS_HOURLY_USER_SQL,
UPSERT_STATS_HOURLY_USER_MODEL_SQL, UPSERT_STATS_HOURLY_USER_SQL,
};
#[derive(Debug, Clone, PartialEq)]
@@ -16,6 +17,7 @@ pub(super) struct StatsHourlyAggregationSummary {
pub(super) hour_utc: DateTime<Utc>,
pub(super) total_requests: i64,
pub(super) user_rows: usize,
pub(super) user_model_rows: usize,
pub(super) model_rows: usize,
pub(super) provider_rows: usize,
}
@@ -31,81 +33,121 @@ pub(super) async fn perform_stats_hourly_aggregation_once(
}
let now_utc = Utc::now();
let hour_utc = stats_hourly_aggregation_target_hour(now_utc);
let target_hour_utc = stats_hourly_aggregation_target_hour(now_utc);
let Some(hour_utc) = next_stats_hourly_bucket(&pool, target_hour_utc)
.await
.map_err(postgres_error)?
else {
return Ok(None);
};
perform_stats_hourly_aggregation_for_hour(&pool, hour_utc, now_utc)
.await
.map(Some)
.map_err(postgres_error)
}
async fn next_stats_hourly_bucket(
pool: &aether_data::postgres::PostgresPool,
target_hour_utc: DateTime<Utc>,
) -> Result<Option<DateTime<Utc>>, sqlx::Error> {
let latest_row = sqlx::query(SELECT_LATEST_STATS_HOURLY_HOUR_SQL)
.fetch_one(pool)
.await?;
let latest_hour = latest_row.try_get::<Option<DateTime<Utc>>, _>("latest_hour")?;
let search_from = latest_hour
.map(|value| value + chrono::Duration::hours(1))
.unwrap_or_else(|| {
DateTime::<Utc>::from_timestamp(0, 0).expect("unix epoch should be valid")
});
let search_until = target_hour_utc + chrono::Duration::hours(1);
if search_from >= search_until {
return Ok(None);
}
let next_row = sqlx::query(SELECT_NEXT_STATS_HOURLY_BUCKET_SQL)
.bind(search_from)
.bind(search_until)
.fetch_one(pool)
.await?;
let next_bucket = next_row.try_get::<Option<DateTime<Utc>>, _>("next_bucket")?;
Ok(next_bucket.filter(|value| *value <= target_hour_utc))
}
async fn perform_stats_hourly_aggregation_for_hour(
pool: &aether_data::postgres::PostgresPool,
hour_utc: DateTime<Utc>,
aggregated_at: DateTime<Utc>,
) -> Result<StatsHourlyAggregationSummary, sqlx::Error> {
let hour_end = hour_utc + chrono::Duration::hours(1);
let aggregated_at = now_utc;
let mut tx = pool.begin().await.map_err(postgres_error)?;
let mut tx = pool.begin().await?;
let row = sqlx::query(SELECT_STATS_HOURLY_AGGREGATE_SQL)
.bind(hour_utc)
.bind(hour_end)
.fetch_one(&mut *tx)
.await
.map_err(postgres_error)?;
let total_requests = row
.try_get::<i64, _>("total_requests")
.map_err(postgres_error)?;
let error_requests = row
.try_get::<i64, _>("error_requests")
.map_err(postgres_error)?;
.await?;
let total_requests = row.try_get::<i64, _>("total_requests")?;
let error_requests = row.try_get::<i64, _>("error_requests")?;
let success_requests = total_requests.saturating_sub(error_requests);
sqlx::query(UPSERT_STATS_HOURLY_SQL)
.bind(Uuid::new_v4().to_string())
.bind(hour_utc)
.bind(total_requests)
.bind(row.try_get::<i64, _>("cache_hit_total_requests")?)
.bind(row.try_get::<i64, _>("cache_hit_requests")?)
.bind(row.try_get::<i64, _>("completed_total_requests")?)
.bind(row.try_get::<i64, _>("completed_cache_hit_requests")?)
.bind(row.try_get::<i64, _>("completed_input_tokens")?)
.bind(row.try_get::<i64, _>("completed_cache_creation_tokens")?)
.bind(row.try_get::<i64, _>("completed_cache_read_tokens")?)
.bind(row.try_get::<i64, _>("completed_total_input_context")?)
.bind(row.try_get::<f64, _>("completed_cache_creation_cost")?)
.bind(row.try_get::<f64, _>("completed_cache_read_cost")?)
.bind(row.try_get::<f64, _>("settled_total_cost")?)
.bind(row.try_get::<i64, _>("settled_total_requests")?)
.bind(row.try_get::<i64, _>("settled_input_tokens")?)
.bind(row.try_get::<i64, _>("settled_output_tokens")?)
.bind(row.try_get::<i64, _>("settled_cache_creation_tokens")?)
.bind(row.try_get::<i64, _>("settled_cache_read_tokens")?)
.bind(row.try_get::<Option<i64>, _>("settled_first_finalized_at_unix_secs")?)
.bind(row.try_get::<Option<i64>, _>("settled_last_finalized_at_unix_secs")?)
.bind(success_requests)
.bind(error_requests)
.bind(
row.try_get::<i64, _>("input_tokens")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<i64, _>("output_tokens")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<i64, _>("cache_creation_tokens")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<i64, _>("cache_read_tokens")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<f64, _>("total_cost")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<f64, _>("actual_total_cost")
.map_err(postgres_error)?,
)
.bind(
row.try_get::<f64, _>("avg_response_time_ms")
.map_err(postgres_error)?,
)
.bind(row.try_get::<i64, _>("input_tokens")?)
.bind(row.try_get::<i64, _>("output_tokens")?)
.bind(row.try_get::<i64, _>("cache_creation_tokens")?)
.bind(row.try_get::<i64, _>("cache_read_tokens")?)
.bind(row.try_get::<f64, _>("total_cost")?)
.bind(row.try_get::<f64, _>("actual_total_cost")?)
.bind(row.try_get::<f64, _>("response_time_sum_ms")?)
.bind(row.try_get::<i64, _>("response_time_samples")?)
.bind(row.try_get::<f64, _>("avg_response_time_ms")?)
.bind(true)
.bind(aggregated_at)
.bind(aggregated_at)
.bind(aggregated_at)
.execute(&mut *tx)
.await
.map_err(postgres_error)?;
.await?;
let user_rows =
upsert_stats_hourly_user_rows(&mut tx, hour_utc, hour_end, aggregated_at).await?;
let user_model_rows =
upsert_stats_hourly_user_model_rows(&mut tx, hour_utc, hour_end, aggregated_at).await?;
let model_rows =
upsert_stats_hourly_model_rows(&mut tx, hour_utc, hour_end, aggregated_at).await?;
let provider_rows =
upsert_stats_hourly_provider_rows(&mut tx, hour_utc, hour_end, aggregated_at).await?;
tx.commit().await.map_err(postgres_error)?;
tx.commit().await?;
Ok(Some(StatsHourlyAggregationSummary {
Ok(StatsHourlyAggregationSummary {
hour_utc,
total_requests,
user_rows,
user_model_rows,
model_rows,
provider_rows,
}))
})
}
async fn upsert_stats_hourly_user_rows(
@@ -113,14 +155,13 @@ async fn upsert_stats_hourly_user_rows(
hour_utc: DateTime<Utc>,
hour_end: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, DataLayerError> {
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_HOURLY_USER_SQL)
.bind(hour_utc)
.bind(hour_end)
.bind(now_utc)
.execute(&mut **tx)
.await
.map_err(postgres_error)?
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
@@ -131,14 +172,30 @@ async fn upsert_stats_hourly_model_rows(
hour_utc: DateTime<Utc>,
hour_end: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, DataLayerError> {
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_HOURLY_MODEL_SQL)
.bind(hour_utc)
.bind(hour_end)
.bind(now_utc)
.execute(&mut **tx)
.await
.map_err(postgres_error)?
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
}
async fn upsert_stats_hourly_user_model_rows(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
hour_utc: DateTime<Utc>,
hour_end: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_HOURLY_USER_MODEL_SQL)
.bind(hour_utc)
.bind(hour_end)
.bind(now_utc)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))
@@ -149,14 +206,13 @@ async fn upsert_stats_hourly_provider_rows(
hour_utc: DateTime<Utc>,
hour_end: DateTime<Utc>,
now_utc: DateTime<Utc>,
) -> Result<usize, DataLayerError> {
) -> Result<usize, sqlx::Error> {
let rows_affected = sqlx::query(UPSERT_STATS_HOURLY_PROVIDER_SQL)
.bind(hour_utc)
.bind(hour_end)
.bind(now_utc)
.execute(&mut **tx)
.await
.map_err(postgres_error)?
.await?
.rows_affected();
Ok(usize::try_from(rows_affected).unwrap_or(usize::MAX))

View File

@@ -21,6 +21,9 @@ use super::{
WALLET_DAILY_USAGE_AGGREGATION_HOUR, WALLET_DAILY_USAGE_AGGREGATION_MINUTE,
};
const STATS_DAILY_CATCH_UP_BURST_LIMIT: usize = 14;
const STATS_HOURLY_CATCH_UP_BURST_LIMIT: usize = 72;
fn log_maintenance_worker_failure(
worker: &'static str,
phase: &'static str,
@@ -110,10 +113,23 @@ pub(crate) fn spawn_stats_aggregation_worker(
Some(tokio::spawn(async move {
loop {
tokio::time::sleep(duration_until_next_stats_aggregation_run(Utc::now())).await;
if let Err(err) = run_stats_aggregation_once(&data).await {
log_maintenance_worker_failure("stats_daily_aggregation", "tick", &err);
let mut processed = 0_usize;
while processed < STATS_DAILY_CATCH_UP_BURST_LIMIT {
match run_stats_aggregation_once(&data).await {
Ok(true) => processed += 1,
Ok(false) => break,
Err(err) => {
log_maintenance_worker_failure("stats_daily_aggregation", "tick", &err);
break;
}
}
}
if processed >= STATS_DAILY_CATCH_UP_BURST_LIMIT {
continue;
}
tokio::time::sleep(duration_until_next_stats_aggregation_run(Utc::now())).await;
}
}))
}
@@ -304,10 +320,23 @@ pub(crate) fn spawn_stats_hourly_aggregation_worker(
Some(tokio::spawn(async move {
loop {
tokio::time::sleep(duration_until_next_stats_hourly_aggregation_run(Utc::now())).await;
if let Err(err) = run_stats_hourly_aggregation_once(&data).await {
log_maintenance_worker_failure("stats_hourly_aggregation", "tick", &err);
let mut processed = 0_usize;
while processed < STATS_HOURLY_CATCH_UP_BURST_LIMIT {
match run_stats_hourly_aggregation_once(&data).await {
Ok(true) => processed += 1,
Ok(false) => break,
Err(err) => {
log_maintenance_worker_failure("stats_hourly_aggregation", "tick", &err);
break;
}
}
}
if processed >= STATS_HOURLY_CATCH_UP_BURST_LIMIT {
continue;
}
tokio::time::sleep(duration_until_next_stats_hourly_aggregation_run(Utc::now())).await;
}
}))
}