diff --git a/backend/tests/fixtures/base.sql b/backend/tests/fixtures/base.sql index 0e38b5c7bd..5e96b75b99 100644 --- a/backend/tests/fixtures/base.sql +++ b/backend/tests/fixtures/base.sql @@ -53,183 +53,74 @@ EXECUTE FUNCTION "notify_queue" (); WHEN (NEW.flow_status IS DISTINCT FROM OLD.flow_status) EXECUTE FUNCTION "notify_queue" (); --- TODO(uael): remove before phase 4 -CREATE OR REPLACE FUNCTION zzz_v2_job_queue_integrity_check() RETURNS TRIGGER AS $$ -DECLARE job v2_job; -DECLARE job_runtime v2_job_runtime; -DECLARE job_status v2_job_status; -BEGIN - IF (OLD.canceled_by IS NOT NULL) IS DISTINCT FROM OLD.__canceled THEN - RAISE EXCEPTION 'canceled mismatch'; - END IF; - -- v2_job: - SELECT * INTO job FROM v2_job WHERE id = OLD.id; - IF job.tag IS DISTINCT FROM OLD.tag THEN - RAISE EXCEPTION 'tag mismatch'; - END IF; - IF job.workspace_id IS DISTINCT FROM OLD.workspace_id THEN - RAISE EXCEPTION 'workspace_id mismatch'; - END IF; - IF job.created_at IS DISTINCT FROM OLD.created_at THEN - RAISE EXCEPTION 'created_at mismatch'; - END IF; - IF job.created_by IS DISTINCT FROM OLD.__created_by THEN - RAISE EXCEPTION 'created_by mismatch'; - END IF; - IF job.permissioned_as IS DISTINCT FROM OLD.__permissioned_as THEN - RAISE EXCEPTION 'permissioned_as mismatch'; - END IF; - IF job.permissioned_as_email IS DISTINCT FROM OLD.__email THEN - RAISE EXCEPTION 'permissioned_as_email mismatch'; - END IF; - IF job.kind IS DISTINCT FROM OLD.__job_kind THEN - RAISE EXCEPTION 'kind mismatch'; - END IF; - IF job.runnable_id IS DISTINCT FROM OLD.__script_hash THEN - RAISE EXCEPTION 'runnable_id mismatch'; - END IF; - IF job.runnable_path IS DISTINCT FROM OLD.__script_path THEN - RAISE EXCEPTION 'runnable_path mismatch'; - END IF; - IF job.parent_job IS DISTINCT FROM OLD.__parent_job THEN - RAISE EXCEPTION 'parent_job mismatch'; - END IF; - IF job.script_lang IS DISTINCT FROM OLD.__language THEN - RAISE EXCEPTION 'script_lang mismatch'; - END IF; - IF job.script_entrypoint_override IS DISTINCT FROM NULLIF(OLD.__args->>'_ENTRYPOINT_OVERRIDE', '__WM_PREPROCESSOR') - AND OLD.__args->>'reason' IS DISTINCT FROM 'PREPROCESSOR_ARGS_ARE_DISCARDED' - THEN - RAISE EXCEPTION 'script_entrypoint_override mismatch'; - END IF; - IF job.flow_step_id IS DISTINCT FROM OLD.__flow_step_id THEN - RAISE EXCEPTION 'flow_step_id mismatch'; - END IF; - IF (job.flow_step_id IS NOT NULL) IS DISTINCT FROM OLD.__is_flow_step THEN - RAISE EXCEPTION 'is_flow_step mismatch'; - END IF; - IF job.flow_innermost_root_job IS DISTINCT FROM OLD.__root_job THEN - RAISE EXCEPTION 'flow_innermost_root_job mismatch'; - END IF; - IF job.trigger IS DISTINCT FROM OLD.__schedule_path THEN - RAISE EXCEPTION 'trigger mismatch'; - END IF; - IF job.same_worker IS DISTINCT FROM OLD.__same_worker THEN - RAISE EXCEPTION 'same_worker mismatch'; - END IF; - IF job.visible_to_owner IS DISTINCT FROM OLD.__visible_to_owner THEN - RAISE EXCEPTION 'visible_to_owner mismatch'; - END IF; - IF job.concurrent_limit IS DISTINCT FROM OLD.__concurrent_limit THEN - RAISE EXCEPTION 'concurrent_limit mismatch'; - END IF; - IF job.concurrency_time_window_s IS DISTINCT FROM OLD.__concurrency_time_window_s THEN - RAISE EXCEPTION 'concurrency_time_window_s mismatch'; - END IF; - IF job.cache_ttl IS DISTINCT FROM OLD.__cache_ttl THEN - RAISE EXCEPTION 'cache_ttl mismatch'; - END IF; - IF job.timeout IS DISTINCT FROM OLD.__timeout THEN - RAISE EXCEPTION 'timeout mismatch'; - END IF; - IF job.priority IS DISTINCT FROM OLD.priority THEN - RAISE EXCEPTION 'priority mismatch'; - END IF; - IF job.args::TEXT IS DISTINCT FROM OLD.__args::TEXT AND OLD.__args->>'_ENTRYPOINT_OVERRIDE' IS DISTINCT FROM '__WM_PREPROCESSOR' THEN - RAISE EXCEPTION 'args mismatch'; - END IF; - IF job.pre_run_error IS DISTINCT FROM OLD.__pre_run_error THEN - RAISE EXCEPTION 'pre_run_error mismatch'; - END IF; - -- v2_job_runtime: - SELECT * INTO job_runtime FROM v2_job_runtime WHERE id = OLD.id; - IF job_runtime.ping IS DISTINCT FROM OLD.__last_ping THEN - RAISE EXCEPTION 'ping mismatch'; - END IF; - IF job_runtime.memory_peak IS DISTINCT FROM OLD.__mem_peak THEN - RAISE EXCEPTION 'memory_peak mismatch'; - END IF; - -- v2_job_status: - IF EXISTS(SELECT 1 FROM v2_job_status WHERE id = OLD.id) THEN - SELECT * INTO job_status FROM v2_job_status WHERE id = OLD.id; - IF COALESCE(job_status.flow_status, job_status.workflow_as_code_status)::TEXT IS DISTINCT FROM OLD.__flow_status::TEXT - THEN - RAISE EXCEPTION 'flow_status mismatch'; - END IF; - IF job_status.flow_leaf_jobs::TEXT IS DISTINCT FROM OLD.__leaf_jobs::TEXT THEN - RAISE EXCEPTION 'leaf_jobs mismatch'; - END IF; - END IF; - RETURN OLD; -END $$ LANGUAGE PLPGSQL; +-- Apply phase 4: +DROP FUNCTION IF EXISTS v2_job_after_update CASCADE; +DROP FUNCTION IF EXISTS v2_job_completed_before_insert CASCADE; +DROP FUNCTION IF EXISTS v2_job_completed_before_update CASCADE; +DROP FUNCTION IF EXISTS v2_job_queue_after_insert CASCADE; +DROP FUNCTION IF EXISTS v2_job_queue_before_insert CASCADE; +DROP FUNCTION IF EXISTS v2_job_queue_before_update CASCADE; +DROP FUNCTION IF EXISTS v2_job_runtime_before_insert CASCADE; +DROP FUNCTION IF EXISTS v2_job_runtime_before_update CASCADE; +DROP FUNCTION IF EXISTS v2_job_status_before_insert CASCADE; +DROP FUNCTION IF EXISTS v2_job_status_before_update CASCADE; -CREATE OR REPLACE TRIGGER zzz_v2_job_queue_integrity_check_before_delete - BEFORE DELETE ON v2_job_queue - FOR EACH ROW -EXECUTE FUNCTION zzz_v2_job_queue_integrity_check(); +DROP VIEW IF EXISTS completed_job, completed_job_view, job, queue, queue_view CASCADE; --- TODO(uael): remove before phase 4 -CREATE OR REPLACE FUNCTION zzz_v2_job_completed_integrity_check() RETURNS TRIGGER AS $$ -DECLARE job v2_job; -BEGIN - IF (NEW.canceled_by IS NOT NULL) IS DISTINCT FROM NEW.__canceled THEN - RAISE EXCEPTION 'canceled mismatch'; - END IF; - SELECT * INTO job FROM v2_job WHERE id = NEW.id; - IF job.tag IS DISTINCT FROM NEW.__tag THEN - RAISE EXCEPTION 'tag mismatch % %', job.tag, NEW.__tag; - END IF; - IF job.workspace_id IS DISTINCT FROM NEW.workspace_id THEN - RAISE EXCEPTION 'workspace_id mismatch'; - END IF; - IF job.created_at IS DISTINCT FROM NEW.__created_at THEN - RAISE EXCEPTION 'created_at mismatch'; - END IF; - IF job.created_by IS DISTINCT FROM NEW.__created_by THEN - RAISE EXCEPTION 'created_by mismatch'; - END IF; - IF job.permissioned_as IS DISTINCT FROM NEW.__permissioned_as THEN - RAISE EXCEPTION 'permissioned_as mismatch'; - END IF; - IF job.permissioned_as_email IS DISTINCT FROM NEW.__email THEN - RAISE EXCEPTION 'permissioned_as_email mismatch'; - END IF; - IF job.kind IS DISTINCT FROM NEW.__job_kind THEN - RAISE EXCEPTION 'kind mismatch'; - END IF; - IF job.runnable_id IS DISTINCT FROM NEW.__script_hash THEN - RAISE EXCEPTION 'runnable_id mismatch'; - END IF; - IF job.runnable_path IS DISTINCT FROM NEW.__script_path THEN - RAISE EXCEPTION 'runnable_path mismatch'; - END IF; - IF job.parent_job IS DISTINCT FROM NEW.__parent_job THEN - RAISE EXCEPTION 'parent_job mismatch'; - END IF; - IF job.script_lang IS DISTINCT FROM NEW.__language THEN - RAISE EXCEPTION 'script_lang mismatch'; - END IF; - IF job.script_entrypoint_override IS DISTINCT FROM NULLIF(NEW.__args->>'_ENTRYPOINT_OVERRIDE', '__WM_PREPROCESSOR') - AND NEW.__args->>'reason' IS DISTINCT FROM 'PREPROCESSOR_ARGS_ARE_DISCARDED' - THEN - RAISE EXCEPTION 'script_entrypoint_override mismatch'; - END IF; - IF (job.flow_step_id IS NOT NULL) IS DISTINCT FROM NEW.__is_flow_step THEN - RAISE EXCEPTION 'is_flow_step mismatch'; - END IF; - IF job.trigger IS DISTINCT FROM NEW.__schedule_path THEN - RAISE EXCEPTION 'trigger mismatch'; - END IF; - IF job.visible_to_owner IS DISTINCT FROM NEW.__visible_to_owner THEN - RAISE EXCEPTION 'visible_to_owner mismatch'; - END IF; - IF job.args::TEXT IS DISTINCT FROM NEW.__args::TEXT AND NEW.__args->>'_ENTRYPOINT_OVERRIDE' IS DISTINCT FROM '__WM_PREPROCESSOR' THEN - RAISE EXCEPTION 'args mismatch'; - END IF; - RETURN NEW; -END $$ LANGUAGE PLPGSQL; +ALTER TABLE v2_job_queue + DROP COLUMN IF EXISTS __parent_job CASCADE, + DROP COLUMN IF EXISTS __created_by CASCADE, + DROP COLUMN IF EXISTS __script_hash CASCADE, + DROP COLUMN IF EXISTS __script_path CASCADE, + DROP COLUMN IF EXISTS __args CASCADE, + DROP COLUMN IF EXISTS __logs CASCADE, + DROP COLUMN IF EXISTS __raw_code CASCADE, + DROP COLUMN IF EXISTS __canceled CASCADE, + DROP COLUMN IF EXISTS __last_ping CASCADE, + DROP COLUMN IF EXISTS __job_kind CASCADE, + DROP COLUMN IF EXISTS __env_id CASCADE, + DROP COLUMN IF EXISTS __schedule_path CASCADE, + DROP COLUMN IF EXISTS __permissioned_as CASCADE, + DROP COLUMN IF EXISTS __flow_status CASCADE, + DROP COLUMN IF EXISTS __raw_flow CASCADE, + DROP COLUMN IF EXISTS __is_flow_step CASCADE, + DROP COLUMN IF EXISTS __language CASCADE, + DROP COLUMN IF EXISTS __same_worker CASCADE, + DROP COLUMN IF EXISTS __raw_lock CASCADE, + DROP COLUMN IF EXISTS __pre_run_error CASCADE, + DROP COLUMN IF EXISTS __email CASCADE, + DROP COLUMN IF EXISTS __visible_to_owner CASCADE, + DROP COLUMN IF EXISTS __mem_peak CASCADE, + DROP COLUMN IF EXISTS __root_job CASCADE, + DROP COLUMN IF EXISTS __leaf_jobs CASCADE, + DROP COLUMN IF EXISTS __concurrent_limit CASCADE, + DROP COLUMN IF EXISTS __concurrency_time_window_s CASCADE, + DROP COLUMN IF EXISTS __timeout CASCADE, + DROP COLUMN IF EXISTS __flow_step_id CASCADE, + DROP COLUMN IF EXISTS __cache_ttl CASCADE; -CREATE OR REPLACE TRIGGER zzz_v2_job_completed_integrity_check_after_insert - AFTER INSERT ON v2_job_completed - FOR EACH ROW -EXECUTE FUNCTION zzz_v2_job_completed_integrity_check(); +LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE; +ALTER TABLE v2_job_completed + DROP COLUMN IF EXISTS __parent_job CASCADE, + DROP COLUMN IF EXISTS __created_by CASCADE, + DROP COLUMN IF EXISTS __created_at CASCADE, + DROP COLUMN IF EXISTS __success CASCADE, + DROP COLUMN IF EXISTS __script_hash CASCADE, + DROP COLUMN IF EXISTS __script_path CASCADE, + DROP COLUMN IF EXISTS __args CASCADE, + DROP COLUMN IF EXISTS __logs CASCADE, + DROP COLUMN IF EXISTS __raw_code CASCADE, + DROP COLUMN IF EXISTS __canceled CASCADE, + DROP COLUMN IF EXISTS __job_kind CASCADE, + DROP COLUMN IF EXISTS __env_id CASCADE, + DROP COLUMN IF EXISTS __schedule_path CASCADE, + DROP COLUMN IF EXISTS __permissioned_as CASCADE, + DROP COLUMN IF EXISTS __raw_flow CASCADE, + DROP COLUMN IF EXISTS __is_flow_step CASCADE, + DROP COLUMN IF EXISTS __language CASCADE, + DROP COLUMN IF EXISTS __is_skipped CASCADE, + DROP COLUMN IF EXISTS __raw_lock CASCADE, + DROP COLUMN IF EXISTS __email CASCADE, + DROP COLUMN IF EXISTS __visible_to_owner CASCADE, + DROP COLUMN IF EXISTS __tag CASCADE, + DROP COLUMN IF EXISTS __priority CASCADE; diff --git a/backend/tests/fixtures/result_format.sql b/backend/tests/fixtures/result_format.sql index 2cf8fa7299..ff99891526 100644 --- a/backend/tests/fixtures/result_format.sql +++ b/backend/tests/fixtures/result_format.sql @@ -4,8 +4,8 @@ INSERT INTO public.v2_job ( '1eecb96a-c8b0-4a3d-b1b6-087878c55e41', 'test-workspace', 'test-user', '2023-01-01 00:00:00', 'script', 'postgresql' ); -INSERT INTO public.completed_job ( - id, workspace_id, created_by, created_at, duration_ms, success, flow_status, result, job_kind, language +INSERT INTO public.v2_job_completed ( + id, workspace_id, duration_ms, status, result_columns, result ) VALUES ( - '1eecb96a-c8b0-4a3d-b1b6-087878c55e41', 'test-workspace', 'test-user', '2023-01-01 00:00:00', 1000, true, '{"_metadata": {"column_order": ["b", "a"]}}', '[{"a": "second", "b": "first"}]', 'script', 'postgresql' + '1eecb96a-c8b0-4a3d-b1b6-087878c55e41', 'test-workspace', 1000, 'success'::job_status, '{b,a}', '[{"a": "second", "b": "first"}]' ) \ No newline at end of file diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 7573e139ee..e1f73f5871 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -4005,11 +4005,6 @@ mod job_payload { let args = job.args.as_ref().unwrap(); assert_eq!(args.get("foo"), Some(&json!("bar"))); assert_eq!(args.get("bar"), Some(&json!("baz"))); - // TODO: remove this check on v2 phase 4 - assert_eq!( - job.flow_status.as_ref().unwrap().get("_metadata"), - Some(&json!({"preprocessed_args": true})) - ); assert_eq!(job.json_result().unwrap(), json!("Hello bar baz")); let job = sqlx::query!("SELECT preprocessed FROM v2_job WHERE id = $1", job.id) .fetch_one(db) diff --git a/backend/windmill-api/src/db.rs b/backend/windmill-api/src/db.rs index c4ba1d5c5d..5bb273c1db 100644 --- a/backend/windmill-api/src/db.rs +++ b/backend/windmill-api/src/db.rs @@ -6,20 +6,21 @@ * LICENSE-AGPL for a copy of the license. */ -use futures::FutureExt; -use sqlx::Executor; +use std::time::Duration; +use futures::FutureExt; use sqlx::{ migrate::{Migrate, MigrateError}, pool::PoolConnection, - PgConnection, Pool, Postgres, + Executor, PgConnection, Pool, Postgres, }; + use windmill_audit::audit_ee::{AuditAuthor, AuditAuthorable}; -use windmill_common::utils::generate_lock_id; use windmill_common::{ db::{Authable, Authed}, error::Error, }; +use windmill_common::{utils::generate_lock_id, worker::MIN_VERSION_IS_AT_LEAST_1_461}; pub type DB = Pool; @@ -224,6 +225,28 @@ pub async fn migrate(db: &DB) -> Result<(), Error> { } }); + if !has_done_migration(db, "v2_finalize_disable_sync").await { + let db2 = db.clone(); + let _ = tokio::task::spawn(async move { + loop { + if !*MIN_VERSION_IS_AT_LEAST_1_461.read().await { + tracing::info!("Waiting for all workers to be at least version 1.461 before applying v2 finalize migration, sleeping for 5s..."); + tokio::time::sleep(Duration::from_secs(5)).await; + continue; + } + if let Err(err) = v2_finalize(&db2).await { + tracing::error!( + "{err:#}: Could not apply v2 finalize migration, retry in 30s.." + ); + tokio::time::sleep(Duration::from_secs(30)).await; + continue; + } + tracing::info!("v2 finalization step successfully applied."); + break; + } + }); + } + Ok(()) } @@ -274,26 +297,32 @@ async fn fix_flow_versioning_migration( Ok(()) } +async fn has_done_migration(db: &DB, migration_job_name: &str) -> bool { + sqlx::query_scalar!( + "SELECT EXISTS(SELECT name FROM windmill_migrations WHERE name = $1)", + migration_job_name + ) + .fetch_one(db) + .await + .ok() + .flatten() + .unwrap_or(false) +} + macro_rules! run_windmill_migration { - ($migration_job_name:expr, $db:expr, $code:block) => { + ($migration_job_name:expr, $db:expr, |$tx:ident| $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 { + let has_done = has_done_migration(db, migration_job_name).await; + if !has_done { tracing::info!("Applying {migration_job_name} migration"); - let mut tx = db.begin().await?; + 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) + .fetch_one(&mut *$tx) .await .map_err(|e| { tracing::error!("Error acquiring {migration_job_name} lock: {e:#}"); @@ -307,15 +336,9 @@ macro_rules! run_windmill_migration { } 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); + let has_done = has_done_migration(db, migration_job_name).await; - if !has_done_migration { + if !has_done { $code @@ -323,7 +346,7 @@ macro_rules! run_windmill_migration { "INSERT INTO windmill_migrations (name) VALUES ($1) ON CONFLICT DO NOTHING", migration_job_name ) - .execute(&mut *tx) + .execute(&mut *$tx) .await?; tracing::info!("Finished applying {migration_job_name} migration"); } else { @@ -331,9 +354,9 @@ macro_rules! run_windmill_migration { } let _ = sqlx::query("SELECT pg_advisory_unlock(4242)") - .execute(&mut *tx) + .execute(&mut *$tx) .await?; - tx.commit().await?; + $tx.commit().await?; tracing::info!("released lock for {migration_job_name}"); } else { tracing::debug!("migration {migration_job_name} already done"); @@ -343,6 +366,99 @@ macro_rules! run_windmill_migration { }; } +async fn v2_finalize(db: &DB) -> Result<(), Error> { + run_windmill_migration!("v2_finalize_disable_sync", db, |tx| { + tx.execute( + r#" + DROP FUNCTION v2_job_after_update CASCADE; + DROP FUNCTION v2_job_completed_before_insert CASCADE; + DROP FUNCTION v2_job_completed_before_update CASCADE; + DROP FUNCTION v2_job_queue_after_insert CASCADE; + DROP FUNCTION v2_job_queue_before_insert CASCADE; + DROP FUNCTION v2_job_queue_before_update CASCADE; + DROP FUNCTION v2_job_runtime_before_insert CASCADE; + DROP FUNCTION v2_job_runtime_before_update CASCADE; + DROP FUNCTION v2_job_status_before_insert CASCADE; + DROP FUNCTION v2_job_status_before_update CASCADE; + + DROP VIEW completed_job, completed_job_view, job, queue, queue_view CASCADE; + "#, + ) + .await?; + }); + run_windmill_migration!("v2_finalize_job_queue", db, |tx| { + tx.execute( + r#" + ALTER TABLE v2_job_queue + DROP COLUMN __parent_job CASCADE, + DROP COLUMN __created_by CASCADE, + DROP COLUMN __script_hash CASCADE, + DROP COLUMN __script_path CASCADE, + DROP COLUMN __args CASCADE, + DROP COLUMN __logs CASCADE, + DROP COLUMN __raw_code CASCADE, + DROP COLUMN __canceled CASCADE, + DROP COLUMN __last_ping CASCADE, + DROP COLUMN __job_kind CASCADE, + DROP COLUMN __env_id CASCADE, + DROP COLUMN __schedule_path CASCADE, + DROP COLUMN __permissioned_as CASCADE, + DROP COLUMN __flow_status CASCADE, + DROP COLUMN __raw_flow CASCADE, + DROP COLUMN __is_flow_step CASCADE, + DROP COLUMN __language CASCADE, + DROP COLUMN __same_worker CASCADE, + DROP COLUMN __raw_lock CASCADE, + DROP COLUMN __pre_run_error CASCADE, + DROP COLUMN __email CASCADE, + DROP COLUMN __visible_to_owner CASCADE, + DROP COLUMN __mem_peak CASCADE, + DROP COLUMN __root_job CASCADE, + DROP COLUMN __leaf_jobs CASCADE, + DROP COLUMN __concurrent_limit CASCADE, + DROP COLUMN __concurrency_time_window_s CASCADE, + DROP COLUMN __timeout CASCADE, + DROP COLUMN __flow_step_id CASCADE, + DROP COLUMN __cache_ttl CASCADE; + "#, + ) + .await?; + }); + run_windmill_migration!("v2_finalize_job_completed", db, |tx| { + tx.execute( + r#" + LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE; + ALTER TABLE v2_job_completed + DROP COLUMN __parent_job CASCADE, + DROP COLUMN __created_by CASCADE, + DROP COLUMN __created_at CASCADE, + DROP COLUMN __success CASCADE, + DROP COLUMN __script_hash CASCADE, + DROP COLUMN __script_path CASCADE, + DROP COLUMN __args CASCADE, + DROP COLUMN __logs CASCADE, + DROP COLUMN __raw_code CASCADE, + DROP COLUMN __canceled CASCADE, + DROP COLUMN __job_kind CASCADE, + DROP COLUMN __env_id CASCADE, + DROP COLUMN __schedule_path CASCADE, + DROP COLUMN __permissioned_as CASCADE, + DROP COLUMN __raw_flow CASCADE, + DROP COLUMN __is_flow_step CASCADE, + DROP COLUMN __language CASCADE, + DROP COLUMN __is_skipped CASCADE, + DROP COLUMN __raw_lock CASCADE, + DROP COLUMN __email CASCADE, + DROP COLUMN __visible_to_owner CASCADE, + DROP COLUMN __tag CASCADE, + DROP COLUMN __priority CASCADE; + "#, + ) + .await?; + }); + Ok(()) +} + 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')" @@ -385,7 +501,7 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { // tx.commit().await?; // } - run_windmill_migration!("fix_job_completed_index_2", &db, { + run_windmill_migration!("fix_job_completed_index_2", &db, |tx| { // 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?; @@ -405,7 +521,7 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { .await?; }); - run_windmill_migration!("fix_job_completed_index_3", &db, { + run_windmill_migration!("fix_job_completed_index_3", &db, |tx| { sqlx::query("DROP INDEX CONCURRENTLY IF EXISTS index_completed_job_on_schedule_path") .execute(db) .await?; @@ -423,7 +539,7 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { .await?; }); - run_windmill_migration!("fix_job_index_1", &db, { + run_windmill_migration!("fix_job_index_1", &db, |tx| { let migration_job_name = "fix_job_completed_index_4"; let mut i = 1; tracing::info!("step {i} of {migration_job_name} migration"); @@ -504,82 +620,7 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { .await?; }); - run_windmill_migration!("add_ix_v2_II", &db, { - sqlx::query!( - "create index concurrently if not exists ix_v2_job_root_by_path - on v2_job (workspace_id, runnable_path, created_at DESC) - where parent_job is null" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_3" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_5" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_6" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_7" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_8" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_9" - ) - .execute(db) - .await?; - - sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_created_at") - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_created_at_new_2" - ) - .execute(db) - .await?; - - sqlx::query!( - "DROP INDEX CONCURRENTLY IF EXISTS ix_completed_job_workspace_id_started_at_new" - ) - .execute(db) - .await?; - - sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path_2") - .execute(db) - .await?; - - sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS scheduled_root_job") - .execute(db) - .await?; - - sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS concurrency_limit_stats_completed_job") - .execute(db) - .await?; - tracing::info!("Finished adding ix_v2_II migration"); - }); - - run_windmill_migration!("fix_labeled_jobs_index", &db, { + run_windmill_migration!("fix_labeled_jobs_index", &db, |tx| { tracing::info!("Special migration to add index concurrently on job labels 2"); sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS labeled_jobs_on_jobs") .execute(db) @@ -589,7 +630,7 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { ).execute(db).await?; }); - run_windmill_migration!("v2_labeled_jobs_index", &db, { + run_windmill_migration!("v2_labeled_jobs_index", &db, |tx| { tracing::info!("Special migration to add index concurrently on job labels"); sqlx::query!( "CREATE INDEX CONCURRENTLY ix_v2_job_labels ON v2_job diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index 23b5ed8f3e..487a217d7f 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -95,6 +95,7 @@ lazy_static::lazy_static! { .unwrap_or(false); pub static ref MIN_VERSION: Arc> = Arc::new(RwLock::new(Version::new(0, 0, 0))); + pub static ref MIN_VERSION_IS_AT_LEAST_1_461: Arc> = Arc::new(RwLock::new(false)); pub static ref MIN_VERSION_IS_AT_LEAST_1_427: Arc> = Arc::new(RwLock::new(false)); pub static ref MIN_VERSION_IS_AT_LEAST_1_432: Arc> = Arc::new(RwLock::new(false)); pub static ref MIN_VERSION_IS_AT_LEAST_1_440: Arc> = Arc::new(RwLock::new(false)); @@ -103,6 +104,8 @@ lazy_static::lazy_static! { pub static ref DISABLE_FLOW_SCRIPT: bool = std::env::var("DISABLE_FLOW_SCRIPT").ok().is_some_and(|x| x == "1" || x == "true"); } +pub static MIN_VERSION_IS_LATEST: AtomicBool = AtomicBool::new(false); + fn format_pull_query(peek: String) -> String { let r = format!( "WITH peek AS ( @@ -662,6 +665,7 @@ pub async fn update_min_version<'c, E: sqlx::Executor<'c, Database = sqlx::Postg tracing::info!("Minimal worker version: {min_version}"); } + *MIN_VERSION_IS_AT_LEAST_1_461.write().await = min_version >= Version::new(1, 461, 0); *MIN_VERSION_IS_AT_LEAST_1_427.write().await = min_version >= Version::new(1, 427, 0); *MIN_VERSION_IS_AT_LEAST_1_432.write().await = min_version >= Version::new(1, 432, 0); *MIN_VERSION_IS_AT_LEAST_1_440.write().await = min_version >= Version::new(1, 440, 0);