From 30d9c4019374121e5d22d83f42226e3e0dd5dc76 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 28 Sep 2024 17:19:52 +0200 Subject: [PATCH] optimize more queries --- ...9a9d035f3eb1d9d833041969fce967daa91a4.json | 23 -- backend/windmill-common/src/flows.rs | 12 + backend/windmill-worker/src/worker.rs | 1 + backend/windmill-worker/src/worker_flow.rs | 259 +++++++++--------- 4 files changed, 136 insertions(+), 159 deletions(-) delete mode 100644 backend/.sqlx/query-ae2f005af8ab4b035a907e0c8fc9a9d035f3eb1d9d833041969fce967daa91a4.json diff --git a/backend/.sqlx/query-ae2f005af8ab4b035a907e0c8fc9a9d035f3eb1d9d833041969fce967daa91a4.json b/backend/.sqlx/query-ae2f005af8ab4b035a907e0c8fc9a9d035f3eb1d9d833041969fce967daa91a4.json deleted file mode 100644 index 74c890726f..0000000000 --- a/backend/.sqlx/query-ae2f005af8ab4b035a907e0c8fc9a9d035f3eb1d9d833041969fce967daa91a4.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT raw_flow->'modules'->$2::int->'retry' FROM queue WHERE id = $1", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "?column?", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Int4" - ] - }, - "nullable": [ - null - ] - }, - "hash": "ae2f005af8ab4b035a907e0c8fc9a9d035f3eb1d9d833041969fce967daa91a4" -} diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index b1ad85da99..054cef7d91 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -285,6 +285,13 @@ pub struct FlowModuleValueWithParallel { pub parallelism: Option, } +#[derive(Deserialize)] +pub struct FlowModuleValueWithSkipFailures { + pub skip_failures: Option, + pub parallel: Option, + pub parallelism: Option, +} + impl FlowModule { pub fn id_append(&mut self, s: &str) { self.id = format!("{}-{}", self.id, s); @@ -293,6 +300,11 @@ impl FlowModule { serde_json::from_str::(self.value.get()).map_err(crate::error::to_anyhow) } + pub fn get_value_with_skip_failures(&self) -> anyhow::Result { + serde_json::from_str::(self.value.get()) + .map_err(crate::error::to_anyhow) + } + pub fn is_flow(&self) -> bool { self.get_type().is_ok_and(|x| x == "flow") } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 8c8856cfff..6c0fce3d3b 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1201,6 +1201,7 @@ pub async fn run_worker() < likelihood_of_suspend || last_suspend_first.elapsed().as_secs_f64() > 5.0; + if suspend_first { last_suspend_first = Instant::now(); } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index b8085e65fe..0c937a28a3 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -168,10 +168,7 @@ pub struct RecUpdateFlowStatusAfterJobCompletion { } #[derive(FromRow)] -pub struct SkipIfStopped { - pub skip_if_stopped: Option, - pub stop_early_expr: Option, - pub continue_on_error: Option, +pub struct RowArgs { pub args: Option>>>, } @@ -183,6 +180,7 @@ struct RecoveryObject { #[derive(sqlx::FromRow, Deserialize)] pub struct RowFlowStatus { pub flow_status: sqlx::types::Json>, + pub current_module: Option>>, } // #[instrument(level = "trace", skip_all)] pub async fn update_flow_status_after_job_completion_internal< @@ -218,7 +216,7 @@ pub async fn update_flow_status_after_job_completion_internal< // tracing::debug!("UPDATE FLOW STATUS: {flow:?} {success} {result:?} {w_id} {depth}"); let old_status_json = sqlx::query_as::<_, RowFlowStatus>( - "SELECT flow_status FROM queue WHERE id = $1 AND workspace_id = $2", + "SELECT flow_status, raw_flow->'modules'->(flow_status->'step')::int as current_module FROM queue WHERE id = $1 AND workspace_id = $2", ) .bind(flow) .bind(w_id) @@ -237,6 +235,16 @@ pub async fn update_flow_status_after_job_completion_internal< ))) })?; + let current_module = if let Some(x) = old_status_json.current_module { + Some(serde_json::from_str::(x.0.get()).or_else(|e| { + Err(Error::InternalErr(format!( + "requiring current module to be parsable as FlowModule: {e:?}" + ))) + })?) + } else { + None + }; + let module_step = Step::from_i32_and_len(old_status.step, old_status.modules.len()); let module_status = match module_step { @@ -259,9 +267,17 @@ pub async fn update_flow_status_after_job_completion_internal< module_status, FlowStatusModule::InProgress { iterator: Some(_), .. } ) { - let (loop_failures, parallelism) = - compute_skip_loop_failures_and_parallelism(flow, old_status.step, db).await?; - (true, loop_failures.unwrap_or(false), parallelism) + let value = current_module + .as_ref() + .and_then(|x| x.get_value_with_skip_failures().ok()); + ( + true, + value + .as_ref() + .and_then(|x| x.skip_failures) + .unwrap_or(false), + value.as_ref().and_then(|x| x.parallelism), + ) } else { (false, false, None) }; @@ -275,76 +291,77 @@ pub async fn update_flow_status_after_job_completion_internal< let is_failure_step = old_status.step >= old_status.modules.len() as i32 && old_status.modules.len() > 0; - let (mut stop_early, mut skip_if_stop_early, continue_on_error) = if let Some(se) = - stop_early_override - { - //do not stop early if module is a flow step - let flow_job = get_queued_job(&flow, w_id, db) - .await? - .ok_or_else(|| Error::InternalErr(format!("requiring flow to be in the queue")))?; - let module = get_module(&flow_job, &module_step); + let (mut stop_early, mut skip_if_stop_early, continue_on_error) = + if let Some(se) = stop_early_override { + //do not stop early if module is a flow step + let flow_job = get_queued_job(&flow, w_id, db).await?.ok_or_else(|| { + Error::InternalErr(format!("requiring flow to be in the queue")) + })?; + let module = get_module(&flow_job, &module_step); - if module.is_some_and(|x| x.is_flow()) { - (false, false, false) - } else { - (true, se, false) - } - } else if is_failure_step || matches!(module_step, Step::PreprocessorStep) { - (false, false, false) - } else { - let r = sqlx::query_as::<_, SkipIfStopped>( - "SELECT - raw_flow->'modules'->$1::int->'stop_after_if'->>'expr' as stop_early_expr, - (raw_flow->'modules'->$1::int->'stop_after_if'->>'skip_if_stopped')::bool as skip_if_stopped, - (raw_flow->'modules'->$1::int->'continue_on_error')::bool as continue_on_error, - args - FROM queue - WHERE id = $2" - ) - .bind(old_status.step) - .bind(flow) - .fetch_one(db) - .await - .map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e:#}")))?; - - let stop_early = success - && !is_branch_all - && if let Some(expr) = r.stop_early_expr.clone() { - let all_iters = match &module_status { - FlowStatusModule::InProgress { flow_jobs: Some(flow_jobs), .. } - if expr.contains("all_iters") => - { - Some(Arc::new( - retrieve_flow_jobs_results(db, w_id, flow_jobs).await?, - )) - } - _ => None, - }; - compute_bool_from_expr( - expr, - Marc::new( - r.args - .map(|x| x.0) - .unwrap_or_else(|| serde_json::from_str("{}").unwrap()) - .to_owned(), - ), - result.clone(), - all_iters, - None, - Some(client), - None, - None, - ) - .await? + if module.is_some_and(|x| x.is_flow()) { + (false, false, false) } else { - false - }; - ( - stop_early, - r.skip_if_stopped.unwrap_or(false), - r.continue_on_error.unwrap_or(false), - ) - }; + (true, se, false) + } + } else if is_failure_step || matches!(module_step, Step::PreprocessorStep) { + (false, false, false) + } else if let Some(current_module) = current_module.as_ref() { + let stop_early = success + && !is_branch_all + && if let Some(ref expr) = current_module + .stop_after_if + .as_ref() + .map(|x| x.expr.clone()) + { + let all_iters = + match &module_status { + FlowStatusModule::InProgress { + flow_jobs: Some(flow_jobs), .. + } if expr.contains("all_iters") => Some(Arc::new( + retrieve_flow_jobs_results(db, w_id, flow_jobs).await?, + )), + _ => None, + }; + let args = sqlx::query_as::<_, RowArgs>( + "SELECT + args + FROM queue + WHERE id = $2", + ) + .bind(old_status.step) + .bind(flow) + .fetch_one(db) + .await + .map_err(|e| { + Error::InternalErr(format!("retrieval of args from state: {e:#}")) + })?; + compute_bool_from_expr( + expr.to_string(), + Marc::new(args.args.unwrap_or_default().0), + result.clone(), + all_iters, + None, + Some(client), + None, + None, + ) + .await? + } else { + false + }; + ( + stop_early, + current_module + .stop_after_if + .as_ref() + .map(|x| x.skip_if_stopped) + .unwrap_or(false), + current_module.continue_on_error.unwrap_or(false), + ) + } else { + (false, false, false) + }; let skip_branch_failure = match module_status { FlowStatusModule::InProgress { @@ -666,25 +683,12 @@ pub async fn update_flow_status_after_job_completion_internal< ) } else { let inc = if continue_on_error { - let retry = sqlx::query_scalar!( - "SELECT raw_flow->'modules'->$2::int->'retry' FROM queue WHERE id = $1", - flow, - old_status.step - ) - .fetch_optional(&mut tx) - .await - .map_err(|e| { - Error::InternalErr(format!( - "error while getting retry from step: {e:#}" - )) - })? - .flatten(); - - let retry = retry - .map(|x| serde_json::from_value::(x).ok()) - .flatten() + let retry = current_module + .as_ref() + .and_then(|x| x.retry.clone()) .unwrap_or_default(); - tracing::info!("update flow status on rety: {retry:#?} "); + + tracing::info!("update flow status on rety: {retry:#?} "); next_retry(&retry, &old_status.retry).is_none() } else { false @@ -816,30 +820,27 @@ pub async fn update_flow_status_after_job_completion_internal< match &new_status { Some(FlowStatusModule::Success { .. }) if is_loop || is_branch_all => { - let r_after_all_iters = sqlx::query_as::<_, SkipIfStopped>( - "SELECT - raw_flow->'modules'->$1::int->'stop_after_all_iters_if'->>'expr' as stop_early_expr, - (raw_flow->'modules'->$1::int->'stop_after_all_iters_if'->>'skip_if_stopped')::bool as skip_if_stopped, - NULL as continue_on_error, + if let Some(ref expr) = current_module + .as_ref() + .and_then(|m| m.stop_after_all_iters_if.as_ref().map(|x| x.expr.clone())) + { + let args = sqlx::query_as::<_, RowArgs>( + "SELECT args FROM queue - WHERE id = $2" + WHERE id = $2", ) .bind(old_status.step) .bind(flow) .fetch_one(db) .await - .map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e:#}")))?; - if let Some(expr) = r_after_all_iters.stop_early_expr { + .map_err(|e| { + Error::InternalErr(format!("retrieval of args from state: {e:#}")) + })?; + let should_stop = compute_bool_from_expr( - expr, - Marc::new( - r_after_all_iters - .args - .map(|x| x.0) - .unwrap_or_else(|| serde_json::from_str("{}").unwrap()) - .to_owned(), - ), + expr.to_string(), + Marc::new(args.args.unwrap_or_default().0), nresult.clone(), None, None, @@ -851,7 +852,14 @@ pub async fn update_flow_status_after_job_completion_internal< if should_stop { stop_early = should_stop; - skip_if_stop_early = r_after_all_iters.skip_if_stopped.unwrap_or(false); + skip_if_stop_early = current_module + .as_ref() + .and_then(|m| { + m.stop_after_all_iters_if + .as_ref() + .map(|x| x.skip_if_stopped) + }) + .unwrap_or(false); } } } @@ -876,6 +884,7 @@ pub async fn update_flow_status_after_job_completion_internal< let flow_job = get_queued_job_tx(flow, w_id, tx.transaction_mut()) .await? .ok_or_else(|| Error::InternalErr(format!("requiring flow to be in the queue")))?; + tx.commit().await?; let job_root = flow_job .root_job @@ -908,14 +917,13 @@ pub async fn update_flow_status_after_job_completion_internal< false if !is_failure_step && !skip_error_handler - && has_failure_module(flow, tx.transaction_mut()).await? => + && has_failure_module(flow, db).await? => { true } false => false, }; - tx.commit().await?; tracing::debug!(id = %flow_job.id, root_id = %job_root, "flow status updated"); ( @@ -1196,24 +1204,6 @@ fn get_module(flow_job: &QueuedJob, module_step: &Step) -> Option { } } -async fn compute_skip_loop_failures_and_parallelism( - flow: Uuid, - step: i32, - db: &DB, -) -> Result<(Option, Option), Error> { - sqlx::query_as( - "SELECT (raw_flow->'modules'->$1->'value'->>'skip_failures')::bool, (raw_flow->'modules'->$1->'value'->>'parallelism')::int - FROM queue - WHERE id = $2", - ) - .bind(step) - .bind(flow) - .fetch_one(db) - .await - .map(|(v, n)| (v,n)) - .map_err(|e| Error::InternalErr(format!("error during retrieval of skip_loop_failures: {e:#}"))) -} - async fn compute_skip_branchall_failure<'c>( flow: Uuid, job: &Uuid, @@ -1262,17 +1252,14 @@ async fn compute_skip_branchall_failure<'c>( }) } -async fn has_failure_module<'c>( - flow: Uuid, - tx: &mut sqlx::Transaction<'c, sqlx::Postgres>, -) -> Result { +async fn has_failure_module<'c>(flow: Uuid, db: &DB) -> Result { sqlx::query_scalar::<_, Option>( "SELECT raw_flow->'failure_module' != 'null'::jsonb FROM queue WHERE id = $1", ) .bind(flow) - .fetch_one(&mut **tx) + .fetch_one(db) .await .map_err(|e| { Error::InternalErr(format!(