84 lines
2.4 KiB
Python
84 lines
2.4 KiB
Python
from abc import ABC, abstractmethod
|
|
from enum import Enum
|
|
from typing import List, Any, Optional
|
|
from radiuma_api.io_port import InPort
|
|
from radiuma_api.io_port import OutPort
|
|
from radiuma_api.execution_context import ExecutionContext
|
|
# Write Codes Here:
|
|
|
|
|
|
class Status(Enum):
|
|
PENDING = "pending"
|
|
READY = "ready"
|
|
RUNNING = "running"
|
|
PAUSED = "paused"
|
|
STOPPED = "stopped"
|
|
COMPLETED = "completed"
|
|
FAILED = "failed"
|
|
|
|
class TaskEvent(Enum):
|
|
BEFORE_RUN = "before_run"
|
|
AFTER_RUN = "after_run"
|
|
ON_ERROR = "on_error"
|
|
STATUS_CHANGED = "status_changed"
|
|
|
|
class TaskEventListener:
|
|
def handle(self, event: TaskEvent, task: "Task", payload: Optional[Any] = None):
|
|
pass
|
|
|
|
class Task(ABC):
|
|
def __init__(self, name: str):
|
|
self.name = name
|
|
self.status = Status.PENDING
|
|
self.inputs: List[InPort] = []
|
|
self.outputs: List[OutPort] = []
|
|
self._listeners: List[TaskEventListener] = []
|
|
|
|
# Events
|
|
def add_listener(self, listener: TaskEventListener):
|
|
self._listeners.append(listener)
|
|
|
|
def _emit(self, event: TaskEvent, payload: Optional[Any] = None):
|
|
for listener in self._listeners:
|
|
listener.handle(event, self, payload)
|
|
|
|
def _set_status(self, new_status: Status):
|
|
old_status = self.status
|
|
self.status = new_status
|
|
if old_status != new_status:
|
|
self._emit(TaskEvent.STATUS_CHANGED, {"from": old_status, "to": new_status})
|
|
|
|
# Lifecycle
|
|
def pause(self):
|
|
if self.status == Status.RUNNING:
|
|
self._set_status(Status.PAUSED)
|
|
|
|
def resume(self):
|
|
if self.status == Status.PAUSED:
|
|
self._set_status(Status.READY)
|
|
|
|
def stop(self):
|
|
if self.status in (Status.RUNNING, Status.PAUSED, Status.READY):
|
|
self._set_status(Status.STOPPED)
|
|
|
|
# Ports
|
|
def _add_input_port(self, port: InPort):
|
|
port.parent_task = self
|
|
self.inputs.append(port)
|
|
|
|
def _add_output_port(self, port: OutPort):
|
|
port.parent_task = self
|
|
self.outputs.append(port)
|
|
|
|
def check_inputs_ready(self, context: ExecutionContext) -> bool:
|
|
"""Check if all required inputs have data in ExecutionContext."""
|
|
for in_port in self.inputs:
|
|
if in_port.required and not context.has_data(in_port):
|
|
return False
|
|
return True
|
|
|
|
# Runnable
|
|
@abstractmethod
|
|
def run(self, context: ExecutionContext):
|
|
pass
|