How Ray Powers Distributed Proof Search in LeanAgent: Architecture and Configuration
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: 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 (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 at line 52, the code explicitly configures:
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 (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, the system initializes the pool:
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:
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:
# 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) 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
BestFirstSearchProverinstances across distributed hardware. - Configuration occurs through environment variables (
RAY_TMPDIRinleanagent.pyline 52) and actor resource options (num_gpusallocation inproof_search.pylines 93-100). - The
ActorPoolabstraction (line 515) manages worker distribution, whilesearch_unorderedhandles asynchronous job dispatch. - Fault tolerance is implemented by catching
RayActorErrorexceptions (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_workersparameter.
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 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 (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 (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.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →