Merge remote-tracking branch 'origin/pr/435' into aether-rust-pioneer

This commit is contained in:
fawney19
2026-05-14 02:02:26 +08:00
20 changed files with 1063 additions and 27 deletions

View File

@@ -18,8 +18,8 @@ pub use types::{
UsageAuditListQuery, UsageAuditSummaryQuery, UsageBodyCaptureResult, UsageBodyCaptureState,
UsageBodyCaptureStorage, UsageBodyField, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupSummary,
UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDailyHeatmapQuery,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupPreviewCounts,
UsageCleanupSummary, UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDailyHeatmapQuery,
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,

View File

@@ -1692,6 +1692,14 @@ pub trait UsageWriteRepository: Send + Sync {
let _ = (window, batch_size, auto_delete_expired_keys);
Ok(UsageCleanupSummary::default())
}
async fn preview_usage_cleanup(
&self,
window: &UsageCleanupWindow,
) -> Result<UsageCleanupPreviewCounts, crate::DataLayerError> {
let _ = window;
Ok(UsageCleanupPreviewCounts::default())
}
}
pub trait UsageRepository: UsageReadRepository + UsageWriteRepository + Send + Sync {}
@@ -1722,6 +1730,14 @@ pub struct UsageCleanupWindow {
pub log_cutoff: DateTime<Utc>,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct UsageCleanupPreviewCounts {
pub detail: u64,
pub compressed: u64,
pub header: u64,
pub log: u64,
}
fn parse_u64(value: i32, field_name: &str) -> Result<u64, crate::DataLayerError> {
u64::try_from(value).map_err(|_| {
crate::DataLayerError::UnexpectedValue(format!("invalid {field_name}: {value}"))

View File

@@ -378,8 +378,8 @@ pub(crate) use aether_data_contracts::repository::usage::{
UsageAuditAggregationGroupBy, UsageAuditAggregationQuery, UsageAuditKeywordSearchQuery,
UsageAuditListQuery, UsageAuditSummaryQuery, UsageBreakdownGroupBy, UsageBreakdownSummaryQuery,
UsageCacheAffinityHitSummaryQuery, UsageCacheAffinityIntervalGroupBy,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupSummary,
UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDailyHeatmapQuery,
UsageCacheAffinityIntervalQuery, UsageCacheHitSummaryQuery, UsageCleanupPreviewCounts,
UsageCleanupSummary, UsageCleanupWindow, UsageCostSavingsSummaryQuery, UsageDailyHeatmapQuery,
UsageDashboardDailyBreakdownQuery, UsageDashboardProviderCountsQuery,
UsageDashboardSummaryQuery, UsageErrorDistributionQuery, UsageLeaderboardGroupBy,
UsageLeaderboardQuery, UsageMonitoringErrorCountQuery, UsageMonitoringErrorListQuery,

View File

@@ -1,7 +1,8 @@
use std::io::Write;
use aether_data_contracts::repository::usage::{
parse_usage_body_ref, usage_body_ref, UsageBodyField, UsageCleanupSummary, UsageCleanupWindow,
parse_usage_body_ref, usage_body_ref, UsageBodyField, UsageCleanupPreviewCounts,
UsageCleanupSummary, UsageCleanupWindow,
};
use chrono::{DateTime, Utc};
use flate2::{write::GzEncoder, Compression};
@@ -473,6 +474,48 @@ impl SqlxUsageReadRepository {
}
}
pub async fn preview_usage_cleanup_impl(
pool: &PostgresPool,
window: &UsageCleanupWindow,
) -> Result<UsageCleanupPreviewCounts, DataLayerError> {
let detail: i64 = sqlx::query_scalar(
"SELECT COUNT(*)::bigint FROM usage WHERE created_at < $1 AND created_at >= $2",
)
.bind(window.detail_cutoff)
.bind(window.compressed_cutoff)
.fetch_one(pool)
.await
.map_err(postgres_error)?;
let compressed: i64 = sqlx::query_scalar(
"SELECT COUNT(*)::bigint FROM usage WHERE created_at < $1 AND created_at >= $2",
)
.bind(window.compressed_cutoff)
.bind(window.log_cutoff)
.fetch_one(pool)
.await
.map_err(postgres_error)?;
let header: i64 = sqlx::query_scalar(
"SELECT COUNT(*)::bigint FROM usage WHERE created_at < $1 AND created_at >= $2",
)
.bind(window.header_cutoff)
.bind(window.log_cutoff)
.fetch_one(pool)
.await
.map_err(postgres_error)?;
let log: i64 = sqlx::query_scalar("SELECT COUNT(*)::bigint FROM usage WHERE created_at < $1")
.bind(window.log_cutoff)
.fetch_one(pool)
.await
.map_err(postgres_error)?;
Ok(UsageCleanupPreviewCounts {
detail: u64::try_from(detail).unwrap_or(0),
compressed: u64::try_from(compressed).unwrap_or(0),
header: u64::try_from(header).unwrap_or(0),
log: u64::try_from(log).unwrap_or(0),
})
}
async fn delete_old_usage_records(
pool: &PostgresPool,
cutoff_time: DateTime<Utc>,

View File

@@ -8221,6 +8221,15 @@ impl UsageWriteRepository for SqlxUsageReadRepository {
) -> Result<UsageCleanupSummary, DataLayerError> {
Self::cleanup_usage(self, window, batch_size, auto_delete_expired_keys).await
}
async fn preview_usage_cleanup(
&self,
window: &UsageCleanupWindow,
) -> Result<aether_data_contracts::repository::usage::UsageCleanupPreviewCounts, DataLayerError>
{
crate::repository::usage::postgres::cleanup::preview_usage_cleanup_impl(&self.pool, window)
.await
}
}
struct StalePendingUsageRow {