From 791c0635d50fc0be283f311e1663565146c018b2 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 10 Nov 2022 08:48:06 +0100 Subject: [PATCH] tweak identity to extract previous_result --- backend/tests/worker.rs | 36 +++++++++++++++++++++++++++ backend/windmill-worker/src/worker.rs | 9 ++++++- 2 files changed, 44 insertions(+), 1 deletion(-) diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 981dd99d88..d236798270 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1070,6 +1070,42 @@ async fn test_deno_flow(db: Pool) { } } +#[sqlx::test(fixtures("base"))] +async fn test_identity(db: Pool) { + initialize_tracing().await; + + let server = ApiServer::start(db.clone()).await; + + let flow: FlowValue = serde_json::from_value(serde_json::json!({ + "modules": [{ + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(): return 42", + }}, { + "value": { + "type": "identity", + }, + }, { + "value": { + "type": "identity", + }, + }, { + "value": { + "type": "identity", + }, + }], + })) + .unwrap(); + + let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) + .run_until_complete(&db, server.addr.port()) + .await + .result + .unwrap(); + assert_eq!(result, serde_json::json!(42)); +} + #[sqlx::test(fixtures("base"))] async fn test_deno_flow_same_worker(db: Pool) { initialize_tracing().await; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 468dd24d4f..d57059cc44 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -621,7 +621,14 @@ async fn handle_queued_job( JobKind::Dependencies => { handle_dependency_job(&job, &mut logs, job_dir, db, timeout, &envs).await } - JobKind::Identity => Ok(job.args.clone().unwrap_or_else(|| Value::Null)), + JobKind::Identity => match job.args.clone() { + Some(Value::Object(args)) + if args.len() == 1 && args.contains_key("previous_result") => + { + Ok(args.get("previous_result").unwrap().clone()) + } + args @ _ => Ok(args.unwrap_or_else(|| Value::Null)), + }, _ => { handle_code_execution_job( &job,