35 lines
1.1 KiB
Python
35 lines
1.1 KiB
Python
from radiuma_api.task import Task, Status, TaskEvent
|
|
from radiuma_api.asset import Asset
|
|
from radiuma_api.execution_context import ExecutionContext
|
|
from radiuma_api.io_port import CompatibilityException
|
|
# Write Codes Here:
|
|
|
|
|
|
class Module(Task):
|
|
"""Leaf node: atomic executable unit."""
|
|
def run(self, context: ExecutionContext):
|
|
if not self.check_inputs_ready(context):
|
|
self._set_status(Status.PENDING)
|
|
return
|
|
|
|
self._set_status(Status.RUNNING)
|
|
|
|
try:
|
|
consumed_data = [
|
|
port.connected_output.produced_asset.data
|
|
for port in self.inputs
|
|
if port.is_ready()
|
|
]
|
|
for output_port in self.outputs:
|
|
asset = Asset(
|
|
data=f"{self.name} processed {consumed_data}",
|
|
type_name=output_port.type_spec
|
|
)
|
|
context.put(output_port, asset)
|
|
|
|
self._set_status(Status.COMPLETED)
|
|
|
|
except CompatibilityException as e:
|
|
self._set_status(Status.FAILED)
|
|
self._emit(TaskEvent.ON_ERROR, {"error": e.reason})
|