Sync to labels
What runs today, from a mailbox sync to a labelled email, and the state each step leaves in Postgres.
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 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
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.
- The scheduler starts a sync of every active mailbox every 30 seconds in development and every 15 minutes in production, unless
MAIL_POLL_INTERVAL_SECONDSsets another interval. On start it returns any email stranded atINGESTION_PROCESSINGtoINGESTION_PENDING. - Enrolling a mailbox, or turning one back on, makes the API add a
sync-mailboxjob 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_messagerow atINGESTION_PENDINGand 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_DAYSstarts it that many days back instead, newest first, andMAIL_SYNC_MAX_MESSAGES_PER_PASSmakes 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 asmail/attachments/<sha256>, in the local folder or S3. It then records the raw message and attachment rows and moves the email toEXTRACTION_PENDING. A raw message missing on every delivery, or a last attempt that fails, leaves itINGESTION_FAILED. - Once those rows are committed, mail adds the email’s
process-emailjob to theemail-processingqueue. 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, orpnpm --filter @flowos/worker devon its own. It takes jobs from theemail-processingqueue, two at a time unlessWORKER_CONCURRENCYsets 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, andemail_messagerecords its key inbody_refwith the sender’s name and the To and Cc addresses. Each stored attachment the body shows as an image is markedembedded; 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 inemail_referencein 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 emptybody_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_ORIGINto the first recorded, which on a newest-first sync can be the later email, its run endsDONE, 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 linkedSUSPECTED. How the origin is found, and the whole copy rule, are on 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 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, orWORKER_MODEL_CALLS=onwhile the switch is unset. Off, it stops after the documents, and the run staysRUNNINGuntilpnpm 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:bodiessaves 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 aspartial, and the job’s log marks the reading partial. TheREADstep names Jev only when neither the categories nor the asks fell back to their rules. Each part’s fallback is on Enquiry reading. - A job for an email that mail has deleted, or is ingesting again, does nothing.
- 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
- 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_PENDINGwith 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 stillRUNNINGand 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_REVIEWwith the errorCould not classify after 3 attemptsand marks the runFAILED. 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.
- 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
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
DONEpipeline 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--allit 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, theemail-processingqueue’s counts and the latest ten emails with their kind. It only reads.
Email status and run state
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
FAILEDstep. 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.
markPipelineRunmirrors a run onto an email atEXTRACTION_PENDINGorPROCESSED:DONEmoves it toPROCESSED, andRUNNINGmoves a processed email back. A run whose latest attempt at any stage failed staysRUNNING, unless its email is a copy. The sweeper is the other writer, and it is what setsNEEDS_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_PENDINGwith 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 stillRUNNING. - 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_PENDINGand 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 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_URLunder that provider’s default terms.
How each layer reads an email is on Processing layers.