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
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,9 +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.
A claimed task whose stored payload cannot be read any more is dropped with the
read error logged. A dropped task records no result, so it leaves the inspector
and cannot be requeued.
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
107 changes: 55 additions & 52 deletions tests/backends/test_redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,38 @@ def _expire_lease(
)


def _assert_dropped(backend: RedisTaskBackend, task_id: str) -> None:
"""Assert a dropped task keeps its running entry without its payload."""
assert (
backend.client.exists(
backend.TASK_KEY.format(prefix=backend.key_prefix, task_id=task_id)
)
== 0
)
running_key = backend._segment_key(TaskResultStatus.RUNNING, "default")
assert backend.client.zscore(running_key, task_id) is not None


_GONE_CALLBACK = "tests.testapp.tasks.gone_from_the_code_base"


def _corrupt_payload(
backend: RedisTaskBackend,
task_id: str,
data: str | None = None,
task: dict | None = None,
**fields: object,
) -> None:
"""Overwrite a leased task's stored payload with raw data or patched fields."""
task_key = backend.TASK_KEY.format(prefix=backend.key_prefix, task_id=task_id)
if task or fields:
payload = json.loads(backend.client.hget(task_key, "data"))
payload["task"].update(task or {})
payload.update(fields)
data = json.dumps(payload)
backend.client.hset(task_key, "data", data)


def _claim_expired(
backend: RedisTaskBackend,
broker: RedisBroker,
Expand Down Expand Up @@ -272,65 +304,36 @@ def test_reap__removes_running_entry_without_task_data(self):
finally:
backend.close()

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."""
@pytest.mark.parametrize(
"corruption",
[
{"data": "{not json"},
{"data": "{}"},
{"task": {"func": _GONE_CALLBACK}},
{"task": {"retry": _GONE_CALLBACK}},
{"status": "CANCELLED"},
],
ids=["malformed", "not-a-result", "gone-func", "gone-retry", "unknown-status"],
)
def test_reap__drops_task_when_payload_cannot_be_read(self, corruption, caplog):
"""A payload the reaper cannot read drops the claimed task."""
backend = _make_backend(
"reap_gone_callback_test", lease_ttl=datetime.timedelta(seconds=1)
"reap_drop_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
)
payload = json.loads(backend.client.hget(task_key, "data"))
payload["task"]["retry"] = "tests.testapp.tasks.gone_from_the_code_base"
backend.client.hset(task_key, "data", json.dumps(payload))
task_result = backend.enqueue(echo, args=[42])
acquired = backend.acquire(
timeout=datetime.timedelta(seconds=1), worker="worker-1"
)
assert acquired is not None
_corrupt_payload(backend, task_result.id, **corruption)
_expire_lease(backend, task_result.id)

with caplog.at_level(logging.ERROR, logger="threadmill.backends.redis"):
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
assert "payload cannot be read; dropping the task" in caplog.text
_assert_dropped(backend, task_result.id)
finally:
backend.close()

Expand Down
59 changes: 17 additions & 42 deletions threadmill/backends/redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@
import collections.abc
import dataclasses
import datetime
import json
import logging
import queue
import random
Expand Down Expand Up @@ -96,49 +95,25 @@ def _reap_running_queue(self, queue_name: str) -> None:
],
)
for task_id in claimed_ids:
try:
self._reap_task(task_id)
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; "
"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 or None,
lease_started_at=_parse_lease_started_at(lease_started_at),
lease_token=lease_token,
)
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),
],
)
)
self._reap_task(task_id)

def _reap_task(self, task_id: str) -> None:
"""Requeue or fail a claimed task."""
if (task_result := self.backend.get_leased_task(task_id)) is None:
"""Requeue or fail it, and drop it when the payload cannot be read."""
try:
task_result = self.backend.get_leased_task(task_id)
except ImportError, ValueError, KeyError, TypeError, AttributeError:
logger.exception(
"Task %r payload cannot be read; dropping the task", task_id
)
# Drop only the payload: the running entry stays for a lease holder's
# late acknowledge, and the reaper script sweeps it once the hash is gone.
self.backend.client.delete(
self.backend.TASK_KEY.format(
prefix=self.backend.key_prefix, task_id=task_id
)
)
return
if task_result is None:
logger.warning("Claimed task %r has no task data; skipping", task_id)
return
now = timezone.now()
Expand Down
Loading