From bcde2e62d7846822cfe7f8ef89a9484938829903 Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Wed, 21 Aug 2024 00:47:34 +0200 Subject: [PATCH] feat: improve early stop (#4257) --- backend/tests/worker.rs | 8 + backend/windmill-api/src/flows.rs | 4 + backend/windmill-common/src/flows.rs | 3 + backend/windmill-queue/src/jobs.rs | 1 + backend/windmill-worker/src/worker_flow.rs | 76 ++++- .../flows/content/FlowModuleComponent.svelte | 7 +- .../flows/content/FlowModuleEarlyStop.svelte | 268 +++++++++++++----- .../flows/content/FlowModuleHeader.svelte | 2 +- .../src/lib/components/flows/flowExplorer.ts | 3 + .../lib/components/flows/map/MapItem.svelte | 2 +- .../flows/propPicker/PropPickerWrapper.svelte | 2 + .../propertyPicker/PropPickerResult.svelte | 7 +- openflow.openapi.yaml | 9 + 13 files changed, 314 insertions(+), 78 deletions(-) diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 66c3ee06a7..6e81aabbd3 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1113,6 +1113,7 @@ async fn test_deno_flow(db: Pool) { } .into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1152,6 +1153,7 @@ async fn test_deno_flow(db: Pool) { } .into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1166,6 +1168,7 @@ async fn test_deno_flow(db: Pool) { } .into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1270,6 +1273,7 @@ async fn test_deno_flow_same_worker(db: Pool) { concurrency_time_window_s: None, }.into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1319,6 +1323,7 @@ async fn test_deno_flow_same_worker(db: Pool) { concurrency_time_window_s: None, }.into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1354,6 +1359,7 @@ async fn test_deno_flow_same_worker(db: Pool) { concurrency_time_window_s: None, }.into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1369,6 +1375,7 @@ async fn test_deno_flow_same_worker(db: Pool) { ], }.into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, @@ -1411,6 +1418,7 @@ async fn test_deno_flow_same_worker(db: Pool) { concurrency_time_window_s: None, }.into(), stop_after_if: Default::default(), + stop_after_all_iters_if: Default::default(), summary: Default::default(), suspend: Default::default(), retry: None, diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index ccec53c4dc..0dbe994040 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -1144,6 +1144,7 @@ mod tests { tag_override: None, }), stop_after_if: None, + stop_after_all_iters_if: None, summary: None, suspend: Default::default(), retry: None, @@ -1172,6 +1173,7 @@ mod tests { expr: "foo = 'bar'".to_string(), skip_if_stopped: false, }), + stop_after_all_iters_if: None, summary: None, suspend: Default::default(), retry: None, @@ -1198,6 +1200,7 @@ mod tests { expr: "previous.isEmpty()".to_string(), skip_if_stopped: false, }), + stop_after_all_iters_if: None, summary: None, suspend: Default::default(), retry: None, @@ -1223,6 +1226,7 @@ mod tests { expr: "previous.isEmpty()".to_string(), skip_if_stopped: false, }), + stop_after_all_iters_if: None, summary: None, suspend: Default::default(), retry: None, diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index dbc44afdd9..849a3afa2c 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -236,6 +236,8 @@ pub struct FlowModule { #[serde(skip_serializing_if = "Option::is_none")] pub stop_after_if: Option, #[serde(skip_serializing_if = "Option::is_none")] + pub stop_after_all_iters_if: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub summary: Option, #[serde(skip_serializing_if = "Option::is_none")] pub suspend: Option, @@ -580,6 +582,7 @@ pub fn add_virtual_items_if_necessary(modules: &mut Vec) { id: format!("{}-v", modules[modules.len() - 1].id), value: crate::worker::to_raw_value(&FlowModuleValue::Identity), stop_after_if: None, + stop_after_all_iters_if: None, summary: Some("Virtual module needed for suspend/sleep when last module".to_string()), mock: None, retry: None, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 3e8e466e21..ad43fb6455 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -3349,6 +3349,7 @@ pub async fn push<'c, 'd, R: rsmq_async::RsmqConnection + Send + 'c>( }, ), stop_after_if: None, + stop_after_all_iters_if: None, summary: None, suspend: None, mock: None, diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index c27400c485..10edf6f86a 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -232,22 +232,27 @@ pub async fn update_flow_status_after_job_completion_internal< // "UPDATE FLOW STATUS 2: {module_index:#?} {module_status:#?} {old_status:#?} " // ); - let (skip_loop_failures, parallelism) = if matches!( + let (is_loop, skip_loop_failures, parallelism) = if matches!( module_status, FlowStatusModule::InProgress { iterator: Some(_), .. } ) { let (loop_failures, parallelism) = compute_skip_loop_failures_and_parallelism(flow, old_status.step, db).await?; - (loop_failures.unwrap_or(false), parallelism) + (true, loop_failures.unwrap_or(false), parallelism) } else { - (false, None) + (false, false, None) }; + let is_branch_all = matches!( + module_status, + FlowStatusModule::InProgress { branchall: Some(_), .. } + ); + // 0 length flows are not failure steps let is_failure_step = old_status.step >= old_status.modules.len() as i32 && old_status.modules.len() > 0; - let (mut stop_early, skip_if_stop_early, continue_on_error) = if let Some(se) = + let (mut stop_early, mut skip_if_stop_early, continue_on_error) = if let Some(se) = stop_early_override { //do not stop early if module is a flow step @@ -280,7 +285,18 @@ pub async fn update_flow_status_after_job_completion_internal< .map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e:#}")))?; let stop_early = success + && !is_branch_all && if let Some(expr) = r.stop_early_expr.clone() { + let all_iters = match &module_status { + FlowStatusModule::InProgress { flow_jobs: Some(flow_jobs), .. } + if expr.contains("all_iters") => + { + Some(Arc::new( + retrieve_flow_jobs_results(db, w_id, flow_jobs).await?, + )) + } + _ => None, + }; compute_bool_from_expr( expr, Marc::new( @@ -290,6 +306,7 @@ pub async fn update_flow_status_after_job_completion_internal< .to_owned(), ), result.clone(), + all_iters, None, Some(client), None, @@ -735,6 +752,50 @@ pub async fn update_flow_status_after_job_completion_internal< _ => result.clone(), }; + match &new_status { + Some(FlowStatusModule::Success { .. }) if is_loop || is_branch_all => { + let r_after_all_iters = sqlx::query_as::<_, SkipIfStopped>( + "SELECT + raw_flow->'modules'->$1::int->'stop_after_all_iters_if'->>'expr' as stop_early_expr, + (raw_flow->'modules'->$1::int->'stop_after_all_iters_if'->>'skip_if_stopped')::bool as skip_if_stopped, + NULL as continue_on_error, + args + FROM queue + WHERE id = $2" + ) + .bind(old_status.step) + .bind(flow) + .fetch_one(db) + .await + .map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e:#}")))?; + if let Some(expr) = r_after_all_iters.stop_early_expr { + let should_stop = compute_bool_from_expr( + expr, + Marc::new( + r_after_all_iters + .args + .map(|x| x.0) + .unwrap_or_else(|| serde_json::from_str("{}").unwrap()) + .to_owned(), + ), + nresult.clone(), + None, + None, + Some(client), + None, + None, + ) + .await?; + + if should_stop { + stop_early = should_stop; + skip_if_stop_early = r_after_all_iters.skip_if_stopped.unwrap_or(false); + } + } + } + _ => {} + } + if old_status.retry.fail_count > 0 && matches!(&new_status, Some(FlowStatusModule::Success { .. })) { @@ -1186,6 +1247,7 @@ async fn compute_bool_from_expr( expr: String, flow_args: Marc>>, result: Arc>, + all_iters: Option>>, by_id: Option, client: Option<&AuthedClient>, resumes: Option<(Arc>, Arc>, Arc>)>, @@ -1193,6 +1255,9 @@ async fn compute_bool_from_expr( ) -> error::Result { let mut context = HashMap::with_capacity(if resumes.is_some() { 7 } else { 3 }); context.insert("result".to_string(), result.clone()); + if let Some(all_iters) = all_iters { + context.insert("all_iters".to_string(), all_iters); + } context.insert("previous_result".to_string(), result.clone()); if let Some(resumes) = resumes { @@ -1568,6 +1633,7 @@ async fn push_next_flow_job arc_flow_job_args.clone(), Arc::new(to_raw_value(&json!("{}"))), None, + None, Some(client), None, Some(vec![( @@ -2982,6 +3048,7 @@ async fn compute_next_flow_transform( b.expr.to_string(), arc_flow_job_args.clone(), arc_last_job_result.clone(), + None, Some(idcontext.clone()), Some(client), Some((resumes.clone(), resume.clone(), approvers.clone())), @@ -3258,6 +3325,7 @@ fn is_simple_modules(modules: &Vec, flow: &FlowValue) -> bool { && modules[0].cache_ttl.is_none() && modules[0].retry.is_none() && modules[0].stop_after_if.is_none() + && modules[0].stop_after_all_iters_if.is_none() && (modules[0].mock.is_none() || modules[0].mock.as_ref().is_some_and(|m| !m.enabled)) && flow.failure_module.is_none(); is_simple diff --git a/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte b/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte index 36f990084c..118e449188 100644 --- a/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte +++ b/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte @@ -402,7 +402,12 @@ {#if !$selectedId.includes('failure')} Runtime Cache - + Early Stop Suspend diff --git a/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte b/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte index e52e71b5d7..9226d915fc 100644 --- a/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte +++ b/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte @@ -2,13 +2,14 @@ import SimpleEditor from '$lib/components/SimpleEditor.svelte' import Toggle from '$lib/components/Toggle.svelte' import PropPickerWrapper from '$lib/components/flows/propPicker/PropPickerWrapper.svelte' - import type { FlowModule } from '$lib/gen' + import type { Flow, FlowModule } from '$lib/gen' import Tooltip from '$lib/components/Tooltip.svelte' import type { FlowEditorContext } from '../types' import { getContext } from 'svelte' import { NEVER_TESTED_THIS_FAR } from '../models' import Section from '$lib/components/Section.svelte' import { getStepPropPicker } from '../previousResults' + import { dfs } from '../previousResults' const { flowStateStore, flowStore, previewArgs } = getContext('FlowEditorContext') @@ -26,80 +27,207 @@ false ) + function checkIfParentLoop(flowStore: typeof $flowStore): string | null { + const flow: Flow = JSON.parse(JSON.stringify(flowStore)) + const parents = dfs(flowModule.id, flow, true) + for (const parent of parents.slice(1)) { + if (parent.value.type === 'forloopflow' || parent.value.type === 'whileloopflow') { + return parent.id + } + } + return null + } + + $: isLoop = flowModule.value.type === 'forloopflow' || flowModule.value.type === 'whileloopflow' + $: isBranchAll = flowModule.value.type === 'branchall' $: isStopAfterIfEnabled = Boolean(flowModule.stop_after_if) + $: isStopAfterAllIterationsEnabled = Boolean(flowModule.stop_after_all_iters_if) $: result = $flowStateStore[flowModule.id]?.previewResult ?? NEVER_TESTED_THIS_FAR + $: parentLoopId = checkIfParentLoop($flowStore)
-
- - - If defined, at the end of the step, the predicate expression will be evaluated to decide if - the flow should stop early. - - - - { - if (isStopAfterIfEnabled && flowModule.stop_after_if) { - flowModule.stop_after_if = undefined - } else { - flowModule.stop_after_if = { - expr: 'result == undefined', - skip_if_stopped: false - } - } - }} - options={{ - right: 'Early stop or Break if condition met' - }} - /> - -
- {#if flowModule.stop_after_if} - - Stop condition expression -
- { - editor?.insertAtCursor(detail) - editor?.focus() - }} - > - + + If defined, at the end of the step, the predicate expression will be evaluated to decide + if the flow should stop early or break if inside a for/while loop. + + + + { + if (isStopAfterIfEnabled && flowModule.stop_after_if) { + flowModule.stop_after_if = undefined + } else { + flowModule.stop_after_if = { + expr: 'result == undefined', + skip_if_stopped: false + } + } + }} + options={{ + right: isLoop + ? 'Break loop' + : parentLoopId + ? 'Break parent loop module' + : 'Stop flow' + ' if condition met' + }} + /> + +
+ {#if flowModule.stop_after_if} + {@const earlyStopResult = isLoop + ? Array.isArray(result) && result.length > 0 + ? result[result.length - 1] + : result === NEVER_TESTED_THIS_FAR + ? result + : undefined + : result} + {#if !parentLoopId && !isLoop} + - -
- {:else} - Stop condition expression -