improve mem handling for update_flow_status

This commit is contained in:
Ruben Fiszel
2023-04-03 01:13:39 +02:00
parent 0b70212c6d
commit 76a809e035
2 changed files with 333 additions and 297 deletions
+6 -3
View File
@@ -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?;
}
+327 -294
View File
@@ -50,58 +50,62 @@ pub async fn update_flow_status_after_job_completion(
worker_dir: &str,
stop_early_override: Option<bool>,
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::<FlowStatus>(old_status_json)
.ok()
.ok_or_else(|| {
Error::InternalErr(format!("requiring status to be parsabled as FlowStatus"))
})?;
let old_status = serde_json::from_value::<FlowStatus>(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::<HashMap<_, _>>();
jobs.as_slice(),
w_id
)
.fetch_all(&mut tx)
.await?
.into_iter()
.map(|r| (r.id, r.result))
.collect::<HashMap<_, _>>();
let results = jobs
.iter()
.map(|j| {
results
.get(j)
.ok_or_else(|| Error::InternalErr(format!("missing job result for {}", j)))
})
.collect::<Result<Vec<_>, _>>()?;
let results = jobs
.iter()
.map(|j| {
results.get(j).ok_or_else(|| {
Error::InternalErr(format!("missing job result for {}", j))
})
})
.collect::<Result<Vec<_>, _>>()?;
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;
}