How to run and shut down a worker pool¶
This guide shows you how to start a pool of workers against a
DurableEngine, shut it down cleanly under a process supervisor, extend a
claim for a long-running step, and cancel a task in flight.
Start the pool¶
from dura import DurableEngine, run_workers
engine = DurableEngine("engine.db")
run_workers(engine, handlers={"send_file": send_file}, worker_count=4)
run_workers starts worker_count threads, each running its own claim
loop against handlers, and blocks the calling thread. Call it from your
main thread: it installs SIGINT/SIGTERM handlers, which only works on
the main thread.
Let it shut down gracefully¶
Send SIGINT (Ctrl-C) or SIGTERM (what docker stop and Kubernetes send
by default). run_workers then:
- stops claiming new work,
- gives in-flight handlers up to
dura.workers.SHUTDOWN_GRACE_SECONDS(10 seconds) to finish, - closes every connection so the database is left without
-wal/-shmsidecar files.
A handler still running when the grace period expires is abandoned, not killed. That's safe: its run's lease simply expires and gets reclaimed on the next start, the same as a crash.
If your handlers can legitimately take longer than that to reach a safe stopping point, design them around checkpoints (see How to checkpoint steps) instead of relying on a longer grace period.
Extend a claim for a long step¶
A run's claim (lease) expires after claim_timeout_secs (120 seconds by
default when going through run_workers). If a handler is about to run a
step you know will take longer than that, extend it first so another
worker doesn't reclaim the run out from under you:
def handler(engine, task):
engine.extend_claim(run_id=task.run_id, by_secs=600)
do_the_slow_thing()
...
Cancel a task in flight¶
This marks the task and any of its non-terminal runs cancelled, and
clears any pending waits. A worker that later tries to settle a cancelled
run gets TaskCancelledError. process_task already handles that for
you and simply drops the run.
Drive the engine yourself instead¶
If you don't want threads, signal handling, or dura.workers' opinions at
all, write your own loop around the same primitives it uses:
while True:
task = engine.claim_task(worker_id="my-worker")
if task is None:
time.sleep(1)
continue
try:
result = handlers[task.name](engine, task)
except WorkflowSuspended:
continue
except Exception as exc:
engine.fail_run(run_id=task.run_id, reason={"type": type(exc).__name__, "message": str(exc)})
continue
engine.complete_run(run_id=task.run_id, result=result)
See dura.workers' reference for the exact
behavior process_task and run_worker add on top of this (settle-race
handling, logging, heartbeats) that you'd otherwise need to reimplement.