import logging from threading import Event logger = logging.getLogger(__name__) class Task: def __init__(self, func, args, kwargs): """Holds the function and parameters, that should be executed by thread worker""" 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 logger.exception(e) finally: self._result_event.set() self._immediate_response_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, exception): if not self.immediate_response_sent(): if isinstance(exception, str): self.immediate_error = RuntimeError(exception) else: self.immediate_error = exception 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 self.immediate_error if self.failed: raise self.result_exception return self.immediate_response or self.result 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