Repository navigation
Lease tasks beside the payload instead of rewriting it - #69
Merged
Merged
Conversation
The acquire script decoded and re-encoded every task payload to stamp it RUNNING, so a fetch cost two full JSON passes on the Redis server — the single largest per-task cost in the backend, and one that grows with the payload. The task hash now keeps the payload exactly as it was enqueued and stores the lease in two fields next to it: the worker holding it and the time it started. acquire() applies that lease to the task result it returns, peek() applies it when reading RUNNING tasks, and the reaper folds it into the failure it writes, so lease visibility is unchanged for the inspector and the reaper. requeue() clears the lease when a task returns to the deferred set. Measured on a dedicated Redis 8.10 instance, one worker process, 20,000 tasks, EVALSHA service time per acquire: payload before after 447 B 10.90 us 6.23 us (-43%) 5.4 KB 36.03 us 12.38 us (-66%) End-to-end drains of 20,000 tasks drop from 30.7 to 21.2 us of server CPU per task (commandstats), and from 172 to 149 us per task of Redis process CPU, with the same 9.0 Redis commands per task. Refs #66
…er-task # Conflicts: # threadmill/backends/lua/reaper.lua
The fetch script popped a task, leased it, and returned the payload for the worker to deserialize. A payload this code cannot read — a document another version wrote, a task module or retry callback gone from the code base — then raised out of acquire(), which the consumer loop does not catch, so the consumer thread died while the id stayed leased where every reaper pass re-claimed and re-skipped it. The reaper's own copy of that read had the same hole, and enumerating exception types twice already missed ImportError. Reading a stored payload is now one rule instead of a list of failures: the acquire script returns the task id beside the payload, and both the fetch path and the reap path catch any read failure, log it with the task ID, and leave the task stored. Nothing is deleted: dropping the hash would turn a deploy or an environment skew into permanent, unaudited task loss, and the fetch path can run before any attempt exists. The reap path distinguishes a payload that is already gone, which a task finishing inside the claim window explains, from one that will not parse. peek() applies the same rule, skipping an entry it cannot read so the inspector still renders the rest of the page, and _decode_text accepts the str replies a REDIS_URL carrying decode_responses=True produces. Refs #66, #69
This was referenced Oct 1, 2026
Comment-only: the rule the guards apply lives in the log messages and the method docstrings, so the paragraphs explaining it at each call site are gone. No behaviour change; the token streams are identical once comments and docstrings are removed.
The guards caught Exception with a BLE001 suppression, which hid the rule behind a linter pragma. They now catch UNREADABLE_PAYLOAD_ERRORS: the failures a stored payload raises when this code cannot read it as a task, each one reachable from a payload and pinned by a test, including the AttributeError a document that parses but is not a TaskResult raises.
A Redis reply is bytes by default, and the URL query overrides any decode_responses kwarg, so the type cannot be pinned from the client construction. The decoder now states the one type it takes, and every payload read decodes before parsing: handing raw bytes to json.loads was the slowest of the three variants measured, because json.loads(bytes) runs its own encoding detection, while decode-then-parse measures the same as a str-configured client. Each decode sits inside the read's guard, so a reply of the wrong shape is reported with its task id and left as stored, rather than raising out of acquire() or the reaper on a path the caller does not catch.
decode_responses=False is the supported client, so a reply is bytes and the one-expression `_decode_text` wrapper only hid a call to `.decode()`. Every reply now decodes at its read site, which removes the last str tolerance in the module: worker_telemetry decodes its pubsub data directly. The unreadable-payload matrix keeps its nine cases; the three that wrapped a literal payload in a one-line function carry the literal instead, applied by one `_apply_poison` for both shapes.
…er-task # Conflicts: # tests/backends/test_redis.py # threadmill/backends/redis.py
EAFP on a field whose absent, cleared and unparseable shapes all mean the same thing here: no start time to apply. Asking permission first only restated the answer.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
The acquire script decoded and re-encoded every task payload to stamp it RUNNING. That is two full JSON passes on the Redis server for each fetch. It was the largest per-task cost in the backend, and it grew with the payload size.
lease_workerandlease_started_at.acquire(),peek()and the reaper apply that lease on read.requeue()clears the lease when a task returns to the deferred set.started_atfrom its lease. The old acquire script intended that and never did it.EVALSHAper acquire drops from 10.30 to 6.01 µs at 447 B, and from 34.14 to 12.11 µs at 5.4 KB. Server CPU per task in a 20,000-task drain drops from 23.01 to 19.70 µs.PUBLISH, the history eviction and the stored result stay as they are. Each was measured before that decision.Fixes #66