# Architecture of the S3-Based Distributed Work Queue in OLMocr

> Explore the S3-based distributed work queue architecture in allenai/olmocr. Discover a fault-tolerant, serverless system for parallel processing using only Amazon S3 object storage.

- Repository: [Ai2/olmocr](https://github.com/allenai/olmocr)
- Tags: architecture
- Published: 2026-07-05

---

**The S3-based distributed work queue system in the allenai/olmocr repository implements a lightweight, serverless coordination mechanism using only Amazon S3 object storage, enabling fault-tolerant parallel processing without external databases or message brokers.**

The `olmocr` project uses a custom **S3-based distributed work queue** to coordinate OCR workloads across multiple workers. Unlike traditional job queues that require Redis, RabbitMQ, or PostgreSQL, this architecture stores all state—pending work, active locks, and completion markers—as objects in a shared S3 bucket. By leveraging atomic S3 operations and object metadata, the system provides strong consistency guarantees while eliminating infrastructure dependencies.

## Core Components of the Work Queue System

### The Backend Abstraction Layer

Located in [`olmocr/work_queue.py`](https://github.com/allenai/olmocr/blob/main/olmocr/work_queue.py), the abstract `Backend` class defines the storage-agnostic API required by the queue. The concrete `S3Backend` implementation uses `boto3` to map queue operations to native S3 calls, handling object creation, deletion, and metadata queries through simple HTTP operations.

### The WorkQueue Orchestrator

The `WorkQueue` class manages the lifecycle of distributed tasks. It handles queue population via `populate_queue()`, distributes work through `get_work()`, and tracks completion using `mark_done()`. This orchestrator maintains an in-memory queue of `WorkItem` objects while delegating persistence to the backend.

### S3 Data Layout and Object Structure

Within a configured `workspace_path` (e.g., `s3://bucket/olmocr-workspace`), the system organizes state into three specific locations:

```

workspace_path/
├── work_index_list.csv.zstd      # Compressed CSV mapping hash → file paths

├── worker_locks/                 # Active worker locks: worker_<hash>.lock

└── done_flags/                   # Completion markers: done_<hash>.flag

```

The directory constants are defined in [`work_queue.py`](https://github.com/allenai/olmocr/blob/main/work_queue.py):

```python
WORKER_LOCKS_DIR = "worker_locks"
DONE_FLAGS_DIR = "done_flags"

```

## How the Distributed Queue Works

### Queue Population and Work Group Indexing

The `populate_queue()` method in `WorkQueue` processes input file lists by grouping them into batches (configurable via `items_per_group`). It computes a deterministic SHA-1 hash for each group using `_compute_workgroup_hash()`, then appends these entries to `work_index_list.csv.zstd`. The helper functions in [`olmocr/s3_utils.py`](https://github.com/allenai/olmocr/blob/main/olmocr/s3_utils.py) handle zstd compression and comma-escaping for paths containing literal commas.

### Worker Coordination and Lock Acquisition

When a worker calls `get_work()`, the system executes a distributed consensus protocol:

1. **Completion Check**: Queries `backend.is_completed()` (a HEAD request for the corresponding file in `done_flags/`)
2. **Lock Verification**: Checks `backend.is_worker_lock_taken()` by examining the `LastModified` timestamp of `worker_<hash>.lock` against a 30-minute timeout
3. **Lock Acquisition**: Creates an empty object via `backend.create_worker_lock()` if the existing lock is stale or absent

### Marking Completion and Releasing Locks

After processing a work group, the worker invokes `mark_done()`, which calls `backend.create_done_flag()` to create an empty object in `done_flags/`. The worker then releases its claim via `backend.delete_worker_lock()`. This atomic write serves as the immutable source of truth for work completion across all nodes.

## Distributed Guarantees and Fault Tolerance

The architecture provides several critical guarantees for distributed systems:

- **Mutual Exclusion**: Only one worker can hold a valid lock for a work group due to S3's atomic PUT semantics
- **Fault Tolerance**: Crashed workers leave stale locks that are automatically reclaimed after the timeout period based on `LastModified` metadata
- **Idempotency**: Completion markers are idempotent; duplicate flag creation is safe and prevents double-processing
- **Scalability**: No central coordinator bottleneck exists; workers communicate only through S3 object operations

## Implementation Example

Here is a complete example showing how to initialize the queue and run a worker loop:

```python
import boto3
from olmocr.work_queue import S3Backend, WorkQueue

# Initialize the S3 backend

s3_client = boto3.client("s3")
backend = S3Backend(s3_client, "s3://my-bucket/olmocr-workspace")
queue = WorkQueue(backend)

# Populate the queue (run once before workers start)

input_paths = [
    "s3://my-bucket/input/doc1.pdf",
    "s3://my-bucket/input/doc2.pdf",
    # ...

]
await queue.populate_queue(input_paths, items_per_group=5)

# Worker processing loop

async def worker_loop():
    await queue.initialize_queue()
    while True:
        item = await queue.get_work()
        if item is None:
            break  # No more work available

        
        # Process files in item.work_paths

        for pdf_path in item.work_paths:
            # Perform OCR processing here

            pass
            
        # Mark completion and release lock

        await queue.mark_done(item)

# Run multiple workers across processes or machines

```

## Summary

- The **S3-based distributed work queue** in `olmocr` uses object storage as the sole coordination mechanism, requiring no external databases
- State persists in three S3 locations: a zstd-compressed CSV index, a `worker_locks/` directory for active claims, and a `done_flags/` directory for completion tracking
- The `S3Backend` class in [`olmocr/work_queue.py`](https://github.com/allenai/olmocr/blob/main/olmocr/work_queue.py) implements distributed locking via atomic empty-object creation with timestamp-based stale detection
- Workers acquire exclusive access using `create_worker_lock()` and finalize work via `create_done_flag()`, ensuring exactly-once processing semantics
- Failed workers are handled automatically through lock timeouts, making the system resilient to individual node failures

## Frequently Asked Questions

### How does the system handle worker crashes?

If a worker crashes mid-processing, its lock object remains in the `worker_locks/` directory indefinitely. Other workers detect this stale lock by comparing the object's `LastModified` timestamp against the configured timeout (default 30 minutes). Once expired, a new worker can safely acquire the lock via `create_worker_lock()` and process the work group.

### What prevents two workers from processing the same work group?

The `S3Backend.is_worker_lock_taken()` method checks for the existence and freshness of lock objects before granting work. Since S3 PUT operations are atomic, only one worker can successfully create `worker_<hash>.lock`. Subsequent attempts fail until the lock is deleted via `delete_worker_lock()` or becomes stale after the timeout period.

### Can this architecture work with S3-compatible storage like MinIO or Wasabi?

Yes. The `S3Backend` class uses standard `boto3` operations. As long as the storage provider supports atomic PUT operations and object metadata queries (HEAD requests), the distributed queue functions correctly. The system has been tested against Amazon S3 and compatible object stores.

### How are work groups defined and hashed?

The `WorkQueue.populate_queue()` method groups input paths into batches of size `items_per_group`. It computes a SHA-1 hash of the sorted paths via `_compute_workgroup_hash()` to create a deterministic, unique identifier for each batch. This hash appears in the CSV index and forms the basis for lock filenames (`worker_<hash>.lock`) and completion flags (`done_<hash>.flag`).