Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions src/conductor/client/workflow/task/lambda_task.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
from __future__ import annotations
from typing import Dict, Optional
from typing_extensions import Self

from conductor.client.workflow.task.task import TaskInterface
from conductor.client.workflow.task.task_type import TaskType


class LambdaTask(TaskInterface):
def __init__(self, task_ref_name: str, script: str, bindings: Optional[Dict[str, str]] = None) -> Self:
super().__init__(
task_reference_name=task_ref_name,
task_type=TaskType.LAMBDA,
input_parameters={
"scriptExpression": script,
}
)
if bindings is not None:
self.input_parameters.update(bindings)
26 changes: 26 additions & 0 deletions tests/unit/workflow/test_lambda_task.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
from conductor.client.workflow.task.lambda_task import LambdaTask
from conductor.client.workflow.task.task_type import TaskType

SCRIPT = "(function(){ return {out: $.x + 1}; })()"


def test_lambda_task_builds_script_expression():
task = LambdaTask(task_ref_name="lambda_ref", script=SCRIPT)
assert task.task_reference_name == "lambda_ref"
assert task.input_parameters == {"scriptExpression": SCRIPT}

workflow_task = task.to_workflow_task()
assert workflow_task.type == TaskType.LAMBDA.value
assert workflow_task.input_parameters["scriptExpression"] == SCRIPT


def test_lambda_task_merges_bindings():
task = LambdaTask(
task_ref_name="lambda_ref",
script=SCRIPT,
bindings={"x": "${workflow.input.x}"},
)
assert task.input_parameters == {
"scriptExpression": SCRIPT,
"x": "${workflow.input.x}",
}