Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ while this project adds the environment-shaped API on top.
- **`cmd/ate-env-api`** — The API service that manages environments and proxies remote guest requests.
- **`cmd/ate-env-guest`** — The daemon server running inside each actor serving command executions, file read/write, and built-in MCP tools.
- **`clients/go`** — The Go client library to manage environments, run commands, and perform file operations.
- **`clients/python`** — The async Python client library ([README](clients/python/README.md)).
- **`clients/python`** — The unified Python client and high-throughput Sandbox Fleet SDK ([README](clients/python/README.md)).
- **`integrations/nemo-gym`** — A [NeMo Gym](https://github.com/NVIDIA-NeMo/Gym) sandbox provider that runs rollout sandboxes as environments, built on the Python client ([README](integrations/nemo-gym/README.md)).

## Installation
Expand Down
9 changes: 5 additions & 4 deletions clients/python/README.md
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
# ate-env-client — Async Python Client
# ate-env-client — Python Client & Sandbox Fleet SDK

> [!WARNING]
> This is an alpha API and is likely to change until v1.0 is released.

Async Python client for the [Agent Substrate Environment](../../README.md)
API (`ate-env-api`): environment lifecycle, remote command execution, and
streaming file I/O over gRPC.
Unified Python client and high-throughput Sandbox Fleet SDK for the [Agent Substrate Environment](../../README.md)
API (`ate-env-api`):
1. **Single Environment Lifecycle & Guest Operations**: `Client` & `Env` for fine-grained gRPC control, process execution, and file streaming.
2. **High-Throughput Fleet & Sandbox Orchestration**: `SandboxFleet` & `AsyncSandboxFleet` for massive RL rollouts (Ray, VeRL, NeMo Gym) with automated pooling, pre-warming, and concurrency control.

Requires Python >= 3.10. Everything is `asyncio`-native: methods are
coroutines, log/file streams are async iterators, and cancellation works
Expand Down
100 changes: 100 additions & 0 deletions clients/python/USER_GUIDE.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
# Reinforcement Learning & Sandbox Fleet Guide

This guide covers how to use high-throughput **Sandbox Fleet Orchestration (`SandboxFleet` & `AsyncSandboxFleet`)** in `ate_env` with distributed Reinforcement Learning frameworks like **VeRL**, **Ray**, and **NeMo-Gym** on top of **Agent Substrate** and **GKE**.

---

## 1. Quickstart

### Synchronous Batch Execution

```python
from ate_env import FleetConfig, SandboxFleet, Task

config = FleetConfig(
backend="substrate",
endpoint="http://localhost:7777",
data_plane="ate_env", # or "router"
)

tasks = [
Task(id="task-1", image="docker.io/library/python:3.11"),
Task(id="task-2", image="docker.io/library/python:3.11"),
]

with SandboxFleet(config) as fleet:
sandboxes = fleet.acquire(tasks)
try:
for task, sb in zip(tasks, sandboxes):
res = sb.exec("python3 -c 'print(\"hello world\")'")
print(f"Task {task.id}: {res.stdout.strip()} (exit_code={res.exit_code})")
finally:
fleet.release(sandboxes)
```

---

## 2. Asynchronous Non-Blocking Rollouts (VeRL / vLLM)

In RL post-training (e.g., GRPO), candidate generation overlaps with sandbox pre-warming to hide provisioning latency:

```python
import asyncio
from ate_env import AsyncSandboxFleet, FleetConfig, Task

async def run_rollouts():
config = FleetConfig(
backend="substrate",
data_plane="ate_env",
endpoint="http://ateapi.ate-system.svc.cluster.local:8080",
grpc_endpoint="substrate-env.ate-system.svc.cluster.local:50051",
batch_size=4,
max_warmpool_replicas=4,
)

tasks = [Task(id=f"t-{i}", image="docker.io/library/python:3.11") for i in range(4)]

fleet = AsyncSandboxFleet(config)
await fleet.setup(tasks)

# 1. Start acquiring pre-warmed sandboxes concurrently while LLM generates tokens
warm_task = asyncio.create_task(fleet.acquire_batch(tasks))
await asyncio.sleep(1.5) # Simulate token generation
sandboxes = await warm_task

# 2. Asynchronously evaluate candidate rollouts
for sb in sandboxes:
async with sb:
await sb.write_file_async("/testbed/calc.py", "def add(a, b): return a + b\n")
res = await sb.exec_async(["python3", "-c", "import calc; assert calc.add(2, 3) == 5"])
print(f"Sandbox {sb.sandbox_id}: ok={res.ok}")

await fleet.teardown()

asyncio.run(run_rollouts())
```

---

## 3. Data Plane Options

| `data_plane` | Protocol | Target | Description |
| :--- | :--- | :--- | :--- |
| **`"ate_env"`** *(default)* | gRPC | `substrate-env:50051` | Native in-guest `ProcessService` and `FileSystemService` streaming. |
| **`"router"`** | HTTP | `atenet-router:8080` | Reverse proxy routing with `ate-target-actor` HTTP headers. |

---

## 4. End-to-End VeRL + Ray Demo on GKE

A full runnable example with Ray actors and GRPO advantage optimization is located in [`examples/verl_swebench`](../../examples/verl_swebench):

* **Local hermetic run**:
```bash
python examples/verl_swebench/async_verl_swebench_pipeline.py --backend mock --num-iters 2 --group-size 2
```
* **Production RayJob on GKE**:
```bash
kubectl apply -f examples/verl_swebench/ray-job.async-verl.yaml
```
See [`examples/verl_swebench/README.md`](../../examples/verl_swebench/README.md) for full deployment instructions.
133 changes: 133 additions & 0 deletions clients/python/examples/poc_rollout_demo.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
#!/usr/bin/env python3
# Copyright 2026 The Kubernetes Authors & Google LLC
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""
Interactive Proof of Concept (PoC) Runner for sandbox-sdk.

Demonstrates:
1. Sizing and double-buffered pipelined prefetching across a task batch.
2. Sandbox acquisition through SandboxFleet (mock backend; Substrate latency not yet measured).
3. Trajectory evaluation for SWE-bench (pytest-5221) with pass/fail candidate patches.
4. Verification of NeMo-Gym provider integration (replacing PR #69).
"""

import sys
import time
from ate_env import FleetConfig, SandboxFleet
from ate_env.adapters.swebench import SWEBENCH_SAMPLE_TASK, SweBenchAdapter
from ate_env.providers.nemo_gym import UnifiedSandboxProvider


def run_poc():
print("=" * 70)
print(" 🚀 SANDBOX-SDK: PROOF OF CONCEPT (PoC) ROLLOUT RUNNER")
print("=" * 70)

# --------------------------------------------------------------------------
# 1. Initialize Fleet with Pipelined Windowing Strategy
# --------------------------------------------------------------------------
print("\n[Step 1] Initializing SandboxFleet with pipelined strategy...")
config = FleetConfig(
backend="mock", # Toggleable to 'substrate' or 'kubernetes'
strategy="pipelined",
batch_size=2,
max_warmpool_replicas=2,
tenancy="rl-genetics-poc",
worker_family="c2",
)
fleet = SandboxFleet(config)

# Prepare batch of tasks
tasks = [
SweBenchAdapter.to_task({
**SWEBENCH_SAMPLE_TASK,
"task_id": f"pytest-dev__pytest-5221-sample-{i}"
})
for i in range(4)
]
fleet.load_tasks(tasks)
fleet.setup()

print(f"✔ Fleet initialized with Run ID: {fleet.run_id}")
print(f"✔ Planned {len(tasks)} tasks across unique images.")

# --------------------------------------------------------------------------
# 2. Simulate RL Policy Rollouts (Passing vs. Buggy Candidates)
# --------------------------------------------------------------------------
print("\n[Step 2] Executing post-training candidate rollout evaluations...")

candidate_solution = (
"--- a/testing/test_helpconfig.py\n"
"+++ b/testing/test_helpconfig.py\n"
"@@ -1,3 +1,4 @@\n"
"+# Fixed fixture formatting\n"
)

def process_task(task, handle):
t0 = time.monotonic()
print(f" -> Acquired sandbox {handle.sandbox_id} for {task.id}")

# In-guest file write
handle.runtime.write_file("/tmp/solution.patch", candidate_solution)

# In-guest execution
handle.exec("git apply /tmp/solution.patch", cwd="/testbed")
res = handle.runtime.exec(task.metadata["test_cmd"], cwd="/testbed")

duration = time.monotonic() - t0
passed = res.ok
return {
"task_id": task.id,
"sandbox_id": handle.sandbox_id,
"passed": passed,
"reward": 1.0 if passed else 0.0,
"duration_s": round(duration, 3)
}

# fleet.run() returns each task's result, or the exception it raised.
results = fleet.run(process_task, concurrency=2)
print("\n✔ Rollout Execution Summary:")
for r in results:
if isinstance(r, Exception):
print(f" ⚠️ ERROR | {type(r).__name__}: {r}")
continue
status_icon = "✅ PASS" if r["passed"] else "❌ FAIL"
print(f" {status_icon} | Task: {r['task_id']:<35} | Duration: {r['duration_s']}s | Reward: {r['reward']}")

# --------------------------------------------------------------------------
# 3. Verify NeMo Gym Provider (Replacing PR #69)
# --------------------------------------------------------------------------
print("\n[Step 3] Verifying NeMo Gym Provider Contract (PR #69 Replacement)...")
nemo_provider = UnifiedSandboxProvider(config={"backend": "mock"})
spec = {
"id": "nemo-sample-ep-1",
"image": "sweb.eval.pytest-5221",
"files": {"/testbed/init.txt": "ready"}
}
handle = nemo_provider.create(spec)
exec_out = nemo_provider.exec(handle, "cat /testbed/init.txt")
print(f"✔ NeMo Gym Provider Exec output: {(exec_out['stdout'] or '').strip()} "
f"(return_code={exec_out['return_code']}, error_type={exec_out['error_type']})")
nemo_provider.close(handle)
print("✔ NeMo Gym Provider close succeeded.")

fleet.teardown()
print("\n" + "=" * 70)
print(" 🎉 ALL PROOF OF CONCEPT STAGES COMPLETED SUCCESSFULLY!")
print("=" * 70)


if __name__ == "__main__":
run_poc()
11 changes: 11 additions & 0 deletions clients/python/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ license = "Apache-2.0"
dependencies = [
"grpcio>=1.83",
"protobuf>=7.35",
"pydantic>=2.0",
]

[project.optional-dependencies]
Expand All @@ -34,6 +35,16 @@ dev = [
"pytest-asyncio>=1.0",
"grpcio-tools>=1.83",
]
kubernetes = [
"kubernetes>=24.0",
"websockets>=11.0",
]
ray = [
"ray>=2.10",
]

[project.entry-points."nemo_gym.sandbox_providers"]
unified_sandbox = "ate_env.providers.nemo_gym:UnifiedSandboxProvider"

[project.urls]
Homepage = "https://github.com/agent-substrate/env"
Expand Down
Loading
Loading