From ccc86ca88b69f9380c6cfe86c3b4c13041722edf Mon Sep 17 00:00:00 2001 From: wendrul <53628737+wendrul@users.noreply.github.com> Date: Wed, 21 Aug 2024 17:15:08 +0200 Subject: [PATCH] refactor windmill concurrent index migrations using a macro (#4275) * Add macro to run migrations * Factor windmill migrations into the new macro * Run formatter --- backend/windmill-api/src/db.rs | 394 +++++++++++++-------------------- 1 file changed, 153 insertions(+), 241 deletions(-) diff --git a/backend/windmill-api/src/db.rs b/backend/windmill-api/src/db.rs index dc8a37149f..e2cc2cfa71 100644 --- a/backend/windmill-api/src/db.rs +++ b/backend/windmill-api/src/db.rs @@ -248,6 +248,78 @@ async fn fix_flow_versioning_migration( Ok(()) } +macro_rules! run_windmill_migration { + ($migration_job_name:expr, $db:expr, $code:block) => { + { + let migration_job_name = $migration_job_name; + let db: &Pool = $db; + + let has_done_migration = sqlx::query_scalar!( + "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", + migration_job_name + ) + .fetch_one(db) + .await? + .unwrap_or(false); + if !has_done_migration { + tracing::info!("Applying {migration_job_name} migration"); + let mut tx = db.begin().await?; + let mut r = false; + while !r { + r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)") + .fetch_one(&mut *tx) + .await + .map_err(|e| { + tracing::error!("Error acquiring {migration_job_name} lock: {e:#}"); + sqlx::migrate::MigrateError::Execute(e) + })? + .unwrap_or(false); + if !r { + tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)"); + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + } + tracing::info!("acquired lock for {migration_job_name}"); + + let has_done_migration = sqlx::query_scalar!( + "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", + migration_job_name + ) + .fetch_one(db) + .await? + .unwrap_or(false); + + if !has_done_migration { + + $code + + sqlx::query!( + "INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING", + migration_job_name + ) + .execute(&mut *tx) + .await?; + tracing::info!("Finished applying {migration_job_name} migration"); + } else { + tracing::info!("migration {migration_job_name} already done"); + } + + let _ = sqlx::query("SELECT pg_advisory_unlock(4242)") + .execute(&mut *tx) + .await?; + tx.commit().await?; + tracing::info!("released lock for {migration_job_name}"); + + + tracing::info!("Finished applying {migration_job_name} migration"); + } else { + tracing::info!("migration {migration_job_name} already done"); + + } + } + }; +} + async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { // let has_done_migration = sqlx::query_scalar!( // "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = 'fix_job_completed_index')" @@ -290,284 +362,124 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { // tx.commit().await?; // } - let migration_job_name = "fix_job_completed_index_2"; - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name - ) - .fetch_one(db) - .await? - .unwrap_or(false); - if !has_done_migration { - tracing::info!("Applying {migration_job_name} migration"); - let mut tx = db.begin().await?; - let mut r = false; - while !r { - r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)") - .fetch_one(&mut *tx) - .await - .map_err(|e| { - tracing::error!("Error acquiring {migration_job_name} lock: {e:#}"); - sqlx::migrate::MigrateError::Execute(e) - })? - .unwrap_or(false); - if !r { - tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)"); - tokio::time::sleep(std::time::Duration::from_secs(5)).await; - } - } + run_windmill_migration!("fix_job_completed_index_2", &db, { + // sqlx::query( + // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new_2 ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, created_at DESC)" + // ).execute(db).await?; - tracing::info!("acquired lock for {migration_job_name}"); - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name + // sqlx::query( + // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_started_at_new ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, started_at DESC)" + // ).execute(db).await?; + + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at") + .execute(db) + .await?; + + sqlx::query( + "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new", ) - .fetch_one(db) - .await? - .unwrap_or(false); + .execute(db) + .await?; + }); - if !has_done_migration { - // sqlx::query( - // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new_2 ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, created_at DESC)" - // ).execute(db).await?; - - // sqlx::query( - // "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_started_at_new ON completed_job (workspace_id, job_kind, success, is_skipped, is_flow_step, started_at DESC)" - // ).execute(db).await?; - - sqlx::query( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at", - ) + run_windmill_migration!("fix_job_completed_index_3", &db, { + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_job_on_schedule_path") .execute(db) .await?; - sqlx::query( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new", - ) + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS concurrency_limit_stats_queue") .execute(db) .await?; - sqlx::query!( - "INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING", - migration_job_name - ) - .execute(&mut *tx) - .await?; - tracing::info!("Finished applying {migration_job_name} migration"); - } else { - tracing::info!("migration {migration_job_name} already done"); - } - - let _ = sqlx::query("SELECT pg_advisory_unlock(4242)") - .execute(&mut *tx) - .await?; - tx.commit().await?; - tracing::info!("released lock for {migration_job_name}"); - } - - let migration_job_name = "fix_job_completed_index_3"; - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name - ) - .fetch_one(db) - .await? - .unwrap_or(false); - if !has_done_migration { - tracing::info!("Applying {migration_job_name} migration"); - let mut tx = db.begin().await?; - let mut r = false; - while !r { - r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)") - .fetch_one(&mut *tx) - .await - .map_err(|e| { - tracing::error!("Error acquiring {migration_job_name} lock: {e:#}"); - sqlx::migrate::MigrateError::Execute(e) - })? - .unwrap_or(false); - if !r { - tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)"); - tokio::time::sleep(std::time::Duration::from_secs(5)).await; - } - } - tracing::info!("acquired lock for {migration_job_name}"); - - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name - ) - .fetch_one(db) - .await? - .unwrap_or(false); - - if !has_done_migration { - sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_job_on_schedule_path") - .execute(db) - .await?; - - sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS concurrency_limit_stats_queue") - .execute(db) - .await?; - - sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index") - .execute(db) - .await?; - - sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_on_created") - .execute(db) - .await?; - - sqlx::query!( - "INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING", - migration_job_name - ) - .execute(&mut *tx) - .await?; - tracing::info!("Finished applying {migration_job_name} migration"); - } else { - tracing::info!("migration {migration_job_name} already done"); - } - - let _ = sqlx::query("SELECT pg_advisory_unlock(4242)") - .execute(&mut *tx) - .await?; - tx.commit().await?; - tracing::info!("released lock for {migration_job_name}"); - } - - let migration_job_name = "fix_job_completed_index_4"; - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name - ) - .fetch_one(db) - .await? - .unwrap_or(false); - if !has_done_migration { - tracing::info!("Applying {migration_job_name} migration"); - let mut tx = db.begin().await?; - let mut r = false; - while !r { - r = sqlx::query_scalar!("SELECT pg_try_advisory_lock(4242)") - .fetch_one(&mut *tx) - .await - .map_err(|e| { - tracing::error!("Error acquiring {migration_job_name} lock: {e:#}"); - sqlx::migrate::MigrateError::Execute(e) - })? - .unwrap_or(false); - if !r { - tracing::info!("PG {migration_job_name} lock already acquired by another server or worker, retrying in 5s. (look for the advisory lock in pg_lock with granted = true)"); - tokio::time::sleep(std::time::Duration::from_secs(5)).await; - } - } - tracing::info!("acquired lock for {migration_job_name}"); - - let has_done_migration = sqlx::query_scalar!( - "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", - migration_job_name - ) - .fetch_one(db) - .await? - .unwrap_or(false); - - if !has_done_migration { - let mut i = 1; - tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_3 ON completed_job (workspace_id, created_at DESC)") - .execute(db) - .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); - - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_8 ON completed_job (workspace_id, created_at DESC) where job_kind in ('deploymentcallback') AND parent_job IS NULL") + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_9 ON completed_job (workspace_id, created_at DESC) where job_kind in ('dependencies', 'flowdependencies', 'appdependencies') AND parent_job IS NULL") + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_on_created") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + }); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_5 ON completed_job (workspace_id, created_at DESC) where job_kind in ('preview', 'flowpreview') AND parent_job IS NULL") + run_windmill_migration!("fix_job_completed_index_4", &db, { + let migration_job_name = "fix_job_completed_index_4"; + let mut i = 1; + tracing::info!("step {i} of {migration_job_name} migration"); + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_3 ON completed_job (workspace_id, created_at DESC)") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_6 ON completed_job (workspace_id, created_at DESC) where job_kind in ('script', 'flow') AND parent_job IS NULL") + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_8 ON completed_job (workspace_id, created_at DESC) where job_kind in ('deploymentcallback') AND parent_job IS NULL") + .execute(db) + .await?; + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); + + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_9 ON completed_job (workspace_id, created_at DESC) where job_kind in ('dependencies', 'flowdependencies', 'appdependencies') AND parent_job IS NULL") + .execute(db) + .await?; + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); + + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_5 ON completed_job (workspace_id, created_at DESC) where job_kind in ('preview', 'flowpreview') AND parent_job IS NULL") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_7 ON completed_job (workspace_id, success, created_at DESC) where job_kind in ('script', 'flow') AND parent_job IS NULL") + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_6 ON completed_job (workspace_id, created_at DESC) where job_kind in ('script', 'flow') AND parent_job IS NULL") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_started_at_new_2 ON completed_job (workspace_id, started_at DESC)") + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_created_at_new_7 ON completed_job (workspace_id, success, created_at DESC) where job_kind in ('script', 'flow') AND parent_job IS NULL") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists root_job_index_by_path_2 ON completed_job (workspace_id, script_path, created_at desc) WHERE parent_job IS NULL") + sqlx::query("create index concurrently if not exists ix_completed_job_workspace_id_started_at_new_2 ON completed_job (workspace_id, started_at DESC)") + .execute(db) + .await?; + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); + + sqlx::query("create index concurrently if not exists root_job_index_by_path_2 ON completed_job (workspace_id, script_path, created_at desc) WHERE parent_job IS NULL") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query("create index concurrently if not exists ix_completed_job_created_at ON completed_job (created_at DESC)") + sqlx::query("create index concurrently if not exists ix_completed_job_created_at ON completed_job (created_at DESC)") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); - sqlx::query( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_2", - ) + sqlx::query( + "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_2", + ) + .execute(db) + .await?; + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); + + sqlx::query( + "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_started_at_new", + ) + .execute(db) + .await?; + i += 1; + tracing::info!("step {i} of {migration_job_name} migration"); + + sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path") .execute(db) .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); - - sqlx::query( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_started_at_new", - ) - .execute(db) - .await?; - i += 1; - tracing::info!("step {i} of {migration_job_name} migration"); - - sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path") - .execute(db) - .await?; - - sqlx::query!( - "INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING", - migration_job_name - ) - .execute(&mut *tx) - .await?; - tracing::info!("Finished applying {migration_job_name} migration"); - } else { - tracing::info!("migration {migration_job_name} already done"); - } - - let _ = sqlx::query("SELECT pg_advisory_unlock(4242)") - .execute(&mut *tx) - .await?; - tx.commit().await?; - tracing::info!("released lock for {migration_job_name}"); - } + }); Ok(()) }