Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,13 @@ curl -sSL https://raw.githubusercontent.com/codingjoe/naming-things/refs/heads/m
- Consistency – We never lose data, even if someone unplugs the power or network.
- Utilization – We keep the CPU saturated with tasks, not with idle time or waiting for locks.

We require a persistent Redis without eviction.

Redis connections use the redis-py default `decode_responses=False`, so all
values read from Redis are bytes. We do not guard against misconfiguration.
We fail loudly instead. The same goes for Redis data altered mid-flight.
These are deliberate design decisions.

## Testing

We have unit tests, integration tests, and benchmarks. Avoid mocking if possible.
Expand Down
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,9 @@ uv run manage.py threadmill inspector

### Redis Backend Options

> [!IMPORTANT]
> Threadmill requires a persistent Redis without eviction.

The `RedisTaskBackend` accepts the following options under `OPTIONS` in your
`TASKS` configuration:

Expand Down
39 changes: 14 additions & 25 deletions tests/backends/test_redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ def _claim_expired(
f"{backend.key_prefix}:task:",
],
)
return [item.decode() if isinstance(item, bytes) else item for item in claimed]
return [item.decode() for item in claimed]


class TestRedisBroker:
Expand Down Expand Up @@ -265,43 +265,32 @@ def test_reap__removes_running_entry_without_task_data(self):
finally:
backend.close()

def test_reap_running_queue__logs_and_continues_after_task_error(self, caplog):
"""A failing reap decision is logged per cause and does not stop the batch."""
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."""
backend = _make_backend(
"reap_error_test", lease_ttl=datetime.timedelta(seconds=1)
"reap_gone_callback_test", lease_ttl=datetime.timedelta(seconds=1)
)
try:
task_ids = []
for _index in range(3):
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
_expire_lease(backend, task_result.id)
task_ids.append(task_result.id)

unreadable_id, gone_id, recovered_id = task_ids
backend.client.hset(
backend.TASK_KEY.format(
prefix=backend.key_prefix, task_id=unreadable_id
),
"data",
"{not json",
)
gone_key = backend.TASK_KEY.format(
prefix=backend.key_prefix, task_id=gone_id
)
payload = json.loads(backend.client.hget(gone_key, "data"))
payload["task"]["retry"] = "tests.testapp.tasks.gone_from_the_code_base"
backend.client.hset(gone_key, "data", json.dumps(payload))
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))

with caplog.at_level(logging.ERROR, logger="threadmill.backends.redis"):
RedisBroker(backend)._reap_running_queue("default")

assert "has an unreadable payload" in caplog.text
assert "gone from the code base" in caplog.text
assert backend.get_result(recovered_id).status == TaskResultStatus.FAILED
assert caplog.text.count("gone from the code base") == 2
finally:
backend.close()

Expand Down
15 changes: 4 additions & 11 deletions threadmill/backends/redis.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ def _reap_running_queue(self, queue_name: str) -> None:
],
)
for member in claimed_ids:
task_id = member.decode() if isinstance(member, bytes) else member
task_id = member.decode()
try:
self._reap_task(task_id)
except ImportError:
Expand All @@ -95,10 +95,6 @@ def _reap_running_queue(self, queue_name: str) -> None:
"skipping the reap",
task_id,
)
except TypeError, ValueError:
logger.exception(
"Task %r has an unreadable payload; skipping the reap", task_id
)

def _reap_task(self, task_id: str) -> None:
"""Requeue or fail a claimed task."""
Expand Down Expand Up @@ -460,7 +456,7 @@ def _peek(
) -> Generator[TaskResult]:
pipe = self.client.pipeline()
for member in self.client.zrange(zset_key, 0, count - 1):
member_id = member.decode() if isinstance(member, bytes) else member
member_id = member.decode()
data_key = data_key_template.format(
prefix=self.key_prefix,
task_id=member_id,
Expand All @@ -472,9 +468,7 @@ def _peek(
pipe.hget(data_key, field)
for data in pipe.execute():
if data:
yield self.deserialize_task_result(
data.decode() if isinstance(data, bytes) else data
)
yield self.deserialize_task_result(data.decode())

def get_result(self, result_id: str) -> TaskResult:
if data := self.client.get(
Expand Down Expand Up @@ -523,8 +517,7 @@ async def worker_telemetry(
try:
async for message in pubsub.listen():
if (data := message.get("data")) is not None:
payload = data.decode() if isinstance(data, bytes) else data
direction, _, queue_name = payload.partition(":")
direction, _, queue_name = data.decode().partition(":")
try:
event = TelemetryEvent(
direction=TelemetryDirection(direction),
Expand Down
Loading