-
Notifications
You must be signed in to change notification settings - Fork 38
Expand file tree
/
Copy pathworker.py
More file actions
148 lines (126 loc) · 5.83 KB
/
worker.py
File metadata and controls
148 lines (126 loc) · 5.83 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
import dataclasses
import inspect
import logging
import time
import traceback
from copy import deepcopy
from typing import Any, Callable, Union
from typing_extensions import Self
from conductor.client.automator import utils
from conductor.client.automator.utils import convert_from_dict_or_list
from conductor.client.configuration.configuration import Configuration
from conductor.client.http.api_client import ApiClient
from conductor.client.http.models import TaskExecLog
from conductor.client.http.models.task import Task
from conductor.client.http.models.task_result import TaskResult
from conductor.client.http.models.task_result_status import TaskResultStatus
from conductor.client.worker.exception import NonRetryableException
from conductor.client.worker.worker_interface import WorkerInterface, DEFAULT_POLLING_INTERVAL
ExecuteTaskFunction = Callable[
[
Union[Task, object]
],
Union[TaskResult, object]
]
logger = logging.getLogger(
Configuration.get_logging_formatted_name(
__name__
)
)
def is_callable_input_parameter_a_task(callable: ExecuteTaskFunction, object_type: Any) -> bool:
parameters = inspect.signature(callable).parameters
if len(parameters) != 1:
return False
parameter = parameters[list(parameters.keys())[0]]
return parameter.annotation == object_type or parameter.annotation == parameter.empty or parameter.annotation is object
def is_callable_return_value_of_type(callable: ExecuteTaskFunction, object_type: Any) -> bool:
return_annotation = inspect.signature(callable).return_annotation
return return_annotation == object_type
class Worker(WorkerInterface):
def __init__(self,
task_definition_name: str,
execute_function: ExecuteTaskFunction,
poll_interval: float = None,
domain: str = None,
worker_id: str = None,
) -> Self:
super().__init__(task_definition_name)
self.api_client = ApiClient()
if poll_interval is None:
self.poll_interval = DEFAULT_POLLING_INTERVAL
else:
self.poll_interval = deepcopy(poll_interval)
self.domain = deepcopy(domain)
if worker_id is None:
self.worker_id = deepcopy(super().get_identity())
else:
self.worker_id = deepcopy(worker_id)
self.execute_function = deepcopy(execute_function)
def execute(self, task: Task) -> TaskResult:
task_input = {}
task_output = None
task_result: TaskResult = self.get_task_result_from_task(task)
try:
if self._is_execute_function_input_parameter_a_task:
task_output = self.execute_function(task)
else:
params = inspect.signature(self.execute_function).parameters
for input_name in params:
typ = params[input_name].annotation
default_value = params[input_name].default
if input_name in task.input_data:
if typ in utils.simple_types:
task_input[input_name] = task.input_data[input_name]
else:
task_input[input_name] = convert_from_dict_or_list(typ, task.input_data[input_name])
else:
if default_value is not inspect.Parameter.empty:
task_input[input_name] = default_value
else:
task_input[input_name] = None
task_output = self.execute_function(**task_input)
if isinstance(task_output, TaskResult):
task_output.task_id = task.task_id
task_output.workflow_instance_id = task.workflow_instance_id
return task_output
else:
task_result.status = TaskResultStatus.COMPLETED
task_result.output_data = task_output
except NonRetryableException as ne:
task_result.status = TaskResultStatus.FAILED_WITH_TERMINAL_ERROR
if len(ne.args) > 0:
task_result.reason_for_incompletion = ne.args[0]
except Exception as ne:
logger.error(
f'Error executing task {task.task_def_name} with id {task.task_id}. error = {traceback.format_exc()}')
task_result.logs = [TaskExecLog(
traceback.format_exc(), task_result.task_id, int(time.time()))]
task_result.status = TaskResultStatus.FAILED
if len(ne.args) > 0:
task_result.reason_for_incompletion = ne.args[0]
if dataclasses.is_dataclass(type(task_result.output_data)):
task_output = dataclasses.asdict(task_result.output_data)
task_result.output_data = task_output
return task_result
if not isinstance(task_result.output_data, dict):
task_output = task_result.output_data
task_result.output_data = self.api_client.sanitize_for_serialization(task_output)
if not isinstance(task_result.output_data, dict):
task_result.output_data = {'result': task_result.output_data}
return task_result
def get_identity(self) -> str:
return self.worker_id
@property
def execute_function(self) -> ExecuteTaskFunction:
return self._execute_function
@execute_function.setter
def execute_function(self, execute_function: ExecuteTaskFunction) -> None:
self._execute_function = execute_function
self._is_execute_function_input_parameter_a_task = is_callable_input_parameter_a_task(
callable=execute_function,
object_type=Task,
)
self._is_execute_function_return_value_a_task_result = is_callable_return_value_of_type(
callable=execute_function,
object_type=TaskResult,
)