from threading import Event from .. import logger log = logger.getChild("Task") class Task: def __init__(self, func, args, kwargs): """Holds the function and parameters, that should be executed by thread worker""" # TODO: .waitfor(), that possibly allows to wait for tasks later, if you remember them self.func = func self.args = args self.kwargs = kwargs self.result = None self.result_exception = None self.immediate_response = None self.immediate_error = None self.done = False self.failed = False self._result_event = Event() self._immediate_response_event = Event() def __call__(self, *args, **kwargs): """Execute task by calling itself. Handles it's own exceptions""" try: self.result = self.func(*self.args, task_handle=self, **self.kwargs) self.done = True except Exception as e: self.failed = True self.result_exception = e log.exception(e) finally: self._result_event.set() def send_immediate_response(self, resp): if not self.immediate_response_sent(): self.immediate_response = resp self._immediate_response_event.set() def send_immediate_error(self, error): if not self.immediate_response_sent(): self.immediate_error = error self._immediate_response_event.set() def immediate_response_sent(self): return self._immediate_response_event.is_set() def wait_for_immediate_response(self, timeout=None): self._immediate_response_event.wait(timeout) if self.immediate_error: raise RuntimeError(self.immediate_error) return self.immediate_response def wait_for_result(self, timeout=None): if self._result_event.wait(timeout): if self.failed: raise self.result_exception return self.result return None def result_sent(self): return self._result_event.is_set() class RatedTask(Task): def __init__(self, func, args, kwargs, rate_tracker): super().__init__(func, args, kwargs) self.rate_tracker = rate_tracker self.rate_tracker.create_wait() def __call__(self, *args, **kwargs): self.rate_tracker.start() super().__call__(*args, **kwargs) self.rate_tracker.done()