diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 61598c368f..c398c96407 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -7,7 +7,6 @@ */ use std::collections::HashMap; -use std::iter; use std::time::Duration; use crate::jobs::{add_completed_job, add_completed_job_error, schedule_again_if_scheduled}; @@ -15,7 +14,6 @@ use crate::js_eval::{eval_timeout, IdContext}; use crate::{worker, AuthedClient, KEEP_JOB_DIR}; use anyhow::Context; use async_recursion::async_recursion; -use dyn_iter::DynIter; use serde_json::{json, Map, Value}; use tokio::sync::mpsc::Sender; use tracing::instrument; @@ -1207,45 +1205,62 @@ async fn push_next_flow_job( let continue_on_same_worker = flow.same_worker && module.suspend.is_none() && module.sleep.is_none(); - let zipped = { - let all_args = match &next_status { - NextStatus::NextLoopIteration(NextIteration { new_args, .. }) => { - let args = args.as_ref().map(|args| { - let mut args = args.clone(); - args.extend(new_args.clone()); - args - }); - DynIter::new(iter::once(args)) - } - - NextStatus::AllFlowJobs { - branchall: Some(BranchAllStatus { len, .. }), - iterator: None, - .. - } => DynIter::new((0..*len).map(|_| args.as_ref().map(|args| args.clone()))), - NextStatus::AllFlowJobs { - branchall: None, - iterator: Some(Iterator { itered, .. }), - .. - } => DynIter::new(itered.into_iter().enumerate().map(|(i, v)| { - args.as_ref().map(|args| { - let mut new_args = args.clone(); - new_args.insert("iter".to_string(), json!({ "index": i, "value": v })); - new_args - }) - })), - - _ => DynIter::new(iter::once(args.as_ref().map(|m| m.clone()))), - }; - - job_payloads.into_iter().zip(all_args).collect::>() - }; - /* Finally, push the job into the queue */ let mut tx = db.begin().await?; let mut uuids = vec![]; - for (payload, args) in zipped { + let len = match &job_payloads { + ContinuePayload::SingleJob(_) => 1, + ContinuePayload::BranchAllJobs(payloads) => payloads.len(), + ContinuePayload::ForloopJobs { n, .. } => *n, + }; + for i in (0..len).into_iter() { + let payload = match &job_payloads { + ContinuePayload::SingleJob(payload) => payload.clone(), + ContinuePayload::BranchAllJobs(payloads) => payloads[i].clone(), + ContinuePayload::ForloopJobs { modules, .. } => { + let mut fm = flow.failure_module.clone(); + if let Some(mut failure_module) = flow.failure_module.clone() { + failure_module.id_append(&format!("{}/{}", status.step, i)); + fm = Some(failure_module); + } + JobPayload::RawFlow { + value: FlowValue { + modules: (*modules).clone(), + failure_module: fm.clone(), + same_worker: flow.same_worker, + }, + path: Some(format!("{}/loop-{}", flow_job.script_path(), i)), + } + } + }; + let args = match &next_status { + NextStatus::AllFlowJobs { + branchall: Some(BranchAllStatus { .. }), + iterator: None, + .. + } => args.as_ref().map(|args| args.clone()), + NextStatus::NextLoopIteration(NextIteration { new_args, .. }) => { + args.as_ref().map(|args| { + let mut args = args.clone(); + args.extend(new_args.clone()); + args + }) + } + NextStatus::AllFlowJobs { + branchall: None, + iterator: Some(Iterator { itered, .. }), + .. + } => args.as_ref().map(|args| { + let mut new_args = args.clone(); + new_args.insert( + "iter".to_string(), + json!({ "index": i, "value": itered[i] }), + ); + new_args + }), + _ => args.as_ref().map(|args| args.clone()), + }; let (ok, err) = match args { Ok(v) => (Some(v), None), Err(e) => (None, Some(e)), @@ -1485,9 +1500,15 @@ enum NextStatus { }, } +enum ContinuePayload { + SingleJob(JobPayload), + ForloopJobs { n: usize, modules: Vec }, + BranchAllJobs(Vec), +} + enum NextFlowTransform { EmptyInnerFlows, - Continue(Vec, NextStatus), + Continue(ContinuePayload, NextStatus), } // a similar function exists on the backend @@ -1521,13 +1542,19 @@ async fn compute_next_flow_transform<'c>( match &module.value { FlowModuleValue::Identity => Ok(( tx, - NextFlowTransform::Continue(vec![JobPayload::Identity], NextStatus::NextStep), + NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayload::Identity), + NextStatus::NextStep, + ), )), FlowModuleValue::Flow { path, .. } => { let payload = JobPayload::Flow(path.to_string()); Ok(( tx, - NextFlowTransform::Continue(vec![payload], NextStatus::NextStep), + NextFlowTransform::Continue( + ContinuePayload::SingleJob(payload), + NextStatus::NextStep, + ), )) } FlowModuleValue::Script { path: script_path, hash: script_hash, .. } => { @@ -1541,7 +1568,10 @@ async fn compute_next_flow_transform<'c>( }; Ok(( tx, - NextFlowTransform::Continue(vec![payload], NextStatus::NextStep), + NextFlowTransform::Continue( + ContinuePayload::SingleJob(payload), + NextStatus::NextStep, + ), )) } FlowModuleValue::RawScript { path, content, language, lock, .. } => { @@ -1551,12 +1581,12 @@ async fn compute_next_flow_transform<'c>( Ok(( tx, NextFlowTransform::Continue( - vec![JobPayload::Code(RawCode { + ContinuePayload::SingleJob(JobPayload::Code(RawCode { path, content: content.clone(), language: language.clone(), lock: lock.clone(), - })], + })), NextStatus::NextStep, ), )) @@ -1654,14 +1684,14 @@ async fn compute_next_flow_transform<'c>( Ok(( tx, NextFlowTransform::Continue( - vec![JobPayload::RawFlow { + ContinuePayload::SingleJob(JobPayload::RawFlow { value: FlowValue { modules: (*modules).clone(), failure_module: fm, same_worker: flow.same_worker, }, path: Some(format!("{}/loop-{}", flow_job.script_path(), ns.index)), - }], + }), NextStatus::NextLoopIteration(ns), ), )) @@ -1669,23 +1699,10 @@ async fn compute_next_flow_transform<'c>( LoopStatus::ParallelIteration { itered } => Ok(( tx, NextFlowTransform::Continue( - (0..itered.len()) - .map(|i| { - let mut fm = flow.failure_module.clone(); - if let Some(mut failure_module) = flow.failure_module.clone() { - failure_module.id_append(&format!("{}/{}", status.step, i)); - fm = Some(failure_module); - } - JobPayload::RawFlow { - value: FlowValue { - modules: (*modules).clone(), - failure_module: fm.clone(), - same_worker: flow.same_worker, - }, - path: Some(format!("{}/loop-{}", flow_job.script_path(), i)), - } - }) - .collect(), + ContinuePayload::ForloopJobs { + n: itered.len(), + modules: (*modules).clone(), + }, NextStatus::AllFlowJobs { branchall: None, iterator: Some(windmill_common::flow_status::Iterator { @@ -1746,7 +1763,7 @@ async fn compute_next_flow_transform<'c>( Ok(( tx, NextFlowTransform::Continue( - vec![JobPayload::RawFlow { + ContinuePayload::SingleJob(JobPayload::RawFlow { value: FlowValue { modules, failure_module: fm, @@ -1757,7 +1774,7 @@ async fn compute_next_flow_transform<'c>( flow_job.script_path(), status.step )), - }], + }), NextStatus::BranchChosen(branch), ), )) @@ -1773,32 +1790,34 @@ async fn compute_next_flow_transform<'c>( return Ok(( tx, NextFlowTransform::Continue( - branches - .iter() - .enumerate() - .map(|(i, b)| { - let mut fm = flow.failure_module.clone(); - if let Some(mut failure_module) = - flow.failure_module.clone() - { - failure_module - .id_append(&format!("{}/{i}", status.step)); - fm = Some(failure_module); - } - JobPayload::RawFlow { - value: FlowValue { - modules: b.modules.clone(), - failure_module: fm.clone(), - same_worker: flow.same_worker, - }, - path: Some(format!( - "{}/branchall-{}", - flow_job.script_path(), - i - )), - } - }) - .collect(), + ContinuePayload::BranchAllJobs( + branches + .iter() + .enumerate() + .map(|(i, b)| { + let mut fm = flow.failure_module.clone(); + if let Some(mut failure_module) = + flow.failure_module.clone() + { + failure_module + .id_append(&format!("{}/{i}", status.step)); + fm = Some(failure_module); + } + JobPayload::RawFlow { + value: FlowValue { + modules: b.modules.clone(), + failure_module: fm.clone(), + same_worker: flow.same_worker, + }, + path: Some(format!( + "{}/branchall-{}", + flow_job.script_path(), + i + )), + } + }) + .collect(), + ), NextStatus::AllFlowJobs { branchall: Some(BranchAllStatus { branch: 0, @@ -1855,7 +1874,7 @@ async fn compute_next_flow_transform<'c>( Ok(( tx, NextFlowTransform::Continue( - vec![JobPayload::RawFlow { + ContinuePayload::SingleJob(JobPayload::RawFlow { value: FlowValue { modules, failure_module: fm.clone(), @@ -1866,7 +1885,7 @@ async fn compute_next_flow_transform<'c>( flow_job.script_path(), branch_status.branch )), - }], + }), NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }), ), ))