---
title: Sync to labels
description: What runs today, from a mailbox sync to a labelled email, and the state each step leaves in Postgres.
sidebar:
  label: Sync to labels
  order: 6
---

This is the path that is built today. Mail syncs and stores each email, then hands it to the worker through the `email-processing` queue, as [Mail to processing](/engineering/mail-to-processing) describes. The worker's listener extracts the email and converts its documents, then classifies it, and labels and reads it when it is an enquiry, unless `WORKER_MODEL_CALLS=off`. Classification and reading send email text to Jev, through TypeSafe's own API with no zero-retention agreement yet, and reading sends the client's contact block to the model at `LLM_BASE_URL`. `pnpm flow` runs the same steps by hand. The job carries only the email id; everything else the two sides share is Postgres rows and stored files.

## The path

```mermaid
flowchart TB
  subgraph mail["apps/mail, one instance"]
    sched["Scheduler\nevery 30 s in development, 15 min in production"]
    sync["Sync a mailbox\nGraph delta pages from its saved delta link"]
    stage["Stage the email\nthread, email_message, delivery\nINGESTION_PENDING"]
    ingest["Ingest job claims it\nINGESTION_PROCESSING"]
    store["Store the raw email and its files\nEXTRACTION_PENDING"]
    handoff["Hand it off\nadd its process-email job"]
    sched --> sync --> stage -->|"mail-ingest queue, email id"| ingest --> store --> handoff
  end

  queue[["email-processing queue\nRedis, one job per email id"]]

  subgraph shared["Shared stores"]
    direction LR
    pg[("Postgres\nemail_message, email_attachment\npipeline_run, pipeline_step, artifact")]
    blobs[("Blob storage\napps/mail/.local/storage or S3")]
  end

  subgraph worker["apps/worker"]
    listener["Listener\ntakes jobs, 2 at a time"]
    sweep["Sweeper, every 5 min\nre-adds a missing or stalled run\nafter 3 stalls, Needs review"]
    extract["Start the run, extract\nregions, text, references, origin"]
    docs["Extract documents\nPREPARE, DOCUMENT_EXTRACTION"]
    classify["Classify with Jev\nCLASSIFY, CLASSIFICATION"]
    label["Label by rules\nLABEL, LABELLED_CONTENT"]
    read["Read an enquiry\nREAD, ENQUIRY_READING"]
    done["Run DONE\nemail PROCESSED"]
    flow["pnpm flow\nthe same steps, by hand"]
    listener --> extract --> docs -->|"while model calls are on"| classify -->|"any other kind"| done
    classify -->|"an enquiry"| label --> read --> done
    extract -->|"a copy"| done
    flow -.-> extract
  end

  handoff -->|"email id"| queue --> listener
  sweep --> queue
  store --> pg
  store --> blobs
  pg --> extract
  blobs --> extract
```

The inbox reads the same rows. It lists emails at `EXTRACTION_PENDING`, `PROCESSED` and `NEEDS_REVIEW`. The first two show as ready, with no status badge. `NEEDS_REVIEW` and `INGESTION_FAILED` show as Needs review. The Needs Review tab lists a conversation when any of its messages is one of those two. Drafts stay off.

### Mail

- The scheduler starts a sync of every active mailbox every 30 seconds in development and every 15 minutes in production, unless `MAIL_POLL_INTERVAL_SECONDS` sets another interval. On start it returns any email stranded at `INGESTION_PROCESSING` to `INGESTION_PENDING`.
- Enrolling a mailbox, or turning one back on, makes the API add a `sync-mailbox` job for it, so it syncs as soon as the sync queue is free rather than at the next scheduled scan. Adding it while one for that mailbox is still queued does nothing. When Redis is unreachable the API logs a warning, the enrolment still succeeds, and the scheduled scan picks the mailbox up.
- A sync reads the mailbox's Graph delta pages from its saved delta link. Each new email is staged in one transaction: its thread, its `email_message` row at `INGESTION_PENDING` and a delivery for the mailbox. The new delta link is saved once the pages are read, and each staged email gets an ingest job carrying its id, tried up to three times.
- With no saved link, a sync starts a new delta round over the Inbox, which ends after 50 emails. `MAIL_SYNC_WINDOW_DAYS` starts it that many days back instead, newest first, and `MAIL_SYNC_MAX_MESSAGES_PER_PASS` makes each sync read one page of that many items and save the link to the next, so a window is walked a page per sync.
- Ingestion claims the email, downloads its raw message and stores it as `mail/raw/<sha256>.eml`, and each file attachment as `mail/attachments/<sha256>`, in the local folder or S3. It then records the raw message and attachment rows and moves the email to `EXTRACTION_PENDING`. A raw message missing on every delivery, or a last attempt that fails, leaves it `INGESTION_FAILED`.
- Once those rows are committed, mail adds the email's `process-email` job to the `email-processing` queue. The job carries only the email id. Adding a job for an email that already has one waiting does nothing, and adding one while it runs queues a single follow-up. A failed add is logged, not retried, and the worker's sweeper finds the email.

### Worker

The worker runs the same steps in two ways, and both save through the same writers.

- **The listener** starts with `pnpm dev`, or `pnpm --filter @flowos/worker dev` on its own. It takes jobs from the `email-processing` queue, two at a time unless `WORKER_CONCURRENCY` sets another number. For each email it starts the pipeline run, extracts the email, saves its display body, records its origin and converts its documents.
  - The display body is what the inbox's reading pane shows: the email's HTML with its embedded images inlined, every other image removed along with its alt text, and anything that could run or submit removed, or its plain text when it has no HTML. Remote images are never loaded, and an embedded image over 512 KB, or past 4 MB for the whole email, is left out. It is stored beside the raw message as `mail/bodies/<sha256>.html`, and `email_message` records its key in `body_ref` with the sender's name and the To and Cc addresses. Each stored attachment the body shows as an image is marked `embedded`; one too large to inline is not, so it stays an attachment. The references extraction found, such as our quotation numbers, purchase orders, the client's enquiry or job numbers and reference-shaped subject tokens to confirm, replace the email's rows in `email_reference` in the same step. It is saved before the origin, so a copy has one too, and with model calls off. A retry keeps a body an earlier attempt saved, with its recipients and references. An email with no body at all keeps an empty `body_ref`.
  - The origin is the client behind any forward. An email that adds nothing to one already recorded from the same client, received within 30 days of it and at the same office, is a copy, unless a forwarding colleague's note adds words of its own. A copy is linked `SAME_ORIGIN` to the first recorded, which on a newest-first sync can be the later email, its run ends `DONE`, and nothing else runs, not even document conversion. Any other email is processed normally; one from the same client under the same subject within seven days is also linked `SUSPECTED`. How the origin is found, and the whole copy rule, are on [Enquiry reading](/engineering/enquiry-reading).
  - It then classifies the email, reusing a classification saved earlier. Any kind but an enquiry is finished there: its classification is saved, again when it was reused, and it is neither labelled nor read. Labelling calls no model: it splits the email and its documents into addressable lines and rows, and tags signatures, disclaimers and firm references by rule. An email classified as an enquiry is then read for its catalogue categories, what the client asks of the offer and by when, and the client's contact, with two Jev calls and one call to the model at `LLM_BASE_URL`, as [Enquiry reading](/engineering/enquiry-reading) describes. Saving the reading, not the labels, finishes an enquiry's run. Model calls are off by default: the mail debug page's switch turns them on, or `WORKER_MODEL_CALLS=on` while the switch is unset. Off, it stops after the documents, and the run stays `RUNNING` until `pnpm flow`, or the job the sweeper re-adds once model calls are on, classifies and labels it, and reads it when it is an enquiry. When TypeSafe or the contact model rate limits a call, the job holds the whole queue for a minute and goes back to wait without using one of its three tries.
  - A raw email it cannot read, or a classification call that fails, fails the job. A display body it cannot write is logged and the email processed on, with its recipients and references saved; `pnpm backfill:bodies` saves the body later. The queue tries each job three times, backing off between tries. A reading call that fails does not fail the job: its part falls back to its rules and is saved as `partial`, and the job's log marks the reading partial. The `READ` step names Jev only when neither the categories nor the asks fell back to their rules. Each part's fallback is on [Enquiry reading](/engineering/enquiry-reading).
  - A job for an email that mail has deleted, or is ingesting again, does nothing.
- **The sweeper** runs every five minutes, once across all worker processes, and once more straight away when the mail debug page switches model calls on. It re-adds a job for each email at `EXTRACTION_PENDING` with no pipeline run, which covers a hand-off mail lost and the backlog when the worker first starts. With model calls on, it also re-adds an email whose run is still `RUNNING` and has no job on the queue: at once when it was never classified, such as one that stopped after the documents while model calls were off, and otherwise once its run has not changed for 30 minutes, such as one whose classification failed all three tries. It does so at most three times per run. An email whose job is waiting, retrying or running is never re-added or counted, however long it waits, and a re-added job that stops after the documents because model calls were switched off does not use up one of the three. An email still being ingested is left to mail, which hands it off once it is stored.
  - On a later sweep, a run already swept three times that is idle again by the same rule, with no job on the queue, is retired. The sweeper sets the email to `NEEDS_REVIEW` with the error `Could not classify after 3 attempts` and marks the run `FAILED`. The email and its artifacts stay. Retirement runs whether or not model calls are on; a run only reaches three sweeps while stalled retries are on. A classification that misses softly, with no probabilities or below the threshold, still finishes as an enquiry and is not retired.
- **`pnpm flow`**, run from the repository root, is a terminal menu. Each step works on the emails picked and saves what it makes, unless a dry run is on.

| Step              | Reads                                                                | Saves                                                                    |
| ----------------- | -------------------------------------------------------------------- | ------------------------------------------------------------------------ |
| Pick emails       | Emails at `EXTRACTION_PENDING` or `PROCESSED`, with their saved classification | Nothing                                                        |
| Extract           | The raw message, through the blob reader                             | With writes on, the email's display body in mail storage, its `body_ref`, sender name and recipients, then its origin on `email_message`, a `CONFIRMED` `SAME_ORIGIN` link for a copy and `SUSPECTED` links to the same client and subject within seven days; later steps skip a copy |
| Extract documents | Each attachment's file, converted to Markdown                        | One `PREPARE` step per attachment set, one `DOCUMENT_EXTRACTION` per file |
| Classify          | The extracted text, sent to Jev after asking                         | A `CLASSIFY` step, a `CLASSIFICATION` artifact, and the email's current kind in `email_classification` |
| Label             | Text and documents of a classified enquiry; rules only              | A `LABEL` step and a `LABELLED_CONTENT` artifact |
| Read enquiries | Labels, classification and the item group master; Jev and the contact model after asking, or by rule alone: the category vocabulary, asks by their keywords, and the contact's rule fields | A `READ` step and an `ENQUIRY_READING` artifact, which finishes an enquiry's run |

- "Newest not finished yet" offers emails with no `DONE` pipeline run.
- A file converted before, for any email, is reused rather than converted again, and an outbound email's files are not converted.
- An image is not converted, since the document reader reads no image format, so a photo pasted into an email is not read yet. The step counts images apart from files the reader rejected: `2 of 21 documents converted (image 19)`.
- Every step reloads the picked emails first and drops any that mail has deleted or is ingesting again. Label rereads the saved classifications and labels only emails classified as an enquiry. Read rereads them too, and reads only emails classified as an enquiry, skipping one whose classification call failed; it counts any with a part read by rule alone because its model call failed, and See results shows a reading's due date, when it has one, and marks it partial when its categories or asks fell back to their rules or its contact call failed. Read sends contact blocks to the contact model only alongside the classifier's model, so without one the setup note says contacts are read by rule. Run everything left asks before reading only when an enquiry may be read, and skips Read when nothing was labelled.
- A dry run saves nothing. Before writes come back on, its results are either saved, in stage order, or discarded, so a saved label never rests on an unsaved classification or document.
- **`pnpm backfill:bodies`**, run from the repository root, saves the display body, recipients and references of every stored email that has no body, such as those processed before bodies were saved; with `--all` it saves every stored email's again, as after the body rendering changes. It reads and extracts each raw message again, 50 at a time, asks before writing unless given `--yes`, calls no model, and ends with how many it saved, how many had no body, and the ids of any that failed.
- **`pnpm status`**, run from the repository root, prints each mailbox's sync, the stored emails by status and how many emails were set aside as copies, what the pipeline saved in the last hour, the `email-processing` queue's counts and the latest ten emails with their kind. It only reads.

## Email status and run state

```mermaid
stateDiagram-v2
  [*] --> INGESTION_PENDING: mail stages it
  INGESTION_PENDING --> INGESTION_PROCESSING: an ingest job claims it
  INGESTION_PROCESSING --> INGESTION_PENDING: attempt failed, retried
  INGESTION_PROCESSING --> INGESTION_FAILED: gone from Graph, or retries used up
  INGESTION_PROCESSING --> EXTRACTION_PENDING: raw email and files stored
  EXTRACTION_PENDING --> PROCESSED: labels or an enquiry's reading saved, or a copy linked, run DONE
  PROCESSED --> EXTRACTION_PENDING: new document, classification or enquiry labels, run RUNNING
  EXTRACTION_PENDING --> NEEDS_REVIEW: swept three times and still stalled, run FAILED
  EXTRACTION_PENDING --> INGESTION_PENDING: sync reset, re-ingest
  PROCESSED --> INGESTION_PENDING: sync reset, re-ingest
```

Mail moves an email up to `EXTRACTION_PENDING`. After that, the email's status mirrors its pipeline run:

| Pipeline run                       | Email status         | When                                                                         |
| ---------------------------------- | -------------------- | ---------------------------------------------------------------------------- |
| None                               | `EXTRACTION_PENDING` | Mail has just stored the email                                               |
| `RUNNING`                          | `EXTRACTION_PENDING` | The listener started it, a document or classification is saved, an enquiry's new labels are saved, or the latest classification failed |
| `DONE`                             | `PROCESSED`          | A classification of anything but an enquiry is saved, or an enquiry's reading is, and no stage's latest attempt failed |
| `DONE`                             | `PROCESSED`          | A copy of an email already recorded                                          |
| `FAILED`, after three stalled sweeps | `NEEDS_REVIEW`     | The sweeper took the email off the queue for review                          |
| Deleted, with its steps and artifacts | Email deleted     | A sync reset purges the email                                                |

- A run is finished once a classification of anything but an enquiry, or an enquiry's reading, is saved after its latest document, or when its email is a copy. A new document, a new classification, or an enquiry's new labels reopen it. A classification of another kind saved again, as when a new document reopened the run, finishes it again; an enquiry's run is finished by reading it again.
- A classification whose model call failed is saved as a `FAILED` step. The run stays open, and the email stays on the "not finished" list, until a later classification succeeds or the sweeper retires it.
- The worker never changes an email that is still being ingested. `markPipelineRun` mirrors a run onto an email at `EXTRACTION_PENDING` or `PROCESSED`: `DONE` moves it to `PROCESSED`, and `RUNNING` moves a processed email back. A run whose latest attempt at any stage failed stays `RUNNING`, unless its email is a copy. The sweeper is the other writer, and it is what sets `NEEDS_REVIEW`.

## Mail debug

`/debug/mail` is the operator view of this path, shown to callers holding `debug:view` (the Developer role).

- **Waiting for worker** is an email at `EXTRACTION_PENDING` with no open pipeline run and no job the worker currently holds. **Worker running** is an email the worker holds, or one whose pipeline run is still `RUNNING`.
- A processed email shows its saved classification when one exists: Enquiry, Revision request, Quotation, Purchase order, Sales order, Follow-up, or Other. The reading pane leads with an enquiry's reading, in sections: who forwarded it and their note; the categories it asks for, solid when present and dashed when unsure, with the client's words that show them and the documents' hints; the due date, marked as counted from when it was sent when the client wrote it relative, and whether the client calls it urgent; the terms the client asks for, grouped by type, with the attachments named as the client's terms; and the client's contact, marked "by rule only" when the model did not read its contact block and noting a PIN the model's address disagrees with, or saying that no client was found, or, for a reading made before contacts were read, that the contact was not read. A reading whose categories were named by the vocabulary alone is marked "vocabulary only", and one whose asks were kept by their keywords "asks by keyword only", each such ask marked "keyword only". The labelled content, with its sections and the lines that were tagged, is folded below the reading, and open when the email has no reading; until the worker has labelled the email, it says so. A finished run of anything but an enquiry shows its Read stage as Not an enquiry; an enquiry processed before reading was built shows it as not run, and an open run labelled again after its reading shows it as not read since. The message ids, deliveries and storage sit under Envelope and storage.
- The **email type** filter lists only emails currently of the kinds picked, and **Not classified** those the worker has not classified yet. It combines with the search, direction, mailbox and tab filters. An email's kind is its current classification, so one classified again moves with it.
- A finished email of any kind but an enquiry shows its Label and Read stages as `Not an enquiry`.
- Emails that share a thread stay together. The thread is placed by its latest activity, and the messages inside it run oldest first. A long thread can still split across pages.
- A copy shows its Documents, Classify, Label and Read stages as done, each naming the first email by the first eight characters of its id: `Copy of 0190aaaa — not processed`. This needs the first email inside the viewer's mailbox scope; a copy of an email delivered to another mailbox shows those stages as not run.
- A forwarded email shows its origin, the client behind the forward, beside the sender: `(Origin: client@example.com)`.
- **Needs review** lists emails at `NEEDS_REVIEW`. That tab is shown while any email is in that state, or while the tab is selected.

## Resets

`pnpm --filter @flowos/mail sync:reset` puts a mailbox back to an earlier state, and settles the worker's runs to match. A pipeline run has no foreign key to its email, so the reset does this itself.

- **Soft reset** and **complete wipe** delete the emails they purge, together with their pipeline runs, steps and artifacts. A copy of a purged email that survives loses its link, so its run is reopened and a processed copy goes back to `EXTRACTION_PENDING`, for the worker to read as an email of its own.
- **Force re-ingest** sends the mailbox's emails back to `INGESTION_PENDING` and reopens their runs. Once each is stored again, mail hands it to the listener again. Their artifacts are kept, so documents converted before are reused.

## Not built yet

- **The child-process pool.** The listener runs each job in its own process, not in the pool that [Worker jobs](/engineering/worker-jobs) describes.
- **Scaling on backlog.** The sweeper reports no backlog metric, and nothing starts more workers.
- **A zero-retention agreement.** The email text the listener sends to Jev is kept under TypeSafe's default terms, and the contact blocks sent to the model at `LLM_BASE_URL` under that provider's default terms.

How each layer reads an email is on [Processing layers](/engineering/processing-layers).
