From 3027272873f46beed8d97deca137ed603b694405 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 23 Jun 2023 16:38:13 +0200 Subject: [PATCH] minor worker_flow refactor --- backend/windmill-worker/src/worker_flow.rs | 299 +++++++++------------ 1 file changed, 130 insertions(+), 169 deletions(-) diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index a070ac4348..62d280cec3 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1291,13 +1291,11 @@ async fn push_next_flow_job } }; - let tx = db.begin().await?; - - let (tx, next_flow_transform) = compute_next_flow_transform( + let next_flow_transform = compute_next_flow_transform( flow_job, &flow, transform_context, - tx, + db, &module, &status, &status_module, @@ -1308,7 +1306,6 @@ async fn push_next_flow_job approvers.clone(), ) .await?; - tx.commit().await?; tracing::info!(id = %flow_job.id, root_id = %job_root, "next flow transform computed"); let (job_payloads, next_status) = match next_flow_transform { @@ -1340,7 +1337,6 @@ async fn push_next_flow_job flow.same_worker && module.suspend.is_none() && module.sleep.is_none(); /* Finally, push the job into the queue */ - let mut tx = (rsmq.clone(), db.begin().await?).into(); let mut uuids = vec![]; let len = match &job_payloads { @@ -1348,6 +1344,9 @@ async fn push_next_flow_job ContinuePayload::BranchAllJobs(payloads) => payloads.len(), ContinuePayload::ForloopJobs { n, .. } => *n, }; + + let mut tx: QueueTransaction<'_, R> = (rsmq.clone(), db.begin().await?).into(); + for i in (0..len).into_iter() { let payload_tag = match &job_payloads { ContinuePayload::SingleJob(payload) => payload.clone(), @@ -1663,11 +1662,11 @@ enum NextFlowTransform { Continue(ContinuePayload, NextStatus), } -async fn compute_next_flow_transform<'c>( +async fn compute_next_flow_transform( flow_job: &QueuedJob, flow: &FlowValue, by_id: Option, - mut tx: sqlx::Transaction<'c, sqlx::Postgres>, + db: &DB, module: &FlowModule, status: &FlowStatus, status_module: &FlowStatusModule, @@ -1676,80 +1675,67 @@ async fn compute_next_flow_transform<'c>( client: &AuthedClient, resumes: &[Value], approvers: Vec, -) -> error::Result<(sqlx::Transaction<'c, sqlx::Postgres>, NextFlowTransform)> { +) -> error::Result { if module.mock.is_some() && module.mock.as_ref().unwrap().enabled { - return Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { - payload: JobPayload::Identity, - tag: None, - }), - NextStatus::NextStep, - ), + return Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { + payload: JobPayload::Identity, + tag: None, + }), + NextStatus::NextStep, )); } - let trivial_next_job = |tx, payload| { - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag: None }), - NextStatus::NextStep, - ), + let trivial_next_job = |payload| { + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag: None }), + NextStatus::NextStep, )) }; match &module.value { - FlowModuleValue::Identity => trivial_next_job(tx, JobPayload::Identity), - FlowModuleValue::Graphql => trivial_next_job(tx, JobPayload::Graphql), - FlowModuleValue::Http => trivial_next_job(tx, JobPayload::Http), - FlowModuleValue::Postgresql => trivial_next_job(tx, JobPayload::Postgresql), + FlowModuleValue::Identity => trivial_next_job(JobPayload::Identity), + FlowModuleValue::Graphql => trivial_next_job(JobPayload::Graphql), + FlowModuleValue::Http => trivial_next_job(JobPayload::Http), + FlowModuleValue::Postgresql => trivial_next_job(JobPayload::Postgresql), FlowModuleValue::Flow { path, .. } => { let payload = JobPayload::Flow(path.to_string()); - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag: None }), - NextStatus::NextStep, - ), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag: None }), + NextStatus::NextStep, )) } FlowModuleValue::Script { path: script_path, hash: script_hash, .. } => { let (payload, tag) = if script_hash.is_none() { + let mut tx: sqlx::Transaction<'_, sqlx::Postgres> = db.begin().await?; script_path_to_payload(script_path, &mut tx, &flow_job.workspace_id).await? } else { let hash = script_hash.clone().unwrap(); + let mut tx: sqlx::Transaction<'_, sqlx::Postgres> = db.begin().await?; let tag = script_hash_to_tag(&hash, &mut tx, &flow_job.workspace_id).await?; ( JobPayload::ScriptHash { hash, path: script_path.to_owned() }, tag, ) }; - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag }), - NextStatus::NextStep, - ), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag }), + NextStatus::NextStep, )) } FlowModuleValue::RawScript { path, content, language, lock, tag, .. } => { let path = path .clone() .or_else(|| Some(format!("{}/step-{}", flow_job.script_path(), status.step))); - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { - payload: JobPayload::Code(RawCode { - path, - content: content.clone(), - language: language.clone(), - lock: lock.clone(), - }), - tag: tag.clone(), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { + payload: JobPayload::Code(RawCode { + path, + content: content.clone(), + language: language.clone(), + lock: lock.clone(), }), - NextStatus::NextStep, - ), + tag: tag.clone(), + }), + NextStatus::NextStep, )) } /* forloop modules are expected set `iter: { value: Value, index: usize }` as job arguments */ @@ -1841,50 +1827,34 @@ async fn compute_next_flow_transform<'c>( }; match next_loop_status { - LoopStatus::EmptyIterator => Ok((tx, NextFlowTransform::EmptyInnerFlows)), + LoopStatus::EmptyIterator => Ok(NextFlowTransform::EmptyInnerFlows), LoopStatus::NextIteration(ns) => { let mut fm = flow.failure_module.clone(); if let Some(mut failure_module) = flow.failure_module.clone() { failure_module.id_append(&format!("{}/{}", status.step, ns.index)); fm = Some(failure_module); } - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { - payload: 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 - )), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { + payload: JobPayload::RawFlow { + value: FlowValue { + modules: (*modules).clone(), + failure_module: fm, + same_worker: flow.same_worker, }, - tag: None, - }), - NextStatus::NextLoopIteration(ns), - ), + path: Some(format!("{}/loop-{}", flow_job.script_path(), ns.index)), + }, + tag: None, + }), + NextStatus::NextLoopIteration(ns), )) } - LoopStatus::ParallelIteration { itered } => Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::ForloopJobs { - n: itered.len(), - modules: (*modules).clone(), - }, - NextStatus::AllFlowJobs { - branchall: None, - iterator: Some(windmill_common::flow_status::Iterator { - index: 0, - itered, - }), - }, - ), + LoopStatus::ParallelIteration { itered } => Ok(NextFlowTransform::Continue( + ContinuePayload::ForloopJobs { n: itered.len(), modules: (*modules).clone() }, + NextStatus::AllFlowJobs { + branchall: None, + iterator: Some(windmill_common::flow_status::Iterator { index: 0, itered }), + }, )), } } @@ -1935,26 +1905,23 @@ async fn compute_next_flow_transform<'c>( failure_module.id_append(&status.step.to_string()); fm = Some(failure_module); } - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { - payload: JobPayload::RawFlow { - value: FlowValue { - modules, - failure_module: fm, - same_worker: flow.same_worker, - }, - path: Some(format!( - "{}/branchone-{}", - flow_job.script_path(), - status.step - )), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { + payload: JobPayload::RawFlow { + value: FlowValue { + modules, + failure_module: fm, + same_worker: flow.same_worker, }, - tag: None, - }), - NextStatus::BranchChosen(branch), - ), + path: Some(format!( + "{}/branchone-{}", + flow_job.script_path(), + status.step + )), + }, + tag: None, + }), + NextStatus::BranchChosen(branch), )) } FlowModuleValue::BranchAll { branches, parallel, .. } => { @@ -1963,51 +1930,48 @@ async fn compute_next_flow_transform<'c>( | FlowStatusModule::WaitingForEvents { .. } | FlowStatusModule::WaitingForExecutor { .. } => { if branches.is_empty() { - return Ok((tx, NextFlowTransform::EmptyInnerFlows)); + return Ok(NextFlowTransform::EmptyInnerFlows); } else if *parallel { - return Ok(( - tx, - NextFlowTransform::Continue( - 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); - } - JobPayloadWithTag { - payload: 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 - )), + return Ok(NextFlowTransform::Continue( + 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); + } + JobPayloadWithTag { + payload: JobPayload::RawFlow { + value: FlowValue { + modules: b.modules.clone(), + failure_module: fm.clone(), + same_worker: flow.same_worker, }, - tag: None, - } - }) - .collect(), - ), - NextStatus::AllFlowJobs { - branchall: Some(BranchAllStatus { - branch: 0, - previous_result: last_result, - len: branches.len(), - }), - iterator: None, - }, + path: Some(format!( + "{}/branchall-{}", + flow_job.script_path(), + i + )), + }, + tag: None, + } + }) + .collect(), ), + NextStatus::AllFlowJobs { + branchall: Some(BranchAllStatus { + branch: 0, + previous_result: last_result, + len: branches.len(), + }), + iterator: None, + }, )); } else { ( @@ -2052,26 +2016,23 @@ async fn compute_next_flow_transform<'c>( failure_module.id_append(&format!("{}/{}", status.step, branch_status.branch)); fm = Some(failure_module); } - Ok(( - tx, - NextFlowTransform::Continue( - ContinuePayload::SingleJob(JobPayloadWithTag { - payload: JobPayload::RawFlow { - value: FlowValue { - modules, - failure_module: fm.clone(), - same_worker: flow.same_worker, - }, - path: Some(format!( - "{}/branchall-{}", - flow_job.script_path(), - branch_status.branch - )), + Ok(NextFlowTransform::Continue( + ContinuePayload::SingleJob(JobPayloadWithTag { + payload: JobPayload::RawFlow { + value: FlowValue { + modules, + failure_module: fm.clone(), + same_worker: flow.same_worker, }, - tag: None, - }), - NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }), - ), + path: Some(format!( + "{}/branchall-{}", + flow_job.script_path(), + branch_status.branch + )), + }, + tag: None, + }), + NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }), )) } }