Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,9 @@ The `RedisTaskBackend` accepts the following options under `OPTIONS` in your

A task whose lease expired reaches the `retry` callback as an
`AcknowledgementTimeout` error, or is marked FAILED when nothing retries it.
A stored retry callback that is gone from the code base fails the task the same
way, recording an `ImportError` beside the timeout error. The failure is stored
as a regular result, so it can be requeued or dropped from the inspector.
Keep `lease_ttl` above your worst-case runtime: a task that outlives its lease
can still be running, so a retry may execute concurrently with it. The
acknowledgement of the lease holder wins: the late result of an expired attempt
Expand Down
41 changes: 37 additions & 4 deletions tests/backends/test_redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -272,20 +272,20 @@ def test_reap__removes_running_entry_without_task_data(self):
finally:
backend.close()

def test_reap_running_queue__logs_and_continues_when_retry_callback_is_gone(
self, caplog
):
"""A batch keeps reaping when a task's retry callback is gone from the code base."""
def test_reap__fails_tasks_when_retry_callback_is_gone(self, caplog):
"""A batch fails each task whose stored retry callback is gone from the code base."""
backend = _make_backend(
"reap_gone_callback_test", lease_ttl=datetime.timedelta(seconds=1)
)
try:
task_ids = []
for _index in range(2):
task_result = backend.enqueue(boom_no_retry, args=[])
acquired = backend.acquire(
timeout=datetime.timedelta(seconds=1), worker="worker-1"
)
assert acquired is not None
task_ids.append(task_result.id)
_expire_lease(backend, task_result.id)
task_key = backend.TASK_KEY.format(
prefix=backend.key_prefix, task_id=task_result.id
Expand All @@ -298,6 +298,39 @@ def test_reap_running_queue__logs_and_continues_when_retry_callback_is_gone(
RedisBroker(backend)._reap_running_queue("default")

assert caplog.text.count("gone from the code base") == 2
running_key = backend._segment_key(TaskResultStatus.RUNNING, "default")
assert backend.client.zcard(running_key) == 0
assert (
backend.client.exists(
backend.TASK_KEY.format(
prefix=backend.key_prefix, task_id=task_ids[0]
)
)
== 0
)
failed = {
result.id: result
for result in backend.peek(
"default", status=TaskResultStatus.FAILED, count=0
)
}
assert set(failed) == set(task_ids)
for result in failed.values():
assert result.status is TaskResultStatus.FAILED
assert result.worker_ids == ["worker-1"]
assert result.task.retry is None
assert result.errors[-1].exception_class_path == "builtins.ImportError"
assert (
failed[task_ids[0]].errors[-2].exception_class_path
== "threadmill.exceptions.AcknowledgementTimeout"
)

# The stored failure is a regular result, so the inspector can requeue it.
backend.requeue(failed[task_ids[0]], timezone.now())
deferred_key = backend.DEFERRED_KEY.format(
prefix=backend.key_prefix, queue_name="default"
)
assert backend.client.zscore(deferred_key, task_ids[0]) is not None
finally:
backend.close()

Expand Down
40 changes: 37 additions & 3 deletions threadmill/backends/redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import collections.abc
import dataclasses
import datetime
import json
import logging
import queue
import random
Expand Down Expand Up @@ -98,11 +99,44 @@ def _reap_running_queue(self, queue_name: str) -> None:
task_id = member.decode()
try:
self._reap_task(task_id)
except ImportError:
logger.exception(
except ImportError as read_error:
task_key = self.backend.TASK_KEY.format(
prefix=self.backend.key_prefix, task_id=task_id
)
data, lease_worker, lease_started_at, lease_token = (
self.backend.client.hmget(
task_key, "data", *self.backend.LEASE_FIELDS
)
)
payload = json.loads(data)
payload["task"].pop("retry", None)
logger.error(
"Task %r retry callback is gone from the code base; "
"skipping the reap",
"failing the task: %s",
task_id,
read_error,
)
task_result = self.backend._apply_lease(
self.backend.deserialize_task_result(json.dumps(payload)),
worker=lease_worker.decode() if lease_worker else None,
lease_started_at=_parse_lease_started_at(lease_started_at),
lease_token=(
lease_token.decode() if lease_token is not None else None
),
)
self.backend.acknowledge(
dataclasses.replace(
task_result,
status=TaskResultStatus.FAILED,
finished_at=timezone.now(),
errors=[
*task_result.errors,
self.backend.create_task_error(
AcknowledgementTimeout("Task processing lease expired.")
),
self.backend.create_task_error(read_error),
],
)
)

def _reap_task(self, task_id: str) -> None:
Expand Down
Loading