For workflow authors
The built-in presets cover standard document ingestion, but workflows are a general, user-extensible pipeline. This page explains the programming model for teams that write their own workflows and processing modules. It is aimed at platform engineers: you write workflows in code inside the ingestion service, not through the HTTP API.
Mental model
A pipeline is built from two kinds of objects:
-
A Workflow is orchestration as data: a Pydantic model whose fields declare the modules it uses, with an async
run(ctx)that dispatches each stage and threads lightweight references between them. The same workflow runs in-process for debugging or distributed across the worker fleet, unchanged. -
A Module is a serializable, pure transform: it receives entities (files, documents, pages, chunks), returns new entities, and never persists anything itself — the framework hydrates inputs and persists outputs at the dispatch boundary.
There is no third "Step" tier. The current model is exactly these two.
The module contract
A module subclasses Module and needs three things:
-
A
kind— aLiteralfield default (for example,kind: Literal["word_count"] = "word_count"). Declaring it auto-registers the class in the module registry and tags the workflow history. -
A
runmethod —run(self, ctx, input)returning one entity or a flat list of entities. Writeasync deffor I/O-bound work or a plaindeffor CPU-bound work. Whether the module is called per item or with batches is inferred from the signature: annotate the input aspage: Pagefor per-item, orpages: list[Page]for batched processing. -
Parameters as Pydantic fields — every field declared with
Field(default=..., description=...)becomes a configurable workflow parameter, and the parameter schema is generated from the model automatically.
Class-level knobs tune execution: compute (threads, processes, or inline), workload (CPU or IO task queue), and per-work-unit batch_size, max_in_flight, and timeout. Two markers change a module’s role in the pipeline: initial = True (seeds the pipeline from the run’s input source) and terminal = True (projects a final result instead of persisting entities).
The runtime context
run receives a RuntimeContext with everything a module may touch:
| Attribute | Purpose |
|---|---|
|
Relational and vector storage — querying existing entities; outputs are persisted for you |
|
File storage: open, read, upload |
|
Model inference: completion, image description, embeddings, transcription |
|
PDF backends: render pages, extract text and images |
|
The run’s scope: search store, workflow, run ID, and input source |
Modules must not open their own connections or widen their scope: every read and write is automatically confined to the search store and workflow run the module executes in.
Constraints and guarantees
-
Idempotency. Entities declare which fields derive their ID, so re-runs and retries produce the same rows instead of duplicates. Enrichment outputs are append-only — saving a duplicate is a no-op — which makes modules safe under the orchestrator’s retries.
-
Error handling. Raise
NonRetryableErrorfor deterministic failures (bad input or configuration) so a doomed run fails fast; any other exception is retried according to the worker’s retry policy. -
Purity. A module is a pure transform of its input plus
ctx. Hydration, persistence, and scope enforcement all happen at the framework boundary, not inside the module.
Composing workflows and presets
A workflow subclasses Workflow, declares its modules as fields (a single module, or a kind-discriminated union to make the module selectable through configuration), and implements async run(ctx) by dispatching stage after stage. Subclassing registers the workflow under its class name — the same name you pass as workflow_name when creating a workflow through the API (see Workflows and presets).
A preset is nothing more than a workflow subclass that overrides field defaults: DocumentAIFull, for example, is DocumentAI with page-level keyword extraction enabled.
Registering a custom workflow
Registration is an import side-effect: the worker imports the preset package at startup, which populates the module and workflow registries and creates one execution activity per module kind. To deploy a new workflow today, its code must live in the ingestion service, or be imported by it. There is no external plugin mechanism yet. Plan custom workflow development together with the platform team.