backend: add raw flow with "restarted from" test (#5251)

This commit is contained in:
Lucas Abel
2025-02-10 11:41:54 +01:00
committed by GitHub
parent 4a744ff313
commit 96cbba43ad
+120
View File
@@ -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<Postgres>) {
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;
}
}