diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 2b9e540226..2c7140a0d1 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -bc904859dd66c55ebad002e8526103c73de841cd +cf96b45aa1183f15b3cc1b971035de5e37a68849 diff --git a/backend/windmill-worker/src/ai_executor.rs b/backend/windmill-worker/src/ai_executor.rs index 76160cfa6b..9e4145274e 100644 --- a/backend/windmill-worker/src/ai_executor.rs +++ b/backend/windmill-worker/src/ai_executor.rs @@ -508,7 +508,7 @@ pub async fn run_agent( .clone() .unwrap_or_else(|| "unknown".to_string()); - Some(get_transform_context(job, &previous_id, flow_status).await?) + Some(get_transform_context(job, &previous_id, flow_status)) } else { None } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index c15c5068b4..63b7c493ef 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -283,6 +283,24 @@ fn get_stop_after_if_data(stop_after_if: Option<&StopAfterIf>) -> (bool, Option< return (false, None); } +async fn get_id_ctx_for_expr( + expr: &str, + flow: uuid::Uuid, + db: &DB, + status: &FlowStatus, +) -> error::Result> { + if expr.contains("results.") || expr.contains("results[") || expr.contains("results?.") { + let flow_job = get_mini_pulled_job(db, &flow).await?; + if let Some(flow_job) = flow_job { + Ok(Some(get_transform_context(&flow_job, "", &status))) + } else { + Ok(None) + } + } else { + Ok(None) + } +} + async fn evaluate_stop_after_all_iters_if( db: &DB, stop_after_all_iters_if: &StopAfterIf, @@ -294,6 +312,8 @@ async fn evaluate_stop_after_all_iters_if( stop_early_err_msg: &mut Option, nresult: &mut Option>>, args: HashMap>, + flow: uuid::Uuid, + status: &FlowStatus, ) -> error::Result<()> { let iters_result = match &module_status { FlowStatusModule::InProgress { flow_jobs: Some(flow_jobs), .. } => { @@ -308,13 +328,15 @@ async fn evaluate_stop_after_all_iters_if( *nresult = Some(iters_result.clone()); // as an optimization, we store the result of all jobs as when stop_early_after_all_iters evaluates to false, it would have to be computed (finished loop/branchall) + let id_ctx = get_id_ctx_for_expr(&stop_after_all_iters_if.expr, flow, db, status).await?; + let stop_early_after_all_iters = compute_bool_from_expr( &stop_after_all_iters_if.expr, Marc::new(args), None, iters_result.clone(), None, - None, + id_ctx.as_ref(), Some(client), None, None, @@ -542,13 +564,15 @@ pub async fn update_flow_status_after_job_completion_internal( }; let args = from_result_to_args(args.as_ref().await.get_ref())?; + let id_ctx = get_id_ctx_for_expr(expr, flow, db, &old_status).await?; + compute_bool_from_expr( &expr, Marc::new(args), None, result.clone(), all_iters, - None, + id_ctx.as_ref(), Some(client), None, None, @@ -816,6 +840,8 @@ pub async fn update_flow_status_after_job_completion_internal( &mut stop_early_err_msg, &mut nresult, args, + flow, + &old_status, ) .await?; } @@ -1023,6 +1049,8 @@ pub async fn update_flow_status_after_job_completion_internal( &mut stop_early_err_msg, &mut nresult, args, + flow, + &old_status, ) .await?; } @@ -2917,9 +2945,7 @@ async fn push_next_flow_job( drop(resume_messages); let is_skipped = if let Some(skip_if) = &module.skip_if { - let idcontext = get_transform_context(&flow_job, previous_id.as_str(), &status) - .warn_after_seconds(3) - .await?; + let idcontext = get_transform_context(&flow_job, previous_id.as_str(), &status); compute_bool_from_expr( &skip_if.expr, arc_flow_job_args.clone(), @@ -2999,9 +3025,7 @@ async fn push_next_flow_job( | FlowModuleValue::Flow { input_transforms, .. } | FlowModuleValue::AIAgent { input_transforms, .. }, ) => { - let ctx = get_transform_context(&flow_job, &previous_id, &status) - .warn_after_seconds(3) - .await?; + let ctx = get_transform_context(&flow_job, &previous_id, &status); transform_context = Some(ctx); let by_id = transform_context.as_ref().unwrap(); // if a failure step, we add flow job id and started_at to the context. This is for error handling of triggers where we wrap scripts into single step flows @@ -3203,9 +3227,7 @@ async fn push_next_flow_job( args.insert("iter".to_string(), to_raw_value(new_args)); if let Some(input_transforms) = simple_input_transforms { //previous id is none because we do not want to use previous id if we are in a for loop - let ctx = get_transform_context(&flow_job, "", &status) - .warn_after_seconds(3) - .await?; + let ctx = get_transform_context(&flow_job, "", &status); let ti = transform_input( Marc::new(args), flow.flow_env.as_ref(), @@ -3260,9 +3282,7 @@ async fn push_next_flow_job( to_raw_value(&json!({ "index": i as i32, "value": itered[i]})), ); if let Some(input_transforms) = simple_input_transforms { - let ctx = get_transform_context(&flow_job, &previous_id, &status) - .warn_after_seconds(3) - .await?; + let ctx = get_transform_context(&flow_job, &previous_id, &status); let ti = transform_input( Marc::new(hm), flow.flow_env.as_ref(), @@ -3377,9 +3397,7 @@ async fn push_next_flow_job( } let evaluated_timeout = if let Some(timeout_transform) = &module.timeout { - let ctx = get_transform_context(&flow_job, &previous_id, &status) - .warn_after_seconds(3) - .await?; + let ctx = get_transform_context(&flow_job, &previous_id, &status); let timeout_value = evaluate_input_transform::( timeout_transform, @@ -3458,9 +3476,7 @@ async fn push_next_flow_job( if let Some(parallelism_transform) = &value_with_parallel.parallelism { tracing::debug!(id = %flow_job.id, root_id = %job_root, "evaluating parallelism expression for forloopflow job {uuid}"); - let ctx = get_transform_context(&flow_job, &previous_id, &status) - .warn_after_seconds(3) - .await?; + let ctx = get_transform_context(&flow_job, &previous_id, &status); let evaluated_parallelism = evaluate_input_transform::( parallelism_transform, @@ -4297,7 +4313,7 @@ async fn compute_next_flow_transform( | FlowStatusModule::WaitingForEvents { .. } | FlowStatusModule::WaitingForExecutor { .. } => { let mut branch_chosen = BranchChosen::Default; - let idcontext = get_transform_context(&flow_job, previous_id, &status).await?; + let idcontext = get_transform_context(&flow_job, previous_id, &status); for (i, b) in branches.iter().enumerate() { let pred = compute_bool_from_expr( &b.expr, @@ -4580,7 +4596,7 @@ async fn next_forloop_status( let by_id = if let Some(x) = by_id { x } else { - get_transform_context(&flow_job, previous_id, &status).await? + get_transform_context(&flow_job, previous_id, &status) }; /* Iterator is an InputTransform, evaluate it into an array. */ let itered_raw = match iterator { @@ -4654,7 +4670,7 @@ async fn next_forloop_status( let by_id = if let Some(x) = by_id { x } else { - get_transform_context(&flow_job, previous_id, &status).await? + get_transform_context(&flow_job, previous_id, &status) }; let itered_raw = match iterator { InputTransform::Static { value } => to_raw_value(value), @@ -4930,18 +4946,18 @@ pub async fn script_to_payload( }) } -pub async fn get_transform_context( +pub fn get_transform_context( flow_job: &MiniPulledJob, previous_id: &str, status: &FlowStatus, -) -> error::Result { +) -> IdContext { let steps_results: HashMap = status .modules .iter() .filter_map(|x| x.job_result().map(|y| (x.id(), y))) .collect(); - Ok(IdContext { flow_job: flow_job.id, steps_results, previous_id: previous_id.to_string() }) + IdContext { flow_job: flow_job.id, steps_results, previous_id: previous_id.to_string() } } // trait IntoArray: Sized { diff --git a/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte b/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte index ea4bce88f2..9cb494bd3e 100644 --- a/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte +++ b/frontend/src/lib/components/flows/content/FlowModuleEarlyStop.svelte @@ -51,10 +51,10 @@ return null } let raise_error_message_stop_after_all_if = $state( - flowModule.stop_after_all_iters_if?.error_message !== undefined + flowModule.stop_after_all_iters_if?.error_message != undefined ) let raise_error_message_stop_after_if = $state( - flowModule.stop_after_if?.error_message !== undefined + flowModule.stop_after_if?.error_message != undefined ) let { isLoop, isParallelLoop } = $derived( flowModule.value.type === 'forloopflow' || flowModule.value.type === 'whileloopflow' @@ -166,8 +166,7 @@ { @@ -180,8 +179,8 @@ lang="javascript" bind:code={flowModule.stop_after_if.expr} class="h-full" - extraLib={`declare const result = ${JSON.stringify(earlyStopResult)};` + - `\n declare const flow_input = ${JSON.stringify(stepPropPicker.pickableProperties.flow_input)};` + + extraLib={`declare const result = ${JSON.stringify(earlyStopResult)};\n` + + stepPropPicker.extraLib + (isLoop ? `\ndeclare const all_iters = ${JSON.stringify(result)};` : '')} /> @@ -303,8 +302,7 @@ { editor?.insertAtCursor(detail) @@ -316,8 +314,8 @@ lang="javascript" bind:code={flowModule.stop_after_all_iters_if.expr} class="h-full" - extraLib={`declare const result = ${JSON.stringify(result)};` + - `\ndeclare const flow_input = ${JSON.stringify(stepPropPicker.pickableProperties.flow_input)};`} + extraLib={`declare const result = ${JSON.stringify(result)};\n` + + stepPropPicker.extraLib} /> diff --git a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte index 5d4761d8d5..a151b31509 100644 --- a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte +++ b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte @@ -144,7 +144,7 @@ : 'transparent'} animationDuration="4s" > - {#if result != undefined} + {#if result != undefined && !pickableProperties} {:else if pickableProperties} | undefined = undefined + export let result: any | undefined = undefined + export let extraResults: any = undefined let variables: Record = {} let resources: Record = {} @@ -69,7 +71,9 @@ : keepByKey(pickableProperties.priorIds, search) flowEnvFiltered = - search === EMPTY_STRING ? pickableProperties.flow_env : keepByKey(pickableProperties.flow_env, search) + search === EMPTY_STRING + ? pickableProperties.flow_env + : keepByKey(pickableProperties.flow_env, search) }, 50) } @@ -221,7 +225,7 @@ await updateCollapsable() } - $: search, $inputMatches, $propPickerConfig, pickableProperties, updateState() + $: (search, $inputMatches, $propPickerConfig, pickableProperties, updateState()) onDestroy(() => { clearTimeout(timeout) @@ -251,6 +255,16 @@ filter: {filteringFlowInputsOrResult} {/if} + {#if result != undefined} + Step Result +
+ +
+ {/if} {#if flowInputsFiltered && (Object.keys(flowInputsFiltered ?? {}).length > 0 || !filterActive)} Flow Input