From ddfb2638908d6a31f66ba79b2f259dec0d64d468 Mon Sep 17 00:00:00 2001 From: Abhinav Rastogi Date: Mon, 5 Oct 2026 16:47:55 +0530 Subject: [PATCH] feat: add exclusive join task builder --- CHANGELOG.md | 2 + docs/WORKFLOW.md | 21 +++++++++++ .../workflow/task/exclusive_join_task.py | 29 +++++++++++++++ .../unit/workflow/test_exclusive_join_task.py | 37 +++++++++++++++++++ 4 files changed, 89 insertions(+) create mode 100644 src/conductor/client/workflow/task/exclusive_join_task.py create mode 100644 tests/unit/workflow/test_exclusive_join_task.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 6a98406ba..fbe6b0600 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- `ExclusiveJoinTask` builder supports exclusive branch joins and optional fallback task references. + - Worker isolation mode: `CONDUCTOR_WORKER_ISOLATION=thread` runs every worker as a thread instead of a `multiprocessing.Process` (default `process` is unchanged — `spawn` remains the start method). For environments where multiprocessing's fork+exec bootstraps (spawn children, `resource_tracker`) fail, e.g. Firecracker microVM guests. Thread-mode tradeoffs: no per-worker force-kill (shutdown is cooperative), CPU-bound workers share the GIL, and `signal.signal` becomes a no-op off the main thread. Implementation: the Windows-only Process→Thread shim moved to `conductor.client.automator.worker_isolation` (the private `worker_manager._patch_conductor_use_threads_on_windows` helper is removed — the Windows gate calls `apply_thread_isolation()` directly) and now also swaps the logging-relay `Queue` for a plain `queue.Queue` - Canonical metrics mode: opt-in harmonized metric surface via `WORKER_CANONICAL_METRICS=true` -- [details](METRICS.md#detailed-technical-notes--unreleased) diff --git a/docs/WORKFLOW.md b/docs/WORKFLOW.md index 2802ffc13..1cd654f54 100644 --- a/docs/WORKFLOW.md +++ b/docs/WORKFLOW.md @@ -17,6 +17,27 @@ configuration = Configuration( workflow_client = OrkesWorkflowClient(configuration) ``` +### Exclusive Join Task + +Use `ExclusiveJoinTask` to merge mutually exclusive branches. The server selects +an executed, non-skipped task from `join_on` and copies its output. If none of +those tasks ran, it checks the optional `default_exclusive_join_task` references. + +```python +from conductor.client.workflow.task.exclusive_join_task import ExclusiveJoinTask + +merge = ExclusiveJoinTask( + task_ref_name="merge_ref", + join_on=["approved_ref", "rejected_ref"], + default_exclusive_join_task=["fallback_ref"], +) +workflow >> merge +``` + +The references must identify tasks in your workflow. This serializes as +`EXCLUSIVE_JOIN`, with `joinOn` and optional `defaultExclusiveJoinTask` fields. +Use `JoinTask` when all listed branches must finish. + ### Start Workflow Execution #### Start using StartWorkflowRequest diff --git a/src/conductor/client/workflow/task/exclusive_join_task.py b/src/conductor/client/workflow/task/exclusive_join_task.py new file mode 100644 index 000000000..967bfcd86 --- /dev/null +++ b/src/conductor/client/workflow/task/exclusive_join_task.py @@ -0,0 +1,29 @@ +from copy import deepcopy +from typing import List, Optional + +from conductor.client.http.models.workflow_task import WorkflowTask +from conductor.client.workflow.task.task import TaskInterface +from conductor.client.workflow.task.task_type import TaskType + + +class ExclusiveJoinTask(TaskInterface): + """Join the first executed branch, with optional fallback task references.""" + + def __init__( + self, + task_ref_name: str, + join_on: List[str], + default_exclusive_join_task: Optional[List[str]] = None, + ) -> None: + super().__init__( + task_reference_name=task_ref_name, + task_type=TaskType.EXCLUSIVE_JOIN, + ) + self._join_on = deepcopy(join_on) + self._default_exclusive_join_task = deepcopy(default_exclusive_join_task) + + def to_workflow_task(self) -> WorkflowTask: + workflow = super().to_workflow_task() + workflow.join_on = self._join_on + workflow.default_exclusive_join_task = self._default_exclusive_join_task + return workflow diff --git a/tests/unit/workflow/test_exclusive_join_task.py b/tests/unit/workflow/test_exclusive_join_task.py new file mode 100644 index 000000000..6e587a50b --- /dev/null +++ b/tests/unit/workflow/test_exclusive_join_task.py @@ -0,0 +1,37 @@ +import unittest + +from conductor.client.http.api_client import ApiClient +from conductor.client.workflow.task.exclusive_join_task import ExclusiveJoinTask +from conductor.client.workflow.task.join_task import JoinTask + + +class TestExclusiveJoinTask(unittest.TestCase): + def test_serializes_join_candidates_and_fallback(self): + task = ExclusiveJoinTask("merge", ["branch_a", "branch_b"], ["fallback"]) + wire = ApiClient().sanitize_for_serialization(task.to_workflow_task()) + self.assertEqual(wire["type"], "EXCLUSIVE_JOIN") + self.assertEqual(wire["taskReferenceName"], "merge") + self.assertEqual(wire["joinOn"], ["branch_a", "branch_b"]) + self.assertEqual(wire["defaultExclusiveJoinTask"], ["fallback"]) + self.assertEqual(JoinTask("normal", ["branch_a"]).to_workflow_task().type, "JOIN") + + def test_fallback_is_optional(self): + for fallback in (None, []): + with self.subTest(fallback=fallback): + task = ExclusiveJoinTask("merge", ["branch_a"], fallback) + wire = ApiClient().sanitize_for_serialization(task.to_workflow_task()) + self.assertEqual(wire["joinOn"], ["branch_a"]) + if fallback is None: + self.assertNotIn("defaultExclusiveJoinTask", wire) + else: + self.assertEqual(wire["defaultExclusiveJoinTask"], []) + + def test_copies_caller_owned_lists(self): + candidates = ["branch_a"] + fallback = ["fallback"] + task = ExclusiveJoinTask("merge", candidates, fallback) + candidates.append("branch_b") + fallback.clear() + workflow_task = task.to_workflow_task() + self.assertEqual(workflow_task.join_on, ["branch_a"]) + self.assertEqual(workflow_task.default_exclusive_join_task, ["fallback"])