From 96cbba43ad70a40dcd741b76ec72d390e97cb448 Mon Sep 17 00:00:00 2001 From: Lucas Abel <22837557+uael@users.noreply.github.com> Date: Mon, 10 Feb 2025 11:41:54 +0100 Subject: [PATCH] backend: add raw flow with "restarted from" test (#5251) --- backend/tests/worker.rs | 120 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 120 insertions(+) diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 1574cf64ca..7573e139ee 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -4498,4 +4498,124 @@ mod job_payload { }; test_for_versions(VERSION_FLAGS.iter().cloned(), test).await; } + + #[sqlx::test(fixtures("base", "hello"))] + async fn test_raw_flow_payload_with_restarted_from(db: Pool) { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await; + let port = server.addr.port(); + + let db = &db; + let test = |restarted_from, arg, result| async move { + let job = RunJob::from(JobPayload::RawFlow { + value: serde_json::from_value(json!({ + "modules": [{ + "id": "a", + "value": { + "type": "rawscript", + "content": r#"export function main(world: string) { + return `Hello ${world}!`; + }"#, + "language": "deno", + "input_transforms": { + "world": { "type": "javascript", "expr": "flow_input.world" } + } + } + }, { + "id": "b", + "value": { + "type": "rawscript", + "content": r#"export function main(world: string, a: string) { + return `${a} ${world}!`; + }"#, + "language": "deno", + "input_transforms": { + "world": { "type": "javascript", "expr": "flow_input.world" }, + "a": { "type": "javascript", "expr": "results.a" } + } + } + }, { + "id": "c", + "value": { + "type": "forloopflow", + "iterator": { "type": "javascript", "expr": "['a', 'b', 'c']" }, + "modules": [{ + "value": { + "input_transforms": { + "world": { "type": "javascript", "expr": "flow_input.world" }, + "b": { "type": "javascript", "expr": "results.b" }, + "x": { "type": "javascript", "expr": "flow_input.iter.value" } + }, + "type": "rawscript", + "language": "deno", + "content": r#"export function main(world: string, b: string, x: string) { + return `${x}: ${b} ${world}!`; + }"#, + }, + }], + } + }], + "schema": { + "$schema": "https://json-schema.org/draft/2020-12/schema", + "properties": { "world": { "type": "string" } }, + "type": "object", + "order": [ "world" ] + } + })) + .unwrap(), + path: None, + restarted_from, + }) + .arg("world", arg) + .run_until_complete(db, port) + .await; + + assert_eq!(job.json_result().unwrap(), result); + job.id + }; + let flow_job_id = test( + None, + json!("foo"), + json!([ + "a: Hello foo! foo! foo!", + "b: Hello foo! foo! foo!", + "c: Hello foo! foo! foo!" + ]), + ) + .await; + let flow_job_id = test( + Some(RestartedFrom { flow_job_id, step_id: "a".into(), branch_or_iteration_n: None }), + json!("foo"), + json!([ + "a: Hello foo! foo! foo!", + "b: Hello foo! foo! foo!", + "c: Hello foo! foo! foo!" + ]), + ) + .await; + let flow_job_id = test( + Some(RestartedFrom { flow_job_id, step_id: "b".into(), branch_or_iteration_n: None }), + json!("bar"), + json!([ + "a: Hello foo! bar! bar!", + "b: Hello foo! bar! bar!", + "c: Hello foo! bar! bar!" + ]), + ) + .await; + let _ = test( + Some(RestartedFrom { + flow_job_id, + step_id: "c".into(), + branch_or_iteration_n: Some(1), + }), + json!("yolo"), + json!([ + "a: Hello foo! bar! bar!", + "b: Hello foo! bar! yolo!", + "c: Hello foo! bar! yolo!" + ]), + ) + .await; + } }