Retry tasks whose processing lease expired - #76
Merged
Merged
Conversation
The reaper finalized expired tasks as FAILED in Lua, so the Python retry callback never saw them and the task hash was deleted, leaving manual inspector requeues as the only mitigation. reaper.lua now claims expired running entries by renewing their lease to an internal claim deadline and keeps the task data. RedisBroker deserializes each claim, appends the AcknowledgementTimeout error, and consults the task's retry callback: a returned delay requeues the task, otherwise it is finalized FAILED. acknowledge() and requeue() accept an optional lease deadline that is checked in Lua, so a worker acknowledgement, another broker, or the inspector cannot be overwritten or resurrected. requeue() is now atomic Lua instead of a Python pipeline, and retry_delay()/create_task_error() moved from WorkerThread to ThreadmillTaskBackend so the broker shares them. Lease expiry still presumes the worker died: lease_ttl must stay above the worst-case task runtime, so a retry may run concurrently with a slow original execution.
The guard made acknowledge()/requeue() conditional on the running entry still holding the reaper's claim, which required a new parameter, a conditional acknowledge.lua, and an atomic requeue.lua. That is more surface than the retry feature needs right now. The claim itself stays: reaper.lua still renews the lease so concurrent broker passes cannot claim the same task and a crashed broker's claim lapses for the next pass. Only the compare-and-set on the decision is gone, leaving the known races between a reap decision and a late worker acknowledgement, a stalled broker, or an inspector action. Issue #71 tracks making the decision conditional again.
The standalone "Retrying lease expiry" subsection repeated what the retry section already explains. Lease expiry now gets one sentence in the retry intro, and the lease_ttl note is compressed to its two facts: the error the callback sees and the concurrent-execution caveat.
List the lease-expiry exception next to HTTPError so the built-in backoff example covers it, and tighten the retry intro sentence it replaces.
The claim script derived "now" and the claim deadline from the broker's clock, which required an epoch-millisecond conversion helper on the Python side. Redis already exposes its clock to scripts, so reaper.lua now reads TIME and takes the claim TTL as a duration, dropping both arguments and the helper.
- inline the bytes decode and read task data through the backend instead of hand-building the key in the broker - drop the two docstrings that only repeat their method name - trim the README lease_ttl row and the duplicated lease-expiry sentence - cut the reaper tests that re-proved covered behavior, drop the fixtures their only consumers used, and assert the delay the retry callback returned instead of overwriting the deferred score with 0 - make the reaper error-path test corrupt a real payload Review issues filed separately: undecodable payloads cycling the reaper, per-task round trips in the reap pass, a late acknowledgement clobbering its retry attempt, and the sleep-based setup in the pre-existing reaper tests.
Bind the stored task data with := in the condition and keep the warn path in the else branch, so _reap_task has no early return.
A claimed task that cannot be decoded is skipped, not swallowed: catch only the payload errors deserialization can raise (a missing import target, a JSON or status ValueError, a shape TypeError) instead of every exception, and drop the noqa that silenced the lint for it.
Split the per-task handler into the two causes it can distinguish: a retry callback that no longer imports, and a stored payload that cannot be read. Each gets its own message, so a missing retry callback reads differently from corrupt data, and the reaper tests pin both, including that the rest of the batch is still finalized.
Cut the header to the facts a reader cannot infer from the call site: what the script claims, that the renewal keeps other passes out, that a lapsed claim returns, that entries without task data are dropped, and that the clock is Redis's.
State what the script does rather than how it does it: claim expired tasks for the broker to decide, take undecided claims up again, and drop running entries without task data.
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.
A task whose worker died was marked FAILED by the reaper in Lua: the Python
retrycallback never saw it and the task data was deleted, so the only mitigation was a manual requeue in the inspector. Lease expiry now goes through the same retry protocol as any other failure.How it works
reaper.luaclaims expired running entries instead of finalizing them. The claim renews the lease so concurrent broker passes leave the task alone, and a claim that is never decided comes back on a later pass.RedisBroker._reap_taskdeserializes each claimed task, appends the existingAcknowledgementTimeouterror, and consults the task'sretrycallback. A returned delay requeues the task; otherwise it is finalized FAILED. ID, error history and attempt count survive, exactly like a worker-observed failure, and a retried task never lands in the failed segment.retry_delayandcreate_task_errormoved fromWorkerThreadtoThreadmillTaskBackendso the worker and the broker share one implementation.TIME) and takes the claim TTL as a duration, so the broker's own clock does not drive expiry.