Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 15 additions & 2 deletions taskqueue.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,16 +11,29 @@
the macOS supervisor launches the bundled ``redis-server`` binary (see
:func:`build_embedded_redis_argv`) and exports its socket URL as ``REDIS_URL``
before the app and workers boot.

``rq`` hardcodes ``UnixSignalDeathPenalty`` as every job registry's default, so
on Windows (no ``signal.SIGALRM``) registry cleanup -- the janitor loop and the
workers' periodic maintenance -- raises ``AttributeError`` whenever an abandoned
job carries an ``on_failure`` callback, and the job is never removed.
``BaseRegistry`` is therefore re-pointed at RQ's own platform dispatcher and
the queues below are built with the same class: a no-op on POSIX, the
timer-based penalty (what RQ's workers already use) on Windows.
"""

from redis import Redis
from rq import Queue, get_current_job
from rq.job import Job
from rq.exceptions import NoSuchJobError
from rq.command import send_stop_job_command
from rq.registry import BaseRegistry
from rq.timeouts import get_default_death_penalty_class

import config

_death_penalty_class = get_default_death_penalty_class()
BaseRegistry.death_penalty_class = _death_penalty_class

__all__ = [
"redis_conn",
"rq_queue_high",
Expand Down Expand Up @@ -54,8 +67,8 @@ def redis_socket_options(url):
**redis_socket_options(config.REDIS_URL),
)

rq_queue_high = Queue('high', connection=redis_conn, default_timeout=-1)
rq_queue_default = Queue('default', connection=redis_conn, default_timeout=-1)
rq_queue_high = Queue('high', connection=redis_conn, default_timeout=-1, death_penalty_class=_death_penalty_class)
rq_queue_default = Queue('default', connection=redis_conn, default_timeout=-1, death_penalty_class=_death_penalty_class)


def build_embedded_redis_argv(server_binary, socket_path, data_dir):
Expand Down
21 changes: 21 additions & 0 deletions tests/unit/test_taskqueue.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import taskqueue
from rq.registry import BaseRegistry
from rq.timeouts import get_default_death_penalty_class


def test_base_registry_uses_platform_death_penalty():
assert BaseRegistry.death_penalty_class is get_default_death_penalty_class()


def test_queues_use_platform_death_penalty():
expected = get_default_death_penalty_class()
assert taskqueue.rq_queue_high.death_penalty_class is expected
assert taskqueue.rq_queue_default.death_penalty_class is expected


def test_registries_built_from_queues_inherit_platform_death_penalty():
expected = get_default_death_penalty_class()
for queue in (taskqueue.rq_queue_high, taskqueue.rq_queue_default):
assert queue.started_job_registry.death_penalty_class is expected
assert queue.finished_job_registry.death_penalty_class is expected
assert queue.failed_job_registry.death_penalty_class is expected
Loading