backend: finalize v2 migration (v2 phase 4 - final) (#5155)

* backend: finalize v2 migration (v2 phase 4 - final)

* update migration

---------

Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
This commit is contained in:
Lucas Abel
2025-02-13 17:10:30 +01:00
committed by GitHub
co-authored by Ruben Fiszel
parent 149d5fb3e1
commit b1f358d4a2
5 changed files with 222 additions and 291 deletions
+68 -177
View File
@@ -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;
+3 -3
View File
@@ -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"}]'
)
-5
View File
@@ -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)
+147 -106
View File
@@ -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<Postgres>;
@@ -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<Postgres> = $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
+4
View File
@@ -95,6 +95,7 @@ lazy_static::lazy_static! {
.unwrap_or(false);
pub static ref MIN_VERSION: Arc<RwLock<Version>> = Arc::new(RwLock::new(Version::new(0, 0, 0)));
pub static ref MIN_VERSION_IS_AT_LEAST_1_461: Arc<RwLock<bool>> = Arc::new(RwLock::new(false));
pub static ref MIN_VERSION_IS_AT_LEAST_1_427: Arc<RwLock<bool>> = Arc::new(RwLock::new(false));
pub static ref MIN_VERSION_IS_AT_LEAST_1_432: Arc<RwLock<bool>> = Arc::new(RwLock::new(false));
pub static ref MIN_VERSION_IS_AT_LEAST_1_440: Arc<RwLock<bool>> = 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);