Skip to content
All work

For an Australian state government client

Document pipeline

A 12-node LangGraph pipeline extracts compliance documents row by row and links each value to its source page.

Status
Status: Production
Role
Pipeline design and engineering
When
2025 to May 2026
Credit
Via Silvatron Pty Ltd
Stack
  • LangGraph
  • Python
  • Pydantic
  • FastAPI
  • SurrealDB
  • Next.js
  • Salesforce

I built a document pipeline for an Australian state government client. It reads regulatory compliance documents, extracts each table row on its own, checks it, and links every value back to the page and cell it came from. The records then load into the client's CRM. The measured results are under Results.

Shown with synthetic data

I built this via Silvatron Pty Ltd under a confidentiality agreement, so I describe the architecture, the decisions and the measured results at a general level. Every table, diagram and film on this page uses made-up data. Nothing here comes from the client's documents, records or schema.

Assurance
  • Every extracted value traceable to its source page
  • Measured on a named benchmark, with precision and recall reported
  • Validation before anything reaches the CRM
  • Human review gates
  • Delivered under a Deed of Confidentiality

The problem

The documents are long PDFs full of tables. Rows carry hierarchical records: a parent row gives context to the child rows beneath it. Every value has to end up in Salesforce custom objects, and it has to pass the CRM's validation rules before it will import.

Until now that meant manual data entry. It was slow, transcription errors crept in, and reporting further down the line waited on it.

Loose pages into one stack. Generated loop, no client material.

A few things made automation hard:

  • Layouts vary. Different authors order columns differently and use different table structures and headings.
  • Tables span pages. Headers repeat mid-document and rows split across page breaks.
  • Records nest. A child row means nothing without its parent, so the hierarchy has to survive extraction.
  • Validation is strict. Some fields only accept values that depend on another field's value. One mismatch blocks the import.
  • Every value needs a source. Reviewers must be able to see the exact page and table cell behind each extracted value.

How it works

The current pipeline is a 12-node LangGraph graph. A deterministic parser reads every table before any model call, and two human review gates sit between the graph and the CRM.

Document to CRM: the 12-node workflow

A deterministic parse runs first, then twelve graph nodes, then two human review gates. Data moves from the document, to structured records, to the CRM.Select a step to see what it does.

Connections: Document in to Deterministic parse; Deterministic parse to 1 · Structure map; 1 · Structure map to 2 · Section inventory; 2 · Section inventory to 3 · Source intelligence; 3 · Source intelligence to 4 · Schema inference; 4 · Schema inference to 5 · Per-section extract; 5 · Per-section extract to 6 · Per-row extract; 6 · Per-row extract to 7 · Normalise; 7 · Normalise to 8 · Validate; 8 · Validate to 9 · Correct; 9 · Correct to 10 · De-duplicate; 10 · De-duplicate to 11 · Missed-row recovery; 11 · Missed-row recovery to 12 · Persist; 12 · Persist to Two review gates; Two review gates to CRM export.

Each node checkpoints its output, so a failed run restarts from the last good step instead of from the top. Every attempt is traced, with tokens and cost per call, so a slow or costly step shows up.

Checkpoint, fail, resume

Each row is one step of a run, each bar one attempt as the trace records it. A diamond is a checkpoint.

Choose a run
  • Deterministic step
  • Model call
  • Failed attempt
  • Checkpoint saved
Time (illustrative)
Fail and resumeThe run stops inside node 9. On restart it loads checkpoint 8 and carries on from there. Nodes 1 to 8 are not run again.

Bar lengths are illustrative, not measured. Measured: one 190-page run took about 31 minutes end to end.

Every value keeps its page

Pick a record to see where it was read from. The coordinates travel with the value from the parse to the review grid.

source.pdfPage 14 of 38
Table 3
SectionField AField BField C
S2A-1B-1C-1
S2A-2B-2C-2
S2A-3
S2A-4B-4C-4
S2A-5B-5C-5
Extracted records
  1. PagePage 14 of 38
  2. CellTable 3, row 4
  3. Positionx 72, y 476, w 396, h 18
  4. RecordRecord 0204, saved with its source
  5. Review gridAny value opens this cell on its page
Record 0204: page 14, table 3, row 4The reviewer clicks any value and sees this cell outlined on the original page before approving it.

The page, table, coordinates and values are synthetic.

One row at a time, checked before it counts

Rows from one table are extracted in parallel, one model call each. Pick a row to see its path.

Rows extracted in parallel
Row R2: corrected
  1. One call to the primary model.
  2. One field fails a rule.
  3. Correct fixes it and logs the change.
  4. The re-check passes and the record is saved.

Rows and outcomes are synthetic.

Each row passes an extractor and two verifiers. A mismatch goes back for a retry, and a person approves at two gates before anything is published. Illustrative data.

A FastAPI backend streams progress to a Next.js review screen. A person checks flagged rows, opens the source page and approves the records at two gates before anything is exported.

Two gates, and a person at each

Play the reviewer. The run cannot reach the CRM until both gates are approved.

The run stops at gate 1 and waits for a person.
Loose cards into a grid. Generated loop, no client material.

The output is a set of structured records like the register below. The data is synthetic.

Extracted records, each with its sourceSynthetic demo data

Scroll sideways for the source column.

RecordGroupClassValueSourceConfidenceStatus
R-1183Group B · Item 04Class 212p. 7 · T2 R40.98Verified
R-1184Group B · Item 05Class 14p. 7 · T2 R50.97Verified
R-1185Group C · Item 02Class 38p. 9 · T1 R20.71Review
R-1190Group C · Item 07Class 220p. 10 · T1 R70.95Verified
R-1192Group A · Item 01Class 16p. 12 · T3 R10.99Verified
R-1201Group A · Item 09Class 33p. 13 · T3 R90.54Flagged
R-1205Group D · Item 03Class 215p. 15 · T4 R30.93Verified

Decisions

One LLM call per row, not one per document

My first prototype sent the whole document to a model and asked for JSON in one pass. It failed three ways. Long documents overflowed the context window and output was silently truncated. The model invented values where a table was ambiguous or split across pages. And there was no intermediate state, so I could not tell where an error came from.

Batching 10 to 20 rows per call came next. It was cheaper, but one malformed response failed the whole batch.

One call per row is slower and costs more. The payoff is that a bad response now costs one row, not twenty, and every row has its own record of what went in and what came out. For a compliance workflow where completeness matters, that trade is worth it.

Earlier design: two parsers and a consensus merge

An earlier version ran two PDF parsers side by side, Docling and MinerU, and merged their output. Where both agreed on a value, it was marked high confidence. Where they disagreed slightly, such as whitespace or punctuation, the merge resolved it automatically. Larger disagreements went to a person.

I record it here as earlier work. The 12-node pipeline above is the current design.

SurrealDB for the record graph

PostgreSQL was the obvious default. I chose SurrealDB because the records are hierarchical and SurrealDB stores graph relations, structured records and nested metadata such as cell coordinates in one place. A query like "every child record under this parent" is a single traversal rather than a recursive query.

The risk was operational: SurrealDB is younger, with fewer managed hosting options. I ran it in Docker next to the FastAPI backend with scheduled exports as a backup.

Python: streaming progress to the review screen

The backend streams pipeline events to the browser with Server-Sent Events. Rows appear in the review grid as they are extracted, instead of after the whole run.

python
import asyncio
import json
from collections.abc import AsyncGenerator
from dataclasses import dataclass, field
from typing import Any, Literal

from fastapi import APIRouter
from fastapi.responses import StreamingResponse

@dataclass
class PipelineEvent:
    event_type: Literal["stage_start", "record_extracted", "stage_complete", "error"]
    job_id: str
    payload: dict[str, Any] = field(default_factory=dict)

class PipelineEventBus:
    """Pipeline nodes publish; the SSE endpoint subscribes per job."""

    def __init__(self) -> None:
        self._queues: dict[str, asyncio.Queue[PipelineEvent | None]] = {}

    def register(self, job_id: str) -> asyncio.Queue[PipelineEvent | None]:
        queue: asyncio.Queue[PipelineEvent | None] = asyncio.Queue()
        self._queues[job_id] = queue
        return queue

    async def publish(self, event: PipelineEvent) -> None:
        queue = self._queues.get(event.job_id)
        if queue:
            await queue.put(event)

    async def close(self, job_id: str) -> None:
        queue = self._queues.pop(job_id, None)
        if queue:
            await queue.put(None)  # sentinel: end of stream

event_bus = PipelineEventBus()
router = APIRouter()

@router.get("/api/jobs/{job_id}/stream")
async def stream_job_events(job_id: str) -> StreamingResponse:
    queue = event_bus.register(job_id)

    async def generate() -> AsyncGenerator[str, None]:
        while (event := await queue.get()) is not None:
            yield f"data: {json.dumps(event.__dict__)}\n\n"

    return StreamingResponse(generate(), media_type="text/event-stream")

Results

F1 on one named benchmark
90%
93.1% precision, 87.1% recall
Source: One benchmark document set. Larger document sets varied.
One 190-page production document
~31 min
About AUD 1.65 for the run
Source: A single dated run, not an average.

Read the F1 score as a benchmark result, not a general accuracy rate: larger document sets varied materially. The 190-page run is one observation and does not describe every document.

Status

I built it for an Australian state government client via Silvatron Pty Ltd, from 2025 to May 2026. It went into production as a managed service, with Langfuse for pipeline traces and Logfire for inspecting Pydantic validation.

Working on something like this?

I build agent workflows, document pipelines and the apps around them. If this looks like your problem, book a call.

Book a call (opens in a new tab)

wihithat@gmail.com

  • Status: ProductionBuilt at Silvatron Pty Ltd

    Claude Code, n8n, Bitbucket, Jira and Discord run first-pass reviews and diagnose failed pipelines. A person approves every fix.

    Evidence: About 20–40 developer hours a month (estimate)

    Stack: Claude · n8n · Node.js · Bitbucket

  • Status: Prototype

    A Monash team prototype that turns receipts and bank statements into matched, tax-ready records using Mistral OCR and GPT-4o-mini. Demo offline.

    Evidence: Prototype with a tested OCR and matching pipeline; public demo offline

    Stack: React · TypeScript · Supabase · n8n