From cc111ba7dcd7fe6c4280d394d4eb0e723e9d8b19 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 5 Aug 2024 12:46:01 +0200 Subject: [PATCH] feat: improve indices of completed_runs for faster load --- backend/windmill-api/src/db.rs | 137 +++++++++++++++++- .../src/lib/components/runs/JobLoader.svelte | 13 +- 2 files changed, 140 insertions(+), 10 deletions(-) diff --git a/backend/windmill-api/src/db.rs b/backend/windmill-api/src/db.rs index 70a9413f87..ffbb018122 100644 --- a/backend/windmill-api/src/db.rs +++ b/backend/windmill-api/src/db.rs @@ -327,13 +327,13 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { .unwrap_or(false); 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_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( + // "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", @@ -436,6 +436,131 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> { 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") + .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"); + + 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"); + + 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"); + + 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"); + + 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?; + + 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(()) } diff --git a/frontend/src/lib/components/runs/JobLoader.svelte b/frontend/src/lib/components/runs/JobLoader.svelte index 57364e0153..dc1cf7d228 100644 --- a/frontend/src/lib/components/runs/JobLoader.svelte +++ b/frontend/src/lib/components/runs/JobLoader.svelte @@ -122,21 +122,26 @@ ): Promise { loadingFetch = true try { + let scriptPathStart = folder === null || folder === '' ? undefined : `f/${folder}/` + let scriptPathExact = path === null || path === '' ? undefined : path return JobService.listJobs({ workspace: $workspaceStore!, createdOrStartedBefore: startedBefore, createdOrStartedAfter: startedAfter, createdOrStartedAfterCompletedJobs: startedAfterCompletedJobs, schedulePath, - scriptPathExact: path === null || path === '' ? undefined : path, + scriptPathExact, createdBy: user === null || user === '' ? undefined : user, - scriptPathStart: folder === null || folder === '' ? undefined : `f/${folder}/`, + scriptPathStart: scriptPathStart, jobKinds, success: success == 'success' ? true : success == 'failure' ? false : undefined, running: success == 'running' ? true : undefined, isSkipped: isSkipped ? undefined : false, - isFlowStep: jobKindsCat != 'all' ? false : undefined, - hasNullParent: jobKindsCat != 'all' ? false : undefined, + // isFlowStep: jobKindsCat != 'all' ? false : undefined, + hasNullParent: + scriptPathExact != undefined || scriptPathStart != undefined || jobKinds != 'all' + ? true + : undefined, label: label === null || label === '' ? undefined : label, isNotSchedule: showSchedules == false ? true : undefined, scheduledForBeforeNow: showFutureJobs == false ? true : undefined,