diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index c8c07b4..fa912fc 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -13,7 +13,7 @@ jobs: - uses: actions/checkout@v3 - name: Lint with ruff run: | - pipx install ruff==0.4.4 + pipx install ruff==0.4.10 ruff format . --check && ruff check . test: diff --git a/django_q/exceptions.py b/django_q/exceptions.py new file mode 100644 index 0000000..eb40677 --- /dev/null +++ b/django_q/exceptions.py @@ -0,0 +1,7 @@ +class TimeoutException(SystemExit): + """ + Exception for when a worker takes too long to complete a task + Raising SystemExit will make sure the function terminates gracefully. + """ + + pass diff --git a/django_q/timeout.py b/django_q/timeout.py new file mode 100644 index 0000000..3d40a06 --- /dev/null +++ b/django_q/timeout.py @@ -0,0 +1,38 @@ +import signal + +from django.utils.translation import gettext_lazy as _ + +from django_q.conf import logger + +from .exceptions import TimeoutException + + +class TimeoutHandler: + def __init__(self, timeout: int): + self._timeout = timeout + + def raise_timeout_exception(self, signum, frame): + raise TimeoutException( + f"Task exceeded maximum timeout value ({self._timeout} seconds)" + ) + + def __enter__(self): + # if the timeout is -1, then there is no timeout and the task will always keep running until it's done or manually killed + if self._timeout == -1: + return + try: + signal.signal(signal.SIGALRM, self.raise_timeout_exception) + except ValueError: # ValueError is raised for Windows users + logger.debug(_("SIGALARM is not available on your platform")) + + signal.alarm(self._timeout) + + def __exit__(self, exc_type, exc_value, traceback): + if self._timeout == -1: + return + """When getting out of the timeout, reset the alarm, so it won't trigger""" + try: + signal.alarm(0) + signal.signal(signal.SIGALRM, signal.SIG_DFL) + except ValueError: # ValueError is raised for Windows users + logger.debug(_("SIGALARM is not available on your platform")) diff --git a/django_q/worker.py b/django_q/worker.py index 47e4fb7..91d45e4 100644 --- a/django_q/worker.py +++ b/django_q/worker.py @@ -17,7 +17,9 @@ except core.exceptions.AppRegistryNotReady: django.setup() from django_q.conf import Conf, error_reporter, logger, resource, setproctitle +from django_q.exceptions import TimeoutException from django_q.signals import post_spawn, pre_execute +from django_q.timeout import TimeoutHandler from django_q.utils import close_old_django_connections, get_func_repr try: @@ -89,25 +91,37 @@ def worker( pre_execute.send(sender="django_q", func=f, task=task) # execute the payload timer.value = timer_value # Busy + if timer.value != -1: + timer.value += 3 # Add buffer so that guard doesn't kill the process on timeout before it gets processed + timeout_error = False try: if f is None: # raise a meaningfull error if task["func"] is not a valid function raise ValueError(f"Function {task['func']} is not defined") - res = f(*task["args"], **task["kwargs"]) + with TimeoutHandler(timer_value): + res = f(*task["args"], **task["kwargs"]) result = (res, True) - except Exception as e: + except (Exception, TimeoutException) as e: + if isinstance(e, TimeoutException): + timeout_error = True result = (f"{e} : {traceback.format_exc()}", False) if error_reporter: error_reporter.report() if task.get("sync", False): raise + with timer.get_lock(): # Process result task["result"] = result[0] task["success"] = result[1] task["stopped"] = timezone.now() result_queue.put(task) + if timeout_error: + # force destroy process due to timeout + timer.value = 0 + break + timer.value = -1 # Idle if setproctitle: setproctitle.setproctitle(f"qcluster {proc_name} idle") diff --git a/docs/configure.rst b/docs/configure.rst index 6e25d0e..864e709 100644 --- a/docs/configure.rst +++ b/docs/configure.rst @@ -79,6 +79,8 @@ timeout The number of seconds a worker is allowed to spend on a task before it's terminated. Defaults to ``None``, meaning it will never time out. Set this to something that makes sense for your project. Can be overridden for individual tasks. +Note: for systems that don't have `SIGALRM` available (e.g. Windows), it will not raise an error properly. It will kill the task, but it will keep retrying until it finishes within the given time. + See :ref:`retry` for details how to set values for timeout and retry. .. _time_zone: