fix(data): preserve API key history end to end

This commit is contained in:
elky
2026-07-18 21:58:21 +08:00
parent 03b7d573e0
commit 8fbda84acb
7 changed files with 462 additions and 90 deletions
@@ -273,35 +273,6 @@ SET is_active = FALSE,
WHERE id = $1
AND is_active IS TRUE
"#;
const NULLIFY_USAGE_API_KEY_BATCH_SQL: &str = r#"
WITH doomed AS (
SELECT id
FROM usage
WHERE api_key_id = $1
ORDER BY created_at ASC, id ASC
LIMIT $2
)
UPDATE usage AS usage_rows
SET api_key_id = NULL,
updated_at = NOW()
FROM doomed
WHERE usage_rows.id = doomed.id
"#;
const NULLIFY_REQUEST_CANDIDATE_API_KEY_BATCH_SQL: &str = r#"
WITH doomed AS (
SELECT id
FROM request_candidates
WHERE api_key_id = $1
ORDER BY created_at ASC, id ASC
LIMIT $2
)
UPDATE request_candidates AS candidate_rows
SET api_key_id = NULL,
updated_at = NOW()
FROM doomed
WHERE candidate_rows.id = doomed.id
"#;
const EXPIRED_API_KEY_PRE_CLEAN_BATCH_SIZE: usize = 2_000;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UsageDetachedBodyBlobWrite {
@@ -1373,8 +1344,6 @@ async fn cleanup_expired_api_keys(
.auto_delete_on_expiry
.unwrap_or(auto_delete_expired_keys);
if should_delete {
nullify_expired_api_key_usage_refs(pool, key.id).await?;
nullify_expired_api_key_candidate_refs(pool, key.id).await?;
sqlx::query(DISABLE_EXPIRED_API_KEY_WALLET_SQL)
.bind(key.id)
.execute(pool)
@@ -1405,46 +1374,6 @@ async fn cleanup_expired_api_keys(
Ok(cleaned)
}
async fn nullify_expired_api_key_usage_refs(
pool: &PostgresPool,
api_key_id: &str,
) -> Result<(), DataLayerError> {
loop {
let updated = sqlx::query(NULLIFY_USAGE_API_KEY_BATCH_SQL)
.bind(api_key_id)
.bind(i64::try_from(EXPIRED_API_KEY_PRE_CLEAN_BATCH_SIZE).unwrap_or(i64::MAX))
.execute(pool)
.await
.map_err(postgres_error)?
.rows_affected();
let updated = usize::try_from(updated).unwrap_or(usize::MAX);
if updated < EXPIRED_API_KEY_PRE_CLEAN_BATCH_SIZE {
break;
}
}
Ok(())
}
async fn nullify_expired_api_key_candidate_refs(
pool: &PostgresPool,
api_key_id: &str,
) -> Result<(), DataLayerError> {
loop {
let updated = sqlx::query(NULLIFY_REQUEST_CANDIDATE_API_KEY_BATCH_SQL)
.bind(api_key_id)
.bind(i64::try_from(EXPIRED_API_KEY_PRE_CLEAN_BATCH_SIZE).unwrap_or(i64::MAX))
.execute(pool)
.await
.map_err(postgres_error)?
.rows_affected();
let updated = usize::try_from(updated).unwrap_or(usize::MAX);
if updated < EXPIRED_API_KEY_PRE_CLEAN_BATCH_SIZE {
break;
}
}
Ok(())
}
fn maybe_externalize_usage_body_field(
plan: &mut UsageBodyExternalizationPlan,
request_id: &str,
@@ -6956,8 +6956,23 @@ WHERE stats_daily_api_key.date >=
.push(" AND stats_daily_api_key.api_key_id IS NOT NULL");
if let Some(user_id) = query.user_id.as_deref() {
builder
.push(" AND api_keys.user_id = ")
.push_bind(user_id.to_string());
.push(" AND (api_keys.user_id = ")
.push_bind(user_id.to_string())
.push(
r#" OR (
api_keys.id IS NULL
AND stats_daily_api_key.api_key_id IN (
SELECT identity_usage.api_key_id
FROM usage AS identity_usage
WHERE identity_usage.user_id = "#,
)
.push_bind(user_id.to_string())
.push(
r#"
AND identity_usage.api_key_id IS NOT NULL
GROUP BY identity_usage.api_key_id"#,
)
.push(")))");
}
builder.push(
" GROUP BY stats_daily_api_key.api_key_id ORDER BY stats_daily_api_key.api_key_id ASC",
@@ -886,6 +886,26 @@ async fn repository_constructs_from_lazy_pool() {
let _ = repository.transaction_runner();
}
#[test]
fn api_key_leaderboard_user_filter_keeps_aggregates_with_historical_identity_evidence() {
let source = include_str!("mod.rs");
let aggregate_path = source
.split("async fn summarize_usage_leaderboard_from_daily_aggregates")
.nth(1)
.and_then(|tail| {
tail.split("pub async fn summarize_usage_leaderboard")
.next()
})
.expect("daily aggregate leaderboard path should be present");
assert!(aggregate_path.contains("api_keys.user_id = "));
assert!(aggregate_path.contains("api_keys.id IS NULL"));
assert!(aggregate_path.contains("stats_daily_api_key.api_key_id IN ("));
assert!(aggregate_path.contains("FROM usage AS identity_usage"));
assert!(aggregate_path.contains("identity_usage.user_id = "));
assert!(aggregate_path.contains("GROUP BY identity_usage.api_key_id"));
}
#[tokio::test]
async fn validates_upsert_before_hitting_database() {
let factory = PostgresPoolFactory::new(PostgresPoolConfig {