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()