KOURCHAL_
Menu
← All insights

Reliable OCR pipelines: transactional outbox and Redis Streams in OCRAgent

Why reliable document processing depends on committed jobs, recoverable queue delivery and explicit failure states. A source-backed look at the OCRAgent worker architecture.

A document pipeline can extract the right fields and still fail its operator. The upload may be stored while the queue is unavailable, or a worker may finish just before its acknowledgment fails. Ironclad OCR addresses these handoff problems with persisted jobs, a transactional outbox and Redis Streams. This article explains the reliability mechanisms in the engineering behind Kaliits.

The analysis uses public OCRAgent commit 647717e, checked on 2 October 2026. It examines source behavior, not a production load test or availability measurement. Kaliits should judge a pilot by whether the operator can account for each submitted case, including failures, rather than by whether one happy-path extraction completed.

The dual-write problem at intake

Consider an API that inserts a job in PostgreSQL and then sends a Redis message. If the database commit succeeds but Redis is unavailable, a stored document has no runnable message. Reversing the order creates another failure window: the worker can receive a message for a job that never committed. Two successful calls are not a shared transaction.

In DossierRepository.create_case_with_sources, the case, source artifacts, processing jobs and corresponding outbox events are written inside conn.transaction(). The ingestion audit event is part of that transaction too. Once it commits, PostgreSQL contains both the work to perform and the intention to publish it. For a Kaliits document workflow, that is a stronger foundation than treating an upload response as proof of processing.

The outbox dispatcher turns intent into a queue message

The dispatcher asks the repository to claim pending outbox events. The claim query selects available events using FOR UPDATE SKIP LOCKED, then marks them PUBLISHING with a lock timestamp and increased attempt count. SKIP LOCKED lets another dispatcher claim other unlocked rows without waiting for the same batch.

For each claimed event, the dispatcher calls queue.enqueue_payload and then mark_outbox_published. If either operation raises, it releases the event with the recorded error. The claim routine also resets PUBLISHING events whose lock is older than five minutes. This gives an abandoned publication attempt a route back to pending work.

Why publication can still happen twice

Suppose Redis accepts the message, then the dispatcher crashes before PostgreSQL records publication. A later dispatcher may publish the same event again. The outbox closes the missing-intent gap; it does not create an atomic transaction spanning PostgreSQL and Redis. Its consumer must tolerate duplicate delivery.

This is an important limit for Kaliits integrations. At-least-once transport is compatible with dependable processing when durable identities and guarded writes are used. It is not evidence of exactly-once execution, and it should never be advertised as a guarantee that every possible external side effect can happen only once.

What Redis Streams does here

RedisQueue publishes serialized payloads with XADD to ironclad:jobs. It creates the ironclad-workers consumer group if needed, reads new messages with XREADGROUP and acknowledges them with XACK. Consumer names include the hostname and process ID. Messages remain pending for the group until acknowledged.

The worker periodically calls claim_idle_messages, implemented with XAUTOCLAIM, to recover messages that have been idle beyond its configured recovery threshold. The inspected worker uses a ten-minute idle threshold and checks recovery about every thirty seconds. Those are implementation settings, not an appropriate latency promise for every document workload.

For Kaliits, reclaim timing needs to be evaluated alongside actual PDF duration. A long-running OCR request and a dead worker are different conditions. Recovery parameters, worker count and provider timeouts require workload testing; the existence of XAUTOCLAIM alone does not establish that concurrency is tuned correctly.

Persist the outcome before acknowledging delivery

The dossier handler calls DossierProcessor.process_source. RetryableProviderError, timeout, connection and PostgreSQL errors go through schedule_retry. Other exceptions go through mark_job_failed. Only after those operations return does the handler acknowledge the Redis message. If storing a retry or failure raises, the handler returns without acknowledgment so recovery remains possible.

schedule_retry locks the job row and records the next attempt in a database transaction. It resets job and artifact state and inserts a delayed outbox event keyed by the job and attempt. Delays start at thirty seconds, then two minutes, then ten minutes; retries stop at the job's attempt limit. A temporary provider problem becomes visible scheduled work rather than an untracked exception.

Duplicate guards live in durable state

Before processing a source, DossierProcessor loads the stored job under its client identity. A COMPLETED or FAILED job returns without repeating extraction. The processor also checks that the queue's case and workflow match the stored job. These checks reject a mismatched handoff and avoid redoing a terminal job after a late redelivery.

save_decision locks the case row and checks its processing state before creating the next decision version. It writes the decision, report payload, pending review and audit event in one transaction. Notifications use a deduplication key. These are concrete guards at selected persistence boundaries, not proof that every stage of the program is immune to concurrent duplicate work.

A fictional failure timeline

An operator uploads a packet, and PostgreSQL commits its case, jobs and outbox entries. Redis is temporarily unavailable, so publication is released for another attempt. Once Redis accepts the message, a worker starts OCR. If the provider times out, the worker stores a delayed retry and acknowledges the original delivery. The dispatcher later publishes that retry. This illustrates code paths a Kaliits pilot should exercise deliberately.

A different test should crash the dispatcher after publication but before marking the event published. Another should fail the acknowledgment after storing results. These tests explore the actual handoff windows. A successful health endpoint cannot substitute for either test, because health says little about a submitted case's history.

Separate technical failure from a business discrepancy

FAILED means processing could not produce the intended outcome. BLOCKED is a business finding generated from document rules. REVIEW_REQUIRED asks a person to resolve uncertainty or policy. Kaliits needs this distinction to keep operational responses sensible: repairing a service timeout is different from asking a supplier to correct a reference.

The LangGraph checkpoint boundary

The worker builds the legacy invoice graph with InMemorySaver. That object is process-local checkpoint memory. PostgreSQL job records and Redis messages have their own persistence, but the graph checkpoint itself is not a durable cross-restart store. The dossier route uses DossierProcessor rather than that graph. A reliability explanation must distinguish these responsibilities.

When this architecture is worth the added parts

The outbox adds a table, a dispatcher and monitoring responsibilities. A small internal workflow may be better served by an existing durable job facility in its platform. Ironclad's design makes sense to study when document work must survive asynchronous handoffs and be accounted for independently of one HTTP request. Kaliits should first check the existing software's import and job capabilities before proposing another service.

Before wider use, measure outbox age, retry volume, pending-message age, terminal failures and time spent awaiting review. Those are suggested pilot measurements, not dashboards claimed to exist in this repository. Recovery paths and retention settings need verification in the intended environment.

Read the reliability code

Transactional intake, retries and decision persistence

Outbox dispatcher

Redis Streams adapter

Worker routing and acknowledgment

Related architecture notes

The full evidence-backed OCR pipeline

Workflow rules and versioned human review

Scope a document-processing pilot with Kaliits