# How Ray Powers Distributed Proof Search in LeanAgent: Architecture and Configuration

> Discover how LeanAgent leverages Ray for distributed proof search. Explore its architecture and configuration for efficient parallel processing using a fault-tolerant actor model.

- Repository: [LeanDojo/leanagent](https://github.com/lean-dojo/leanagent)
- Tags: architecture
- Published: 2026-03-05

---

**LeanAgent uses Ray as the backbone of its distributed architecture to parallelize expensive proof-search workloads across multiple CPUs or GPUs using a fault-tolerant actor model.**

The `lean-dojo/leanagent` repository implements a distributed theorem-proving system that scales beyond single-process limitations. By integrating Ray, LeanAgent transforms the `BestFirstSearchProver` algorithm into a resilient, parallel computation engine capable of exploring multiple Lean theorems simultaneously. This article examines the specific role Ray plays in the architecture and how the codebase configures it for high-performance proof search.

## The Role of Ray in Distributed Proof Search

Ray provides the **actor model** foundation that enables LeanAgent to distribute proof-search tasks across heterogeneous hardware. Rather than managing threads or processes manually, the system delegates orchestration to Ray's distributed runtime, which handles scheduling, resource allocation, and failure recovery automatically.

### Actor-Based Parallelism

At the core of the architecture are two specialized Ray actors defined in [`prover/proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/prover/proof_search.py): `ProverActor` and `VllmActor`. Each actor encapsulates an independent instance of `BestFirstSearchProver` or a vLLM inference engine, respectively. These actors are decorated with `@ray.remote`, allowing them to execute concurrently on available cluster resources. The `DistributedProver` class creates a pool of these remote actors and dispatches proof-search jobs asynchronously, enabling parallel exploration of different theorems without blocking the main process.

### Fault Tolerance and Resilience

Ray's built-in fault tolerance ensures that hardware failures or actor crashes do not corrupt entire proof-search runs. When an actor crashes, Ray raises a `ray.exceptions.RayActorError`, which LeanAgent catches and handles explicitly. According to the source code in [`prover/proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/prover/proof_search.py) (lines 36-38), the system logs the error and terminates the entire run with `sys.exit(1)` to prevent silent data corruption. This design guarantees that partial failures are detected immediately rather than allowing inconsistent proof states to propagate.

## Configuring Ray for LeanAgent

LeanAgent configures Ray through environment variables, resource allocation strategies, and actor pool initialization. These settings determine how the system utilizes available compute resources, whether running on a single multi-core machine or a multi-node cluster.

### Environment Variables and Temporary Storage

Before initializing Ray, LeanAgent sets the `RAY_TMPDIR` environment variable to a high-performance workspace directory. In [`leanagent.py`](https://github.com/lean-dojo/leanagent/blob/main/leanagent.py) at line 52, the code explicitly configures:

```python
os.environ['RAY_TMPDIR'] = f"{RAID_DIR}/tmp"

```

This ensures that Ray's temporary files—including object store data and spillover storage—reside on fast RAID storage rather than default system temporary directories, which is critical for I/O-intensive proof search operations.

### GPU Resource Allocation

When GPUs are available, LeanAgent calculates `num_gpus_per_worker` by dividing the total GPU count by the number of workers. In [`prover/proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/prover/proof_search.py) (lines 93-100), the system creates actors with `ProverActor.options(num_gpus=...)` or `VllmActor.options(num_gpus=...)` depending on whether GPU-accelerated vLLM inference is enabled. For vLLM configurations, GPU handling is delegated entirely to the `VllmActor`, while standard CPU-based workers receive `num_gpus=0`.

### Actor Pool Initialization

The `DistributedProver` aggregates remote actors into an `ActorPool` for efficient job distribution. At line 515 of [`proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/proof_search.py), the system initializes the pool:

```python
self.prover_pool = ActorPool(provers)

```

This pool manages a collection of `ProverActor` or `VllmActor` instances, enabling the `search_unordered` method to dispatch theorem positions to available workers without manual load balancing.

## Distributed Proof Search Workflow

The implementation follows a clear workflow that separates resource initialization from execution. When `num_workers` exceeds 1, LeanAgent enters distributed mode; otherwise, it falls back to single-process execution.

### CPU-Based Configuration

For CPU-only environments, instantiate `DistributedProver` with `use_vllm=False` and specify the number of worker processes:

```python
from prover.proof_search import DistributedProver

# Configure 4 CPU workers without GPU acceleration

prover = DistributedProver(
    use_vllm=False,
    ckpt_path=None,
    indexed_corpus_path=None,
    tactic="my_tactic",
    module=None,
    num_workers=4,
    num_gpus=0,
    timeout=300,
    max_expansions=None,
    num_sampled_tactics=5,
    raid_dir="/path/to/raid",
    checkpoint_dir="checkpoints",
    debug=False,
)

# Execute proof search across multiple theorems

results = prover.search_unordered(
    repo=my_lean_repo,
    theorems=[thm1, thm2, thm3],
    positions=[pos1, pos2, pos3],
)

```

### GPU-Accelerated Configuration

For vLLM-based inference, configure the prover to distribute GPU resources among workers:

```python

# Configure 2 workers sharing 4 GPUs for vLLM inference

prover = DistributedProver(
    use_vllm=True,
    ckpt_path="model.ckpt",
    indexed_corpus_path=None,
    tactic=None,
    module=None,
    num_workers=2,
    num_gpus=4,               # Divided as 2 GPUs per worker

    timeout=300,
    max_expansions=None,
    num_sampled_tactics=5,
    raid_dir="/path/to/raid",
    checkpoint_dir="checkpoints",
    debug=False,
)

```

The `search_unordered` method (lines 518-536 in [`proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/proof_search.py)) submits jobs to the actor pool via `p.search.remote(repo, thm, pos)`, allowing Ray to schedule these calls across the available CPU cores or GPUs while handling retries and actor failures automatically.

## Summary

- **Ray provides the actor model** that enables LeanAgent to parallelize `BestFirstSearchProver` instances across distributed hardware.
- **Configuration occurs through environment variables** (`RAY_TMPDIR` in [`leanagent.py`](https://github.com/lean-dojo/leanagent/blob/main/leanagent.py) line 52) and actor resource options (`num_gpus` allocation in [`proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/proof_search.py) lines 93-100).
- **The `ActorPool` abstraction** (line 515) manages worker distribution, while `search_unordered` handles asynchronous job dispatch.
- **Fault tolerance is implemented** by catching `RayActorError` exceptions (lines 36-38) to prevent silent corruption of proof-search results.
- **Resource allocation scales dynamically** from single-process fallback to multi-GPU distributed execution based on the `num_workers` parameter.

## Frequently Asked Questions

### What is the role of Ray in LeanAgent's distributed architecture?

Ray serves as the distributed computing engine that manages the actor-based parallelism in LeanAgent. It provides the `@ray.remote` infrastructure for `ProverActor` and `VllmActor` classes, handles resource allocation across CPUs and GPUs, and offers automatic fault tolerance through its exception handling mechanisms. Without Ray, LeanAgent would require manual process management and custom load balancing logic.

### How does LeanAgent configure Ray's temporary directory?

LeanAgent explicitly sets the `RAY_TMPDIR` environment variable before Ray initialization to ensure temporary files are stored on high-performance RAID storage. In [`leanagent.py`](https://github.com/lean-dojo/leanagent/blob/main/leanagent.py) at line 52, the code assigns `os.environ['RAY_TMPDIR'] = f"{RAID_DIR}/tmp"`, preventing I/O bottlenecks that could occur if Ray used default system temporary directories during large-scale proof search operations.

### How are GPU resources allocated in LeanAgent's Ray setup?

GPU resources are allocated by dividing the total available GPUs by the number of workers (`num_gpus_per_worker = num_gpus / num_workers`). When creating actors in [`proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/proof_search.py) (lines 93-100), the system passes this calculated value to `ProverActor.options(num_gpus=...)` or `VllmActor.options(num_gpus=...)`. For vLLM configurations, each actor manages its own GPU resources independently, while standard CPU workers receive zero GPU allocation.

### How does LeanAgent handle actor failures in Ray?

LeanAgent implements strict fault tolerance by catching `ray.exceptions.RayActorError` when remote actors crash during execution. According to [`prover/proof_search.py`](https://github.com/lean-dojo/leanagent/blob/main/prover/proof_search.py) (lines 36-38), the system logs the specific error and immediately terminates the entire run using `sys.exit(1)`. This prevents partial proof results from corrupting the dataset, ensuring that any worker failure results in a complete, detectable failure rather than silent data loss.