fix: improve error handler behavior wrt to parrallel branchall & forloops (#6273)
* error handler improvement * fix: improve error handler behavior with parallel for loops * all * Error handler
This commit is contained in:
@@ -98,7 +98,7 @@ pub async fn update_flow_status_after_job_completion(
|
||||
success,
|
||||
result,
|
||||
stop_early_override,
|
||||
skip_error_handler: false,
|
||||
has_triggered_error_handler: false,
|
||||
};
|
||||
let mut unrecoverable = unrecoverable;
|
||||
loop {
|
||||
@@ -115,7 +115,7 @@ pub async fn update_flow_status_after_job_completion(
|
||||
same_worker_tx,
|
||||
worker_dir,
|
||||
rec.stop_early_override,
|
||||
rec.skip_error_handler,
|
||||
rec.has_triggered_error_handler,
|
||||
worker_name,
|
||||
job_completed_tx.clone(),
|
||||
#[cfg(feature = "benchmark")]
|
||||
@@ -140,7 +140,7 @@ pub async fn update_flow_status_after_job_completion(
|
||||
same_worker_tx,
|
||||
worker_dir,
|
||||
rec.stop_early_override,
|
||||
rec.skip_error_handler,
|
||||
rec.has_triggered_error_handler,
|
||||
worker_name,
|
||||
job_completed_tx.clone(),
|
||||
#[cfg(feature = "benchmark")]
|
||||
@@ -188,7 +188,7 @@ pub struct RecUpdateFlowStatusAfterJobCompletion {
|
||||
success: bool,
|
||||
result: Arc<Box<RawValue>>,
|
||||
stop_early_override: Option<bool>,
|
||||
skip_error_handler: bool,
|
||||
has_triggered_error_handler: bool,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
@@ -233,11 +233,12 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
same_worker_tx: &SameWorkerSender,
|
||||
worker_dir: &str,
|
||||
stop_early_override: Option<bool>,
|
||||
skip_error_handler: bool,
|
||||
has_triggered_error_handler: bool,
|
||||
worker_name: &str,
|
||||
job_completed_tx: JobCompletedSender,
|
||||
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
|
||||
) -> error::Result<UpdateFlowStatusAfterJobCompletion> {
|
||||
let mut has_triggered_error_handler = has_triggered_error_handler;
|
||||
add_time!(bench, "update flow status internal START");
|
||||
let (
|
||||
should_continue_flow,
|
||||
@@ -292,6 +293,11 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
Step::Step(i) => flow_value.modules.get(i),
|
||||
_ => None,
|
||||
};
|
||||
|
||||
if current_module.is_some_and(|x| x.is_flow()) {
|
||||
has_triggered_error_handler = false;
|
||||
}
|
||||
|
||||
let module_status = match module_step {
|
||||
Step::PreprocessorStep => old_status
|
||||
.preprocessor_module
|
||||
@@ -565,10 +571,39 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
let branches = current_module
|
||||
.and_then(|x| x.get_branches_skip_failures().ok())
|
||||
.map(|x| {
|
||||
x.branches
|
||||
.iter()
|
||||
.map(|b| b.skip_failure.unwrap_or(false))
|
||||
.collect::<Vec<_>>()
|
||||
});
|
||||
|
||||
let mut njobs = Vec::new();
|
||||
|
||||
let jobs_filtered = if branchall.is_some() {
|
||||
if let Some(branches) = branches {
|
||||
for (branch, job) in branches.iter().zip(jobs.iter()) {
|
||||
if !branch {
|
||||
njobs.push(job.clone());
|
||||
}
|
||||
}
|
||||
njobs.as_slice()
|
||||
} else {
|
||||
jobs.as_slice()
|
||||
}
|
||||
} else {
|
||||
jobs.as_slice()
|
||||
};
|
||||
|
||||
|
||||
|
||||
let new_status = if skip_loop_failures
|
||||
|| sqlx::query_scalar!(
|
||||
"SELECT status = 'success' OR status = 'skipped' AS \"success!\" FROM v2_job_completed WHERE id = ANY($1)",
|
||||
jobs.as_slice()
|
||||
jobs_filtered
|
||||
)
|
||||
.fetch_all(&mut *tx)
|
||||
.await
|
||||
@@ -620,7 +655,15 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
tracing::info!(
|
||||
"parallel iteration {job_id_for_status} of flow {flow} has finished",
|
||||
);
|
||||
(true, Some(new_status))
|
||||
|
||||
// for parallel branchall and forloop, we do not want to trigger the error handler again at the forloop/branchall node since it was already triggered at the leaf level
|
||||
// so we want to ignore the has_triggered_error_handler flag at the forloop/branchall node and reset it based on if the node is a success or failure
|
||||
if !success && flow_value.failure_module.is_some() {
|
||||
has_triggered_error_handler = true;
|
||||
} else {
|
||||
has_triggered_error_handler = false;
|
||||
}
|
||||
(success, Some(new_status))
|
||||
} else {
|
||||
add_time!(bench, "handle parallel flow start");
|
||||
tx.commit().await?;
|
||||
@@ -791,17 +834,9 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
}
|
||||
};
|
||||
|
||||
let skip_parallel_branchall_failure = match (module_status, new_status.as_ref()) {
|
||||
(
|
||||
FlowStatusModule::InProgress { branchall: Some(_), parallel: true, .. },
|
||||
Some(FlowStatusModule::Success { flow_jobs_success, .. }),
|
||||
) => compute_skip_branchall_failure(0, true, current_module, flow_jobs_success).await?,
|
||||
(
|
||||
FlowStatusModule::InProgress { branchall: Some(_), parallel: true, .. },
|
||||
Some(FlowStatusModule::Failure { flow_jobs_success, .. }),
|
||||
) => compute_skip_branchall_failure(0, true, current_module, flow_jobs_success).await?,
|
||||
_ => false,
|
||||
};
|
||||
|
||||
|
||||
|
||||
let step_counter = if inc_step_counter {
|
||||
sqlx::query!(
|
||||
"UPDATE v2_job_status
|
||||
@@ -1113,6 +1148,7 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
.map(|x| x.to_string())
|
||||
.unwrap_or_else(|| "none".to_string());
|
||||
|
||||
|
||||
let should_continue_flow = match success {
|
||||
_ if stop_early => false,
|
||||
_ if flow_job.is_canceled() => false,
|
||||
@@ -1120,10 +1156,10 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
false if unrecoverable => false,
|
||||
false
|
||||
if skip_seq_branch_failure
|
||||
|| skip_parallel_branchall_failure
|
||||
|| skip_loop_failures
|
||||
|| continue_on_error =>
|
||||
{
|
||||
|
||||
!is_last_step
|
||||
}
|
||||
false
|
||||
@@ -1152,7 +1188,7 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
}
|
||||
false
|
||||
if !is_failure_step
|
||||
&& !skip_error_handler
|
||||
&& !has_triggered_error_handler
|
||||
&& flow_value.failure_module.is_some() =>
|
||||
{
|
||||
true
|
||||
@@ -1160,7 +1196,10 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
false => false,
|
||||
};
|
||||
|
||||
tracing::info!(id = %flow_job.id, root_id = %job_root, "flow should continue: {should_continue_flow}");
|
||||
tracing::info!(id = %flow_job.id, root_id = %job_root, success = %success, stop_early = %stop_early, is_last_step = %is_last_step, unrecoverable = %unrecoverable,
|
||||
skip_seq_branch_failure = %skip_seq_branch_failure, skip_loop_failures = %skip_loop_failures,
|
||||
current_module_id = %current_module.map(|x| x.id.clone()).unwrap_or_default(),
|
||||
continue_on_error = %continue_on_error, should_continue_flow = %should_continue_flow, "computed if flow should continue");
|
||||
|
||||
(
|
||||
should_continue_flow,
|
||||
@@ -1264,7 +1303,6 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
}
|
||||
let success = success
|
||||
&& (!is_failure_step || result_has_recover_true(nresult.clone()))
|
||||
&& !skip_error_handler
|
||||
&& stop_early_err_msg.is_none();
|
||||
|
||||
add_time!(bench, "flow status update 1");
|
||||
@@ -1356,7 +1394,7 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
} else {
|
||||
None
|
||||
},
|
||||
skip_error_handler: skip_error_handler || is_failure_step,
|
||||
has_triggered_error_handler: has_triggered_error_handler || is_failure_step,
|
||||
},
|
||||
));
|
||||
}
|
||||
@@ -1779,10 +1817,11 @@ async fn push_next_flow_job(
|
||||
.flow_innermost_root_job
|
||||
.map(|x| x.to_string())
|
||||
.unwrap_or_else(|| "none".to_string());
|
||||
tracing::info!(id = %flow_job.id, root_id = %job_root, "pushing next flow job");
|
||||
|
||||
let mut step = Step::from_i32_and_len(status.step, flow.modules.len());
|
||||
|
||||
tracing::info!(id = %flow_job.id, root_id = %job_root, step = ?step, "pushing next flow job");
|
||||
|
||||
let mut status_module = match step {
|
||||
Step::Step(i) => status
|
||||
.modules
|
||||
@@ -1796,6 +1835,8 @@ async fn push_next_flow_job(
|
||||
Step::FailureStep => status.failure_module.module_status.clone(),
|
||||
};
|
||||
|
||||
// tracing::error!("status_module: {status_module:#?}");
|
||||
|
||||
let fj: mappable_rc::Marc<MiniPulledJob> = flow_job.clone().into();
|
||||
let arc_flow_job_args: Marc<HashMap<String, Box<RawValue>>> = Marc::map(fj, |x| {
|
||||
if let Some(args) = &x.args {
|
||||
|
||||
@@ -88,7 +88,10 @@ shouldUsePortal={true} -->
|
||||
{#if kind === 'trigger'}
|
||||
<SchedulePollIcon size={14} />
|
||||
{:else if kind === 'failure'}
|
||||
<Bug size={14} />
|
||||
<div class="flex items-center gap-1">
|
||||
<Bug size={14} />
|
||||
<span class="text-xs w-20">Error Handler</span>
|
||||
</div>
|
||||
{:else}
|
||||
<Cross size={iconSize} />
|
||||
{/if}
|
||||
|
||||
@@ -821,7 +821,6 @@ export function graphBuilder(
|
||||
type: 'branchOneStart'
|
||||
}
|
||||
|
||||
|
||||
nodes.push(defaultBranch)
|
||||
|
||||
addEdge(module.id, defaultBranch.id, { rootId: module.id, branch: 0 }, prefix, {
|
||||
@@ -884,7 +883,7 @@ export function graphBuilder(
|
||||
let expanded = expandedSubflows[module.id]
|
||||
if (expanded) {
|
||||
expanded = $state.snapshot(expanded)
|
||||
const startId = `${module.id}-subflow-start`
|
||||
const startId = `${module.id}`
|
||||
const idWithoutPrefix = module.id.startsWith('subflow:')
|
||||
? module.id.substring(8)
|
||||
: module.id
|
||||
|
||||
Reference in New Issue
Block a user