From 79141a7bec086ef12d797b598c2d4946306cb94f Mon Sep 17 00:00:00 2001 From: nthmost-orkes Date: Wed, 22 Jul 2026 20:22:40 -0700 Subject: [PATCH] fix: use dynamic_fork_tasks_param instead of dynamic_fork_join_tasks_param in DynamicForkTask Fixes the regression introduced in baadd201 that undid PR #352. Setting dynamic_fork_join_tasks_param alongside dynamic_fork_tasks_input_param_name violates Conductor's validation (only one approach is valid at a time). Adds unit tests to prevent this from silently regressing again. Fixes conductor-oss/python-sdk#377 --- .../client/workflow/task/dynamic_fork_task.py | 2 +- tests/unit/workflow/test_dynamic_fork_task.py | 32 +++++++++++++++++++ 2 files changed, 33 insertions(+), 1 deletion(-) create mode 100644 tests/unit/workflow/test_dynamic_fork_task.py diff --git a/src/conductor/client/workflow/task/dynamic_fork_task.py b/src/conductor/client/workflow/task/dynamic_fork_task.py index 259a27bdd..9866cd502 100644 --- a/src/conductor/client/workflow/task/dynamic_fork_task.py +++ b/src/conductor/client/workflow/task/dynamic_fork_task.py @@ -20,7 +20,7 @@ def __init__(self, task_ref_name: str, tasks_param: str = "dynamicTasks", tasks_ def to_workflow_task(self) -> WorkflowTask: wf_task = super().to_workflow_task() - wf_task.dynamic_fork_join_tasks_param = self.tasks_param + wf_task.dynamic_fork_tasks_param = self.tasks_param wf_task.dynamic_fork_tasks_input_param_name = self.tasks_input_param_name tasks = [ wf_task, diff --git a/tests/unit/workflow/test_dynamic_fork_task.py b/tests/unit/workflow/test_dynamic_fork_task.py new file mode 100644 index 000000000..a61e70a0b --- /dev/null +++ b/tests/unit/workflow/test_dynamic_fork_task.py @@ -0,0 +1,32 @@ +import unittest + +from conductor.client.workflow.task.dynamic_fork_task import DynamicForkTask +from conductor.client.workflow.task.join_task import JoinTask + + +class TestDynamicForkTask(unittest.TestCase): + def test_to_workflow_task_uses_dynamic_fork_tasks_param(self): + task = DynamicForkTask( + task_ref_name="fork", + tasks_param="myTasks", + tasks_input_param_name="myTasksInput", + ) + tasks = task.to_workflow_task() + wf_task = tasks[0] + self.assertEqual(wf_task.dynamic_fork_tasks_param, "myTasks") + self.assertEqual(wf_task.dynamic_fork_tasks_input_param_name, "myTasksInput") + self.assertIsNone(wf_task.dynamic_fork_join_tasks_param) + + def test_to_workflow_task_with_join(self): + join = JoinTask("join", join_on=[]) + task = DynamicForkTask(task_ref_name="fork", join_task=join) + tasks = task.to_workflow_task() + self.assertEqual(len(tasks), 2) + self.assertEqual(tasks[0].dynamic_fork_tasks_param, "dynamicTasks") + self.assertEqual(tasks[0].dynamic_fork_tasks_input_param_name, "dynamicTasksInputs") + self.assertIsNone(tasks[0].dynamic_fork_join_tasks_param) + self.assertEqual(tasks[1].task_reference_name, "join") + + +if __name__ == "__main__": + unittest.main()