diff --git a/README.md b/README.md index 128e84f..bade689 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index 08a8284..a113c84 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -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 @@ -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() diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index f6a0630..442c4db 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -3,6 +3,7 @@ import collections.abc import dataclasses import datetime +import json import logging import queue import random @@ -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: