From 47d67b8e7a8adeca47ec0f86a4b9d00c8c7c99f2 Mon Sep 17 00:00:00 2001 From: yufoo1 Date: Wed, 7 Oct 2026 12:33:22 +0800 Subject: [PATCH] feat: support shared-GPU simulation evaluation Pin MuJoCo rendering independently of HTTP inference replicas, allowing GPU0 to render while already-running services on GPU0 and GPU1 both serve inference. Statically shard episodes across process-isolated workers and keep each worker's sessions on one replica. Reuse the existing SimulationService lifecycle and original inference batcher. Add an evaluation CLI, paired placement configurations, result accounting, worker timeout and failure handling. Leave EmbodiInfer unchanged. This integrates the deployment used by the earlier external-driver 1.762x warmed-loop result (1.485x including reset). This native runner has not been GPU-benchmarked. Validation before removing the added test file from this commit: 1044 tests passed, 45 optional-dependency tests skipped; Ruff checks passed. Runtime code is unchanged by the documentation/test removal. --- .../simulation/offline-eval-dedicated.json | 44 ++++ configs/simulation/offline-eval-shared.json | 45 ++++ pyproject.toml | 1 + .../application/simulation/evaluation.py | 221 ++++++++++++++++++ src/embodirun/services/simulation/evaluate.py | 31 +++ 5 files changed, 342 insertions(+) create mode 100644 configs/simulation/offline-eval-dedicated.json create mode 100644 configs/simulation/offline-eval-shared.json create mode 100644 src/embodirun/application/simulation/evaluation.py create mode 100644 src/embodirun/services/simulation/evaluate.py diff --git a/configs/simulation/offline-eval-dedicated.json b/configs/simulation/offline-eval-dedicated.json new file mode 100644 index 0000000..fff9969 --- /dev/null +++ b/configs/simulation/offline-eval-dedicated.json @@ -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} + ] +} diff --git a/configs/simulation/offline-eval-shared.json b/configs/simulation/offline-eval-shared.json new file mode 100644 index 0000000..0e572bb --- /dev/null +++ b/configs/simulation/offline-eval-shared.json @@ -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} + ] +} diff --git a/pyproject.toml b/pyproject.toml index 869ddff..a9481a5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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. diff --git a/src/embodirun/application/simulation/evaluation.py b/src/embodirun/application/simulation/evaluation.py new file mode 100644 index 0000000..3d3599a --- /dev/null +++ b/src/embodirun/application/simulation/evaluation.py @@ -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), + } diff --git a/src/embodirun/services/simulation/evaluate.py b/src/embodirun/services/simulation/evaluate.py new file mode 100644 index 0000000..3368161 --- /dev/null +++ b/src/embodirun/services/simulation/evaluate.py @@ -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())