perf: retire the legacy pre-partitioning audit table (#11240)

* perf: index audit logs by operation and limit each audit table on its own

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf: key the audit operation index on id with a trailing timestamp

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf: index the legacy audit table by operation too

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf: drop the legacy audit index and per-table limit from this PR

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf: build the audit operation index in a background migration

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf: retire the legacy pre-partitioning audit table

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: bound the legacy audit retirement's lock waits and keep the backfill guard

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: bound the audit index parent's lock wait and retry background steps

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* chore: update ee-repo-ref to 5370ae3a95a3dc10171f72da74286a654996eee6

This commit updates the EE repository reference after PR #818 was merged in windmill-ee-private.

Previous ee-repo-ref: aa367619a970de231f5c10cad470d9d42169fff8

New ee-repo-ref: 5370ae3a95a3dc10171f72da74286a654996eee6

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
Ruben Fiszel
2026-09-21 15:24:07 +00:00
committed by GitHub
co-authored by Claude Opus 5 windmill-internal-app[bot]
parent 3b8c3264ec
commit 9f2685d788
12 changed files with 218 additions and 120 deletions
@@ -1,38 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "WITH all_audit AS (SELECT username, operation, timestamp FROM audit_partitioned UNION ALL SELECT username, operation, timestamp FROM audit),\n active_users as (SELECT distinct username as email FROM all_audit WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),\n active_authors as (SELECT distinct email FROM usr WHERE usr.operator IS false AND email IN (SELECT email FROM active_users)),\n active_authors_agg as (SELECT array_agg(email) as authors FROM active_authors),\n active_ops_agg as (SELECT array_agg(email) as operators from active_users WHERE email NOT IN (SELECT email FROM active_authors))\n SELECT active_authors_agg.authors, active_ops_agg.operators, array_length(active_authors_agg.authors, 1) as author_count, array_length(active_ops_agg.operators, 1) as operator_count FROM active_authors_agg, active_ops_agg",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "authors",
"type_info": "VarcharArray"
},
{
"ordinal": 1,
"name": "operators",
"type_info": "VarcharArray"
},
{
"ordinal": 2,
"name": "author_count",
"type_info": "Int4"
},
{
"ordinal": 3,
"name": "operator_count",
"type_info": "Int4"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null,
null,
null
]
},
"hash": "359cd29f531d263a8cf7205e0869229a610767087f01e0154be8da0620fa114b"
}
@@ -1,38 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH capped AS (\n SELECT timestamp, operation, resource, parameters\n FROM audit_partitioned\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.capped'\n UNION ALL\n SELECT timestamp, operation, resource, parameters\n FROM audit\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.capped'\n ORDER BY timestamp DESC\n LIMIT 200\n ), uncapped AS (\n SELECT timestamp, operation, resource, parameters\n FROM audit_partitioned\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.uncapped'\n UNION ALL\n SELECT timestamp, operation, resource, parameters\n FROM audit\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.uncapped'\n ORDER BY timestamp DESC\n LIMIT 200\n )\n SELECT timestamp AS \"timestamp!\",\n operation::text AS \"operation!\",\n resource AS workspace_id,\n parameters\n FROM (SELECT * FROM capped UNION ALL SELECT * FROM uncapped) e\n ORDER BY timestamp DESC\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "timestamp!",
"type_info": "Timestamptz"
},
{
"ordinal": 1,
"name": "operation!",
"type_info": "Text"
},
{
"ordinal": 2,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "parameters",
"type_info": "Jsonb"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null,
null,
null
]
},
"hash": "6fcdb09cdd7fedd7e54fdc0e49203f453fc1b85272fe212c0e1bf0ebd28bf58a"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "WITH active_users AS (SELECT distinct username as email FROM (SELECT username, timestamp, operation FROM audit_partitioned UNION ALL SELECT username, timestamp, operation FROM audit) AS a WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),\n authors as (SELECT distinct email FROM usr WHERE usr.operator IS false)\n SELECT email as \"email!\", (email NOT IN (SELECT email FROM authors)) as operator_only, NULL::bool as is_workspace_admin, login_type::text, verified as \"verified!\", super_admin as \"super_admin!\", devops as \"devops!\", name, company, username, first_time_user as \"first_time_user!\", role_source as \"role_source!\", disabled as \"disabled!\", NULL::text as workspace_id\n FROM password\n WHERE email IN (SELECT email FROM active_users)\n UNION ALL\n SELECT email as \"email!\", operator as operator_only, is_admin as is_workspace_admin, 'service_account'::text as login_type, true as \"verified!\", false as \"super_admin!\", false as \"devops!\", NULL::text as name, NULL::text as company, username, false as \"first_time_user!\", 'service_account'::text as \"role_source!\", disabled as \"disabled!\", workspace_id\n FROM usr\n WHERE is_service_account IS true\n ORDER BY \"super_admin!\" DESC, \"devops!\" DESC\n LIMIT $1 OFFSET $2",
"query": "WITH active_users AS (SELECT distinct username as email FROM audit_partitioned WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),\n authors as (SELECT distinct email FROM usr WHERE usr.operator IS false)\n SELECT email as \"email!\", (email NOT IN (SELECT email FROM authors)) as operator_only, NULL::bool as is_workspace_admin, login_type::text, verified as \"verified!\", super_admin as \"super_admin!\", devops as \"devops!\", name, company, username, first_time_user as \"first_time_user!\", role_source as \"role_source!\", disabled as \"disabled!\", NULL::text as workspace_id\n FROM password\n WHERE email IN (SELECT email FROM active_users)\n UNION ALL\n SELECT email as \"email!\", operator as operator_only, is_admin as is_workspace_admin, 'service_account'::text as login_type, true as \"verified!\", false as \"super_admin!\", false as \"devops!\", NULL::text as name, NULL::text as company, username, false as \"first_time_user!\", 'service_account'::text as \"role_source!\", disabled as \"disabled!\", workspace_id\n FROM usr\n WHERE is_service_account IS true\n ORDER BY \"super_admin!\" DESC, \"devops!\" DESC\n LIMIT $1 OFFSET $2",
"describe": {
"columns": [
{
@@ -97,5 +97,5 @@
null
]
},
"hash": "8aae160c589adf02e20b7e6ba860b66fb08f156776b96b7f3701aed330a7d000"
"hash": "d50d7a93f5c97570e3ef51ef542de41e9c1a9230911a2331d7127572722e340e"
}
@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM audit WHERE timestamp <= now() - ($1::bigint::text || ' s')::interval",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Int8"
]
},
"nullable": []
},
"hash": "d8186f0cee285aa50db7626409aec7e0504b068ffa8bb185d9384ce1422fd3d1"
}
@@ -0,0 +1,38 @@
{
"db_name": "PostgreSQL",
"query": "WITH active_users as (SELECT distinct username as email FROM audit_partitioned WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),\n active_authors as (SELECT distinct email FROM usr WHERE usr.operator IS false AND email IN (SELECT email FROM active_users)),\n active_authors_agg as (SELECT array_agg(email) as authors FROM active_authors),\n active_ops_agg as (SELECT array_agg(email) as operators from active_users WHERE email NOT IN (SELECT email FROM active_authors))\n SELECT active_authors_agg.authors, active_ops_agg.operators, array_length(active_authors_agg.authors, 1) as author_count, array_length(active_ops_agg.operators, 1) as operator_count FROM active_authors_agg, active_ops_agg",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "authors",
"type_info": "VarcharArray"
},
{
"ordinal": 1,
"name": "operators",
"type_info": "VarcharArray"
},
{
"ordinal": 2,
"name": "author_count",
"type_info": "Int4"
},
{
"ordinal": 3,
"name": "operator_count",
"type_info": "Int4"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null,
null,
null
]
},
"hash": "dc37f9ad685a2f2bd781be8678e0548a860815f30e866362344434f23aa7c83c"
}
@@ -0,0 +1,38 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH capped AS (\n SELECT timestamp, operation, resource, parameters\n FROM audit_partitioned\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.capped'\n ORDER BY timestamp DESC\n LIMIT 200\n ), uncapped AS (\n SELECT timestamp, operation, resource, parameters\n FROM audit_partitioned\n WHERE workspace_id = 'admins'\n AND operation = 'workspace_fairness.uncapped'\n ORDER BY timestamp DESC\n LIMIT 200\n )\n SELECT timestamp AS \"timestamp!\",\n operation::text AS \"operation!\",\n resource AS workspace_id,\n parameters\n FROM (SELECT * FROM capped UNION ALL SELECT * FROM uncapped) e\n ORDER BY timestamp DESC\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "timestamp!",
"type_info": "Timestamptz"
},
{
"ordinal": 1,
"name": "operation!",
"type_info": "Text"
},
{
"ordinal": 2,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "parameters",
"type_info": "Jsonb"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null,
null,
null
]
},
"hash": "dfa54341a687431ea46ec842ab4167bcfd875cb4646ef859a7b855428bf48290"
}
+1 -1
View File
@@ -1 +1 @@
6f3161192ab0eba36a8630b92e6765a36a3326fc
5370ae3a95a3dc10171f72da74286a654996eee6
-14
View File
@@ -1953,20 +1953,6 @@ pub async fn delete_expired_items(db: &DB) -> () {
tracing::error!("Error deleting expired shared AI artifacts: {:?}", e);
}
let audit_retention_days = audit_log_retention_days().await;
let audit_retention_secs: i64 = audit_retention_days * 60 * 60 * 24;
// Clean up old (non-partitioned) audit table — will eventually be empty and dropped
if let Err(e) = sqlx::query_scalar!(
"DELETE FROM audit WHERE timestamp <= now() - ($1::bigint::text || ' s')::interval",
audit_retention_secs,
)
.fetch_all(db)
.await
{
tracing::error!("Error deleting audit log: {:?}", e);
}
if let Err(e) = sqlx::query_scalar!(
"DELETE FROM autoscaling_event WHERE applied_at <= now() - ($1::bigint::text || ' s')::interval",
30 * 24 * 60 * 60, // 30 days
+1 -1
View File
@@ -56,7 +56,7 @@ app_version_lite: id(bigint), value(jsonb)
FK: (id) -> app_version(id)
asset: workspace_id(char), path(char), kind(asset_kind), usage_access_type(asset_access_type), usage_path(char), usage_kind(asset_usage_kind), created_at(ts), id(bigint), columns(jsonb)
FK: (workspace_id) -> workspace(id)
audit: workspace_id(char), id(bigint), timestamp(ts), username(char), operation(char), action_kind(action_kind), resource(char), parameters(jsonb), email(char), span(char)
audit_partitioned: workspace_id(char), id(bigint), timestamp(ts), username(char), operation(char), action_kind(action_kind), resource(char), parameters(jsonb), email(char), span(char)
autoscaling_event: id(int), worker_group(text), event_type(autoscaling_event_type), desired_workers(int), applied_at(ts), reason(text)
capture: workspace_id(char), path(char), created_at(ts), created_by(char), main_args(jsonb), is_flow(bool), trigger_kind(trigger_kind), preprocessor_args(jsonb), id(bigint)
FK: (workspace_id) -> workspace(id)
+1 -1
View File
@@ -530,7 +530,7 @@ async fn list_users_as_super_admin(
let rows = if active_only.is_some_and(|x| x) {
sqlx::query_as!(
GlobalUserInfo,
r#"WITH active_users AS (SELECT distinct username as email FROM (SELECT username, timestamp, operation FROM audit_partitioned UNION ALL SELECT username, timestamp, operation FROM audit) AS a WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),
r#"WITH active_users AS (SELECT distinct username as email FROM audit_partitioned WHERE timestamp > NOW() - INTERVAL '1 month' AND (operation = 'users.login' OR operation = 'oauth.login' OR operation = 'users.token.refresh')),
authors as (SELECT distinct email FROM usr WHERE usr.operator IS false)
SELECT email as "email!", (email NOT IN (SELECT email FROM authors)) as operator_only, NULL::bool as is_workspace_admin, login_type::text, verified as "verified!", super_admin as "super_admin!", devops as "devops!", name, company, username, first_time_user as "first_time_user!", role_source as "role_source!", disabled as "disabled!", NULL::text as workspace_id
FROM password
-10
View File
@@ -436,11 +436,6 @@ async fn get_workspace_fairness_events(
WITH capped AS (
SELECT timestamp, operation, resource, parameters
FROM audit_partitioned
WHERE workspace_id = 'admins'
AND operation = 'workspace_fairness.capped'
UNION ALL
SELECT timestamp, operation, resource, parameters
FROM audit
WHERE workspace_id = 'admins'
AND operation = 'workspace_fairness.capped'
ORDER BY timestamp DESC
@@ -448,11 +443,6 @@ async fn get_workspace_fairness_events(
), uncapped AS (
SELECT timestamp, operation, resource, parameters
FROM audit_partitioned
WHERE workspace_id = 'admins'
AND operation = 'workspace_fairness.uncapped'
UNION ALL
SELECT timestamp, operation, resource, parameters
FROM audit
WHERE workspace_id = 'admins'
AND operation = 'workspace_fairness.uncapped'
ORDER BY timestamp DESC
+137 -1
View File
@@ -160,7 +160,7 @@ async fn run_background_migrations(db: &DB) -> Result<(), Error> {
return Ok(());
}
for step in [AUDIT_OPERATION_INDEX] {
for step in [AUDIT_OPERATION_INDEX, RETIRE_LEGACY_AUDIT] {
if background_migration_done(&mut conn, step).await? {
continue;
}
@@ -170,6 +170,7 @@ async fn run_background_migrations(db: &DB) -> Result<(), Error> {
loop {
let run = match step {
AUDIT_OPERATION_INDEX => create_audit_operation_index(&mut conn).await,
RETIRE_LEGACY_AUDIT => retire_legacy_audit_table(&mut conn).await,
_ => unreachable!("background migration {step} has no step function"),
};
match run {
@@ -190,6 +191,7 @@ async fn run_background_migrations(db: &DB) -> Result<(), Error> {
}
const AUDIT_OPERATION_INDEX: &str = "audit_partitioned_workspace_operation_index";
const RETIRE_LEGACY_AUDIT: &str = "retire_legacy_audit_table";
// Short on purpose: a statement waiting for its lock is also a wait for every audit insert
// queued behind it.
@@ -277,3 +279,137 @@ async fn create_audit_operation_index(conn: &mut PgConnection) -> Result<(), Err
}
Ok(())
}
const AUDIT_COLUMNS: &str =
"workspace_id, id, timestamp, username, operation, action_kind, resource, parameters, email, span";
/// Moves the last 30 days of the pre-partitioning `audit` table into daily partitions and drops
/// it with anything older. An empty `audit` view takes its place for servers still running an
/// older version, which read it in a `UNION ALL` with `audit_partitioned`.
async fn retire_legacy_audit_table(conn: &mut PgConnection) -> Result<(), Error> {
let is_table = sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM pg_class WHERE oid = to_regclass('audit') AND relkind = 'r')",
)
.fetch_one(&mut *conn)
.await?;
if !is_table {
return Ok(());
}
// Creating a partition and re-owning the sequence queue every audit insert (so every job
// push) behind them while they wait for their own lock, e.g. on a long-running reader, and
// the DROP queues older servers' reads the same way. Give up early and retry instead.
conn.execute(STEP_LOCK_TIMEOUT).await?;
let days: Vec<chrono::NaiveDate> = sqlx::query_scalar(
"SELECT DISTINCT timestamp::date FROM audit
WHERE timestamp > now() - interval '30 days' ORDER BY 1",
)
.fetch_all(&mut *conn)
.await?;
for day in days {
let next = day + chrono::Duration::days(1);
// Created outside any transaction: a partition created inside one keeps a lock on
// audit_partitioned that blocks every job push's audit insert until it commits.
conn.execute(
format!(
"CREATE TABLE IF NOT EXISTS \"audit_{}\" PARTITION OF audit_partitioned \
FOR VALUES FROM ('{day}') TO ('{next}')",
day.format("%Y%m%d")
)
.as_str(),
)
.await?;
// One statement per day, so a row is always in exactly one of the two tables.
let moved = sqlx::query(&format!(
"WITH moved AS (
DELETE FROM audit WHERE timestamp >= $1::date AND timestamp < $1::date + 1
RETURNING {AUDIT_COLUMNS}
)
INSERT INTO audit_partitioned ({AUDIT_COLUMNS}) SELECT {AUDIT_COLUMNS} FROM moved"
))
.bind(day)
.execute(&mut *conn)
.await?;
tracing::info!(
"Moved {} legacy audit rows of {day} into audit_partitioned",
moved.rows_affected()
);
}
// audit_partitioned draws its ids from this sequence, which dropping its owner would drop.
// Not in the transaction below: its lock stops every insert's nextval until commit, and the
// DROP can wait there on readers of the old table.
conn.execute("ALTER SEQUENCE audit_id_seq OWNED BY audit_partitioned.id")
.await?;
let mut tx = conn.begin().await?;
// Servers older than audit partitioning still write to the table during a rolling upgrade:
// rows they added after their day was moved above are swept up under the lock.
tx.execute("LOCK TABLE audit IN ACCESS EXCLUSIVE MODE")
.await?;
tx.execute(
format!(
"WITH moved AS (
DELETE FROM audit WHERE timestamp > now() - interval '30 days'
RETURNING {AUDIT_COLUMNS}
)
INSERT INTO audit_partitioned ({AUDIT_COLUMNS}) SELECT {AUDIT_COLUMNS} FROM moved"
)
.as_str(),
)
.await?;
tx.execute("DROP TABLE audit").await?;
tx.execute(
format!("CREATE VIEW audit AS SELECT {AUDIT_COLUMNS} FROM audit_partitioned WHERE false")
.as_str(),
)
.await?;
tx.execute("GRANT ALL ON audit TO windmill_user, windmill_admin")
.await?;
tx.commit().await?;
conn.execute("RESET lock_timeout").await?;
tracing::info!("Retired the legacy audit table");
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[sqlx::test(migrations = "../migrations")]
async fn retire_legacy_audit_table_keeps_the_last_30_days(db: DB) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO audit (workspace_id, id, timestamp, username, operation, action_kind)
VALUES ('w', -2, now() - interval '3 days', 'u', 'recent', 'execute'),
('w', -1, now() - interval '60 days', 'u', 'old', 'execute')",
)
.execute(&db)
.await?;
retire_legacy_audit_table(&mut *db.acquire().await?).await?;
let moved: Vec<(i64, String)> = sqlx::query_as(
"SELECT id, operation::text FROM audit_partitioned WHERE workspace_id = 'w'",
)
.fetch_all(&db)
.await?;
assert_eq!(moved, vec![(-2, "recent".to_string())]);
// The id sequence outlives the table that owned it.
sqlx::query(
"INSERT INTO audit_partitioned (workspace_id, username, operation, action_kind)
VALUES ('w', 'u', 'after', 'execute')",
)
.execute(&db)
.await?;
// Servers on an older version still read `audit` in a UNION with audit_partitioned.
let seen: i64 = sqlx::query_scalar(
"SELECT count(*) FROM (SELECT * FROM audit_partitioned UNION ALL SELECT * FROM audit) a
WHERE workspace_id = 'w'",
)
.fetch_one(&db)
.await?;
assert_eq!(seen, 2);
Ok(())
}
}