From e7d458278969897aa7312dcd20a8091aaad772d7 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 14 Feb 2025 04:12:55 +0100 Subject: [PATCH] fix: improve handling of empty branches and loops --- backend/windmill-worker/src/worker_flow.rs | 83 +++++++++++++--------- 1 file changed, 51 insertions(+), 32 deletions(-) diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index f265adaf93..e21df1133f 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1589,21 +1589,24 @@ pub async fn handle_flow( ); } } + let mut rec = Some(PushNextFlowJobRec { flow_job: flow_job, status: status }); + while let Some(nrec) = rec { + rec = push_next_flow_job( + nrec.flow_job, + nrec.status, + flow, + db, + client, + last_result.clone(), + same_worker_tx.clone(), + worker_dir, + job_completed_tx.clone(), + worker_name, + ) + .warn_after_seconds(10) + .await?; + } - push_next_flow_job( - flow_job, - status, - flow, - db, - client, - last_result, - same_worker_tx, - worker_dir, - job_completed_tx, - worker_name, - ) - .warn_after_seconds(10) - .await?; Ok(()) } @@ -1646,6 +1649,11 @@ fn potentially_crash_for_testing() { lazy_static::lazy_static! { pub static ref EHM: HashMap> = HashMap::new(); } + +struct PushNextFlowJobRec { + flow_job: Arc, + status: FlowStatus, +} // #[async_recursion] // #[instrument(level = "trace", skip_all)] async fn push_next_flow_job( @@ -1659,7 +1667,7 @@ async fn push_next_flow_job( worker_dir: &str, job_completed_tx: Sender, worker_name: &str, -) -> error::Result<()> { +) -> error::Result> { let job_root = flow_job .root_job .map(|x| x.to_string()) @@ -1714,7 +1722,7 @@ async fn push_next_flow_job( )) })?; - return Ok(()); + return Ok(None); } if matches!(step, Step::Step(0)) { @@ -1772,7 +1780,7 @@ async fn push_next_flow_job( )) })?; - return Ok(()); + return Ok(None); } } } @@ -1809,7 +1817,7 @@ async fn push_next_flow_job( )) })?; - return Ok(()); + return Ok(None); } } } @@ -2055,7 +2063,7 @@ async fn push_next_flow_job( .await?; tx.commit().await?; - return Ok(()); + return Ok(None); /* cancelled or we're WaitingForEvents but we don't have enough messages (timed out) */ } else { @@ -2112,7 +2120,7 @@ async fn push_next_flow_job( )) })?; - return Ok(()); + return Ok(None); } } } @@ -2412,10 +2420,11 @@ async fn push_next_flow_job( let (job_payloads, next_status) = match next_flow_transform { NextFlowTransform::Continue(job_payload, next_state) => (job_payload, next_state), NextFlowTransform::EmptyInnerFlows => { - sqlx::query!( + let raw_status = sqlx::query_scalar!( "UPDATE v2_job_status SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2) - WHERE id = $3", + WHERE id = $3 + RETURNING flow_status AS \"flow_status: Json>\"", status.step.to_string(), json!(FlowStatusModule::Success { id: status_module.id(), @@ -2429,15 +2438,25 @@ async fn push_next_flow_job( }), flow_job.id ) - .execute(db) - .await?; - // flow is reprocessed by the worker in a state where the module has completed successfully. - // The next steps are pull -> handle flow -> push next flow job -> update flow status since module status is success - same_worker_tx - .send(SameWorkerPayload { job_id: flow_job.id, recoverable: true }) - .await - .expect("send to same worker"); - return Ok(()); + .fetch_optional(db) + .await? + .flatten(); + + let status = raw_status + .as_ref() + .and_then(|v| serde_json::from_str::((**v).get()).ok()); + + if let Some(status) = status { + // // flow is reprocessed by the worker in a state where the module has completed successfully. + return Ok(Some(PushNextFlowJobRec { + flow_job: flow_job, + status: status, + })); + } else { + return Err(Error::BadRequest( + "impossible to parse new flow status after applying innr flows".to_string(), + )); + } } }; @@ -2923,7 +2942,7 @@ async fn push_next_flow_job( .await .map_err(to_anyhow)?; } - return Ok(()); + return Ok(None); } // async fn jump_to_next_step(