From a59e8e7ea76c2045ace2cdb169c72d580a80a675 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 15:47:52 +0200 Subject: [PATCH] Drop a claimed task whose stored payload cannot be read `_reap_task` deserialized the stored payload before it decided the fate of the task. A payload with a task function, a retry callback, or a status the code base does not know raised on every pass. The claim script had already renewed the lease, so the task stayed in the running set and the same error repeated every `CLAIM_TTL`. - Drop a claimed task whose payload cannot be read any more: log the read error with its traceback, and delete the stored payload. The running entry stays, so a lease holder's late acknowledge still writes its result, and the reaper script sweeps the entry once the hash is gone. - Document the behavior in the README. Fix #73 --- README.md | 6 +- tests/backends/test_redis.py | 107 ++++++++++++++++++----------------- threadmill/backends/redis.py | 59 ++++++------------- 3 files changed, 75 insertions(+), 97 deletions(-) diff --git a/README.md b/README.md index bade689..f2156de 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index a5eea12..3d94e95 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -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, @@ -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() diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index 3078031..b270a2e 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -3,7 +3,6 @@ import collections.abc import dataclasses import datetime -import json import logging import queue import random @@ -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()