feat: handle $var: and $res: in arrays for transform_json_value (#7949)

* feat: handle $var: and $res: in arrays for transform_json_value

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: resolve $var/$res arrays inside objects in transform_json_value

Replace top_level: bool with depth: u8 to fix the main bug where arrays
nested inside objects were never resolved. The Object arm recursed with
top_level: false, which caused the Array arm guard to reject the value.

With depth <= 2, arrays inside objects (depth 1-2) are now processed.
Also: transform all array elements unconditionally (handles $encrypted:
in worker, reserved $vars, nested structures), and warn on >1000 items.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-02-15 09:35:06 +01:00
committed by GitHub
parent 43218c6285
commit e4a34d031b
2 changed files with 267 additions and 6 deletions
+119 -3
View File
@@ -517,6 +517,7 @@ pub async fn get_resource_value_interpolated_internal<'a>(
value,
&job_id,
token_for_context,
0,
)
.await?;
if allow_cache {
@@ -535,6 +536,7 @@ pub async fn transform_json_value(
v: Value,
job_id: &Option<Uuid>,
token: Option<&str>,
depth: u8,
) -> Result<Value> {
match v {
Value::String(y) if y.starts_with("$var:") => {
@@ -563,7 +565,8 @@ pub async fn transform_json_value(
tx.commit().await?;
let v = not_found_if_none(v, "Resource", path)?;
if let Some(v) = v {
transform_json_value(db_with_opt_authed, workspace, v, job_id, token).await
transform_json_value(db_with_opt_authed, workspace, v, job_id, token, depth + 1)
.await
} else {
Ok(Value::Null)
}
@@ -636,10 +639,41 @@ pub async fn transform_json_value(
.unwrap_or_else(|| y);
Ok(serde_json::json!(value))
}
Value::Array(mut arr) if depth <= 2 && arr.len() <= 1000 => {
for i in 0..arr.len() {
let val = std::mem::take(&mut arr[i]);
arr[i] = transform_json_value(
db_with_opt_authed,
workspace,
val,
job_id,
token,
depth + 1,
)
.await?;
}
Ok(Value::Array(arr))
}
Value::Array(arr) => {
if arr.len() > 1000 {
tracing::warn!(
"Array with {} items exceeds 1000 item limit for variable/resource resolution, skipping",
arr.len()
);
}
Ok(Value::Array(arr))
}
Value::Object(mut m) => {
for (a, b) in m.clone().into_iter() {
let v =
transform_json_value(db_with_opt_authed, workspace, b, job_id, token).await?;
let v = transform_json_value(
db_with_opt_authed,
workspace,
b,
job_id,
token,
depth + 1,
)
.await?;
m.insert(a.clone(), v);
}
Ok(Value::Object(m))
@@ -1822,6 +1856,7 @@ pub async fn interpolate(
value,
&None,
None,
0,
)
.await?
{
@@ -1829,3 +1864,84 @@ pub async fn interpolate(
v => Err(anyhow::anyhow!("Expected string, got {:?}", v)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use windmill_common::audit::AuditAuthor;
use windmill_common::db::DbWithOptAuthed;
fn test_db_with_opt_authed(db: DB) -> DbWithOptAuthed<'static, ApiAuthed> {
DbWithOptAuthed::DB {
db,
audit_author: AuditAuthor {
username: "test".to_string(),
email: "test@test.com".to_string(),
username_override: None,
token_prefix: None,
},
}
}
#[tokio::test]
async fn test_transform_array_over_1000_passthrough() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let dba = test_db_with_opt_authed(pool);
let arr: Vec<Value> = (0..1001).map(|i| json!(format!("$var:x/{i}"))).collect();
let input = Value::Array(arr.clone());
let result = transform_json_value(&dba, "test", input, &None, None, 0)
.await
.unwrap();
assert_eq!(result, Value::Array(arr));
}
#[tokio::test]
async fn test_transform_array_non_matching_strings_passthrough() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let dba = test_db_with_opt_authed(pool);
let input = json!(["hello", "world", 42, true, null, {"key": "val"}]);
let result = transform_json_value(&dba, "test", input.clone(), &None, None, 0)
.await
.unwrap();
assert_eq!(result, input);
}
#[tokio::test]
async fn test_transform_array_resolved_inside_object() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let dba = test_db_with_opt_authed(pool);
let input = json!({"urls": ["$var:u/test/nonexistent", "plain"]});
let result = transform_json_value(&dba, "test", input, &None, None, 0).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_transform_array_attempts_matching_items() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let dba = test_db_with_opt_authed(pool);
let input = json!(["$var:u/test/nonexistent", "plain"]);
let result = transform_json_value(&dba, "test", input, &None, None, 0).await;
assert!(result.is_err());
}
}
+148 -3
View File
@@ -170,7 +170,8 @@ pub async fn transform_json<'a>(
let value = serde_json::from_str(inner_vs).map_err(|e| {
error::Error::internal_err(format!("Error while parsing inner arg: {e:#}"))
})?;
let transformed = transform_json_value(&k, &client, workspace, value, job, db).await?;
let transformed =
transform_json_value(&k, &client, workspace, value, job, db, 0).await?;
let as_raw = serde_json::from_value(transformed).map_err(|e| {
error::Error::internal_err(format!("Error while parsing inner arg: {e:#}"))
})?;
@@ -196,7 +197,8 @@ pub async fn transform_json_as_values<'a>(
let value = serde_json::from_str(inner_vs).map_err(|e| {
error::Error::internal_err(format!("Error while parsing inner arg: {e:#}"))
})?;
let transformed = transform_json_value(&k, &client, workspace, value, job, db).await?;
let transformed =
transform_json_value(&k, &client, workspace, value, job, db, 0).await?;
let as_raw = serde_json::from_value(transformed).map_err(|e| {
error::Error::internal_err(format!("Error while parsing inner arg: {e:#}"))
})?;
@@ -237,6 +239,7 @@ pub async fn transform_json_value(
v: Value,
job: &MiniPulledJob,
conn: &Connection,
depth: u8,
) -> error::Result<Value> {
match v {
Value::String(y) if y.starts_with("$var:") => {
@@ -306,11 +309,28 @@ pub async fn transform_json_value(
.unwrap_or_else(|| y);
Ok(json!(value))
}
Value::Array(mut arr) if depth <= 2 && arr.len() <= 1000 => {
for i in 0..arr.len() {
let val = std::mem::take(&mut arr[i]);
arr[i] = transform_json_value(name, client, workspace, val, job, conn, depth + 1)
.await?;
}
Ok(Value::Array(arr))
}
Value::Array(arr) => {
if arr.len() > 1000 {
tracing::warn!(
"Array with {} items exceeds 1000 item limit for variable/resource resolution, skipping",
arr.len()
);
}
Ok(Value::Array(arr))
}
Value::Object(mut m) => {
for (a, b) in m.clone().into_iter() {
m.insert(
a.clone(),
transform_json_value(&a, client, workspace, b, job, conn).await?,
transform_json_value(&a, client, workspace, b, job, conn, depth + 1).await?,
);
}
Ok(Value::Object(m))
@@ -1408,3 +1428,128 @@ impl MaybeLock {
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use windmill_common::client::AuthedClient;
fn test_client() -> AuthedClient {
AuthedClient::new(
"http://localhost:0".to_string(),
"test".to_string(),
"test-token".to_string(),
None,
)
}
fn test_job() -> MiniPulledJob {
MiniPulledJob {
workspace_id: "test".to_string(),
id: uuid::Uuid::nil(),
args: None,
parent_job: None,
created_by: "test".to_string(),
scheduled_for: chrono::Utc::now(),
started_at: None,
runnable_path: None,
kind: windmill_common::jobs::JobKind::Noop,
runnable_id: None,
canceled_reason: None,
canceled_by: None,
permissioned_as: "test".to_string(),
permissioned_as_email: "test@test.com".to_string(),
flow_status: None,
tag: "test".to_string(),
script_lang: None,
same_worker: false,
pre_run_error: None,
concurrent_limit: None,
concurrency_time_window_s: None,
flow_innermost_root_job: None,
root_job: None,
timeout: None,
flow_step_id: None,
cache_ttl: None,
cache_ignore_s3_path: None,
priority: None,
preprocessed: None,
script_entrypoint_override: None,
trigger: None,
trigger_kind: None,
visible_to_owner: false,
permissioned_as_end_user_email: None,
runnable_settings_handle: None,
}
}
#[tokio::test]
async fn test_transform_array_over_1000_passthrough() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let conn = Connection::Sql(pool);
let client = test_client();
let job = test_job();
let arr: Vec<Value> = (0..1001).map(|i| json!(format!("$var:x/{i}"))).collect();
let input = Value::Array(arr.clone());
let result = transform_json_value("test", &client, "test", input, &job, &conn, 0)
.await
.unwrap();
assert_eq!(result, Value::Array(arr));
}
#[tokio::test]
async fn test_transform_array_non_matching_strings_passthrough() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let conn = Connection::Sql(pool);
let client = test_client();
let job = test_job();
let input = json!(["hello", "world", 42, true, null, {"key": "val"}]);
let result = transform_json_value("test", &client, "test", input.clone(), &job, &conn, 0)
.await
.unwrap();
assert_eq!(result, input);
}
#[tokio::test]
async fn test_transform_array_resolved_inside_object() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let conn = Connection::Sql(pool);
let client = test_client();
let job = test_job();
let input = json!({"urls": ["$var:u/test/nonexistent", "plain"]});
let result = transform_json_value("test", &client, "test", input, &job, &conn, 0).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_transform_array_attempts_matching_items() {
let db_url = std::env::var("DATABASE_URL")
.unwrap_or("postgres://postgres:changeme@localhost:5432/windmill".to_string());
let pool = sqlx::PgPool::connect(&db_url).await.unwrap();
let conn = Connection::Sql(pool);
let client = test_client();
let job = test_job();
let input = json!(["$var:u/test/nonexistent", "plain"]);
let result = transform_json_value("test", &client, "test", input, &job, &conn, 0).await;
assert!(result.is_err());
}
}