minor worker_flow refactor

This commit is contained in:
Ruben Fiszel
2023-06-23 16:38:21 +02:00
parent dd9526d8bf
commit 3027272873
+130 -169
View File
@@ -1291,13 +1291,11 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
}
};
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<IdContext>,
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<String>,
) -> error::Result<(sqlx::Transaction<'c, sqlx::Postgres>, NextFlowTransform)> {
) -> error::Result<NextFlowTransform> {
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 }),
))
}
}