diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 23c7d9234e..4559829a5e 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -804,7 +804,8 @@ async fn handle_job_error( same_worker_tx, worker_dir, None, - base_internal_url + base_internal_url, + 0 ) .await; @@ -965,7 +966,8 @@ async fn handle_queued_job( same_worker_tx.clone(), worker_dir, None, - base_internal_url + base_internal_url, + 0 ) .await?; } @@ -1017,7 +1019,8 @@ async fn handle_queued_job( same_worker_tx, worker_dir, None, - base_internal_url + base_internal_url, + 0 ) .await?; } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index f14b93b71e..17a86b077e 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -50,58 +50,62 @@ pub async fn update_flow_status_after_job_completion( worker_dir: &str, stop_early_override: Option, base_internal_url: &str, + depth: u8, ) -> error::Result<()> { - tracing::debug!("UPDATE FLOW STATUS: {flow:?} {success} {result:?} {w_id}"); + let (should_continue_flow, flow_job, stop_early, skip_if_stop_early, nresult) = { + tracing::debug!("UPDATE FLOW STATUS: {flow:?} {success} {result:?} {w_id} {depth}"); - let mut tx = db.begin().await?; + let mut tx = db.begin().await?; - let old_status_json = sqlx::query_scalar!( - "SELECT flow_status FROM queue WHERE id = $1 AND workspace_id = $2", - flow, - w_id - ) - .fetch_one(&mut tx) - .await - .map_err(|e| { - Error::InternalErr(format!( - "fetching flow status {flow} while reporting {success} {result:?}: {e}" - )) - })? - .ok_or_else(|| Error::InternalErr(format!("requiring a previous status")))?; + let old_status_json = sqlx::query_scalar!( + "SELECT flow_status FROM queue WHERE id = $1 AND workspace_id = $2", + flow, + w_id + ) + .fetch_one(&mut tx) + .await + .map_err(|e| { + Error::InternalErr(format!( + "fetching flow status {flow} while reporting {success} {result:?}: {e}" + )) + })? + .ok_or_else(|| Error::InternalErr(format!("requiring a previous status")))?; - let old_status = serde_json::from_value::(old_status_json) - .ok() - .ok_or_else(|| { - Error::InternalErr(format!("requiring status to be parsabled as FlowStatus")) - })?; + let old_status = serde_json::from_value::(old_status_json) + .ok() + .ok_or_else(|| { + Error::InternalErr(format!("requiring status to be parsabled as FlowStatus")) + })?; - let module_index = usize::try_from(old_status.step).ok(); + let module_index = usize::try_from(old_status.step).ok(); - let module_status = module_index - .and_then(|i| old_status.modules.get(i)) - .unwrap_or(&old_status.failure_module.module_status); + let module_status = module_index + .and_then(|i| old_status.modules.get(i)) + .unwrap_or(&old_status.failure_module.module_status); - tracing::debug!("UPDATE FLOW STATUS 2: {module_index:#?} {module_status:#?} {old_status:#?} "); + tracing::debug!( + "UPDATE FLOW STATUS 2: {module_index:#?} {module_status:#?} {old_status:#?} " + ); - let skip_loop_failures = if matches!( - module_status, - FlowStatusModule::InProgress { iterator: Some(_), .. } - ) { - compute_skip_loop_failures(flow, old_status.step, &mut tx) - .await? - .unwrap_or(false) - } else { - false - }; + let skip_loop_failures = if matches!( + module_status, + FlowStatusModule::InProgress { iterator: Some(_), .. } + ) { + compute_skip_loop_failures(flow, old_status.step, &mut tx) + .await? + .unwrap_or(false) + } else { + false + }; - let is_failure_step = old_status.step >= old_status.modules.len() as i32; + let is_failure_step = old_status.step >= old_status.modules.len() as i32; - let (mut stop_early, skip_if_stop_early) = if let Some(se) = stop_early_override { - (true, se) - } else if is_failure_step { - (false, false) - } else { - let r = sqlx::query!( + let (mut stop_early, skip_if_stop_early) = if let Some(se) = stop_early_override { + (true, se) + } else if is_failure_step { + (false, false) + } else { + let r = sqlx::query!( " 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, @@ -116,37 +120,39 @@ pub async fn update_flow_status_after_job_completion( .await .map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e}")))?; - let stop_early = success - && if let Some(expr) = r.stop_early_expr.clone() { - compute_bool_from_expr(expr, &r.args, result.clone(), None, Some(client)).await? - } else { - false - }; - (stop_early, r.skip_if_stopped.unwrap_or(false)) - }; + let stop_early = success + && if let Some(expr) = r.stop_early_expr.clone() { + compute_bool_from_expr(expr, &r.args, result.clone(), None, Some(client)) + .await? + } else { + false + }; + (stop_early, r.skip_if_stopped.unwrap_or(false)) + }; - let skip_branch_failure = match module_status { - FlowStatusModule::InProgress { - branchall: Some(BranchAllStatus { branch, .. }), .. - } => compute_skip_branchall_failure(flow, old_status.step, *branch, &mut tx) - .await? - .unwrap_or(false), - _ => false, - }; + let skip_branch_failure = match module_status { + FlowStatusModule::InProgress { + branchall: Some(BranchAllStatus { branch, .. }), + .. + } => compute_skip_branchall_failure(flow, old_status.step, *branch, &mut tx) + .await? + .unwrap_or(false), + _ => false, + }; - let skip_failure = skip_branch_failure || skip_loop_failures; + let skip_failure = skip_branch_failure || skip_loop_failures; - let (inc_step_counter, new_status) = match module_status { - FlowStatusModule::InProgress { - iterator, - branchall, - parallel, - flow_jobs: Some(jobs), - .. - } if *parallel => { - let (nindex, len) = match (iterator, branchall) { - (Some(Iterator { itered, .. }), _) => { - let nindex = sqlx::query_scalar!( + let (inc_step_counter, new_status) = match module_status { + FlowStatusModule::InProgress { + iterator, + branchall, + parallel, + flow_jobs: Some(jobs), + .. + } if *parallel => { + let (nindex, len) = match (iterator, branchall) { + (Some(Iterator { itered, .. }), _) => { + let nindex = sqlx::query_scalar!( " UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'iterator', 'index'], ((flow_status->'modules'->$1::int->'iterator'->>'index')::int + 1)::text::jsonb) @@ -159,10 +165,10 @@ pub async fn update_flow_status_after_job_completion( .fetch_one(&mut tx) .await? .ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?; - (nindex, itered.len() as i32) - } - (_, Some(BranchAllStatus { len, .. })) => { - let nindex = sqlx::query_scalar!( + (nindex, itered.len() as i32) + } + (_, Some(BranchAllStatus { len, .. })) => { + let nindex = sqlx::query_scalar!( " UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'branchall', 'branch'], ((flow_status->'modules'->$1::int->'branchall'->>'branch')::int + 1)::text::jsonb) @@ -175,269 +181,292 @@ pub async fn update_flow_status_after_job_completion( .fetch_one(&mut tx) .await? .ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?; - (nindex, *len as i32) - } - _ => Err(Error::InternalErr(format!( - "unexpected status for parallel module" - )))?, - }; - if nindex == len { - let new_status = if skip_loop_failures - || sqlx::query_scalar!( - " + (nindex, *len as i32) + } + _ => Err(Error::InternalErr(format!( + "unexpected status for parallel module" + )))?, + }; + if nindex == len { + let new_status = if skip_loop_failures + || sqlx::query_scalar!( + " SELECT success FROM completed_job WHERE id = ANY($1) ", - jobs.as_slice(), - ) - .fetch_all(&mut tx) - .await? - .into_iter() - .all(|x| x) - { - FlowStatusModule::Success { - id: module_status.id(), - job: job_id_for_status.clone(), - flow_jobs: Some(jobs.clone()), - branch_chosen: None, - approvers: vec![], - } + jobs.as_slice(), + ) + .fetch_all(&mut tx) + .await? + .into_iter() + .all(|x| x) + { + FlowStatusModule::Success { + id: module_status.id(), + job: job_id_for_status.clone(), + flow_jobs: Some(jobs.clone()), + branch_chosen: None, + approvers: vec![], + } + } else { + FlowStatusModule::Failure { + id: module_status.id(), + job: job_id_for_status.clone(), + flow_jobs: Some(jobs.clone()), + branch_chosen: None, + } + }; + (true, Some(new_status)) } else { - FlowStatusModule::Failure { - id: module_status.id(), - job: job_id_for_status.clone(), - flow_jobs: Some(jobs.clone()), - branch_chosen: None, - } - }; - (true, Some(new_status)) - } else { - tx.commit().await?; - return Ok(()); - } - } - FlowStatusModule::InProgress { - iterator: Some(windmill_common::flow_status::Iterator { index, itered, .. }), - .. - } if (*index + 1 < itered.len() && (success || skip_loop_failures)) && !stop_early => { - (false, None) - } - FlowStatusModule::InProgress { - branchall: Some(BranchAllStatus { branch, len, .. }), - .. - } if branch.to_owned() < len - 1 && (success || skip_branch_failure) => (false, None), - _ => { - if stop_early - && matches!( - module_status, - FlowStatusModule::InProgress { iterator: Some(_), .. } - ) - { - // if we're stopping early inside a loop, we just want to break the loop instead - stop_early = false; - } - let (flow_jobs, branch_chosen) = match module_status { - FlowStatusModule::InProgress { flow_jobs, branch_chosen, .. } => { - (flow_jobs.clone(), branch_chosen.clone()) + tx.commit().await?; + return Ok(()); } - _ => (None, None), - }; - if success || (flow_jobs.is_some() && (skip_loop_failures || skip_branch_failure)) { - ( - true, - Some(FlowStatusModule::Success { - id: module_status.id(), - job: job_id_for_status.clone(), - flow_jobs, - branch_chosen, - approvers: vec![], - }), - ) - } else { - ( - false, - Some(FlowStatusModule::Failure { - id: module_status.id(), - job: job_id_for_status.clone(), - flow_jobs, - branch_chosen, - }), - ) } - } - }; + FlowStatusModule::InProgress { + iterator: Some(windmill_common::flow_status::Iterator { index, itered, .. }), + .. + } if (*index + 1 < itered.len() && (success || skip_loop_failures)) && !stop_early => { + (false, None) + } + FlowStatusModule::InProgress { + branchall: Some(BranchAllStatus { branch, len, .. }), + .. + } if branch.to_owned() < len - 1 && (success || skip_branch_failure) => (false, None), + _ => { + if stop_early + && matches!( + module_status, + FlowStatusModule::InProgress { iterator: Some(_), .. } + ) + { + // if we're stopping early inside a loop, we just want to break the loop instead + stop_early = false; + } + let (flow_jobs, branch_chosen) = match module_status { + FlowStatusModule::InProgress { flow_jobs, branch_chosen, .. } => { + (flow_jobs.clone(), branch_chosen.clone()) + } + _ => (None, None), + }; + if success || (flow_jobs.is_some() && (skip_loop_failures || skip_branch_failure)) { + ( + true, + Some(FlowStatusModule::Success { + id: module_status.id(), + job: job_id_for_status.clone(), + flow_jobs, + branch_chosen, + approvers: vec![], + }), + ) + } else { + ( + false, + Some(FlowStatusModule::Failure { + id: module_status.id(), + job: job_id_for_status.clone(), + flow_jobs, + branch_chosen, + }), + ) + } + } + }; - let step_counter = if inc_step_counter { - sqlx::query!( - " + let step_counter = if inc_step_counter { + sqlx::query!( + " UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1) WHERE id = $2 ", - json!(old_status.step + 1), - flow - ) - .execute(&mut tx) - .await?; - old_status.step + 1 - } else { - old_status.step - }; + json!(old_status.step + 1), + flow + ) + .execute(&mut tx) + .await?; + old_status.step + 1 + } else { + old_status.step + }; - /* is_last_step is true when the step_counter (the next step index) is an invalid index */ - let is_last_step = usize::try_from(step_counter) - .map(|i| !(..old_status.modules.len()).contains(&i)) - .unwrap_or(true); + /* is_last_step is true when the step_counter (the next step index) is an invalid index */ + let is_last_step = usize::try_from(step_counter) + .map(|i| !(..old_status.modules.len()).contains(&i)) + .unwrap_or(true); - if let Some(new_status) = new_status.as_ref() { - if is_failure_step { - let parent_module = sqlx::query_scalar!( + if let Some(new_status) = new_status.as_ref() { + if is_failure_step { + let parent_module = sqlx::query_scalar!( "SELECT flow_status->'failure_module'->>'parent_module' FROM queue WHERE id = $1", flow ) - .fetch_one(&mut tx) - .await?; + .fetch_one(&mut tx) + .await?; - sqlx::query!( - " + sqlx::query!( + " UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1) WHERE id = $2 ", - json!(FlowStatusModuleWParent { parent_module, module_status: new_status.clone() }), - flow - ) - .execute(&mut tx) - .await?; - } else { - sqlx::query!( - " - UPDATE queue - SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2) - WHERE id = $3 - ", - old_status.step.to_string(), - json!(new_status), - flow - ) - .execute(&mut tx) - .await?; - - if let Some(job_result) = new_status.job_result() { - sqlx::query!( - " - UPDATE queue - SET leaf_jobs = JSONB_SET(coalesce(leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2) - WHERE COALESCE((SELECT root_job FROM queue WHERE id = $3), $3) = id - ", - new_status.id(), - json!(job_result), + json!(FlowStatusModuleWParent { + parent_module, + module_status: new_status.clone() + }), flow ) .execute(&mut tx) .await?; + } else { + sqlx::query!( + " + UPDATE queue + SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2) + WHERE id = $3 + ", + old_status.step.to_string(), + json!(new_status), + flow + ) + .execute(&mut tx) + .await?; + + if let Some(job_result) = new_status.job_result() { + sqlx::query!( + " + UPDATE queue + SET leaf_jobs = JSONB_SET(coalesce(leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2) + WHERE COALESCE((SELECT root_job FROM queue WHERE id = $3), $3) = id + ", + new_status.id(), + json!(job_result), + flow + ) + .execute(&mut tx) + .await?; + } } } - } - let result = match &new_status { - Some(FlowStatusModule::Success { flow_jobs: Some(jobs), .. }) - | Some(FlowStatusModule::Failure { flow_jobs: Some(jobs), .. }) => { - let results = sqlx::query!( - " + let nresult = match &new_status { + Some(FlowStatusModule::Success { flow_jobs: Some(jobs), .. }) + | Some(FlowStatusModule::Failure { flow_jobs: Some(jobs), .. }) => { + let results = sqlx::query!( + " SELECT result, id FROM completed_job WHERE id = ANY($1) AND workspace_id = $2 ", - jobs.as_slice(), - w_id - ) - .fetch_all(&mut tx) - .await? - .into_iter() - .map(|r| (r.id, r.result)) - .collect::>(); + jobs.as_slice(), + w_id + ) + .fetch_all(&mut tx) + .await? + .into_iter() + .map(|r| (r.id, r.result)) + .collect::>(); - let results = jobs - .iter() - .map(|j| { - results - .get(j) - .ok_or_else(|| Error::InternalErr(format!("missing job result for {}", j))) - }) - .collect::, _>>()?; + let results = jobs + .iter() + .map(|j| { + results.get(j).ok_or_else(|| { + Error::InternalErr(format!("missing job result for {}", j)) + }) + }) + .collect::, _>>()?; - json!(results) - } - _ => result, - }; + json!(results) + } + _ => result, + }; - if matches!(&new_status, Some(FlowStatusModule::Success { .. })) { - sqlx::query( - " + if matches!(&new_status, Some(FlowStatusModule::Success { .. })) { + sqlx::query( + " UPDATE queue SET flow_status = flow_status - 'retry' WHERE id = $1 RETURNING flow_status ", - ) - .bind(flow) - .execute(&mut tx) - .await - .context("remove flow status retry")?; - } - - let flow_job = get_queued_job(flow, w_id, &mut tx) - .await? - .ok_or_else(|| Error::InternalErr(format!("requiring flow to be in the queue")))?; - - let job_root = flow_job - .root_job - .map(|x| x.to_string()) - .unwrap_or_else(|| "none".to_string()); - tracing::info!(id = %flow_job.id, root_id = %job_root, "update flow status"); - - let raw_flow = flow_job.parse_raw_flow(); - let module = raw_flow.as_ref().and_then(|module| { - module_index.and_then(|i| module.modules.get(i).or(module.failure_module.as_ref())) - }); - - let should_continue_flow = match success { - _ if stop_early => false, - _ if flow_job.canceled => false, - true => !is_last_step, - false if unrecoverable => false, - false if skip_failure => !is_last_step, - false - if next_retry( - &module.and_then(|m| m.retry.clone()).unwrap_or_default(), - &old_status.retry, ) - .is_some() => - { - true + .bind(flow) + .execute(&mut tx) + .await + .context("remove flow status retry")?; } - false if has_failure_module(flow, &mut tx).await? && !is_failure_step => true, - false => false, - }; - if old_status.step == 0 - && !flow_job.is_flow_step - && flow_job.schedule_path.is_some() - && flow_job.script_path.is_some() - { - tx = schedule_again_if_scheduled( - tx, - db, - flow_job.schedule_path.as_ref().unwrap(), - flow_job.script_path.as_ref().unwrap(), - &w_id, + let flow_job = get_queued_job(flow, w_id, &mut tx) + .await? + .ok_or_else(|| Error::InternalErr(format!("requiring flow to be in the queue")))?; + + let job_root = flow_job + .root_job + .map(|x| x.to_string()) + .unwrap_or_else(|| "none".to_string()); + tracing::info!(id = %flow_job.id, root_id = %job_root, "update flow status"); + + let module = { + let raw_flow = flow_job.parse_raw_flow(); + if let Some(raw_flow) = raw_flow { + if let Some(i) = module_index { + if let Some(module) = raw_flow.modules.get(i) { + Some(module.clone()) + } else { + raw_flow.failure_module + } + } else { + None + } + } else { + None + } + }; + + let should_continue_flow = match success { + _ if stop_early => false, + _ if flow_job.canceled => false, + true => !is_last_step, + false if unrecoverable => false, + false if skip_failure => !is_last_step, + false + if next_retry( + &module.and_then(|m| m.retry.clone()).unwrap_or_default(), + &old_status.retry, + ) + .is_some() => + { + true + } + false if has_failure_module(flow, &mut tx).await? && !is_failure_step => true, + false => false, + }; + + if old_status.step == 0 + && !flow_job.is_flow_step + && flow_job.schedule_path.is_some() + && flow_job.script_path.is_some() + { + tx = schedule_again_if_scheduled( + tx, + db, + flow_job.schedule_path.as_ref().unwrap(), + flow_job.script_path.as_ref().unwrap(), + &w_id, + ) + .await?; + } + tx.commit().await?; + ( + should_continue_flow, + flow_job, + stop_early, + skip_if_stop_early, + nresult, ) - .await?; - } - tx.commit().await?; + }; let done = if !should_continue_flow { let logs = if flow_job.canceled { @@ -464,7 +493,7 @@ pub async fn update_flow_status_after_job_completion( &flow_job, success, stop_early && skip_if_stop_early, - result.clone(), + nresult.clone(), logs, ) .await?; @@ -475,7 +504,7 @@ pub async fn update_flow_status_after_job_completion( &flow_job, db, client, - result.clone(), + nresult.clone(), same_worker_tx.clone(), worker_dir, base_internal_url, @@ -503,6 +532,8 @@ pub async fn update_flow_status_after_job_completion( } if let Some(parent_job) = flow_job.parent_job { + drop(flow_job); + tracing::error!("update flow status after job completion"); update_flow_status_after_job_completion( db, client, @@ -510,7 +541,7 @@ pub async fn update_flow_status_after_job_completion( &flow, w_id, success, - result, + nresult, metrics, false, same_worker_tx.clone(), @@ -521,6 +552,7 @@ pub async fn update_flow_status_after_job_completion( None }, base_internal_url, + depth + 1, ) .await?; } @@ -832,6 +864,7 @@ async fn push_next_flow_job( worker_dir, None, base_internal_url, + 0, ) .await; }