from abc import ABCMeta, abstractmethod from multiprocessing import Event, Process from multiprocessing.connection import Connection from typing import Any, Dict, Optional, Set, Type, Union from smelog.factory import BoundLogger from logger import get_logger from utils.key import Key from utils.perf_counter import PerfCounter, timeit __all__ = ["BaseStep", "StepError"] class StepError(Exception): pass class BaseStepMeta(ABCMeta): __step_name_prefix: str = f"{__name__.split('.')[0]}." @property def step_name(cls) -> str: return cls.__module__.replace(cls.__step_name_prefix, "") class BaseStep(Process, metaclass=BaseStepMeta): depends_on: Set[Union[Type["BaseStep"], str]] = set() priority: int = 10 # less means higher priority def __init__(self, key: Key, conn: Connection): super().__init__(name=self.__class__.step_name) self.__conn = conn self.key = key self.perf_counter = PerfCounter(self.__class__.step_name) self.logger: Optional[BoundLogger] = None self.__is_done = Event() def _pre_run(self): self.logger = get_logger(str(self.key), pid=self.pid, step=self.__class__.step_name) def _post_run(self): pass @property def timestamp(self) -> int: return self.key.timestamp.to_int() @property def uuid(self) -> str: return self.key.uuid def is_done(self) -> bool: return self.__is_done.is_set() def run(self) -> None: self._pre_run() self.logger.info(f"Executing step '{self.__class__.step_name}' with key '{self.key}'") # type: ignore with timeit(self.perf_counter, "total"): try: self.process() finally: self._post_run() # TODO: fix this self.__conn.send(({}, self.perf_counter.report)) # self.__conn.send((result, self.perf_counter.report)) <- payload is to heavy, python can't manage it # self.__conn.send(({}, {})) self.__is_done.set() @abstractmethod def process(self) -> Dict[str, Any]: pass