---
title: Worker jobs
description: How a processing job is claimed, run and retried, and how the number of workers grows with the backlog.
sidebar:
  label: Worker jobs
  order: 4
---

Mail adds one job per email to the processing queue, as described on [Mail to processing](/engineering/mail-to-processing). This page covers what happens next: which worker takes the job, what it reads and writes, what happens when it fails, and how the number of workers grows with the backlog. The listener, its retries and the sweeper are built. The child-process pool and scaling on backlog are not: each job runs in the listener's own process, and one worker process runs until it is stopped. What a job does, step by step, is on [Sync to labels](/engineering/sync-to-labels).

## One job, start to finish

[![One job from hand-off through claim, run and retry, with the stores it reads and writes](/blume-assets/content/content/engineering/worker-takes-a-job-preview.png)](/engineering/worker-takes-a-job-diagram)

<a href="/engineering/worker-takes-a-job-diagram">Open the interactive diagram</a>

- Every worker task runs one listener that waits on the queue. Redis hands a new job to the listener that has waited longest, and the listener claims it with a lock that it keeps renewing while the job runs.
- The listener runs the job, two at a time by default. The job loads the `email_message` row by id, reads the raw email from blob storage, processes it, and saves its artifacts to the pipeline tables. It records the email's origin first; a copy of an email already recorded is linked to it and finishes there, with no documents converted and no model call (see [Sync to labels](/engineering/sync-to-labels)). Running each job in a child process from a small pool is planned, so that a crash fails one job instead of the task.
- If the job fails, or the whole task dies and its lock expires, the job waits out a backoff and returns to the waiting list. After three failed tries it stays failed, in the queue's failed list on mail's `/queues` page.
- Every five minutes, and at once when model calls are switched on, a sweeper re-adds jobs for stored emails that have no pipeline run and, with model calls on, for those whose run is open with no job on the queue: at once when the email was never classified, after 30 minutes when it was. It does so at most three times per run. Jobs are deduplicated by the email's id, so a re-added job cannot create a duplicate, and an email waiting in the queue is never counted as stalled. A run swept three times that is still stalled is taken off the queue: the email is marked for review and the run is failed. The detail is on [Sync to labels](/engineering/sync-to-labels).

## The worker fleet

This part is planned; no worker runs on ECS yet.

[![Worker tasks pulling from one queue, the shared stores, and autoscaling on backlog](/blume-assets/content/content/engineering/worker-fleet-scaling-preview.png)](/engineering/worker-fleet-scaling-diagram)

<a href="/engineering/worker-fleet-scaling-diagram">Open the interactive diagram</a>

- Workers run as tasks of one ECS service, and every task pulls from the same queue. Adding a task adds capacity with no other change.
- Each task runs a set number of jobs at once, one per child process. Capacity is the number of tasks times that number.
- The sweeper reports the backlog per task to CloudWatch. When it passes the target, Auto Scaling starts another task, and that task's listener starts claiming jobs.
- When the queue drains, Auto Scaling stops tasks. A stopping task finishes its running jobs first; a job cut off after 120 seconds is retried on another task.
