Task queues: workers pull, so redeploys are safe
The engine never pushes into your process — it places tasks on a queue that stateless workers long-poll, so 'no worker is running my workflow right now' is a normal, recoverable state and a mid-workflow redeploy just leaves the next task waiting on the queue.
Stateless workers pulling tasks made redeploys safe: kill them all and the task just waits. But surviving the worker dying is different from surviving the WORK failing — the payment API returns a 503, the warehouse is briefly unreachable. We can't fail ORDER #1001 over a transient blip, so how does the engine re-attempt a failed activity without re-running the whole workflow?
Scene 05
Task queues: workers pull, so redeploys are safe
- Watch
- Try it
- Predict
- Capture
Here's the surprise that makes everything else possible: the engine never reaches into your process to run your code. It can't — your workers might be redeploying, scaled to zero, or on fire. So instead the engine does something humbler. It writes down the next unit of work — "run ShipPackage for ORDER #1001" — and drops it somewhere safe to wait. That somewhere is the task queue: a queue where the engine parks the next step of a workflow until something is ready to do it. Notice nothing has run yet, and nothing is lost. On the right sit your workers: plain, stateless processes that sit and long-poll the queue — "anything for me? anything for me?" — and when a task appears, one of them claims it, replays ORDER #1001's history to the current step, runs ShipPackage for real, and reports the result back as a new history event. Watch it happen. The load-bearing detail: when the worker is done it keeps NO state of its own. Everything that matters — the charge, the reservation — already lives in the engine's history.
Highlighted lines are the ones running in the diagram right now.
def dispatchNextStep(workflow_id):history = db.load(workflow_id) # source of truthnext_step = fold(history).next # ShipPackageif next_step is None:return # workflow done# place a task; the engine holds NO worker connectiontask_queue.put(Task(workflow_id, next_step))# task now waits here until SOME worker polls
def pollLoop(task_queue):while True:task = task_queue.poll() # long-poll; blocks if emptyhistory = engine.getHistory(task.workflow_id)state = replay(history) # hands back recorded resultsresult = run(task.step) # ShipPackage, for realengine.report(task, result) # becomes a history event# worker keeps nothing — loops back to poll
def poll():task = take_visible() # claim ittask.invisible_until = now() + visibility_timeoutreturn task # held, not yet ackeddef sweep(): # runs continuouslyfor t in claimed:if now() > t.invisible_until: # worker died before reportmake_visible(t) # a peer will grab it
Where this sits in Build a workflow engine (Temporal / Airflow / Cadence style)
Scene 05 of 13, in the Runtime act — Workers pull; retries with backoff; idempotency.. The engine never pushes work; stateless workers pull tasks from a queue, so a redeploy is just 'no worker for a moment' and the task simply waits to be picked up.
Up next. A worker dying is recoverable; a flaky downstream is the next problem. When ChargeCard's payment API returns a 503, the engine doesn't fail the order — it retries just that activity, waiting longer and longer between attempts so it doesn't hammer an already-sick service. That spacing is exponential backoff.
All 13 scenes in Build a workflow engine (Temporal / Airflow / Cadence style) · Every curriculum