mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 00:02:08 +00:00
fix: allow results access inside nested functions in input transforms (#11358)
* fix: allow results access inside nested functions in input transforms Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: decode escaped bracket step ids and test deferred fetch errors Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: let quickjs decode bracket step ids and match quoted forms Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: prefetch results read through spread syntax Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: keep prefetched bracket literals on a single line Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: decode prefetch step literals as data and skip unparsable ones Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: run the results prefetch outside the expression scope Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: keep the transform expression a zero-arg iife after prefetch Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
da866c5eff
commit
ca8a04a869
@@ -57,13 +57,16 @@ const END_BRACKET_PATTERN: &str = "\"]";
|
||||
// ── Regex statics ─────────────────────────────────────────────────────
|
||||
|
||||
lazy_static! {
|
||||
// `results` is fetched lazily via the `__getResult` async proxy, so we
|
||||
// wrap each `results.X` access with `(await ...)` to drive the proxy.
|
||||
// `flow_env` used to be wrapped here too (it was an async Deno op-backed
|
||||
// proxy in the deno_core era); QuickJS now exposes flow_env as a plain
|
||||
// in-memory object, so no await is needed.
|
||||
// Step ids statically referenced as `results.X` / `results["X"]`. They are
|
||||
// prefetched before the expression runs so `results` reads synchronously:
|
||||
// rewriting accesses to `await` breaks inside non-async inner functions.
|
||||
// Over-matching (e.g. inside a string literal) only costs a spurious fetch.
|
||||
// Bracket keys are captured as the whole JS string literal and decoded by
|
||||
// QuickJS (`__loadResults`), dropping any that fail to parse: a match may
|
||||
// come from a comment. The boundary must allow `.`, or `...results.a` is
|
||||
// never prefetched.
|
||||
static ref RE: Regex = Regex::new(
|
||||
r#"(?m)(?P<r>results(?:\?)?(?:(?:\.[a-zA-Z_0-9]+)|(?:\[\".*?\"\])))"#
|
||||
r#"(?:^|[^a-zA-Z0-9_$])results(?:\??\.([a-zA-Z_0-9]+)|(?:\?\.)?\[("(?:[^"\\\r\n]|\\.)*"|'(?:[^'\\\r\n]|\\.)*')\])"#
|
||||
)
|
||||
.unwrap();
|
||||
// SQL fast-path: simple `results.X.Y[i]...` accesses are dispatched to
|
||||
@@ -89,9 +92,20 @@ pub fn replace_with_await(expr: String, fn_name: &str) -> String {
|
||||
s
|
||||
}
|
||||
|
||||
/// JS string literals naming the step ids statically read from `results`.
|
||||
#[cfg(feature = "quickjs")]
|
||||
pub fn replace_with_await_result(expr: String) -> String {
|
||||
RE.replace_all(&expr, "(await $r)").to_string()
|
||||
fn referenced_step_ids(expr: &str) -> Vec<String> {
|
||||
let mut ids: Vec<String> = RE
|
||||
.captures_iter(expr)
|
||||
.filter_map(|c| match (c.get(1), c.get(2)) {
|
||||
(Some(ident), _) => Some(format!("\"{}\"", ident.as_str())),
|
||||
(_, Some(literal)) => Some(literal.as_str().to_string()),
|
||||
_ => None,
|
||||
})
|
||||
.collect();
|
||||
ids.sort();
|
||||
ids.dedup();
|
||||
ids
|
||||
}
|
||||
|
||||
#[cfg(feature = "quickjs")]
|
||||
@@ -442,11 +456,18 @@ async fn eval_quickjs_inner(
|
||||
let op_state_clone = op_state.clone();
|
||||
let by_id_clone = by_id.clone();
|
||||
|
||||
// Transform expression to add await for variable/resource/results access
|
||||
let expr_with_funcs = ["variable", "resource", "flow_user_state"]
|
||||
let transformed_expr = ["variable", "resource", "flow_user_state"]
|
||||
.into_iter()
|
||||
.fold(expr.to_string(), replace_with_await);
|
||||
let transformed_expr = replace_with_await_result(expr_with_funcs);
|
||||
let step_ids = by_id
|
||||
.as_ref()
|
||||
.map(|_| referenced_step_ids(expr))
|
||||
.unwrap_or_default();
|
||||
let prefetch = if step_ids.is_empty() {
|
||||
"Promise.resolve()".to_string()
|
||||
} else {
|
||||
format!("__loadResults({})", serde_json::to_string(&step_ids)?)
|
||||
};
|
||||
|
||||
async_with!(context => |ctx| {
|
||||
let globals = ctx.globals();
|
||||
@@ -552,10 +573,12 @@ async fn eval_quickjs_inner(
|
||||
.map_err(quickjs_error_to_anyhow)?;
|
||||
|
||||
// Determine if we need to add return statement.
|
||||
// The prefetch runs outside the expression's function so user declarations
|
||||
// can't shadow `__loadResults`; the expression stays a zero-arg IIFE.
|
||||
let code = if should_add_return_quickjs(&transformed_expr) {
|
||||
format!("(async function() {{ return {}; }})().then((x) => JSON.stringify(x ?? null))", transformed_expr)
|
||||
format!("{}.then(() => (async function() {{ return {}; }})()).then((x) => JSON.stringify(x ?? null))", prefetch, transformed_expr)
|
||||
} else {
|
||||
format!("(async function() {{ {} }})().then((x) => JSON.stringify(x ?? null))", transformed_expr)
|
||||
format!("{}.then(() => (async function() {{ {} }})()).then((x) => JSON.stringify(x ?? null))", prefetch, transformed_expr)
|
||||
};
|
||||
|
||||
// Evaluate the expression (returns a Promise that resolves to a JSON string)
|
||||
@@ -802,16 +825,37 @@ fn setup_results_proxy<'js>(
|
||||
.map_err(quickjs_error_to_anyhow)?;
|
||||
}
|
||||
|
||||
// Prefetched steps read synchronously; a fetch error is rethrown only when
|
||||
// the step is actually read. Other (dynamic) names still return a Promise.
|
||||
let proxy_setup = r#"
|
||||
function __resolveResult(name) {
|
||||
if (name === __previous_id && typeof previous_result !== 'undefined') {
|
||||
return Promise.resolve(previous_result);
|
||||
}
|
||||
return __getResult(name);
|
||||
}
|
||||
const __loaded = new Map();
|
||||
async function __loadResults(literals) {
|
||||
const ids = new Set();
|
||||
for (const lit of literals) {
|
||||
try { ids.add((0, eval)(lit)); } catch {}
|
||||
}
|
||||
await Promise.all([...ids].map((id) => __resolveResult(id).then(
|
||||
(v) => __loaded.set(id, { v }),
|
||||
(e) => __loaded.set(id, { e }),
|
||||
)));
|
||||
}
|
||||
const results = new Proxy({}, {
|
||||
get: function(target, name, receiver) {
|
||||
if (typeof name === 'symbol') {
|
||||
return undefined;
|
||||
}
|
||||
if (name === __previous_id && typeof previous_result !== 'undefined') {
|
||||
return Promise.resolve(previous_result);
|
||||
const r = __loaded.get(name);
|
||||
if (r) {
|
||||
if ('e' in r) throw r.e;
|
||||
return r.v;
|
||||
}
|
||||
return __getResult(name);
|
||||
return __resolveResult(name);
|
||||
}
|
||||
});
|
||||
"#;
|
||||
@@ -1116,19 +1160,94 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_replace_with_await_result() {
|
||||
fn test_referenced_step_ids() {
|
||||
assert_eq!(
|
||||
replace_with_await_result("results.step_a".to_string()),
|
||||
"(await results.step_a)"
|
||||
referenced_step_ids(r#"results.b + results?.["a"] + results['c'] + my_results.e"#),
|
||||
vec![r#""a""#, r#""b""#, "'c'"]
|
||||
);
|
||||
assert!(referenced_step_ids("// results[\"\n\"]").is_empty());
|
||||
assert!(referenced_step_ids("no_results_here").is_empty());
|
||||
assert_eq!(referenced_step_ids("{...results.a}"), vec![r#""a""#]);
|
||||
}
|
||||
|
||||
/// `results.a` resolves to `previous_result` (step `a` is the previous
|
||||
/// step), so no client is needed.
|
||||
async fn eval_with_results(expr: &str) -> anyhow::Result<serde_json::Value> {
|
||||
let mut ctx = HashMap::new();
|
||||
ctx.insert(
|
||||
"previous_result".to_string(),
|
||||
Arc::new(to_raw_value(&json!({"x": 1, "y": 2}))),
|
||||
);
|
||||
let mut fi = HashMap::new();
|
||||
fi.insert("ids".to_string(), to_raw_value(&json!(["y", "x"])));
|
||||
let by_id = IdContext {
|
||||
flow_job: Uuid::nil(),
|
||||
root_flow_job: Uuid::nil(),
|
||||
steps_results: HashMap::new(),
|
||||
previous_id: "a".to_string(),
|
||||
};
|
||||
let result = eval_timeout_quickjs(
|
||||
expr.to_string(),
|
||||
ctx,
|
||||
Some(mappable_rc::Marc::new(fi)),
|
||||
None,
|
||||
None,
|
||||
Some(&by_id),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
Ok(serde_json::from_str(result.get())?)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_results_in_nested_functions() {
|
||||
assert_eq!(
|
||||
eval_with_results("return (() => { const tm = results.a; return tm.x + tm.y; })();")
|
||||
.await
|
||||
.unwrap(),
|
||||
json!(3)
|
||||
);
|
||||
assert_eq!(
|
||||
replace_with_await_result("results.a + results.b".to_string()),
|
||||
"(await results.a) + (await results.b)"
|
||||
eval_with_results("flow_input.ids.map(id => results.a[id])")
|
||||
.await
|
||||
.unwrap(),
|
||||
json!([2, 1])
|
||||
);
|
||||
assert_eq!(
|
||||
replace_with_await_result("no_results_here".to_string()),
|
||||
"no_results_here"
|
||||
eval_with_results(r#""results.a is " + JSON.stringify({...results.a})"#)
|
||||
.await
|
||||
.unwrap(),
|
||||
json!(r#"results.a is {"x":1,"y":2}"#)
|
||||
);
|
||||
assert_eq!(
|
||||
eval_with_results(
|
||||
r#"[0].map(() => results["\x61"].x + results?.['a'].y + results.a.y)"#
|
||||
)
|
||||
.await
|
||||
.unwrap(),
|
||||
json!([5])
|
||||
);
|
||||
assert_eq!(
|
||||
eval_with_results(r#"results.a.x /* results["\xZZ"] */"#)
|
||||
.await
|
||||
.unwrap(),
|
||||
json!(1)
|
||||
);
|
||||
assert_eq!(
|
||||
eval_with_results(
|
||||
"function __loadResults() {}\nreturn [results.a.x, arguments.length];"
|
||||
)
|
||||
.await
|
||||
.unwrap(),
|
||||
json!([1, 0])
|
||||
);
|
||||
// A failed prefetch only throws when that step is actually read.
|
||||
assert_eq!(
|
||||
eval_with_results("true ? 1 : results.other").await.unwrap(),
|
||||
json!(1)
|
||||
);
|
||||
let err = eval_with_results("results.other").await.unwrap_err();
|
||||
assert!(err.to_string().contains("Result fetching not available"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
Reference in New Issue
Block a user