Skip to content
Merged
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
44 changes: 44 additions & 0 deletions configs/simulation/offline-eval-dedicated.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
{
"simulation": {
"schema": "rlinf.simulation.config.v1",
"runtime_id": "libero-eval",
"binding": "franka.panda.pi05",
"server": {
"bind": "127.0.0.1",
"port": 8200
},
"inference": {
"backend": "embodiinfer",
"transport": "http",
"endpoint": "http://127.0.0.1:8101",
"options": {}
},
"simulator": {
"id": "libero",
"type": "libero",
"options": {
"suite": "libero_10",
"task_id": 0,
"max_episode_steps": 480
}
}
},
"evaluation": {
"workers": 8,
"inference_endpoints": [
"http://127.0.0.1:8101"
],
"render_gpu_id": 0,
"worker_timeout_s": 600
},
"episodes": [
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-0", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 0, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-1", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 1, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-2", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 2, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-3", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 3, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-4", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 4, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-5", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 5, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-6", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 6, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-7", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 7, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120}
]
}
45 changes: 45 additions & 0 deletions configs/simulation/offline-eval-shared.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
{
"simulation": {
"schema": "rlinf.simulation.config.v1",
"runtime_id": "libero-eval",
"binding": "franka.panda.pi05",
"server": {
"bind": "127.0.0.1",
"port": 8200
},
"inference": {
"backend": "embodiinfer",
"transport": "http",
"endpoint": "http://127.0.0.1:8101",
"options": {}
},
"simulator": {
"id": "libero",
"type": "libero",
"options": {
"suite": "libero_10",
"task_id": 0,
"max_episode_steps": 480
}
}
},
"evaluation": {
"workers": 8,
"inference_endpoints": [
"http://127.0.0.1:8100",
"http://127.0.0.1:8101"
],
"render_gpu_id": 0,
"worker_timeout_s": 600
},
"episodes": [
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-0", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 0, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-1", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 1, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-2", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 2, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-3", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 3, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-4", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 4, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-5", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 5, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-6", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 6, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120},
{"schema": "rlinf.simulation.episode.v1", "request_id": "episode-7", "runtime_id": "libero-eval", "prompt": null, "task": "libero_10:0", "seed": 7, "chunk_steps": 5, "max_policy_steps": 96, "inference_timeout_s": 120}
]
}
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,7 @@ sim-vlabench = { requires-python = ">=3.12" }
embodirun = "embodirun.services.host.cli.cli:main"
embodirun-control-serve = "embodirun.services.control.server:main"
embodirun-simulation-serve = "embodirun.services.simulation.server:main"
embodirun-evaluate = "embodirun.services.simulation.evaluate:main"
embodirun-sglang-serve = "embodirun.services.inference.adapters.sglang.pi05:main"
embodirun-go2-streamvln = "embodirun.bindings.unitree.go2.streamvln.cli:main"
# Compatibility aliases for existing deployments.
Expand Down
221 changes: 221 additions & 0 deletions src/embodirun/application/simulation/evaluation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,221 @@
"""Parallel evaluation with explicit rendering placement and HTTP replicas."""

from __future__ import annotations

import dataclasses
import math
import multiprocessing as mp
import os
import time
import traceback
from collections.abc import Callable, Sequence
from contextlib import suppress
from dataclasses import dataclass
from multiprocessing.connection import Connection, wait
from typing import Any
from urllib.parse import urlsplit

from .contracts import EpisodeRequest, EpisodeResult, SimulationServiceConfig
from .service import SimulationService


@dataclass(frozen=True)
class EvaluationConfig:
"""Place simulator workers separately from already-running inference services.

Workers and their statically assigned episodes stay on one HTTP replica.
All endpoints must serve identical weights and adapter contracts. An empty
endpoint list uses the simulation service's existing endpoint. Rendering
placement applies to MuJoCo EGL; None preserves the simulator's own setup.
worker_timeout_s bounds startup and the time between completed episodes for
each worker, including reset and cleanup, not just an individual HTTP call.
"""

workers: int = 8
inference_endpoints: tuple[str, ...] = ()
render_gpu_id: int | None = None
worker_timeout_s: float = 600.0

def __post_init__(self) -> None:
if type(self.workers) is not int or self.workers < 1:
raise ValueError("workers must be a positive integer")
if self.render_gpu_id is not None and (type(self.render_gpu_id) is not int or self.render_gpu_id < 0):
raise ValueError("render_gpu_id must be a nonnegative integer or None")
if (
isinstance(self.worker_timeout_s, bool)
or not isinstance(self.worker_timeout_s, (int, float))
or not math.isfinite(self.worker_timeout_s)
or self.worker_timeout_s <= 0
):
raise ValueError("worker_timeout_s must be finite and positive")
if not isinstance(self.inference_endpoints, (list, tuple)):
raise ValueError("inference_endpoints must be a list or tuple of HTTP URLs")
endpoints = tuple(self.inference_endpoints)
for endpoint in endpoints:
if not isinstance(endpoint, str) or endpoint != endpoint.strip():
raise ValueError("inference_endpoints must contain HTTP URLs")
parsed = urlsplit(endpoint)
if (
parsed.scheme not in {"http", "https"}
or not parsed.hostname
or parsed.username is not None
or parsed.password is not None
or parsed.query
or parsed.fragment
):
raise ValueError("inference_endpoints must contain credential-free HTTP URLs")
_ = parsed.port
if len(set(endpoints)) != len(endpoints):
raise ValueError("inference_endpoints must be unique")
object.__setattr__(self, "inference_endpoints", endpoints)


def _worker(
service_config: SimulationServiceConfig,
render_gpu_id: int | None,
jobs: Sequence[tuple[int, EpisodeRequest]],
connection: Connection,
service_factory: Callable[[SimulationServiceConfig], SimulationService],
) -> None:
try:
if render_gpu_id is not None:
os.environ.update(MUJOCO_EGL_DEVICE_ID=str(render_gpu_id), MUJOCO_GL="egl", PYOPENGL_PLATFORM="egl")
service = service_factory(service_config)
try:
for index, request in jobs:
started = time.monotonic()
result = service.execute(request)
connection.send(("result", index, result.to_payload(), time.monotonic() - started))
finally:
service.close()
connection.send(("done",))
except BaseException:
with suppress(BrokenPipeError, OSError):
connection.send(("error", traceback.format_exc()))
finally:
connection.close()


class EvaluationRunner:
"""Run finite simulation jobs on a fixed pool of process-isolated workers.

Reuses SimulationService's episode/session lifecycle and original inference
clients. There is no per-step coordinator, CPU quota, or inference batching
implementation here. Failure aborts the run without retrying an episode on
another replica. Results are returned in input order only after all workers
close successfully. Custom service factories must be spawn-picklable.
"""

def __init__(
self,
service: SimulationServiceConfig,
config: EvaluationConfig,
*,
service_factory: Callable[[SimulationServiceConfig], SimulationService] = SimulationService,
) -> None:
if service.inference_transport != "http":
raise ValueError("parallel evaluation currently requires HTTP inference")
if (
config.render_gpu_id is not None
and service_factory is SimulationService
and service.simulator_kind not in {"libero", "vlabench"}
):
raise ValueError("render_gpu_id currently supports the LIBERO and VLABench MuJoCo adapters")
self.service, self.config, self._service_factory = service, config, service_factory

def run(self, requests: Sequence[EpisodeRequest]) -> dict[str, Any]:
"""Run each request once; reject invalid identities before starting workers."""
requests = list(requests)
if any(request.runtime_id != self.service.runtime_id for request in requests):
raise ValueError("all requests must target the configured runtime")
if len({request.request_id for request in requests}) != len(requests):
raise ValueError("episode request IDs must be unique")
endpoints = self.config.inference_endpoints or (self.service.inference_endpoint,)
count = min(self.config.workers, len(requests))
ctx = mp.get_context("spawn")
processes = []
readers: dict[Connection, int] = {}
activity: dict[int, float] = {}
results: dict[int, dict[str, Any]] = {}
completed = False
started = time.monotonic()
try:
for worker_id in range(count):
config = dataclasses.replace(
self.service,
simulator_id=f"{self.service.simulator_id}-{worker_id}",
inference_endpoint=endpoints[worker_id % len(endpoints)],
)
jobs = [(index, requests[index]) for index in range(worker_id, len(requests), count)]
reader, writer = ctx.Pipe(duplex=False)
process = ctx.Process(
target=_worker,
args=(config, self.config.render_gpu_id, jobs, writer, self._service_factory),
name=f"embodirun-eval-{worker_id}",
)
try:
process.start()
except BaseException:
reader.close()
raise
finally:
writer.close()
processes.append(process)
readers[reader] = worker_id
activity[worker_id] = time.monotonic()
while readers:
oldest = min(readers.values(), key=activity.__getitem__)
remaining = self.config.worker_timeout_s - (time.monotonic() - activity[oldest])
if remaining <= 0:
raise TimeoutError(f"evaluation worker {oldest} exceeded worker_timeout_s")
for reader in wait(list(readers), timeout=remaining):
worker_id = readers[reader]
try:
message = reader.recv()
except EOFError:
raise RuntimeError(f"evaluation worker {worker_id} exited without completion") from None
activity[worker_id] = time.monotonic()
if message[0] == "error":
raise RuntimeError(f"evaluation worker {worker_id} failed:\n{message[1]}")
if message[0] == "done":
reader.close()
del readers[reader]
elif message[0] == "result":
_, index, payload, duration = message
if index in results or index not in range(worker_id, len(requests), count):
raise RuntimeError("unexpected or duplicate episode completion")
result = EpisodeResult.from_payload(payload)
if (result.request_id, result.runtime_id) != (
requests[index].request_id,
self.service.runtime_id,
):
raise RuntimeError("episode result identity does not match request")
results[index] = {
**payload,
"worker": worker_id,
"inference_endpoint": endpoints[worker_id % len(endpoints)],
"wall_s": duration,
}
else:
raise RuntimeError("unexpected evaluation worker message")
if len(results) != len(requests):
raise RuntimeError("evaluation completed with missing results")
completed = True
finally:
for process in processes:
if not completed and process.is_alive():
process.terminate()
for process in processes:
process.join(timeout=2)
if process.is_alive():
process.kill()
process.join(timeout=2)
for reader in readers:
reader.close()
if any(process.exitcode != 0 for process in processes):
raise RuntimeError("evaluation worker failed during cleanup")
return {
"results": [results[index] for index in range(len(requests))],
"wall_s": time.monotonic() - started,
"evaluation": dataclasses.asdict(self.config),
}
31 changes: 31 additions & 0 deletions src/embodirun/services/simulation/evaluate.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
"""Run a finite evaluation manifest against existing HTTP inference replicas."""

from __future__ import annotations

import argparse
import json
from pathlib import Path

from embodirun.application.simulation.contracts import EpisodeRequest, SimulationServiceConfig
from embodirun.application.simulation.evaluation import EvaluationConfig, EvaluationRunner


def main(argv: list[str] | None = None) -> int:
"""Parse jobs, run process-isolated simulators and save the complete report."""
parser = argparse.ArgumentParser(prog="embodirun-evaluate")
parser.add_argument("--config", required=True, type=Path)
parser.add_argument("--output", required=True, type=Path)
args = parser.parse_args(argv)
manifest = json.loads(args.config.read_text())
service = SimulationServiceConfig.from_json(json.dumps(manifest["simulation"]))
config = EvaluationConfig(**manifest.get("evaluation", {}))
requests = [EpisodeRequest.from_payload(item) for item in manifest["episodes"]]
report = EvaluationRunner(service, config).run(requests)
report["manifest"] = manifest
args.output.parent.mkdir(parents=True, exist_ok=True)
args.output.write_text(json.dumps(report, indent=2, allow_nan=False) + "\n")
return 0


if __name__ == "__main__":
raise SystemExit(main())
Loading