# Hands-on: build a queue-based document pipeline with pipes and filters and claim check

> A PDF pushed straight into a queue, a processing step that runs twice, a client that resends a POST: all three familiar failures can be handled in a few dozen lines of Python, provided the design is right from the start.

Original: https://fdetimes.net/en/guides/document-pipeline-pipes-filters-claim-check/

The detail most often overlooked in Microsoft's reference documentation on pipes and filters is that the messages passing between steps do not contain the image being processed. Each message holds only a reference (a claim check) pointing to the image sitting in storage.

This keeps messages small, but it does not remove a separate risk: a queue can deliver the same message more than once. So, alongside claim check, every filter must also be idempotent, meaning that receiving a message twice still produces exactly one result. Step 3 deals with this.

On FDE engagements this problem turns up under many names: loan applications, scanned contracts, insurance claims. The customer uploads a heavy file, the file has to pass through several validation and extraction steps, and someone has to know when processing is finished. This guide builds such a pipeline on a laptop using only the Python standard library.

## What you will build, and what you need

The pipeline has three steps, validate, extract and store, connected by queues. This is the approach suggested in Microsoft's best-practice guidance on background jobs: a job that passes through several stages should be split into separate filters connected by queues. Outside the pipeline sits an API that accepts documents and reports processing status.

All you need is Python 3 and an empty directory. Everything below is a **simplified simulation**: `queue.Queue` stands in for the message broker, a local directory stands in for blob storage, and the API is written as plain functions with no framework. On a customer site you swap each piece for their real services; the logic stays the same. (In the code, `docs` is Vietnamese for "application file".)

## Step 1: write the file first, send the ticket second

As Microsoft describes it, claim check means storing a large payload in an external store and sending only a reference through the messaging system. The consumer receives the reference and fetches the payload itself. The order of these two operations matters: the documentation requires that the reference be published only after the payload has been written successfully.

```python
import hashlib, queue, uuid
from pathlib import Path

STORE = Path("blobs"); STORE.mkdir(exist_ok=True)
q_validate, q_extract, q_store, q_dead = (queue.Queue() for _ in range(4))

def put_payload(data: bytes) -> str:
    key = hashlib.sha256(data).hexdigest()
    (STORE / key).write_bytes(data)
    return key

def submit(data: bytes) -> str:
    job_id = str(uuid.uuid4())
    ref = put_payload(data)                      # 1. write payload
    q_validate.put({"job_id": job_id, "blob": ref, "schema": "docs.v1"})  # 2. publish
    return job_id
```

**Check:** call `submit(b"%PDF-1.4 ...")` and look in the `blobs/` directory. It should contain exactly one file, while the message in `q_validate` is only a few dozen bytes.

Using a hash of the content as the key is deliberate: the same file always produces the same key. So if the write is retried, no extra copy is created.

Microsoft notes that the write and the publish do not happen atomically, so you have to account for both duplicate messages and orphaned payloads, meaning files that were written but whose message was never sent.

## Step 2: write a common filter template, so the pipe only carries messages

The Azure documentation is explicit: a pipe does no routing and holds no logic; it simply turns the output of one filter into the input of the next. All processing, including error handling, therefore lives in the filters.

```python
def run_filter(inbox, outbox, work, max_attempts=3):
    msg = inbox.get()
    try:
        result = work(msg)
        if outbox is not None:
            outbox.put({**msg, **result})
    except Exception:
        msg["attempts"] = msg.get("attempts", 0) + 1
        (q_dead if msg["attempts"] >= max_attempts else inbox).put(msg)
```

The `except` branch handles poison messages. A message that fails repeatedly is not requeued forever but moved to a dead-letter queue, exactly as the background-job guidance recommends.

The `schema` field attached to the message has a purpose too: each filter knows only its own input and output schemas, so once schemas are standardised you can reorder the steps without rewriting them.

## Step 3: three filters, and the last one must be idempotent

```python
DONE = {}

def validate(msg):
    data = (STORE / msg["blob"]).read_bytes()
    if not data.startswith(b"%PDF"):
        raise ValueError("not a PDF")
    return {}

def extract(msg):
    data = (STORE / msg["blob"]).read_bytes()
    return {"size_bytes": len(data)}   # simplified: replace with real OCR/extraction

def store(msg):
    if msg["job_id"] in DONE:          # already processed, skip
        return {}
    DONE[msg["job_id"]] = {"blob": msg["blob"], "size_bytes": msg["size_bytes"]}
    return {}
```

(The comments read, in order: "simplified: replace with real OCR/extraction" and "already processed, skip"; the error message means "not a PDF".)

The reason for the `if` line in `store`: a queue guarantees only *at-least-once* delivery, which means the same message can arrive several times.

Microsoft describes a specific scenario: a filter posts its result and then fails, the message is rerun on another instance, and the result is duplicated.

Any filter that writes to the outside world therefore needs a key to recognise work already done.

**Check:** run the three filters in turn with `run_filter(q_validate, q_extract, validate)`, `run_filter(q_extract, q_store, extract)` and `run_filter(q_store, None, store)`. Then manually put the same message into `q_store` once more and run it again. `len(DONE)` must still equal 1.

**Key point:** Write the data before sending the ticket, and make every filter able to handle receiving a message twice.

## Step 4: find the slowest filter before optimising

The Azure documentation says that the time taken to process a request depends on the slowest filters in the pipeline. Imagine validate takes 1 second, extract 8 seconds and store 1 second.

Optimising validate changes almost nothing. The right move is to run several instances of extract in parallel, all reading from `q_extract`, and that is a direct benefit of separating the steps with queues.

In the simulation you can create a few `threading.Thread` instances that all call `run_filter` on `q_extract`. Because `store` is already idempotent, even if two workers happen to process the same message, the result is still correct.

## Step 5: the API returns 202, blocks duplicates, and returns 303 when done

The customer should not have to hold a connection open while the pipeline runs. Under the async request-reply pattern, the API validates the request and immediately returns HTTP 202 with `Location` and `Retry-After` headers; the client then polls the status endpoint.

```python
JOBS_BY_KEY = {}

def post_docs(body: bytes, idem_key: str):
    job_id = JOBS_BY_KEY.get(idem_key)
    if job_id is None:
        job_id = submit(body)
        JOBS_BY_KEY[idem_key] = job_id
    return 202, {"Location": f"/docs/status/{job_id}", "Retry-After": "5"}

def get_status(job_id: str):
    if job_id in DONE:
        return 303, {"Location": f"/docs/{job_id}"}
    return 200, {"status": "processing"}
```

Two details in this code prevent two different failures. The `Idempotency-Key` header lets the backend return the existing status resource when a client resubmits, instead of putting another work item on the queue.

And when the job is finished, Microsoft recommends 303 See Other rather than 302, because with some clients a 302 can cause the original POST to be resent.

**Check:** call `post_docs` twice with the same key. Both calls must return the same `Location`, and `q_validate` must grow by only one message.

## Final exercise: break validate on purpose

This exercise ties all five steps together. Call `post_docs(b"hello", "key-bad")`: the file is still written to `blobs/` and the API still returns 202, but the content does not start with `%PDF`, so `validate` will raise an error.

Now run `run_filter(q_validate, q_extract, validate)` three times. On the first two runs the message returns to `q_validate` with `attempts` incrementing; on the third it lands in `q_dead`. Confirm that `q_dead.qsize()` equals 1, `q_extract` is still empty, and `get_status` for this job still returns `processing`.

That last result exposes a gap in the simulation: the client will poll forever without learning that the document failed. The natural next step is to have the status endpoint also read the dead-letter queue and return a failed status.

## Common mistakes

The first mistake is publishing the message before the file has been fully written, reversing the order the claim check documentation requires. The first filter may then read a reference that points at nothing.

The second is that nobody owns payload deletion. The claim check documentation requires you to state clearly who owns deletion and retention, and to align the payload's retention period with the message's lifetime, so that a still-valid reference never points to data that has already been deleted.

The third concerns security: putting a security token into the claim check for convenience. Microsoft advises against this. The last mistake is applying the pattern in the wrong place. If processing must complete within the original request, request/response style, pipes and filters is not a fit.

## How this skill shows up on a customer site

On a customer site the first task is not writing code but asking one question: will users accept receiving the result later, through a status endpoint, or must they have it immediately?

The answer decides whether you use a pipeline at all. Next, map the current steps into filters, mark the slowest one, and ask where files are stored and who is allowed to delete them.

On your CV, do not just write "used a message queue". Spell out how you handled duplicate delivery, where you placed the dead-letter queue and who was responsible for cleaning up payloads.

The book Enterprise Integration Patterns frames this pattern around a question: how do you process a message through several complex steps while keeping those steps independent and flexible? The exercise above shows the answer: each filter knows only its own input and output schemas, so you can reorder steps or add instances without rewriting any filter.

**Try this week:**

- Run the simulation code from this guide, push the same message into the store queue twice, and check whether the result is duplicated.
- Open a batch system you maintain, redraw it as filters and mark the slowest step. That is the first place to add instances.
- Add a line to your CV describing a pipeline you have built, stating how you handled duplicate delivery and the dead-letter queue.

## Sources

- [Pipes and Filters pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/pipes-and-filters)

- [Claim Check pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/claim-check)

- [Best Practices for Background Jobs - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/best-practices/background-jobs)

- [Asynchronous Request-Reply Pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/async-request-reply)

- [Pipes and Filters - Enterprise Integration Patterns](http://www.enterpriseintegrationpatterns.com/patterns/messaging/PipesAndFilters.html)
