From 0d38752671253a931adb6a937ce562dcbb576fce Mon Sep 17 00:00:00 2001 From: hugocasa Date: Wed, 16 Sep 2026 16:09:01 +0200 Subject: [PATCH] refactor: resolve $flow_expr tags by path lookup instead of expression evaluation Co-Authored-By: Claude Opus 5 --- backend/tests/flow_engine_parity.rs | 2 +- backend/windmill-worker/src/worker_flow.rs | 120 ++++++++++-------- .../lib/components/AssignableTagsInner.svelte | 8 +- 3 files changed, 69 insertions(+), 61 deletions(-) diff --git a/backend/tests/flow_engine_parity.rs b/backend/tests/flow_engine_parity.rs index 5f17c98a06..b6a943a23b 100644 --- a/backend/tests/flow_engine_parity.rs +++ b/backend/tests/flow_engine_parity.rs @@ -3314,7 +3314,7 @@ export function main(i: number) { } // A `$flow_expr[...]` step tag is resolved from the flow's state before the step is pushed, and -// one that cannot be evaluated fails the step instead of queueing it on a tag no worker serves. +// one that cannot be resolved fails the step instead of queueing it on a tag no worker serves. #[cfg(feature = "deno_core")] #[sqlx::test(fixtures("base"))] async fn test_flow_expr_step_tag(db: Pool) -> anyhow::Result<()> { diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index e58f0357b3..d25ed15582 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -3114,53 +3114,77 @@ fn resolve_flow_step_tag( } } -/// Resolves each `$flow_expr[path]` of a step tag by evaluating `path` as a flow expression -/// (`results.a.foo`, `flow_input.region`, `flow_env.pool`, `resume.region`). A string renders -/// bare, `null` (also what a missing field evaluates to) renders empty, any other value as its -/// JSON text. +/// Resolves each `$flow_expr[root.path]` of a step tag by reading `path` from `results` (keyed by +/// step id), `flow_input` or `flow_env`. As for `$args[...]`, a path that reaches nothing renders +/// empty, a string renders bare and any other value as its JSON text. async fn interpolate_flow_expr_tag( tag: &str, - env: HashMap>>, - flow_args: Marc>>, + db: &DB, + flow_job: &MiniPulledJob, + flow_input: &HashMap>, flow_env: Option<&HashMap>>, - client: &AuthedClient, - by_id: &IdContext, ) -> error::Result { let mut rendered: HashMap<&str, String> = HashMap::new(); for cap in RE_FLOW_EXPR_TAG.captures_iter(tag) { - let expr = cap.get(1).unwrap().as_str(); - if rendered.contains_key(expr) { + let path = cap.get(1).unwrap().as_str(); + if rendered.contains_key(path) { continue; } - let value = eval_timeout( - expr.to_string(), - env.clone(), - Some(flow_args.clone()), - flow_env, - Some(client), - Some(by_id), - None, - ) - .await - .map_err(|e| { - Error::ExecutionErr(format!( - "Could not resolve the step tag `{tag}`: error during evaluation of `{expr}`:\n{e:#}" - )) - })?; - let value = serde_json::from_str::(value.get()).map_err(to_anyhow)?; - rendered.insert(expr, render_flow_expr_tag_value(value)); + let mut segments = path.split('.'); + let root = segments.next().unwrap_or_default(); + let segments = segments.collect::>(); + let value = match (root, segments.split_first()) { + ("flow_input", _) => read_flow_expr_path(Some(flow_input), &segments), + ("flow_env", _) => read_flow_expr_path(flow_env, &segments), + ("results", Some((step_id, rest))) => { + match windmill_queue::get_result_by_id( + db.clone(), + flow_job.workspace_id.clone(), + flow_job.id, + step_id.to_string(), + (!rest.is_empty()).then(|| rest.join(".")), + ) + .await + { + Ok(result) => serde_json::from_str(result.get()).unwrap_or_default(), + Err(Error::NotFound(_)) => Value::Null, + Err(e) => { + return Err(Error::ExecutionErr(format!( + "Could not resolve the step tag `{tag}`: {e}" + ))) + } + } + } + _ => { + return Err(Error::ExecutionErr(format!( + "Could not resolve the step tag `{tag}`: `{path}` must start with \ + `results.`, `flow_input` or `flow_env`" + ))) + } + }; + let value = match value { + Value::String(s) => s, + Value::Null => String::new(), + v => v.to_string(), + }; + rendered.insert(path, value); } Ok(RE_FLOW_EXPR_TAG .replace_all(tag, |cap: ®ex::Captures| rendered[&cap[1]].clone()) .into_owned()) } -fn render_flow_expr_tag_value(value: Value) -> String { - match value { - Value::String(s) => s, - Value::Null => String::new(), - v => v.to_string(), - } +fn read_flow_expr_path(map: Option<&HashMap>>, segments: &[&str]) -> Value { + let Some((key, rest)) = segments.split_first() else { + return Value::Null; + }; + map.and_then(|m| m.get(*key)) + .and_then(|raw| serde_json::from_str::(raw.get()).ok()) + .and_then(|v| { + v.pointer(&rest.iter().map(|s| format!("/{s}")).collect::()) + .cloned() + }) + .unwrap_or_default() } #[cfg(test)] @@ -4241,10 +4265,9 @@ async fn push_next_flow_job( None }; - // The scope the step's input transforms were evaluated in, which a `$flow_expr[...]` tag - // must read too. - let mut expr_flow_args = arc_flow_job_args.clone(); - let mut expr_previous_id = previous_id.as_str(); + // The `flow_input` the step's input transforms read, which a `$flow_expr[flow_input...]` + // tag must read too: the body of a simple for-loop also sees `iter` there. + let mut step_flow_input = arc_flow_job_args.clone(); let marc; let me; @@ -4271,8 +4294,7 @@ async fn push_next_flow_job( //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); let args = Marc::new(args); - expr_flow_args = args.clone(); - expr_previous_id = ""; + step_flow_input = args.clone(); let ti = transform_input( args, flow_env, @@ -4514,23 +4536,9 @@ async fn push_next_flow_job( let mut tag_err = None; let tag = match tag { Some(t) if err.is_none() && tag_reads_flow_expr(&t) => { - let ctx = get_transform_context(&flow_job, expr_previous_id, &status); - let env = HashMap::from([ - ("previous_result".to_string(), arc_last_job_result.clone()), - ("resume".to_string(), resume.clone()), - ("resumes".to_string(), resumes.clone()), - ("approvers".to_string(), approvers.clone()), - ]); - match interpolate_flow_expr_tag( - &t, - env, - expr_flow_args.clone(), - flow_env, - client, - &ctx, - ) - .warn_after_seconds(3) - .await + match interpolate_flow_expr_tag(&t, db, &flow_job, &step_flow_input, flow_env) + .warn_after_seconds(3) + .await { Ok(resolved) => Some(resolved), Err(e) => { diff --git a/frontend/src/lib/components/AssignableTagsInner.svelte b/frontend/src/lib/components/AssignableTagsInner.svelte index f014dfeed5..351a95c782 100644 --- a/frontend/src/lib/components/AssignableTagsInner.svelte +++ b/frontend/src/lib/components/AssignableTagsInner.svelte @@ -210,8 +210,8 @@ {/if} {#if dynamicTag?.kind == 'flow_expr'}
- Interpolated tag based on flow expression {dynamicTag.path}, resolved when the - flow step starts + Interpolated tag based on the flow value at {dynamicTag.path}, resolved when + the flow step starts
{:else if dynamicTag}
Interpolated tag based on args input of {dynamicTag.path}
@@ -259,8 +259,8 @@ based on args input, use
$args[a.b.c]
where
a.b.c
is the path to the value in the args object.
{#if variant !== 'drawer'}
{/if} - On flow steps, use a flow expression path to base the tag on earlier results, flow inputs or flow - env, e.g.
$flow_expr[results.a.region]
or + On flow steps, use a path into earlier step results, flow inputs or flow env, e.g. +
$flow_expr[results.a.region]
or
$flow_expr[flow_input.region]
. {/if}