From d323ffb18a67e77eae38b5ae325c7346fcffe97a Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 14:14:33 +0200 Subject: [PATCH 1/4] Fix #78 -- Fail tasks whose stored retry callback is gone The running reaper skipped a payload it could not deserialize. The claim script had already renewed the lease, so the task stayed in the running set and the broker logged the same ImportError every CLAIM_TTL. The inspector could not list or clear it. The reaper now fails such a task. It drops the unimportable retry callback, records an ImportError beside the lease timeout, and acknowledges the failure through the regular result path. The task leaves the running set, shows on the Failed tab, and can be requeued or dropped. --- README.md | 3 ++ tests/backends/test_redis.py | 53 +++++++++++++++++++++++++++++++++--- threadmill/backends/redis.py | 53 ++++++++++++++++++++++++++++++++---- 3 files changed, 99 insertions(+), 10 deletions(-) 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..8b9ccab 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,51 @@ 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() + + def test_fail_unreadable_task__skips_when_task_data_is_missing(self, caplog): + """Failing an unreadable task logs and skips when its task data is gone.""" + backend = _make_backend("reap_gone_callback_missing_data_test") + try: + with caplog.at_level(logging.WARNING, logger="threadmill.backends.redis"): + RedisBroker(backend)._fail_unreadable_task( + "missing-task-id", ImportError("gone") + ) + assert "has no task data" in caplog.text finally: backend.close() diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index f6a0630..c067e84 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,12 +99,52 @@ def _reap_running_queue(self, queue_name: str) -> None: task_id = member.decode() try: self._reap_task(task_id) - except ImportError: - logger.exception( - "Task %r retry callback is gone from the code base; " - "skipping the reap", - task_id, - ) + except ImportError as read_error: + self._fail_unreadable_task(task_id, read_error) + + def _fail_unreadable_task(self, task_id: str, read_error: ImportError) -> None: + """Fail a claimed task whose stored retry callback is gone from the code base. + + The unimportable callback is dropped from the payload, so the payload + stays readable and the failure is stored as a regular result. The + inspector can then requeue or drop the task. + """ + 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 + ) + if data is None: + logger.warning("Claimed task %r has no task data; skipping", task_id) + return + logger.error( + "Task %r retry callback is gone from the code base; failing the task: %s", + task_id, + read_error, + ) + payload = json.loads(data) + payload["task"].pop("retry", None) + 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: """Requeue or fail a claimed task.""" From bfbe79a479709f554dfc743b65d9ba3b6e3892fb Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 14:19:36 +0200 Subject: [PATCH 2/4] Drop the missing task data guard from the reaper CONTRIBUTING.md says we do not guard against Redis data altered mid-flight, we fail loudly instead. The helper reads the task hash right after the claim returned that same key, so the guard is dead code that hides a failure instead of surfacing it. --- tests/backends/test_redis.py | 12 ------------ threadmill/backends/redis.py | 3 --- 2 files changed, 15 deletions(-) diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index 8b9ccab..a113c84 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -334,18 +334,6 @@ def test_reap__fails_tasks_when_retry_callback_is_gone(self, caplog): finally: backend.close() - def test_fail_unreadable_task__skips_when_task_data_is_missing(self, caplog): - """Failing an unreadable task logs and skips when its task data is gone.""" - backend = _make_backend("reap_gone_callback_missing_data_test") - try: - with caplog.at_level(logging.WARNING, logger="threadmill.backends.redis"): - RedisBroker(backend)._fail_unreadable_task( - "missing-task-id", ImportError("gone") - ) - assert "has no task data" in caplog.text - finally: - backend.close() - class TestRedisTaskBackend: """Tests for the RedisTaskBackend update and lease functionality.""" diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index c067e84..03905ab 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -115,9 +115,6 @@ def _fail_unreadable_task(self, task_id: str, read_error: ImportError) -> None: data, lease_worker, lease_started_at, lease_token = self.backend.client.hmget( task_key, "data", *self.backend.LEASE_FIELDS ) - if data is None: - logger.warning("Claimed task %r has no task data; skipping", task_id) - return logger.error( "Task %r retry callback is gone from the code base; failing the task: %s", task_id, From 63302477dca75e817d5e9b1f39af967af0bd1b03 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 14:21:57 +0200 Subject: [PATCH 3/4] Trim the reaper helper docstring to its contract Private methods in this module carry a single summary line. The explanation of the payload rewrite belongs in the user docs, which already describe the terminal failure. --- threadmill/backends/redis.py | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index 03905ab..c3a2b8f 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -103,12 +103,7 @@ def _reap_running_queue(self, queue_name: str) -> None: self._fail_unreadable_task(task_id, read_error) def _fail_unreadable_task(self, task_id: str, read_error: ImportError) -> None: - """Fail a claimed task whose stored retry callback is gone from the code base. - - The unimportable callback is dropped from the payload, so the payload - stays readable and the failure is stored as a regular result. The - inspector can then requeue or drop the task. - """ + """Fail a claimed task whose stored retry callback is gone from the code base.""" task_key = self.backend.TASK_KEY.format( prefix=self.backend.key_prefix, task_id=task_id ) From 04c86ea53b7d82d9150e668d21f9ab0fbaab1d95 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 14:27:39 +0200 Subject: [PATCH 4/4] Inline the unreadable payload failure into the reaper loop The handler only had one caller and passing the ImportError as an argument read awkwardly. It now runs where the exception is caught. --- threadmill/backends/redis.py | 73 ++++++++++++++++++------------------ 1 file changed, 37 insertions(+), 36 deletions(-) diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index c3a2b8f..442c4db 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -100,43 +100,44 @@ def _reap_running_queue(self, queue_name: str) -> None: try: self._reap_task(task_id) except ImportError as read_error: - self._fail_unreadable_task(task_id, read_error) - - def _fail_unreadable_task(self, task_id: str, read_error: ImportError) -> None: - """Fail a claimed task whose stored retry callback is gone from the code base.""" - 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 - ) - logger.error( - "Task %r retry callback is gone from the code base; failing the task: %s", - task_id, - read_error, - ) - payload = json.loads(data) - payload["task"].pop("retry", None) - 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.") + 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.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.create_task_error(read_error), - ], - ) - ) + ) + 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: """Requeue or fail a claimed task."""