diff --git a/CHANGELOG.md b/CHANGELOG.md index 6a98406b..071cbc50 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 +- `NoopTask` builder and `TaskType.NOOP` for explicit no-op workflow steps, including SWITCH default branches. + - 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 2802ffc1..8e60124c 100644 --- a/docs/WORKFLOW.md +++ b/docs/WORKFLOW.md @@ -1,5 +1,23 @@ # Workflow Management +## No-op workflow steps + +`NoopTask` creates a built-in `NOOP` task that completes immediately on the +server. It needs no worker or input parameters and can serve as an explicit +placeholder when a workflow branch has nothing to do: + +```python +from conductor.client.workflow.task.noop_task import NoopTask +from conductor.client.workflow.task.switch_task import SwitchTask + +route = SwitchTask("route", "${workflow.input.action}").default_case( + [NoopTask(task_ref_name="skip_action")] +) +``` + +Add `route` to a workflow like any other task. The Conductor server must support +the `NOOP` system task type. + ## Workflow Client ### Initialization diff --git a/src/conductor/client/workflow/task/noop_task.py b/src/conductor/client/workflow/task/noop_task.py new file mode 100644 index 00000000..1b3a0cf7 --- /dev/null +++ b/src/conductor/client/workflow/task/noop_task.py @@ -0,0 +1,12 @@ +from conductor.client.workflow.task.task import TaskInterface +from conductor.client.workflow.task.task_type import TaskType + + +class NoopTask(TaskInterface): + """A system task that completes immediately without running a worker.""" + + def __init__(self, task_ref_name: str) -> None: + super().__init__( + task_reference_name=task_ref_name, + task_type=TaskType.NOOP, + ) diff --git a/src/conductor/client/workflow/task/task_type.py b/src/conductor/client/workflow/task/task_type.py index 36108f72..a056e68b 100644 --- a/src/conductor/client/workflow/task/task_type.py +++ b/src/conductor/client/workflow/task/task_type.py @@ -3,6 +3,7 @@ class TaskType(str, Enum): SIMPLE = "SIMPLE" + NOOP = "NOOP" DYNAMIC = "DYNAMIC" FORK_JOIN = "FORK_JOIN" FORK_JOIN_DYNAMIC = "FORK_JOIN_DYNAMIC" diff --git a/tests/unit/workflow/test_noop_task.py b/tests/unit/workflow/test_noop_task.py new file mode 100644 index 00000000..6ddd348e --- /dev/null +++ b/tests/unit/workflow/test_noop_task.py @@ -0,0 +1,33 @@ +import unittest + +from conductor.client.configuration.configuration import Configuration +from conductor.client.http.api_client import ApiClient +from conductor.client.workflow.task.noop_task import NoopTask +from conductor.client.workflow.task.switch_task import SwitchTask +from conductor.client.workflow.task.task_type import TaskType + + +class TestNoopTask(unittest.TestCase): + def setUp(self): + self.client = ApiClient(configuration=Configuration()) + + def test_serializes_as_builtin_task_without_inputs(self): + task = NoopTask(task_ref_name="do_nothing") + self.assertEqual(task.task_type, TaskType.NOOP) + wire = self.client.sanitize_for_serialization(task.to_workflow_task()) + self.assertEqual(wire["name"], "do_nothing") + self.assertEqual(wire["taskReferenceName"], "do_nothing") + self.assertEqual(wire["type"], "NOOP") + self.assertEqual(wire["inputParameters"], {}) + + def test_serializes_in_switch_default_branch(self): + switch = SwitchTask("route", "${workflow.input.action}").default_case( + [NoopTask("skip_action")] + ) + wire = self.client.sanitize_for_serialization(switch.to_workflow_task()) + self.assertEqual(wire["type"], "SWITCH") + self.assertEqual(len(wire["defaultCase"]), 1) + fallback = wire["defaultCase"][0] + self.assertEqual(fallback["type"], "NOOP") + self.assertEqual(fallback["taskReferenceName"], "skip_action") + self.assertEqual(fallback["inputParameters"], {})