File: execution.py

package info (click to toggle)
python-moto 5.1.18-3
  • links: PTS, VCS
  • area: main
  • in suites: forky, sid
  • size: 116,520 kB
  • sloc: python: 636,725; javascript: 181; makefile: 39; sh: 3
file content (358 lines) | stat: -rw-r--r-- 13,325 bytes parent folder | download
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
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
from __future__ import annotations

import datetime
import json
import logging
from typing import Optional

from moto.stepfunctions.parser.api import (
    Arn,
    CloudWatchEventsExecutionDataDetails,
    DescribeExecutionOutput,
    DescribeStateMachineForExecutionOutput,
    ExecutionListItem,
    ExecutionStatus,
    GetExecutionHistoryOutput,
    HistoryEventList,
    InvalidName,
    SensitiveCause,
    SensitiveError,
    StartExecutionOutput,
    StartSyncExecutionOutput,
    StateMachineType,
    SyncExecutionStatus,
    Timestamp,
    TraceHeader,
    VariableReferences,
)
from moto.stepfunctions.parser.asl.eval.evaluation_details import (
    AWSExecutionDetails,
    EvaluationDetails,
    ExecutionDetails,
    StateMachineDetails,
)
from moto.stepfunctions.parser.asl.eval.event.logging import (
    CloudWatchLoggingSession,
)
from moto.stepfunctions.parser.asl.eval.program_state import (
    ProgramEnded,
    ProgramError,
    ProgramState,
    ProgramStopped,
    ProgramTimedOut,
)
from moto.stepfunctions.parser.asl.static_analyser.variable_references_static_analyser import (
    VariableReferencesStaticAnalyser,
)
from moto.stepfunctions.parser.asl.utils.encoding import to_json_str
from moto.stepfunctions.parser.backend.activity import Activity
from moto.stepfunctions.parser.backend.execution_worker import (
    ExecutionWorker,
    SyncExecutionWorker,
)
from moto.stepfunctions.parser.backend.execution_worker_comm import (
    ExecutionWorkerCommunication,
)
from moto.stepfunctions.parser.backend.state_machine import (
    StateMachineInstance,
    StateMachineVersion,
)

LOG = logging.getLogger(__name__)


class BaseExecutionWorkerCommunication(ExecutionWorkerCommunication):
    execution: Execution

    def __init__(self, execution: Execution):
        self.execution = execution

    def _reflect_execution_status(self):
        exit_program_state: ProgramState = (
            self.execution.exec_worker.env.program_state()
        )
        self.execution.stop_date = datetime.datetime.now(tz=datetime.timezone.utc)
        if isinstance(exit_program_state, ProgramEnded):
            self.execution.exec_status = ExecutionStatus.SUCCEEDED
            self.execution.output = self.execution.exec_worker.env.states.get_input()
        elif isinstance(exit_program_state, ProgramStopped):
            self.execution.exec_status = ExecutionStatus.ABORTED
        elif isinstance(exit_program_state, ProgramError):
            self.execution.exec_status = ExecutionStatus.FAILED
            self.execution.error = exit_program_state.error.get("error")
            self.execution.cause = exit_program_state.error.get("cause")
        elif isinstance(exit_program_state, ProgramTimedOut):
            self.execution.exec_status = ExecutionStatus.TIMED_OUT
        else:
            raise RuntimeWarning(
                f"Execution ended with unsupported ProgramState type '{type(exit_program_state)}'."
            )

    def terminated(self) -> None:
        self._reflect_execution_status()


class Execution:
    name: str
    sm_type: StateMachineType
    role_arn: Arn
    exec_arn: Arn

    account_id: str
    region_name: str

    state_machine: StateMachineInstance
    start_date: Timestamp
    input_data: Optional[json]
    input_details: Optional[CloudWatchEventsExecutionDataDetails]
    trace_header: Optional[TraceHeader]
    _cloud_watch_logging_session: Optional[CloudWatchLoggingSession]

    exec_status: Optional[ExecutionStatus]
    stop_date: Optional[Timestamp]

    output: Optional[json]
    output_details: Optional[CloudWatchEventsExecutionDataDetails]

    error: Optional[SensitiveError]
    cause: Optional[SensitiveCause]

    exec_worker: Optional[ExecutionWorker]

    _activity_store: dict[Arn, Activity]

    def __init__(
        self,
        name: str,
        sm_type: StateMachineType,
        role_arn: Arn,
        exec_arn: Arn,
        account_id: str,
        region_name: str,
        state_machine: StateMachineInstance,
        start_date: Timestamp,
        cloud_watch_logging_session: Optional[CloudWatchLoggingSession],
        activity_store: dict[Arn, Activity],
        input_data: str,
        trace_header: Optional[TraceHeader] = None,
    ):
        self.name = name
        self.sm_type = sm_type
        self.role_arn = role_arn
        self.exec_arn = exec_arn
        self.execution_arn = exec_arn
        self.account_id = account_id
        self.region_name = region_name
        self.state_machine = state_machine
        self._cloud_watch_logging_session = cloud_watch_logging_session
        self.input_data = json.loads(input_data)
        self.input_details = CloudWatchEventsExecutionDataDetails(included=True)
        self.trace_header = trace_header
        self.exec_status = None
        self.stop_date = None
        self.output = None
        self.output_details = CloudWatchEventsExecutionDataDetails(included=True)
        self.exec_worker = None
        self.error = None
        self.cause = None
        self._activity_store = activity_store

        # Compatibility with mock SFN
        self.state_machine_arn = state_machine.arn
        self.start_date = start_date
        self.execution_input = input_data

    @property
    def status(self):
        return self.exec_status.value

    def to_start_output(self) -> StartExecutionOutput:
        return StartExecutionOutput(
            executionArn=self.exec_arn, startDate=self.start_date
        )

    def to_describe_output(self) -> DescribeExecutionOutput:
        describe_output = DescribeExecutionOutput(
            executionArn=self.exec_arn,
            stateMachineArn=self.state_machine.arn,
            name=self.name,
            status=self.exec_status,
            startDate=self.start_date,
            stopDate=self.stop_date,
            input=to_json_str(self.input_data, separators=(",", ":")),
            inputDetails=self.input_details,
            traceHeader=self.trace_header,
        )
        if describe_output["status"] == ExecutionStatus.SUCCEEDED:
            describe_output["output"] = to_json_str(self.output, separators=(",", ":"))
            describe_output["outputDetails"] = self.output_details
        if self.error is not None:
            describe_output["error"] = self.error
        if self.cause is not None:
            describe_output["cause"] = self.cause
        return describe_output

    def to_describe_state_machine_for_execution_output(
        self,
    ) -> DescribeStateMachineForExecutionOutput:
        state_machine: StateMachineInstance = self.state_machine
        state_machine_arn = (
            state_machine.source_arn
            if isinstance(state_machine, StateMachineVersion)
            else state_machine.arn
        )
        out = DescribeStateMachineForExecutionOutput(
            stateMachineArn=state_machine_arn,
            name=state_machine.name,
            definition=state_machine.definition,
            roleArn=self.role_arn,
            # The date and time the state machine associated with an execution was updated.
            updateDate=state_machine.create_date,
            loggingConfiguration=state_machine.logging_config,
        )
        revision_id = self.state_machine.revision_id
        if self.state_machine.revision_id:
            out["revisionId"] = revision_id
        variable_references: VariableReferences = (
            VariableReferencesStaticAnalyser.process_and_get(
                definition=self.state_machine.definition
            )
        )
        if variable_references:
            out["variableReferences"] = variable_references
        return out

    def to_execution_list_item(self) -> ExecutionListItem:
        if isinstance(self.state_machine, StateMachineVersion):
            state_machine_arn = self.state_machine.source_arn
            state_machine_version_arn = self.state_machine.arn
        else:
            state_machine_arn = self.state_machine.arn
            state_machine_version_arn = None

        item = ExecutionListItem(
            executionArn=self.exec_arn,
            stateMachineArn=state_machine_arn,
            name=self.name,
            status=self.exec_status,
            startDate=self.start_date,
            stopDate=self.stop_date,
        )
        if state_machine_version_arn is not None:
            item["stateMachineVersionArn"] = state_machine_version_arn
        return item

    def to_history_output(self) -> GetExecutionHistoryOutput:
        env = self.exec_worker.env
        event_history: HistoryEventList = []
        if env is not None:
            # The execution has not started yet.
            event_history: HistoryEventList = env.event_manager.get_event_history()
        return GetExecutionHistoryOutput(events=event_history)

    def _get_start_execution_worker_comm(self) -> BaseExecutionWorkerCommunication:
        return BaseExecutionWorkerCommunication(self)

    def _get_start_aws_execution_details(self) -> AWSExecutionDetails:
        return AWSExecutionDetails(
            account=self.account_id, region=self.region_name, role_arn=self.role_arn
        )

    def get_start_execution_details(self) -> ExecutionDetails:
        return ExecutionDetails(
            arn=self.exec_arn,
            name=self.name,
            role_arn=self.role_arn,
            inpt=self.input_data,
            start_time=self.start_date,
        )

    def get_start_state_machine_details(self) -> StateMachineDetails:
        return StateMachineDetails(
            arn=self.state_machine.arn,
            name=self.state_machine.name,
            typ=self.state_machine.sm_type,
            definition=self.state_machine.definition,
        )

    def _get_start_execution_worker(self) -> ExecutionWorker:
        return ExecutionWorker(
            evaluation_details=EvaluationDetails(
                aws_execution_details=self._get_start_aws_execution_details(),
                execution_details=self.get_start_execution_details(),
                state_machine_details=self.get_start_state_machine_details(),
            ),
            exec_comm=self._get_start_execution_worker_comm(),
            cloud_watch_logging_session=self._cloud_watch_logging_session,
            activity_store=self._activity_store,
        )

    def start(self) -> None:
        # TODO: checks exec_worker does not exists already?
        if self.exec_worker:
            raise InvalidName()  # TODO.
        self.exec_worker = self._get_start_execution_worker()
        self.exec_status = ExecutionStatus.RUNNING
        self.exec_worker.start()

    def stop(
        self, stop_date: datetime.datetime, error: Optional[str], cause: Optional[str]
    ):
        exec_worker: Optional[ExecutionWorker] = self.exec_worker
        if exec_worker:
            exec_worker.stop(stop_date=stop_date, cause=cause, error=error)


class SyncExecutionWorkerCommunication(BaseExecutionWorkerCommunication):
    execution: SyncExecution

    def _reflect_execution_status(self) -> None:
        super()._reflect_execution_status()
        exit_status: ExecutionStatus = self.execution.exec_status
        if exit_status == ExecutionStatus.SUCCEEDED:
            self.execution.sync_execution_status = SyncExecutionStatus.SUCCEEDED
        elif exit_status == ExecutionStatus.TIMED_OUT:
            self.execution.sync_execution_status = SyncExecutionStatus.TIMED_OUT
        else:
            self.execution.sync_execution_status = SyncExecutionStatus.FAILED


class SyncExecution(Execution):
    sync_execution_status: Optional[SyncExecutionStatus] = None

    def _get_start_execution_worker(self) -> SyncExecutionWorker:
        return SyncExecutionWorker(
            evaluation_details=EvaluationDetails(
                aws_execution_details=self._get_start_aws_execution_details(),
                execution_details=self.get_start_execution_details(),
                state_machine_details=self.get_start_state_machine_details(),
            ),
            exec_comm=self._get_start_execution_worker_comm(),
            cloud_watch_logging_session=self._cloud_watch_logging_session,
            activity_store=self._activity_store,
        )

    def _get_start_execution_worker_comm(self) -> BaseExecutionWorkerCommunication:
        return SyncExecutionWorkerCommunication(self)

    def to_start_sync_execution_output(self) -> StartSyncExecutionOutput:
        start_output = StartSyncExecutionOutput(
            executionArn=self.exec_arn,
            stateMachineArn=self.state_machine.arn,
            name=self.name,
            status=self.sync_execution_status,
            startDate=self.start_date,
            stopDate=self.stop_date,
            input=to_json_str(self.input_data, separators=(",", ":")),
            inputDetails=self.input_details,
            traceHeader=self.trace_header,
        )
        if self.sync_execution_status == SyncExecutionStatus.SUCCEEDED:
            start_output["output"] = to_json_str(self.output, separators=(",", ":"))
        if self.output_details:
            start_output["outputDetails"] = self.output_details
        if self.error is not None:
            start_output["error"] = self.error
        if self.cause is not None:
            start_output["cause"] = self.cause
        return start_output