From 614dd68b7b44300697953774d97020058421b95c Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 10:16:11 +0200 Subject: [PATCH 1/8] Fix #75 -- Issue and require secure lease To avoid a worker form publishing results after their lease has expired, they receive a signed lease which they will need to return. If a new lease was issue the invalid lease will block publishing. --- README.md | 4 +- tests/backends/test_redis.py | 108 ++++++++++++++++++++++-- threadmill/backends/base.py | 47 ++++++++++- threadmill/backends/lua/acknowledge.lua | 14 ++- threadmill/backends/lua/acquire.lua | 7 +- threadmill/backends/redis.py | 48 +++++++---- 6 files changed, 198 insertions(+), 30 deletions(-) diff --git a/README.md b/README.md index 893ef6f..128e84f 100644 --- a/README.md +++ b/README.md @@ -137,7 +137,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. 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. +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 +is discarded. All keys for one backend alias share a Redis Cluster hash tag (`{alias}`), so every multi-key operation — including the cross-queue acquire — runs on a single diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index 1a3f1f9..f4a1f6e 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -13,6 +13,7 @@ import pytest from django.tasks import default_task_backend from django.tasks.base import TaskResultStatus +from django.tasks.exceptions import TaskResultDoesNotExist from django.utils import timezone from tests.testapp.tasks import ( @@ -105,6 +106,11 @@ def _measure_wait_deltas(calls: list[float]) -> list[float]: return [calls[index + 1] - calls[index] for index in range(len(calls) - 1)] +def _rotation_offset(sent_args: list[str]) -> str: + """Return the rotation offset from recorded acquire script arguments.""" + return sent_args[-2] + + def _now_ms() -> float: """Return the current time in milliseconds since the UNIX epoch.""" return timezone.now().timestamp() * 1000 @@ -335,7 +341,7 @@ def test_acquire__moves_to_running_set(self): backend.close() def test_acquire__stamps_lease_on_task_hash(self): - """Record the acquiring worker and lease start on the task hash.""" + """Record the acquiring worker, lease issue time, and token on the task hash.""" backend = RedisTaskBackend( "acquire_lease_test", { @@ -356,6 +362,7 @@ def test_acquire__stamps_lease_on_task_hash(self): assert acquired.last_attempted_at is not None assert acquired.started_at == acquired.last_attempted_at assert acquired.worker_ids == ["test-worker"] + assert acquired.lease_token is not None # Verify the lease is persisted, not only applied in memory. restored = backend.get_leased_task(task_result.id) @@ -363,6 +370,7 @@ def test_acquire__stamps_lease_on_task_hash(self): assert restored.worker_ids == acquired.worker_ids assert restored.started_at == acquired.started_at assert restored.last_attempted_at == acquired.last_attempted_at + assert restored.lease_token == acquired.lease_token finally: backend.close() @@ -478,6 +486,32 @@ def test_acknowledge__stores_the_leased_attempt(self): finally: backend.close() + def test_acknowledge__keeps_lease_token_out_of_the_result(self): + """Persist no lease token: it is attempt state, not result state.""" + backend = _make_backend("acknowledge_token_test") + try: + task_result = backend.enqueue(echo, args=[42]) + acquired = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="token-test" + ) + assert acquired is not None + assert acquired.lease_token is not None + + backend.acknowledge( + dataclasses.replace( + acquired, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + + result_key = backend.RESULT_KEY.format( + prefix=backend.key_prefix, result_id=task_result.id + ) + assert "lease_token" not in json.loads(backend.client.get(result_key)) + finally: + backend.close() + async def test_running_reaper__fails_expired_tasks(self): """Running reaper creates FAILED results for tasks with expired lease.""" backend = RedisTaskBackend( @@ -565,6 +599,63 @@ def test_stale_acknowledge__is_noop(self): finally: backend.close() + def test_stale_acknowledge__keeps_retry_attempt_after_requeue(self): + """Discard a late acknowledgement while a retry attempt holds the lease.""" + backend = _make_backend( + "stale_ack_retry_test", lease_ttl=datetime.timedelta(seconds=1) + ) + try: + task_result = backend.enqueue(echo_retry_on_lease_expiry, args=[42]) + expired = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="expired-worker" + ) + assert expired is not None + _expire_lease(backend, task_result.id) + RedisBroker(backend).main() + + # Make the scheduled retry due and lease it to the next attempt. + deferred_key = backend.DEFERRED_KEY.format( + prefix=backend.key_prefix, queue_name="default" + ) + backend.client.zadd(deferred_key, {task_result.id: 0}) + RedisBroker(backend).main() + retry = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="retry-worker" + ) + assert retry is not None + assert retry.id == task_result.id + + backend.acknowledge( + dataclasses.replace( + expired, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + + with pytest.raises(TaskResultDoesNotExist): + backend.get_result(task_result.id) + running_key = backend._segment_key(TaskResultStatus.RUNNING, "default") + task_key = backend.TASK_KEY.format( + prefix=backend.key_prefix, task_id=task_result.id + ) + assert backend.client.zscore(running_key, task_result.id) is not None + assert backend.client.exists(task_key) + + backend.acknowledge( + dataclasses.replace( + retry, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + assert backend.get_result(task_result.id).worker_ids == [ + "expired-worker", + "retry-worker", + ] + finally: + backend.close() + async def test_queue_stats__empty_backend(self): """queue_stats returns zero counts for an empty backend.""" backend = RedisTaskBackend( @@ -910,6 +1001,7 @@ def test_peek__running_tasks(self): assert results[0].worker_ids == ["peek-test"] assert results[0].last_attempted_at == acquired.last_attempted_at assert results[0].started_at == acquired.started_at + assert results[0].lease_token == acquired.lease_token def test_peek__running_task_without_lease(self): """Peek RUNNING marks a task whose hash carries no lease as RUNNING.""" @@ -931,8 +1023,8 @@ def test_peek__running_task_without_lease(self): assert result.last_attempted_at is None assert result.started_at is None - def test_peek__running_task_with_unparseable_lease_start(self): - """Read a lease start that is not a timestamp as no start time.""" + def test_peek__running_task_with_unparseable_lease_issued_at(self): + """Read a lease issue time that is not a timestamp as no start time.""" task_result = default_task_backend.enqueue(echo, args=[1]) default_task_backend.acquire( timeout=datetime.timedelta(seconds=1), worker="malformed-test" @@ -940,7 +1032,7 @@ def test_peek__running_task_with_unparseable_lease_start(self): task_key = default_task_backend.TASK_KEY.format( prefix=default_task_backend.key_prefix, task_id=task_result.id ) - default_task_backend.client.hset(task_key, "lease_started_at", "not-a-time") + default_task_backend.client.hset(task_key, "lease_issued_at", "not-a-time") (result,) = default_task_backend.peek( queue_name="default", status=TaskResultStatus.RUNNING, count=10 @@ -1593,7 +1685,9 @@ def test_acquire__wraps_rotation_offset_through_every_queue(self): "default", "compute", ] * 2 - assert [sent_args[-1] for sent_args in recorder.sent_args] == [ + assert [ + _rotation_offset(sent_args) for sent_args in recorder.sent_args + ] == [ "2", "0", "1", @@ -1618,7 +1712,9 @@ def test_acquire__keep_rotation_on_idle_polls(self): timeout=datetime.timedelta(seconds=0.3), ) assert len(recorder.sent_args) > 1 - assert {sent_args[-1] for sent_args in recorder.sent_args} == {"5"} + assert { + _rotation_offset(sent_args) for sent_args in recorder.sent_args + } == {"5"} assert backend._rotation_offset == 5 finally: backend.close() diff --git a/threadmill/backends/base.py b/threadmill/backends/base.py index 959a85f..a789578 100644 --- a/threadmill/backends/base.py +++ b/threadmill/backends/base.py @@ -65,6 +65,34 @@ def __reduce__(self): return (reconstructor, (kwargs,)) +@dataclasses.dataclass(frozen=True, slots=True, kw_only=True) +class ThreadmillTaskResult(TaskResult): + """Task result with threadmill-specific fields. + + A backend that leases tasks returns one from `acquire`. Its `lease_token` + is stamped beside the stored payload, and a retry attempt stamps a fresh + one, so `acknowledge` can discard the late result of an expired attempt. + The token is attempt state: it is never serialized, so it reaches neither + a stored payload nor a published result. + """ + + lease_token: str | None = None + + @classmethod + def from_result( + cls, task_result: TaskResult, *, lease_token: str | None + ) -> ThreadmillTaskResult: + """Return the task result as a leased attempt.""" + return cls( + **{ + field.name: getattr(task_result, field.name) + for field in dataclasses.fields(TaskResult) + if field.init + }, + lease_token=lease_token, + ) + + @dataclasses.dataclass(kw_only=True, slots=True) class QueueCounts: """Point-in-time cardinality of each queue segment.""" @@ -152,13 +180,23 @@ def _parse_datetime(value: object) -> object: class TaskResultEncoder(DjangoJSONEncoder): - """JSON encoder for TaskResult and TaskError objects.""" + """JSON encoder for TaskResult and TaskError objects. + + Only the fields of the base types are written, so threadmill-specific + fields such as the lease token stay out of stored data and published + results. + """ def default(self, o): - if isinstance(o, (TaskResult, TaskError)): + if isinstance(o, TaskResult): return { field.name: getattr(o, field.name) - for field in dataclasses.fields(type(o)) + for field in dataclasses.fields(TaskResult) + } + if isinstance(o, TaskError): + return { + field.name: getattr(o, field.name) + for field in dataclasses.fields(TaskError) } if isinstance(o, RetryTask): data = { @@ -241,6 +279,9 @@ def acquire( """ Return and lock the next task to be processed without removing it from the queue. + A backend that leases tasks returns a `ThreadmillTaskResult` carrying + the lease token of the attempt, so `acknowledge` can prove the lease. + Args: queue_names: The names of the queues to acquire tasks from. timeout: The maximum time to wait for a task. If None, wait indefinitely. diff --git a/threadmill/backends/lua/acknowledge.lua b/threadmill/backends/lua/acknowledge.lua index 896f02f..f8ff19a 100644 --- a/threadmill/backends/lua/acknowledge.lua +++ b/threadmill/backends/lua/acknowledge.lua @@ -18,7 +18,19 @@ -- ARGV[5] -- status (SUCCESSFUL or FAILED) -- ARGV[6] -- telemetry pub/sub channel -- ARGV[7] -- queue name --- Returns: 1 on success, 0 if task was not in the running set +-- ARGV[8] -- lease token of the acknowledging attempt (empty when the attempt +-- holds no lease) +-- Returns: 1 on success, 0 when the attempt no longer holds the lease or the +-- task is not in the running set + +-- Fence off a late acknowledgement: once a lease expires the reaper requeues +-- the task, and the retry attempt stamps a fresh lease token. Only the attempt +-- holding the stored token may remove the lease, publish its result, and drop +-- the task data. A hash without a token predates leases and is acknowledged. +local lease_token = redis.call('HGET', KEYS[3], 'lease_token') +if lease_token and lease_token ~= ARGV[8] then + return 0 -- Leased to a retry attempt, skip +end local removed = redis.call('ZREM', KEYS[1], ARGV[1]) if removed == 0 then diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 5352e53..9a248c3 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -3,19 +3,20 @@ -- a backlogged queue cannot starve its neighbours. -- -- The payload comes back as it was enqueued. Apply the lease stamped beside it --- (lease_worker, lease_started_at) to report the task as RUNNING. +-- (lease_worker, lease_issued_at, lease_token) to report the task as RUNNING. -- -- KEYS[1..N] -- interleaved running keys and queue keys, one pair per queue: -- KEYS[1] = running set, KEYS[2] = queue set, KEYS[3] = running, -- KEYS[4] = queue, etc. -- ARGV[1] -- current time in milliseconds --- ARGV[2] -- current time as ISO-8601 string, stamped as lease_started_at +-- ARGV[2] -- current time as ISO-8601 string, stamped as lease_issued_at -- ARGV[3] -- task key prefix (e.g. "threadmill:task:") -- ARGV[4] -- number of queue pairs (N/2) -- ARGV[5] -- worker name, stamped as lease_worker -- ARGV[6] -- lease TTL in milliseconds -- ARGV[7] -- start_index; 0-based index of the queue pair to scan first, so -- start_index 0 is the pair at KEYS[1] and KEYS[2] +-- ARGV[8] -- random lease token that proves this attempt at acknowledgement -- Returns: the stored payload, or nil when no queue yields a task. An entry -- whose hash holds no data returns nil too, leaving its queue unleased. @@ -32,7 +33,7 @@ for offset = 0, num_queues - 1 do if data then local deadline = tonumber(ARGV[1]) + lease_ttl_ms redis.call('ZADD', KEYS[queue_index * 2 - 1], deadline, task_id) - redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'lease_started_at', ARGV[2]) + redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'lease_issued_at', ARGV[2], 'lease_token', ARGV[8]) return data end end diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index 1becfcc..b8db69c 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -27,6 +27,7 @@ TelemetryDirection, TelemetryEvent, ThreadmillTaskBackend, + ThreadmillTaskResult, ) from threadmill.exceptions import AcknowledgementTimeout @@ -40,8 +41,8 @@ def _load_lua(name: str) -> str: return (_LUA_DIR / f"{name}.lua").read_text() -def _parse_lease_started_at(value: bytes | None) -> datetime.datetime | None: - """Return the lease start stamped beside a task, or None when the hash holds no timestamp.""" +def _parse_lease_issued_at(value: bytes | None) -> datetime.datetime | None: + """Return the lease issue time stamped beside a task, or None when the hash holds no timestamp.""" try: return datetime.datetime.fromisoformat(value.decode()) except AttributeError, ValueError: @@ -166,7 +167,7 @@ class RedisTaskBackend(ThreadmillTaskBackend): SEGMENT_KEY = "{prefix}:{queue_name}:{status}" DEFERRED_KEY = "{prefix}:{queue_name}:deferred" - LEASE_FIELDS = ("lease_worker", "lease_started_at") + LEASE_FIELDS = ("lease_worker", "lease_issued_at", "lease_token") TELEMETRY_CHANNEL = "{prefix}:telemetry" @@ -182,10 +183,10 @@ def _segment_key(self, status: TaskResultStatus, queue_name: str) -> str: status=status.value.lower(), ) - def get_leased_task(self, task_id: str) -> TaskResult | None: + def get_leased_task(self, task_id: str) -> ThreadmillTaskResult | None: """Return a running task as its lease holds it, or None when its hash is gone.""" task_key = self.TASK_KEY.format(prefix=self.key_prefix, task_id=task_id) - data, lease_worker, lease_started_at = self.client.hmget( + data, lease_worker, lease_issued_at, lease_token = self.client.hmget( task_key, "data", *self.LEASE_FIELDS ) if data is None: @@ -193,7 +194,8 @@ def get_leased_task(self, task_id: str) -> TaskResult | None: return self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - lease_started_at=_parse_lease_started_at(lease_started_at), + lease_issued_at=_parse_lease_issued_at(lease_issued_at), + lease_token=lease_token.decode() if lease_token is not None else None, ) def __init__(self, alias: str, params: dict) -> None: @@ -323,6 +325,7 @@ def acquire( now = timezone.now() now_ms = now.timestamp() * 1000 now_iso = now.isoformat() + lease_token = str(uuid.uuid7()) if data := self._acquire_script( keys=keys, @@ -334,6 +337,7 @@ def acquire( worker, str(int(self.lease_ttl.total_seconds() * 1000)), str(self._rotation_offset), + lease_token, ], ): self._miss_count = 0 @@ -341,7 +345,8 @@ def acquire( return self._apply_lease( self.deserialize_task_result(data.decode()), worker=worker, - lease_started_at=now, + lease_issued_at=now, + lease_token=lease_token, ) try: @@ -366,22 +371,24 @@ def _apply_lease( task_result: TaskResult, *, worker: str | None, - lease_started_at: datetime.datetime | None, - ) -> TaskResult: + lease_issued_at: datetime.datetime | None, + lease_token: str | None, + ) -> ThreadmillTaskResult: """Return a stored task result as a running attempt. `worker=None` means the lease records no worker; an empty string still counts as an attempt, and a task that already records a start keeps it. """ + leased = ThreadmillTaskResult.from_result(task_result, lease_token=lease_token) return dataclasses.replace( - task_result, + leased, status=TaskResultStatus.RUNNING, - started_at=task_result.started_at or lease_started_at, - last_attempted_at=lease_started_at or task_result.last_attempted_at, + started_at=leased.started_at or lease_issued_at, + last_attempted_at=lease_issued_at or leased.last_attempted_at, worker_ids=( - [*task_result.worker_ids, worker] + [*leased.worker_ids, worker] if worker is not None - else task_result.worker_ids + else leased.worker_ids ), ) @@ -402,6 +409,11 @@ def acknowledge(self, task_result: TaskResult) -> None: ) finished_at = task_result.finished_at or timezone.now() finish_score = finished_at.timestamp() * 1000 + lease_token = ( + task_result.lease_token + if isinstance(task_result, ThreadmillTaskResult) + else None + ) self._acknowledge_script( keys=[ @@ -419,6 +431,7 @@ def acknowledge(self, task_result: TaskResult) -> None: task_result.status.name, self.telemetry_channel, task_result.task.queue_name, + lease_token or "", ], ) @@ -518,13 +531,16 @@ def _peek_tasks( if not stored: continue if leased: - data, lease_worker, lease_started_at = stored + data, lease_worker, lease_issued_at, lease_token = stored if not data: continue yield self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - lease_started_at=_parse_lease_started_at(lease_started_at), + lease_issued_at=_parse_lease_issued_at(lease_issued_at), + lease_token=( + lease_token.decode() if lease_token is not None else None + ), ) else: yield self.deserialize_task_result(stored.decode()) From 80b2ed6e14d9143ffd2b13ad1fe77157394e71d1 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 10:34:56 +0200 Subject: [PATCH 2/8] Drop test, docstring, and comment changes for now --- tests/backends/test_redis.py | 108 ++---------------------- threadmill/backends/base.py | 20 +---- threadmill/backends/lua/acknowledge.lua | 11 +-- threadmill/backends/lua/acquire.lua | 5 +- threadmill/backends/redis.py | 2 +- 5 files changed, 12 insertions(+), 134 deletions(-) diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index f4a1f6e..1a3f1f9 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -13,7 +13,6 @@ import pytest from django.tasks import default_task_backend from django.tasks.base import TaskResultStatus -from django.tasks.exceptions import TaskResultDoesNotExist from django.utils import timezone from tests.testapp.tasks import ( @@ -106,11 +105,6 @@ def _measure_wait_deltas(calls: list[float]) -> list[float]: return [calls[index + 1] - calls[index] for index in range(len(calls) - 1)] -def _rotation_offset(sent_args: list[str]) -> str: - """Return the rotation offset from recorded acquire script arguments.""" - return sent_args[-2] - - def _now_ms() -> float: """Return the current time in milliseconds since the UNIX epoch.""" return timezone.now().timestamp() * 1000 @@ -341,7 +335,7 @@ def test_acquire__moves_to_running_set(self): backend.close() def test_acquire__stamps_lease_on_task_hash(self): - """Record the acquiring worker, lease issue time, and token on the task hash.""" + """Record the acquiring worker and lease start on the task hash.""" backend = RedisTaskBackend( "acquire_lease_test", { @@ -362,7 +356,6 @@ def test_acquire__stamps_lease_on_task_hash(self): assert acquired.last_attempted_at is not None assert acquired.started_at == acquired.last_attempted_at assert acquired.worker_ids == ["test-worker"] - assert acquired.lease_token is not None # Verify the lease is persisted, not only applied in memory. restored = backend.get_leased_task(task_result.id) @@ -370,7 +363,6 @@ def test_acquire__stamps_lease_on_task_hash(self): assert restored.worker_ids == acquired.worker_ids assert restored.started_at == acquired.started_at assert restored.last_attempted_at == acquired.last_attempted_at - assert restored.lease_token == acquired.lease_token finally: backend.close() @@ -486,32 +478,6 @@ def test_acknowledge__stores_the_leased_attempt(self): finally: backend.close() - def test_acknowledge__keeps_lease_token_out_of_the_result(self): - """Persist no lease token: it is attempt state, not result state.""" - backend = _make_backend("acknowledge_token_test") - try: - task_result = backend.enqueue(echo, args=[42]) - acquired = backend.acquire( - timeout=datetime.timedelta(seconds=1), worker="token-test" - ) - assert acquired is not None - assert acquired.lease_token is not None - - backend.acknowledge( - dataclasses.replace( - acquired, - status=TaskResultStatus.SUCCESSFUL, - finished_at=timezone.now(), - ) - ) - - result_key = backend.RESULT_KEY.format( - prefix=backend.key_prefix, result_id=task_result.id - ) - assert "lease_token" not in json.loads(backend.client.get(result_key)) - finally: - backend.close() - async def test_running_reaper__fails_expired_tasks(self): """Running reaper creates FAILED results for tasks with expired lease.""" backend = RedisTaskBackend( @@ -599,63 +565,6 @@ def test_stale_acknowledge__is_noop(self): finally: backend.close() - def test_stale_acknowledge__keeps_retry_attempt_after_requeue(self): - """Discard a late acknowledgement while a retry attempt holds the lease.""" - backend = _make_backend( - "stale_ack_retry_test", lease_ttl=datetime.timedelta(seconds=1) - ) - try: - task_result = backend.enqueue(echo_retry_on_lease_expiry, args=[42]) - expired = backend.acquire( - timeout=datetime.timedelta(seconds=1), worker="expired-worker" - ) - assert expired is not None - _expire_lease(backend, task_result.id) - RedisBroker(backend).main() - - # Make the scheduled retry due and lease it to the next attempt. - deferred_key = backend.DEFERRED_KEY.format( - prefix=backend.key_prefix, queue_name="default" - ) - backend.client.zadd(deferred_key, {task_result.id: 0}) - RedisBroker(backend).main() - retry = backend.acquire( - timeout=datetime.timedelta(seconds=1), worker="retry-worker" - ) - assert retry is not None - assert retry.id == task_result.id - - backend.acknowledge( - dataclasses.replace( - expired, - status=TaskResultStatus.SUCCESSFUL, - finished_at=timezone.now(), - ) - ) - - with pytest.raises(TaskResultDoesNotExist): - backend.get_result(task_result.id) - running_key = backend._segment_key(TaskResultStatus.RUNNING, "default") - task_key = backend.TASK_KEY.format( - prefix=backend.key_prefix, task_id=task_result.id - ) - assert backend.client.zscore(running_key, task_result.id) is not None - assert backend.client.exists(task_key) - - backend.acknowledge( - dataclasses.replace( - retry, - status=TaskResultStatus.SUCCESSFUL, - finished_at=timezone.now(), - ) - ) - assert backend.get_result(task_result.id).worker_ids == [ - "expired-worker", - "retry-worker", - ] - finally: - backend.close() - async def test_queue_stats__empty_backend(self): """queue_stats returns zero counts for an empty backend.""" backend = RedisTaskBackend( @@ -1001,7 +910,6 @@ def test_peek__running_tasks(self): assert results[0].worker_ids == ["peek-test"] assert results[0].last_attempted_at == acquired.last_attempted_at assert results[0].started_at == acquired.started_at - assert results[0].lease_token == acquired.lease_token def test_peek__running_task_without_lease(self): """Peek RUNNING marks a task whose hash carries no lease as RUNNING.""" @@ -1023,8 +931,8 @@ def test_peek__running_task_without_lease(self): assert result.last_attempted_at is None assert result.started_at is None - def test_peek__running_task_with_unparseable_lease_issued_at(self): - """Read a lease issue time that is not a timestamp as no start time.""" + def test_peek__running_task_with_unparseable_lease_start(self): + """Read a lease start that is not a timestamp as no start time.""" task_result = default_task_backend.enqueue(echo, args=[1]) default_task_backend.acquire( timeout=datetime.timedelta(seconds=1), worker="malformed-test" @@ -1032,7 +940,7 @@ def test_peek__running_task_with_unparseable_lease_issued_at(self): task_key = default_task_backend.TASK_KEY.format( prefix=default_task_backend.key_prefix, task_id=task_result.id ) - default_task_backend.client.hset(task_key, "lease_issued_at", "not-a-time") + default_task_backend.client.hset(task_key, "lease_started_at", "not-a-time") (result,) = default_task_backend.peek( queue_name="default", status=TaskResultStatus.RUNNING, count=10 @@ -1685,9 +1593,7 @@ def test_acquire__wraps_rotation_offset_through_every_queue(self): "default", "compute", ] * 2 - assert [ - _rotation_offset(sent_args) for sent_args in recorder.sent_args - ] == [ + assert [sent_args[-1] for sent_args in recorder.sent_args] == [ "2", "0", "1", @@ -1712,9 +1618,7 @@ def test_acquire__keep_rotation_on_idle_polls(self): timeout=datetime.timedelta(seconds=0.3), ) assert len(recorder.sent_args) > 1 - assert { - _rotation_offset(sent_args) for sent_args in recorder.sent_args - } == {"5"} + assert {sent_args[-1] for sent_args in recorder.sent_args} == {"5"} assert backend._rotation_offset == 5 finally: backend.close() diff --git a/threadmill/backends/base.py b/threadmill/backends/base.py index a789578..8d48307 100644 --- a/threadmill/backends/base.py +++ b/threadmill/backends/base.py @@ -67,22 +67,12 @@ def __reduce__(self): @dataclasses.dataclass(frozen=True, slots=True, kw_only=True) class ThreadmillTaskResult(TaskResult): - """Task result with threadmill-specific fields. - - A backend that leases tasks returns one from `acquire`. Its `lease_token` - is stamped beside the stored payload, and a retry attempt stamps a fresh - one, so `acknowledge` can discard the late result of an expired attempt. - The token is attempt state: it is never serialized, so it reaches neither - a stored payload nor a published result. - """ - lease_token: str | None = None @classmethod def from_result( cls, task_result: TaskResult, *, lease_token: str | None ) -> ThreadmillTaskResult: - """Return the task result as a leased attempt.""" return cls( **{ field.name: getattr(task_result, field.name) @@ -180,12 +170,7 @@ def _parse_datetime(value: object) -> object: class TaskResultEncoder(DjangoJSONEncoder): - """JSON encoder for TaskResult and TaskError objects. - - Only the fields of the base types are written, so threadmill-specific - fields such as the lease token stay out of stored data and published - results. - """ + """JSON encoder for TaskResult and TaskError objects.""" def default(self, o): if isinstance(o, TaskResult): @@ -279,9 +264,6 @@ def acquire( """ Return and lock the next task to be processed without removing it from the queue. - A backend that leases tasks returns a `ThreadmillTaskResult` carrying - the lease token of the attempt, so `acknowledge` can prove the lease. - Args: queue_names: The names of the queues to acquire tasks from. timeout: The maximum time to wait for a task. If None, wait indefinitely. diff --git a/threadmill/backends/lua/acknowledge.lua b/threadmill/backends/lua/acknowledge.lua index f8ff19a..f112600 100644 --- a/threadmill/backends/lua/acknowledge.lua +++ b/threadmill/backends/lua/acknowledge.lua @@ -18,18 +18,11 @@ -- ARGV[5] -- status (SUCCESSFUL or FAILED) -- ARGV[6] -- telemetry pub/sub channel -- ARGV[7] -- queue name --- ARGV[8] -- lease token of the acknowledging attempt (empty when the attempt --- holds no lease) --- Returns: 1 on success, 0 when the attempt no longer holds the lease or the --- task is not in the running set +-- Returns: 1 on success, 0 if task was not in the running set --- Fence off a late acknowledgement: once a lease expires the reaper requeues --- the task, and the retry attempt stamps a fresh lease token. Only the attempt --- holding the stored token may remove the lease, publish its result, and drop --- the task data. A hash without a token predates leases and is acknowledged. local lease_token = redis.call('HGET', KEYS[3], 'lease_token') if lease_token and lease_token ~= ARGV[8] then - return 0 -- Leased to a retry attempt, skip + return 0 end local removed = redis.call('ZREM', KEYS[1], ARGV[1]) diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 9a248c3..7a7eb95 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -3,20 +3,19 @@ -- a backlogged queue cannot starve its neighbours. -- -- The payload comes back as it was enqueued. Apply the lease stamped beside it --- (lease_worker, lease_issued_at, lease_token) to report the task as RUNNING. +-- (lease_worker, lease_started_at) to report the task as RUNNING. -- -- KEYS[1..N] -- interleaved running keys and queue keys, one pair per queue: -- KEYS[1] = running set, KEYS[2] = queue set, KEYS[3] = running, -- KEYS[4] = queue, etc. -- ARGV[1] -- current time in milliseconds --- ARGV[2] -- current time as ISO-8601 string, stamped as lease_issued_at +-- ARGV[2] -- current time as ISO-8601 string, stamped as lease_started_at -- ARGV[3] -- task key prefix (e.g. "threadmill:task:") -- ARGV[4] -- number of queue pairs (N/2) -- ARGV[5] -- worker name, stamped as lease_worker -- ARGV[6] -- lease TTL in milliseconds -- ARGV[7] -- start_index; 0-based index of the queue pair to scan first, so -- start_index 0 is the pair at KEYS[1] and KEYS[2] --- ARGV[8] -- random lease token that proves this attempt at acknowledgement -- Returns: the stored payload, or nil when no queue yields a task. An entry -- whose hash holds no data returns nil too, leaving its queue unleased. diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index b8db69c..6a4b474 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -42,7 +42,7 @@ def _load_lua(name: str) -> str: def _parse_lease_issued_at(value: bytes | None) -> datetime.datetime | None: - """Return the lease issue time stamped beside a task, or None when the hash holds no timestamp.""" + """Return the lease start stamped beside a task, or None when the hash holds no timestamp.""" try: return datetime.datetime.fromisoformat(value.decode()) except AttributeError, ValueError: From fd336cc11d2ea4c4ab59001459c18ec9732c88d4 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 11:28:37 +0200 Subject: [PATCH 3/8] Name the lease timestamp last_attempted_at --- threadmill/backends/lua/acquire.lua | 2 +- threadmill/backends/redis.py | 20 ++++++++++---------- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 7a7eb95..1c187bd 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -32,7 +32,7 @@ for offset = 0, num_queues - 1 do if data then local deadline = tonumber(ARGV[1]) + lease_ttl_ms redis.call('ZADD', KEYS[queue_index * 2 - 1], deadline, task_id) - redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'lease_issued_at', ARGV[2], 'lease_token', ARGV[8]) + redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'last_attempted_at', ARGV[2], 'lease_token', ARGV[8]) return data end end diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index 6a4b474..8d4595e 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -41,7 +41,7 @@ def _load_lua(name: str) -> str: return (_LUA_DIR / f"{name}.lua").read_text() -def _parse_lease_issued_at(value: bytes | None) -> datetime.datetime | None: +def _parse_last_attempted_at(value: bytes | None) -> datetime.datetime | None: """Return the lease start stamped beside a task, or None when the hash holds no timestamp.""" try: return datetime.datetime.fromisoformat(value.decode()) @@ -167,7 +167,7 @@ class RedisTaskBackend(ThreadmillTaskBackend): SEGMENT_KEY = "{prefix}:{queue_name}:{status}" DEFERRED_KEY = "{prefix}:{queue_name}:deferred" - LEASE_FIELDS = ("lease_worker", "lease_issued_at", "lease_token") + LEASE_FIELDS = ("lease_worker", "last_attempted_at", "lease_token") TELEMETRY_CHANNEL = "{prefix}:telemetry" @@ -186,7 +186,7 @@ def _segment_key(self, status: TaskResultStatus, queue_name: str) -> str: def get_leased_task(self, task_id: str) -> ThreadmillTaskResult | None: """Return a running task as its lease holds it, or None when its hash is gone.""" task_key = self.TASK_KEY.format(prefix=self.key_prefix, task_id=task_id) - data, lease_worker, lease_issued_at, lease_token = self.client.hmget( + data, lease_worker, last_attempted_at, lease_token = self.client.hmget( task_key, "data", *self.LEASE_FIELDS ) if data is None: @@ -194,7 +194,7 @@ def get_leased_task(self, task_id: str) -> ThreadmillTaskResult | None: return self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - lease_issued_at=_parse_lease_issued_at(lease_issued_at), + last_attempted_at=_parse_last_attempted_at(last_attempted_at), lease_token=lease_token.decode() if lease_token is not None else None, ) @@ -345,7 +345,7 @@ def acquire( return self._apply_lease( self.deserialize_task_result(data.decode()), worker=worker, - lease_issued_at=now, + last_attempted_at=now, lease_token=lease_token, ) @@ -371,7 +371,7 @@ def _apply_lease( task_result: TaskResult, *, worker: str | None, - lease_issued_at: datetime.datetime | None, + last_attempted_at: datetime.datetime | None, lease_token: str | None, ) -> ThreadmillTaskResult: """Return a stored task result as a running attempt. @@ -383,8 +383,8 @@ def _apply_lease( return dataclasses.replace( leased, status=TaskResultStatus.RUNNING, - started_at=leased.started_at or lease_issued_at, - last_attempted_at=lease_issued_at or leased.last_attempted_at, + started_at=leased.started_at or last_attempted_at, + last_attempted_at=last_attempted_at or leased.last_attempted_at, worker_ids=( [*leased.worker_ids, worker] if worker is not None @@ -531,13 +531,13 @@ def _peek_tasks( if not stored: continue if leased: - data, lease_worker, lease_issued_at, lease_token = stored + data, lease_worker, last_attempted_at, lease_token = stored if not data: continue yield self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - lease_issued_at=_parse_lease_issued_at(lease_issued_at), + last_attempted_at=_parse_last_attempted_at(last_attempted_at), lease_token=( lease_token.decode() if lease_token is not None else None ), From 2fbe024fb0ea63fc146af6244fb875816e91a084 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 11:44:23 +0200 Subject: [PATCH 4/8] Mint the lease token in the acquire script --- threadmill/backends/lua/acquire.lua | 11 +++++++++-- threadmill/backends/redis.py | 7 +++---- 2 files changed, 12 insertions(+), 6 deletions(-) diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 1c187bd..7365a96 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -22,6 +22,9 @@ local num_queues = tonumber(ARGV[4]) local lease_ttl_ms = tonumber(ARGV[6]) local start_index = tonumber(ARGV[7]) +if not redis.REDIS_VERSION_NUM or redis.REDIS_VERSION_NUM < 0x070000 then + return redis.error_reply('ERR threadmill requires Redis 7.0 or later') +end for offset = 0, num_queues - 1 do local queue_index = (start_index + offset) % num_queues + 1 local result = redis.call('ZPOPMIN', KEYS[queue_index * 2]) @@ -31,9 +34,13 @@ for offset = 0, num_queues - 1 do local data = redis.call('HGET', task_key, 'data') if data then local deadline = tonumber(ARGV[1]) + lease_ttl_ms + local lease_token = string.format( + '%06x%06x%06x%06x', + math.random(0, 0xffffff), math.random(0, 0xffffff), + math.random(0, 0xffffff), math.random(0, 0xffffff)) redis.call('ZADD', KEYS[queue_index * 2 - 1], deadline, task_id) - redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'last_attempted_at', ARGV[2], 'lease_token', ARGV[8]) - return data + redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'last_attempted_at', ARGV[2], 'lease_token', lease_token) + return { data, lease_token } end end end diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index 8d4595e..afd7799 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -325,9 +325,8 @@ def acquire( now = timezone.now() now_ms = now.timestamp() * 1000 now_iso = now.isoformat() - lease_token = str(uuid.uuid7()) - if data := self._acquire_script( + if result := self._acquire_script( keys=keys, args=[ str(now_ms), @@ -337,16 +336,16 @@ def acquire( worker, str(int(self.lease_ttl.total_seconds() * 1000)), str(self._rotation_offset), - lease_token, ], ): + data, lease_token = result self._miss_count = 0 self._rotation_offset = (self._rotation_offset + 1) % len(queue_names) return self._apply_lease( self.deserialize_task_result(data.decode()), worker=worker, last_attempted_at=now, - lease_token=lease_token, + lease_token=lease_token.decode(), ) try: From fe907172785f03bf92ca6b401c48c1195e806963 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 12:05:50 +0200 Subject: [PATCH 5/8] Revert the lease timestamp rename --- threadmill/backends/lua/acquire.lua | 2 +- threadmill/backends/redis.py | 20 ++++++++++---------- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 7365a96..9396c44 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -39,7 +39,7 @@ for offset = 0, num_queues - 1 do math.random(0, 0xffffff), math.random(0, 0xffffff), math.random(0, 0xffffff), math.random(0, 0xffffff)) redis.call('ZADD', KEYS[queue_index * 2 - 1], deadline, task_id) - redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'last_attempted_at', ARGV[2], 'lease_token', lease_token) + redis.call('HSET', task_key, 'lease_worker', ARGV[5], 'lease_started_at', ARGV[2], 'lease_token', lease_token) return { data, lease_token } end end diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index afd7799..e1a4369 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -41,7 +41,7 @@ def _load_lua(name: str) -> str: return (_LUA_DIR / f"{name}.lua").read_text() -def _parse_last_attempted_at(value: bytes | None) -> datetime.datetime | None: +def _parse_lease_started_at(value: bytes | None) -> datetime.datetime | None: """Return the lease start stamped beside a task, or None when the hash holds no timestamp.""" try: return datetime.datetime.fromisoformat(value.decode()) @@ -167,7 +167,7 @@ class RedisTaskBackend(ThreadmillTaskBackend): SEGMENT_KEY = "{prefix}:{queue_name}:{status}" DEFERRED_KEY = "{prefix}:{queue_name}:deferred" - LEASE_FIELDS = ("lease_worker", "last_attempted_at", "lease_token") + LEASE_FIELDS = ("lease_worker", "lease_started_at", "lease_token") TELEMETRY_CHANNEL = "{prefix}:telemetry" @@ -186,7 +186,7 @@ def _segment_key(self, status: TaskResultStatus, queue_name: str) -> str: def get_leased_task(self, task_id: str) -> ThreadmillTaskResult | None: """Return a running task as its lease holds it, or None when its hash is gone.""" task_key = self.TASK_KEY.format(prefix=self.key_prefix, task_id=task_id) - data, lease_worker, last_attempted_at, lease_token = self.client.hmget( + data, lease_worker, lease_started_at, lease_token = self.client.hmget( task_key, "data", *self.LEASE_FIELDS ) if data is None: @@ -194,7 +194,7 @@ def get_leased_task(self, task_id: str) -> ThreadmillTaskResult | None: return self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - last_attempted_at=_parse_last_attempted_at(last_attempted_at), + lease_started_at=_parse_lease_started_at(lease_started_at), lease_token=lease_token.decode() if lease_token is not None else None, ) @@ -344,7 +344,7 @@ def acquire( return self._apply_lease( self.deserialize_task_result(data.decode()), worker=worker, - last_attempted_at=now, + lease_started_at=now, lease_token=lease_token.decode(), ) @@ -370,7 +370,7 @@ def _apply_lease( task_result: TaskResult, *, worker: str | None, - last_attempted_at: datetime.datetime | None, + lease_started_at: datetime.datetime | None, lease_token: str | None, ) -> ThreadmillTaskResult: """Return a stored task result as a running attempt. @@ -382,8 +382,8 @@ def _apply_lease( return dataclasses.replace( leased, status=TaskResultStatus.RUNNING, - started_at=leased.started_at or last_attempted_at, - last_attempted_at=last_attempted_at or leased.last_attempted_at, + started_at=leased.started_at or lease_started_at, + last_attempted_at=lease_started_at or leased.last_attempted_at, worker_ids=( [*leased.worker_ids, worker] if worker is not None @@ -530,13 +530,13 @@ def _peek_tasks( if not stored: continue if leased: - data, lease_worker, last_attempted_at, lease_token = stored + data, lease_worker, lease_started_at, lease_token = stored if not data: continue yield self._apply_lease( self.deserialize_task_result(data.decode()), worker=lease_worker.decode() if lease_worker else None, - last_attempted_at=_parse_last_attempted_at(last_attempted_at), + lease_started_at=_parse_lease_started_at(lease_started_at), lease_token=( lease_token.decode() if lease_token is not None else None ), From dc868ec18d3a70e2b8d6a2720227a3025ad2d98f Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 12:06:27 +0200 Subject: [PATCH 6/8] Override task_result in _apply_lease --- threadmill/backends/redis.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index e1a4369..f6a0630 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -378,16 +378,18 @@ def _apply_lease( `worker=None` means the lease records no worker; an empty string still counts as an attempt, and a task that already records a start keeps it. """ - leased = ThreadmillTaskResult.from_result(task_result, lease_token=lease_token) + task_result = ThreadmillTaskResult.from_result( + task_result, lease_token=lease_token + ) return dataclasses.replace( - leased, + task_result, status=TaskResultStatus.RUNNING, - started_at=leased.started_at or lease_started_at, - last_attempted_at=lease_started_at or leased.last_attempted_at, + started_at=task_result.started_at or lease_started_at, + last_attempted_at=lease_started_at or task_result.last_attempted_at, worker_ids=( - [*leased.worker_ids, worker] + [*task_result.worker_ids, worker] if worker is not None - else leased.worker_ids + else task_result.worker_ids ), ) From d46b26e1b77869bcb88afe5b4b8b8d60bddb8235 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 12:36:22 +0200 Subject: [PATCH 7/8] Add minimal lease token regression tests --- tests/backends/test_redis.py | 84 ++++++++++++++++++++++++++++++++++++ 1 file changed, 84 insertions(+) diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index 1a3f1f9..08a8284 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -13,6 +13,7 @@ import pytest from django.tasks import default_task_backend from django.tasks.base import TaskResultStatus +from django.tasks.exceptions import TaskResultDoesNotExist from django.utils import timezone from tests.testapp.tasks import ( @@ -356,6 +357,7 @@ def test_acquire__stamps_lease_on_task_hash(self): assert acquired.last_attempted_at is not None assert acquired.started_at == acquired.last_attempted_at assert acquired.worker_ids == ["test-worker"] + assert acquired.lease_token is not None # Verify the lease is persisted, not only applied in memory. restored = backend.get_leased_task(task_result.id) @@ -363,6 +365,7 @@ def test_acquire__stamps_lease_on_task_hash(self): assert restored.worker_ids == acquired.worker_ids assert restored.started_at == acquired.started_at assert restored.last_attempted_at == acquired.last_attempted_at + assert restored.lease_token == acquired.lease_token finally: backend.close() @@ -478,6 +481,31 @@ def test_acknowledge__stores_the_leased_attempt(self): finally: backend.close() + def test_acknowledge__keeps_lease_token_out_of_the_result(self): + """Persist no lease token: it is attempt state, not result state.""" + backend = _make_backend("acknowledge_token_test") + try: + task_result = backend.enqueue(echo, args=[42]) + acquired = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="token-test" + ) + assert acquired is not None + + backend.acknowledge( + replace( + acquired, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + + result_key = backend.RESULT_KEY.format( + prefix=backend.key_prefix, result_id=task_result.id + ) + assert "lease_token" not in json.loads(backend.client.get(result_key)) + finally: + backend.close() + async def test_running_reaper__fails_expired_tasks(self): """Running reaper creates FAILED results for tasks with expired lease.""" backend = RedisTaskBackend( @@ -565,6 +593,62 @@ def test_stale_acknowledge__is_noop(self): finally: backend.close() + def test_stale_acknowledge__keeps_retry_attempt_after_requeue(self): + """Discard a late acknowledgement while a retry attempt holds the lease.""" + backend = _make_backend( + "stale_ack_retry_test", lease_ttl=datetime.timedelta(seconds=1) + ) + try: + task_result = backend.enqueue(echo_retry_on_lease_expiry, args=[42]) + expired = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="expired-worker" + ) + assert expired is not None + _expire_lease(backend, task_result.id) + RedisBroker(backend).main() + + deferred_key = backend.DEFERRED_KEY.format( + prefix=backend.key_prefix, queue_name="default" + ) + backend.client.zadd(deferred_key, {task_result.id: 0}) + RedisBroker(backend).main() + retry = backend.acquire( + timeout=datetime.timedelta(seconds=1), worker="retry-worker" + ) + assert retry is not None + assert retry.id == task_result.id + + backend.acknowledge( + replace( + expired, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + + with pytest.raises(TaskResultDoesNotExist): + backend.get_result(task_result.id) + running_key = backend._segment_key(TaskResultStatus.RUNNING, "default") + task_key = backend.TASK_KEY.format( + prefix=backend.key_prefix, task_id=task_result.id + ) + assert backend.client.zscore(running_key, task_result.id) is not None + assert backend.client.exists(task_key) + + backend.acknowledge( + replace( + retry, + status=TaskResultStatus.SUCCESSFUL, + finished_at=timezone.now(), + ) + ) + assert backend.get_result(task_result.id).worker_ids == [ + "expired-worker", + "retry-worker", + ] + finally: + backend.close() + async def test_queue_stats__empty_backend(self): """queue_stats returns zero counts for an empty backend.""" backend = RedisTaskBackend( From 9dffd3c592cc6cc649b9c60ef14025aad8542025 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 12:43:02 +0200 Subject: [PATCH 8/8] Drop the Redis version guard from the acquire script --- threadmill/backends/lua/acquire.lua | 3 --- 1 file changed, 3 deletions(-) diff --git a/threadmill/backends/lua/acquire.lua b/threadmill/backends/lua/acquire.lua index 9396c44..fda2c0b 100644 --- a/threadmill/backends/lua/acquire.lua +++ b/threadmill/backends/lua/acquire.lua @@ -22,9 +22,6 @@ local num_queues = tonumber(ARGV[4]) local lease_ttl_ms = tonumber(ARGV[6]) local start_index = tonumber(ARGV[7]) -if not redis.REDIS_VERSION_NUM or redis.REDIS_VERSION_NUM < 0x070000 then - return redis.error_reply('ERR threadmill requires Redis 7.0 or later') -end for offset = 0, num_queues - 1 do local queue_index = (start_index + offset) % num_queues + 1 local result = redis.call('ZPOPMIN', KEYS[queue_index * 2])