fix(admin): run destructive purges as cleanup tasks

This commit is contained in:
fawney19
2026-05-10 00:41:21 +08:00
parent 391c2fbe5d
commit a81053e6ff
9 changed files with 682 additions and 126 deletions

View File

@@ -434,6 +434,99 @@ VALUES
assert_eq!(admin_exists, 1);
}
#[tokio::test]
async fn admin_system_request_bodies_purge_clears_inline_usage_body_fields() {
let config = SqlDatabaseConfig {
driver: DatabaseDriver::Sqlite,
url: "sqlite::memory:".to_string(),
pool: SqlPoolConfig {
max_connections: 1,
..SqlPoolConfig::default()
},
};
let backend = SqliteBackend::from_config(config).expect("backend should build");
run_sqlite_migrations(backend.pool())
.await
.expect("sqlite migrations should run");
for (column, ty) in [
("request_body", "TEXT"),
("response_body", "TEXT"),
("provider_request_body", "TEXT"),
("client_response_body", "TEXT"),
("request_body_compressed", "BLOB"),
("response_body_compressed", "BLOB"),
("provider_request_body_compressed", "BLOB"),
("client_response_body_compressed", "BLOB"),
] {
sqlx::query(&format!(r#"ALTER TABLE "usage" ADD COLUMN {column} {ty}"#))
.execute(backend.pool())
.await
.expect("legacy body column should be added");
}
sqlx::query(
r#"
INSERT INTO "usage" (
request_id,
provider_name,
model,
request_body,
response_body,
provider_request_body,
client_response_body,
request_body_compressed,
response_body_compressed,
provider_request_body_compressed,
client_response_body_compressed,
created_at_unix_ms
)
VALUES (
'request-1',
'openai',
'gpt-4.1',
'client request',
'provider response',
'provider request',
'client response',
X'01',
X'02',
X'03',
X'04',
1
)
"#,
)
.execute(backend.pool())
.await
.expect("usage row should insert");
let summary = backend
.purge_admin_system_data(AdminSystemPurgeTarget::RequestBodies)
.await
.expect("request body purge should run");
assert_eq!(summary.affected.get("usage_body_fields_cleaned"), Some(&1));
let remaining: i64 = sqlx::query_scalar(
r#"
SELECT COUNT(*)
FROM "usage"
WHERE request_body IS NOT NULL
OR response_body IS NOT NULL
OR provider_request_body IS NOT NULL
OR client_response_body IS NOT NULL
OR request_body_compressed IS NOT NULL
OR response_body_compressed IS NOT NULL
OR provider_request_body_compressed IS NOT NULL
OR client_response_body_compressed IS NOT NULL
"#,
)
.fetch_one(backend.pool())
.await
.expect("remaining body count should load");
assert_eq!(remaining, 0);
}
async fn sqlite_count(pool: &sqlx::SqlitePool, table: &str) -> i64 {
let sql = format!("SELECT COUNT(*) FROM \"{table}\"");
sqlx::query_scalar::<_, i64>(&sql)

View File

@@ -561,6 +561,17 @@ const ADMIN_USAGE_CHILD_TABLES: &[&str] = &[
"usage_settlement_snapshots",
];
const USAGE_BODY_FIELD_COLUMNS: &[&str] = &[
"request_body",
"response_body",
"provider_request_body",
"client_response_body",
"request_body_compressed",
"response_body_compressed",
"provider_request_body_compressed",
"client_response_body_compressed",
];
const ADMIN_USER_SCOPED_TABLES: &[&str] = &[
"stats_user_daily_cost_savings_model_provider",
"stats_user_daily_cost_savings_model",
@@ -720,9 +731,10 @@ WHERE request_count <> 0
}
AdminSystemPurgeTarget::RequestBodies => {
pg_delete_table(tx, "usage_body_blobs", summary).await?;
pg_execute_if_table(
pg_execute_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
UPDATE public.usage
@@ -1080,6 +1092,33 @@ WHERE request_count <> 0
}
AdminSystemPurgeTarget::RequestBodies => {
mysql_delete_table(tx, "usage_body_blobs", summary).await?;
mysql_execute_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
UPDATE `usage`
SET request_body = NULL,
response_body = NULL,
provider_request_body = NULL,
client_response_body = NULL,
request_body_compressed = NULL,
response_body_compressed = NULL,
provider_request_body_compressed = NULL,
client_response_body_compressed = NULL
WHERE request_body IS NOT NULL
OR response_body IS NOT NULL
OR provider_request_body IS NOT NULL
OR client_response_body IS NOT NULL
OR request_body_compressed IS NOT NULL
OR response_body_compressed IS NOT NULL
OR provider_request_body_compressed IS NOT NULL
OR client_response_body_compressed IS NOT NULL
"#,
summary,
)
.await?;
mysql_execute_if_table(
tx,
"usage_http_audits",
@@ -1375,6 +1414,33 @@ WHERE request_count <> 0
}
AdminSystemPurgeTarget::RequestBodies => {
sqlite_delete_table(tx, "usage_body_blobs", summary).await?;
sqlite_execute_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
UPDATE "usage"
SET request_body = NULL,
response_body = NULL,
provider_request_body = NULL,
client_response_body = NULL,
request_body_compressed = NULL,
response_body_compressed = NULL,
provider_request_body_compressed = NULL,
client_response_body_compressed = NULL
WHERE request_body IS NOT NULL
OR response_body IS NOT NULL
OR provider_request_body IS NOT NULL
OR client_response_body IS NOT NULL
OR request_body_compressed IS NOT NULL
OR response_body_compressed IS NOT NULL
OR provider_request_body_compressed IS NOT NULL
OR client_response_body_compressed IS NOT NULL
"#,
summary,
)
.await?;
sqlite_execute_if_table(
tx,
"usage_http_audits",
@@ -1573,9 +1639,10 @@ WHERE blobs.body_ref = doomed.body_ref
limit,
)
.await?;
pg_execute_batch_if_table(
pg_execute_batch_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
WITH batch AS (
@@ -1674,9 +1741,10 @@ WHERE body_ref IN (
limit,
)
.await?;
mysql_execute_batch_if_table(
mysql_execute_batch_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
UPDATE `usage`
@@ -1772,9 +1840,10 @@ WHERE body_ref IN (
limit,
)
.await?;
sqlite_execute_batch_if_table(
sqlite_execute_batch_if_table_has_columns(
tx,
"usage",
USAGE_BODY_FIELD_COLUMNS,
"usage_body_fields_cleaned",
r#"
UPDATE "usage"
@@ -1923,6 +1992,26 @@ async fn pg_execute_if_table(
Ok(())
}
async fn pg_execute_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
) -> Result<(), DataLayerError> {
if !pg_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.execute(&mut **tx)
.await
.map_postgres_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn pg_execute_batch_if_table(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
table: &str,
@@ -1944,6 +2033,28 @@ async fn pg_execute_batch_if_table(
Ok(())
}
async fn pg_execute_batch_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
limit: i64,
) -> Result<(), DataLayerError> {
if !pg_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.bind(limit)
.execute(&mut **tx)
.await
.map_postgres_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn pg_table_exists(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
table: &str,
@@ -1956,6 +2067,38 @@ async fn pg_table_exists(
.map_postgres_err()
}
async fn pg_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
table: &str,
columns: &[&str],
) -> Result<bool, DataLayerError> {
let table = checked_sql_identifier(table)?;
if !pg_table_exists(tx, table).await? {
return Ok(false);
}
for column in columns {
let column = checked_sql_identifier(column)?;
let exists: bool = sqlx::query_scalar(
"SELECT EXISTS (
SELECT 1
FROM information_schema.columns
WHERE table_schema = 'public'
AND table_name = $1
AND column_name = $2
)",
)
.bind(table)
.bind(column)
.fetch_one(&mut **tx)
.await
.map_postgres_err()?;
if !exists {
return Ok(false);
}
}
Ok(true)
}
async fn mysql_delete_table(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
@@ -2036,6 +2179,26 @@ async fn mysql_execute_if_table(
Ok(())
}
async fn mysql_execute_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
) -> Result<(), DataLayerError> {
if !mysql_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.execute(&mut **tx)
.await
.map_sql_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn mysql_execute_batch_if_table(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
@@ -2057,6 +2220,28 @@ async fn mysql_execute_batch_if_table(
Ok(())
}
async fn mysql_execute_batch_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
limit: i64,
) -> Result<(), DataLayerError> {
if !mysql_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.bind(limit)
.execute(&mut **tx)
.await
.map_sql_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn mysql_table_exists(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
@@ -2072,6 +2257,36 @@ async fn mysql_table_exists(
Ok(total > 0)
}
async fn mysql_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::MySql>,
table: &str,
columns: &[&str],
) -> Result<bool, DataLayerError> {
let table = checked_sql_identifier(table)?;
if !mysql_table_exists(tx, table).await? {
return Ok(false);
}
for column in columns {
let column = checked_sql_identifier(column)?;
let total: i64 = sqlx::query_scalar(
"SELECT COUNT(*)
FROM information_schema.columns
WHERE table_schema = DATABASE()
AND table_name = ?
AND column_name = ?",
)
.bind(table)
.bind(column)
.fetch_one(&mut **tx)
.await
.map_sql_err()?;
if total == 0 {
return Ok(false);
}
}
Ok(true)
}
async fn sqlite_delete_table(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
@@ -2152,6 +2367,26 @@ async fn sqlite_execute_if_table(
Ok(())
}
async fn sqlite_execute_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
) -> Result<(), DataLayerError> {
if !sqlite_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.execute(&mut **tx)
.await
.map_sql_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn sqlite_execute_batch_if_table(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
@@ -2173,6 +2408,28 @@ async fn sqlite_execute_batch_if_table(
Ok(())
}
async fn sqlite_execute_batch_if_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
columns: &[&str],
key: &str,
sql: &str,
summary: &mut AdminSystemPurgeSummary,
limit: i64,
) -> Result<(), DataLayerError> {
if !sqlite_table_has_columns(tx, checked_sql_identifier(table)?, columns).await? {
return Ok(());
}
let rows = sqlx::query(sql)
.bind(limit)
.execute(&mut **tx)
.await
.map_sql_err()?
.rows_affected();
summary.add(key, rows);
Ok(())
}
async fn sqlite_table_exists(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
@@ -2187,6 +2444,31 @@ async fn sqlite_table_exists(
Ok(total > 0)
}
async fn sqlite_table_has_columns(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
table: &str,
columns: &[&str],
) -> Result<bool, DataLayerError> {
let table = checked_sql_identifier(table)?;
if !sqlite_table_exists(tx, table).await? {
return Ok(false);
}
for column in columns {
let column = checked_sql_identifier(column)?;
let total: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM pragma_table_info(?) WHERE name = ?")
.bind(table)
.bind(column)
.fetch_one(&mut **tx)
.await
.map_sql_err()?;
if total == 0 {
return Ok(false);
}
}
Ok(true)
}
fn current_unix_secs() -> u64 {
chrono::Utc::now().timestamp().max(0) as u64
}