# How the Fully Async RL Pipeline in Miles Decouples Rollout and Training Workers

> Discover how the Miles RL pipeline decouples rollout and training workers using a persistent background Ray actor and queue based buffer for efficient asynchronous data fetching and training.

- Repository: [RadixArk/miles](https://github.com/radixark/miles)
- Tags: internals
- Published: 2026-09-06

---

**The Fully Async RL pipeline in Miles decouples rollout and training workers by running environment rollouts in a persistent background Ray actor with a queue-based data buffer, while the training loop fetches batches asynchronously and overlaps the next rollout preparation with current training steps.**

Miles (`radixark/miles`) is an open-source RL framework that implements a **fully-asynchronous** reinforcement learning architecture. Unlike synchronous pipelines where generation and training block each other, Miles separates these concerns through three interconnected components: a dedicated rollout executor, an inference controller, and independent training models. This design eliminates global barriers and maximizes GPU utilization.

## Core Components of the Fully Async RL Pipeline

### Rollout Executor with Persistent Background Worker

The `rollout_executor` runs as a **Ray actor** that continuously generates rollout data—token sequences, rewards, and other trajectory information—without pausing for the training loop.

Key implementation details from [`miles/rollout/fully_async_rollout.py`](https://github.com/radixark/miles/blob/main/miles/rollout/fully_async_rollout.py):

- The executor owns a **`FullyAsyncDataBuffer`** that queues generated rollout groups
- This buffer smooths production-consumption mismatches, allowing the trainer to fetch data at its own pace
- Metrics including queue length, staleness, and dropped groups are exposed for monitoring under `rollout/fully_async/...`

```python

# From fully_async_rollout.py — persistent worker with data buffer

class FullyAsyncRolloutFn:
    def __init__(self):
        self.data_buffer = FullyAsyncDataBuffer()  # Decouples producer/consumer

        # ... worker runs continuously in background Ray actor

```

### Inference Controller for Coordination

The **`inference_controller`** orchestrates between the trainer and executor without forcing synchronization points:

| Function | Purpose |
|----------|---------|
| `prepare_rollout()` | Prepares the model state for the next rollout generation |
| `update_weights()` primitive | Pushes policy updates asynchronously to the executor |

The controller exposes a thin RPC surface. The trainer requests rollouts via `rollout_executor.get.remote()` without blocking on weight synchronization.

### Training Models with Async Weight Updates

The **`actor_model`** and **`critic_model`** perform RL training steps on rollout data fetched from the executor. They receive policy weights through the **same asynchronous channel** used by the executor, ensuring approximate synchronization without global barriers.

## Pipeline Execution Flow in train_async.py

The driver script [`train_async.py`](https://github.com/radixark/miles/blob/main/train_async.py) implements the following sequence:

### 1. Component Initialization (Lines 35-38)

```python
inference_controller, rollout_executor, num_rollout = await create_rollout_components(args)
actor_model, critic_model = await create_training_models(
    args, inference_controller, rollout_executor
)

```

### 2. Initial Weight Synchronization (Line 53)

```python
await update_weights(actor_model, rollout_executor)  # Non-blocking push to executor

```

### 3. Main Async Loop with Overlapped Operations (Lines 76-89)

```python

# Eagerly start first rollout preparation

rollout_future = await eager_create_task(prepare_and_generate(start_rollout_id))

for rid in range(start_rollout_id, args.num_rollout):
    # Await the rollout produced by background worker

    rollout_data = await rollout_future
    
    # Immediately kick off next rollout (overlaps with training)

    if rid + 1 < args.num_rollout:
        rollout_future = await eager_create_task(prepare_and_generate(rid + 1))
    
    # Train on current rollout while next one generates

    await actor_model.train(rid, rollout_data)

```

### 4. Periodic Weight Updates (Line 119)

```python
if (rid + 1) % args.update_weights_interval == 0:
    await update_weights(actor_model, rollout_executor, rollout_id=rid)

```

The `eager_create_task` primitive is critical: it ensures the **next rollout preparation begins before training completes**, achieving true pipelining.

## Key Decoupling Mechanisms

| Mechanism | How Decoupling Works |
|-----------|----------------------|
| **Ray actor isolation** | Rollout executor runs in separate process/thread from trainer |
| **`FullyAsyncDataBuffer`** | Queue-based buffer absorbs rate mismatches between producer and consumer |
| **Asynchronous `get.remote()`** | Training loop pulls batches without blocking on rollout generation |
| **`update_weights` channel** | Weight pushes happen independently; no global synchronization barrier |
| **Eager task creation** | Next rollout starts early, overlapping generation with training |

## Complete Implementation Example

```python
import asyncio
from miles.rollout.fully_async_rollout import FullyAsyncRolloutFn
from miles.train_async import (
    create_rollout_components,
    create_training_models,
    update_weights,
    eager_create_task
)

async def fully_async_training_loop(args):
    # Build async components

    inference_controller, rollout_executor, num_rollout = await create_rollout_components(args)
    
    # Create models with references to executor for data fetching

    actor_model, critic_model = await create_training_models(
        args, inference_controller, rollout_executor
    )
    
    # Initial policy sync to executor

    await update_weights(actor_model, rollout_executor)
    
    # Overlapped rollout generation and training

    rollout_future = await eager_create_task(
        inference_controller.prepare_rollout(start_rollout_id)
        and rollout_executor.get.remote(start_rollout_id)
    )
    
    for rid in range(start_rollout_id, args.num_rollout):
        # Consume rollout from buffer

        rollout_data = await rollout_future
        
        # Start next rollout immediately (non-blocking)

        if rid + 1 < args.num_rollout:
            rollout_future = await eager_create_task(
                inference_controller.prepare_rollout(rid + 1)
                and rollout_executor.get.remote(rid + 1)
            )
        
        # Train while next rollout generates in background

        await actor_model.train(rid, rollout_data)
        await critic_model.train(rid, rollout_data)
        
        # Periodic weight push without stopping executor

        if (rid + 1) % args.update_weights_interval == 0:
            await update_weights(actor_model, rollout_executor, rollout_id=rid)
    
    return actor_model

```

## Relevant Source Files

| File Path | Description |
|-----------|-------------|
| [`miles/rollout/fully_async_rollout.py`](https://github.com/radixark/miles/blob/main/miles/rollout/fully_async_rollout.py) | `FullyAsyncRolloutFn` class with persistent worker and data buffer logic |
| [`miles/rollout/fully_async_data_buffer.py`](https://github.com/radixark/miles/blob/main/miles/rollout/fully_async_data_buffer.py) | Queue implementation tracking staleness, queue length, dropped groups |
| [`train_async.py`](https://github.com/radixark/miles/blob/main/train_async.py) | Driver orchestrating inference controller, executor, and training models |
| [`miles/utils/arguments.py`](https://github.com/radixark/miles/blob/main/miles/utils/arguments.py) | `--fully-async` flag parsing and rollout implementation selection |
| [`tests/fast/rollout/test_fully_async_rollout.py`](https://github.com/radixark/miles/blob/main/tests/fast/rollout/test_fully_async_rollout.py) | Unit tests for buffer metrics and staleness handling |

## Summary

- **Rollout executor** runs as persistent Ray actor with dedicated `FullyAsyncDataBuffer` queue
- **Inference controller** coordinates without blocking through `prepare_rollout` and `update_weights` primitives
- **Training loop** fetches batches asynchronously via `get.remote()` and overlaps next rollout with current training
- **Weight updates** flow through async channel, avoiding global barriers while keeping policy roughly synchronized
- **`eager_create_task`** enables true pipelining by starting rollouts before training completes

## Frequently Asked Questions

### How does Miles prevent the training loop from starving when rollout generation is slow?

The `FullyAsyncDataBuffer` maintains a queue of generated rollouts. The training loop consumes from this buffer asynchronously via `await rollout_executor.get.remote()`. If generation lags, the trainer simply waits; if generation runs ahead, excess rollouts queue up. The buffer exposes metrics (`rollout/fully_async/queue_length`) to monitor this balance.

### What happens to rollouts generated with stale policy weights?

Miles tracks **staleness** explicitly in [`fully_async_data_buffer.py`](https://github.com/radixark/miles/blob/main/fully_async_data_buffer.py). Each rollout records its generation policy version, and the buffer calculates the gap between consumed rollout version and current trainer version. While some staleness is inherent to async training, the `update_weights_interval` parameter bounds it—weights push periodically to refresh the executor's policy.

### Can the Fully Async pipeline run on a single GPU?

Yes, though the design targets multi-GPU scenarios. On single GPU, the Ray actor executor and trainer share the device with time-slicing. The `prepare_rollout` call in the inference controller handles resource allocation. The async structure still provides latency hiding benefits even without true parallelism.

### How does `eager_create_task` differ from standard `asyncio.create_task`?

`eager_create_task` (implemented in Miles utilities) ensures the task begins execution immediately rather than waiting for the next event loop iteration. This minimizes the gap between requesting and starting the next rollout, tightening the overlap between generation and training.