From 9e27ee8b576aaa2ecaf9c0f47cd5e145db589338 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 11 Jul 2024 01:15:12 +0200 Subject: [PATCH] improve try migration lock --- ...bf9e96a48e45daf4a3c54f7b81a7b92d39431.json | 20 +++++++++++++++++++ backend/windmill-api/src/db.rs | 19 ++++++++++++++---- 2 files changed, 35 insertions(+), 4 deletions(-) create mode 100644 backend/.sqlx/query-9dd0cf1627f9e767af8758f26d8bf9e96a48e45daf4a3c54f7b81a7b92d39431.json diff --git a/backend/.sqlx/query-9dd0cf1627f9e767af8758f26d8bf9e96a48e45daf4a3c54f7b81a7b92d39431.json b/backend/.sqlx/query-9dd0cf1627f9e767af8758f26d8bf9e96a48e45daf4a3c54f7b81a7b92d39431.json new file mode 100644 index 0000000000..ff9881bb4e --- /dev/null +++ b/backend/.sqlx/query-9dd0cf1627f9e767af8758f26d8bf9e96a48e45daf4a3c54f7b81a7b92d39431.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT pg_try_advisory_lock(4242)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "pg_try_advisory_lock", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + null + ] + }, + "hash": "9dd0cf1627f9e767af8758f26d8bf9e96a48e45daf4a3c54f7b81a7b92d39431" +} diff --git a/backend/windmill-api/src/db.rs b/backend/windmill-api/src/db.rs index bc4045cb7f..1ee7705034 100644 --- a/backend/windmill-api/src/db.rs +++ b/backend/windmill-api/src/db.rs @@ -248,10 +248,21 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { if !has_done_migration { tracing::info!("Applying fix_job_completed_index migration"); let mut tx = db.begin().await?; - let _ = sqlx::query("SELECT pg_advisory_lock(4242)") - .execute(&mut *tx) - .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 fix_job_completed_index lock: {e:#}"); + sqlx::migrate::MigrateError::Execute(e) + })? + .unwrap_or(false); + if !r { + tracing::info!("PG fix_job_completed_index_migration 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; + } + } sqlx::query( "CREATE INDEX CONCURRENTLY IF NOT EXISTS ix_completed_job_workspace_id_created_at_new ON completed_job (workspace_id, job_kind, is_skipped, is_flow_step, created_at DESC, started_at DESC)" ).execute(db).await?;