import collections import dataclasses import io import logging import math import pathlib import subprocess import imageio from datetime import datetime import numpy as np import random from openpi_client import image_tools from openpi_client import websocket_client_policy as _websocket_client_policy import tqdm import tyro import json import os import robocasa.utils.robomimic.robomimic_dataset_utils as FileUtils import robocasa.utils.robomimic.robomimic_env_utils as EnvUtils import robocasa.utils.robomimic.robomimic_obs_utils as ObsUtils import robocasa from robocasa.utils.dataset_registry import TASK_SET_REGISTRY from robocasa.utils.dataset_registry_utils import get_task_horizon import gymnasium as gym from robocasa.utils.env_utils import convert_action # noqa: E402 import gzip as _gzip import pickle as _pickle # ============================================================================= # EXACT-REPLAY SUPPORT (added 2026-09-22) # # The negmesh8 "hard set" is 8 objects x 20 episodes whose initial conditions were # snapshotted (MuJoCo XML + flattened sim state + robocasa ep_meta) rather than # re-derived from a seed. Replaying them reproduces the exact scenes the base 60k # policy was screened on, which a seeded env.reset() cannot do. # # The restore sequence below mirrors robocasa's own reference implementation, # robocasa/scripts/dataset_scripts/playback_dataset.py::reset_to() -- set ep_meta, # hard reset, edit + reload the XML, soft reset, restore flattened state, forward, # update state. Deviating from that order does not work: reset_from_xml_string only # soft-resets, so the ep_meta-driven env.reset() before it is required. # ============================================================================= def load_replay_bank(replay_root: str) -> list: """Return the bank's episode dirs in a stable, reproducible order. Layout: /episodes//episode_NNN/. Sorted by (object, episode index) so a shard slice always names the same episodes. """ root = pathlib.Path(replay_root) base = root / "episodes" if (root / "episodes").is_dir() else root eps = [] for obj_dir in sorted(p for p in base.iterdir() if p.is_dir()): for ep_dir in sorted(obj_dir.glob("episode_*")): if (ep_dir / "initial_model.xml.gz").exists(): eps.append(ep_dir) if not eps: raise FileNotFoundError(f"no replay episodes found under {base}") return eps def reset_from_replay(env, ep_dir: pathlib.Path, seed=None): """Restore one snapshotted episode. Returns (obs, info) like env.reset().""" # gym.make() wraps RoboCasaGymEnv; .unwrapped is that wrapper, .env the kitchen env. gym_env = env.unwrapped base = gym_env.env with open(ep_dir / "episode_metadata.pkl", "rb") as f: ep_meta = _pickle.load(f) if hasattr(base, "set_attrs_from_ep_meta"): base.set_attrs_from_ep_meta(ep_meta) elif hasattr(base, "set_ep_meta"): base.set_ep_meta(ep_meta) else: raise RuntimeError("env exposes neither set_attrs_from_ep_meta nor set_ep_meta") # Reset through the WRAPPER, not base.reset(). # # Two things have to happen here and only this call does both: # 1. the hard reset that reset_from_xml_string needs (it only soft-resets and will # not reload the model), and # 2. satisfying gymnasium's OrderEnforcing wrapper, which gym.make() puts around the # env and which raises # gymnasium.error.ResetNeeded: Cannot call env.step() before calling env.reset() # until it has seen a reset() of its own. base.reset() performs the physics reset # but leaves OrderEnforcing._has_reset False, so the FIRST env.step() of the # episode dies. (This is exactly how jobs 20087/20088 failed: the smoke test missed # it because it called env.reset() before exercising the replay path, which had # already flipped that flag.) # The scene this builds is thrown away by reset_from_xml_string below; seed is passed # only to keep env.rng deterministic. env.reset(seed=seed) with _gzip.open(ep_dir / "initial_model.xml.gz", "rt") as f: xml = f.read() xml = base.edit_model_xml(xml) # repoint recorded asset paths at the local robocasa base.reset_from_xml_string(xml) base.sim.reset() with np.load(ep_dir / "initial_sim_state.npz") as st: flat = np.asarray(st["states"]).ravel() base.sim.set_state_from_flattened(flat) base.sim.forward() if hasattr(base, "update_sites"): base.update_sites() if hasattr(base, "update_state"): base.update_state() raw_obs = ( base.viewer._get_observations(force_update=True) if base.viewer_get_obs else base._get_observations(force_update=True) ) obs = gym_env.get_observation(raw_obs) return obs, {"success": False} def gpu_model() -> str: """The GPU this eval actually ran on, recorded INTO stats.json. Measured 2026-09-22: eval results are not comparable across GPU architectures. The same checkpoint, seeds, config and code path gave 5/8 on A100 (det_A, det_B -- two different physical A100s, agreeing on every episode) and 3/8 on L40S (det_C), with EVERY cross-architecture episode pair diverging at step 0 / query 0 -- i.e. in the very first rendered frame, before the policy acts. EGL rasterisation differs between architectures and the closed loop amplifies it. So the GPU is a controlled variable, and it has to live in the artifact rather than only in a job log that gets rotated away -- otherwise two runs on different hardware look comparable and silently are not. """ try: out = subprocess.run(["nvidia-smi", "--query-gpu=name", "--format=csv,noheader"], capture_output=True, text=True, timeout=30, check=False) names = sorted({ln.strip() for ln in out.stdout.splitlines() if ln.strip()}) return ",".join(names) if names else "unknown" except (OSError, subprocess.SubprocessError): return "unknown" def mp4_roundtrip_image(img_uint8: np.ndarray, fps: int = 20, crf: int = 23, n_frames: int = 5, take_idx: int = 2) -> np.ndarray: """Encode a single image into an h264/yuv420p mp4 in memory and decode it back. This mimics the lerobot training pipeline so eval-time pixels match the compression artifacts the policy was trained on. """ buf = io.BytesIO() writer = imageio.get_writer( buf, format="mp4", fps=fps, codec="h264", pixelformat="yuv420p", output_params=["-crf", str(crf)], ) for _ in range(n_frames): writer.append_data(img_uint8.astype(np.uint8)) writer.close() buf.seek(0) reader = imageio.get_reader(buf.getvalue(), format="mp4") return np.asarray(reader.get_data(take_idx)).astype(np.uint8) @dataclasses.dataclass class Args: ################################################################################################################# # Model server parameters ################################################################################################################# host: str = "0.0.0.0" port: int = 8000 resize_size: int = 224 replan_steps: int = 5 split: str = "pretrain" num_trials: int = 50 # Number of rollouts per task task_set: list = None # REPRO FIX (optional): evaluate exactly these env names instead of expanding task_set. env_names: list = None # If true, encode/decode each eval image through h264/yuv420p mp4 to match # the compression artifacts present in the lerobot training pipeline. This # matches what the policy saw during training and is required for non-zero # eval success on RoboCasa365. mp4_roundtrip: bool = True # Skip the first N episodes (calls env.reset() + sleeps) to advance the # robocasa RNG without running the policy. Combined with the same seed, this # lets multiple shards reproduce a different slice of the same 50-episode # sequence (shard 1 starts at episode 0, shard 2 at episode 25, etc). start_episode_idx: int = 0 # Optional suffix appended to the log dir, useful to avoid colliding with # an in-progress single-job run when sharding (e.g. "shard_21_30"). log_dir_suffix: str = "" ################################################################################################################# # Utils ################################################################################################################# log_dir: str = None seed: int = 7 # Random Seed (for reproducibility) # REPRO FIX: per-episode seeds (see EVAL_SEEDS.md). Every episode is seeded from its own # index, so results no longer depend on how many episodes/resets ran before it. env_seed_offset: int = 0 # env seed for episode i = env_seed_offset + i action_seed: int = 42 # flow-matching noise seed base (per episode, per query) # ---- exact-replay mode (negmesh8 hard set) -------------------------------- # When set, episodes come from a snapshot bank instead of seeded env.reset(). # Unset => behaviour is byte-identical to before this flag existed. replay_root: str = "" # "start:count" slice of the (sorted) bank, for sharding. Empty = the whole bank. replay_episodes: str = "" # Flow-matching noise convention: # "per_query_epid" rng = default_rng(action_seed + ep_id*1_000_003 + query_idx) # -- this repo's scheme; noise differs per episode AND per query. # "per_episode_reset" rng = default_rng(action_seed) once per episode, one draw per # query -- upstream's scheme in examples/robocasa/ # eval_single_task_seeded.py ("NumPy PCG64 is reset to # action_seed at the start of every episode"). This is what the # negmesh8 replay bank's recorded outcomes were produced under, # so use it when comparing against those. noise_scheme: str = "per_query_epid" # PERF (2026-09-22): main.py ran the mp4 round-trip on EVERY step, but the round-tripped # image is only ever read inside the `if not action_plan:` replan branch -- on the other # (replan_steps-1) steps out of every replan_steps the result was computed and thrown # away. Each call is an imageio/ffmpeg SUBPROCESS at ~479 ms measured on this cluster, and # there are two per step, so this dominated everything: episode 0 of job 20077 spent # ~268 s on 277 steps == 0.97 s/step == exactly 2 x 479 ms. At that rate a 750-step # episode costs ~12 min and a 50-episode shard ~8 h, well past the 4 h eval QoS cap. # # Doing the round-trip only on replan steps feeds the policy the SAME pixels -- the # skipped results were unused -- while cutting the calls by replan_steps (5x here). # This is NOT the same as main_optimized.py, which additionally swaps imageio for PyAV; # that DOES change the pixels (measured max|d|=49, 85% of pixels differ, because PyAV # encodes one I-frame while this path encodes 5 frames and takes the P-frame at index 2), # so the encoder is deliberately left alone here. # # Set mp4_eager=True to restore the old every-step behaviour (slow; kept so the # equivalence can be re-checked). mp4_eager: bool = False def eval_main(args: Args) -> None: # Set random seed np.random.seed(args.seed) split = args.split log_dir = args.log_dir num_trials = args.num_trials resize_size = args.resize_size replan_steps = args.replan_steps host = args.host port = args.port all_env_names = [] if args.env_names: # REPRO FIX: explicit env list overrides task_set expansion all_env_names = list(args.env_names) else: for task in args.task_set: env_names = TASK_SET_REGISTRY[task] all_env_names.extend(env_names) for env_name in all_env_names: # try: eval_env(env_name, split, log_dir, num_trials, resize_size, replan_steps, host, port, args.seed, mp4_roundtrip=args.mp4_roundtrip, start_episode_idx=args.start_episode_idx, log_dir_suffix=args.log_dir_suffix, env_seed_offset=args.env_seed_offset, action_seed=args.action_seed, replay_root=args.replay_root, replay_episodes=args.replay_episodes, noise_scheme=args.noise_scheme, mp4_eager=args.mp4_eager) # except Exception as e: # print("Exception!") # print(e) def eval_env(env_name, split, log_dir, num_trials, resize_size, replan_steps, host, port, seed, mp4_roundtrip: bool = False, start_episode_idx: int = 0, log_dir_suffix: str = "", env_seed_offset: int = 0, action_seed: int = 42, replay_root: str = "", replay_episodes: str = "", noise_scheme: str = "per_query_epid", mp4_eager: bool = False): # set args based on task assert split in ["pretrain", "target"] task_horizon = get_task_horizon(env_name) # set dataset path and horizon horizon = int(task_horizon * 1.5) # the policy moves slow so give the policy extra time now = datetime.now() now_formatted = now.strftime("%Y-%m-%d-%H-%M") suffix = f"_{log_dir_suffix}" if log_dir_suffix else "" log_path = f"{log_dir}/evals/{split}/{env_name}/{now_formatted}{suffix}" # Only skip if THIS exact dir (with suffix) already has stats — different # shards must not mutually skip. if log_dir_suffix: if os.path.exists(os.path.join(log_path, "stats.json")): print(f"{env_name}/{split}/{log_dir_suffix}, stats exists, skipping.") return else: for root, dirs, files in os.walk(os.path.dirname(log_path)): if "stats.json" in files: print(f"{env_name}/{split}, stats path exists, skipping.") return pathlib.Path(log_path).mkdir(parents=True, exist_ok=True) client = _websocket_client_policy.WebsocketClientPolicy(host, port) # REPRO FIX: size the per-query noise from the server's advertised action shape _md = client.get_server_metadata() or {} noise_shape = (int(_md.get('action_horizon') or 50), int(_md.get('action_dim') or 32)) logging.info(f"[repro] per-query noise shape from server metadata: {noise_shape}") # Start evaluation total_episodes, total_successes = 0, 0 # Get task env = gym.make(f"robocasa/{env_name}", split=split, seed=seed) # REPLAY: resolve the bank slice up front so num_trials and the manifest agree. replay_eps = None if replay_root: replay_eps = load_replay_bank(replay_root) if replay_episodes: _start, _count = (int(x) for x in replay_episodes.split(":")) replay_eps = replay_eps[_start:_start + _count] num_trials = len(replay_eps) logging.info(f"[replay] {num_trials} episode(s) from {replay_root}" + (f" slice {replay_episodes}" if replay_episodes else " (full bank)")) manifest = [ {"episode_idx": i, "object": d.parent.name, "bank_episode": d.name, "path": str(d)} for i, d in enumerate(replay_eps) ] with open(os.path.join(log_path, "replay_manifest.json"), "w") as f: json.dump(manifest, f, indent=2) if start_episode_idx > 0: # REPRO FIX: warm-up resets are no longer needed -- episodes are seeded from their own # index (env_seed_offset + episode_idx), so a shard can start anywhere. Kept as a no-op # for CLI compatibility; use --env_seed_offset to select the episode range instead. logging.info(f"[shard] start_episode_idx={start_episode_idx}: warm-up resets skipped (per-episode seeding)") # Start episodes task_episodes, task_successes = 0, 0 per_episode_log = [] for episode_idx in tqdm.tqdm(range(num_trials)): # Reset environment ep_id = env_seed_offset + episode_idx # REPRO FIX: global episode id; every RNG below keys on it np.random.seed(ep_id); random.seed(ep_id) # REPRO FIX: global RNGs (robosuite/robocasa use np.random in places) if replay_eps is None: obs, info = env.reset(seed=ep_id) # REPRO FIX: per-episode env seed else: # REPLAY: the scene comes from the snapshot, not from a seed. ep_id still keys # the action noise (under "per_query_epid"), so replays stay reproducible. _ep_dir = replay_eps[episode_idx] logging.info(f"[replay] episode {episode_idx}: {_ep_dir.parent.name}/{_ep_dir.name}") obs, info = reset_from_replay(env, _ep_dir, seed=ep_id) query_idx = 0 # REPRO FIX: policy-query counter within the episode # Upstream's convention: one PCG64 reset to action_seed per episode, one draw per # query. Used for the negmesh8 bank comparison; default keeps this repo's scheme. ep_noise_rng = np.random.default_rng(action_seed) if noise_scheme == "per_episode_reset" else None ep_actions = [] # REPRO FIX: executed actions, saved per episode for exact-repro checks task_lang = obs["annotation.human.task_description"] action_plan = collections.deque() # Setup t = 0 replay_images = [] logging.info(f"Starting episode {task_episodes+1}...") while t < horizon: # Pass the raw 256x256 uint8 frames straight to the policy. At training time, # GrootOpenpiSingleDataset returns 256x256 uint8 images and the only resize is # applied later by `ModelTransformFactory.ResizeImages(224, 224)` inside the # policy itself. Doing an extra resize_with_pad here was double-resizing with a # different interpolation kernel and shifting pixel statistics enough to push # the policy out of distribution. img = np.ascontiguousarray(obs["video.robot0_agentview_left"]) wrist_img = np.ascontiguousarray(obs["video.robot0_eye_in_hand"]) img = image_tools.convert_to_uint8(img) wrist_img = image_tools.convert_to_uint8(wrist_img) if mp4_roundtrip and mp4_eager: # Legacy path: round-trip every step (see Args.mp4_eager). img = mp4_roundtrip_image(img) wrist_img = mp4_roundtrip_image(wrist_img) # Save preprocessed image for replay video # replay_images.append(img) if not action_plan: if mp4_roundtrip and not mp4_eager: # Match lerobot's mp4-decoded pixel statistics so the policy sees the # same compression artifacts as during training. Only needed here -- # this is the only place img/wrist_img reach the model. img = mp4_roundtrip_image(img) wrist_img = mp4_roundtrip_image(wrist_img) # Match the training-time concat order in # `openpi.groot_utils.groot_openpi_dataset.GrootOpenpiSingleDataset.__getitem__`: # eef_pos_rel, eef_rot_rel, base_pos, base_rot, gripper_qpos. state = np.concatenate( ( obs["state.end_effector_position_relative"], obs["state.end_effector_rotation_relative"], obs["state.base_position"], obs["state.base_rotation"], obs["state.gripper_qpos"], ), axis=0 ) # state = np.ascontiguousarray(state) # Finished executing previous action chunk -- compute new chunk # Prepare observations dict element = { "observation/image": img, "observation/wrist_image": wrist_img, "observation/state": state, "prompt": task_lang, } # Query model to get action # REPRO FIX: client-supplied noise, seeded per (episode, query). Same convention # as the GR00T eval: base + episode*1_000_003 + query. Removes the dependence on # the policy server's cumulative RNG (shared across all clients). if ep_noise_rng is not None: element["noise"] = ep_noise_rng.standard_normal(noise_shape).astype(np.float32) else: _nrng = np.random.default_rng(action_seed + ep_id * 1_000_003 + query_idx) element["noise"] = _nrng.standard_normal(noise_shape).astype(np.float32) query_idx += 1 action_chunk = client.infer(element)["actions"] assert ( len(action_chunk) >= replan_steps ), f"We want to replan every {replan_steps} steps, but policy only predicts {len(action_chunk)} steps." action_plan.extend(action_chunk[: replan_steps]) action = action_plan.popleft() ep_actions.append(np.asarray(action, dtype=np.float32)) # raw policy output, before convert_action action = convert_action(action) # Execute action in environment obs, reward, done, truncated, info = env.step(action) done = info["success"] # for robocasa, usuccess entry in info replay_img = env.render() replay_img = np.ascontiguousarray(replay_img) replay_img = image_tools.convert_to_uint8( replay_img ) if t % 2 == 0 or t == horizon - 1 or done: replay_images.append(replay_img) if done: task_successes += 1 total_successes += 1 break t += 1 task_episodes += 1 total_episodes += 1 per_episode_log.append({ "episode_idx": episode_idx, "ep_id": ep_id, "success": bool(done), "steps": len(ep_actions), "queries": query_idx, **({"object": replay_eps[episode_idx].parent.name, "bank_episode": replay_eps[episode_idx].name} if replay_eps is not None else {}), }) # Save a replay video of the episode suffix = "success" if done else "failure" np.save(pathlib.Path(log_path) / f"actions_{episode_idx}_{suffix}.npy", np.stack(ep_actions) if ep_actions else np.zeros((0,))) imageio.mimwrite( pathlib.Path(log_path) / f"rollout_{episode_idx}_{suffix}.mp4", [np.asarray(x) for x in replay_images], fps=20, ) # Log current results logging.info(f"Success: {done}") logging.info(f"# episodes completed so far: {total_episodes}") logging.info(f"# successes: {total_successes} ({total_successes / total_episodes * 100:.1f}%)") # Log final results logging.info(f"Current task success rate: {float(task_successes) / float(task_episodes)}") logging.info(f"Current total success rate: {float(total_successes) / float(total_episodes)}") logging.info(f"[{env_name}] Total success rate: {float(total_successes) / float(total_episodes)}") logging.info(f"[{env_name}] Total episodes: {total_episodes}") print() with open(os.path.join(log_path, "stats.json"), "w") as f: stats = { "num_episodes": total_episodes, "num_successes": total_successes, "success_rate": float(total_successes) / float(total_episodes), "env_name": env_name, "gpu_model": gpu_model(), "split": split, "env_seed_offset": env_seed_offset, "action_seed": action_seed, "noise_scheme": noise_scheme, "replan_steps": replan_steps, "mp4_roundtrip": bool(mp4_roundtrip), "replay_root": replay_root or None, "replay_episodes": replay_episodes or None, "per_episode": per_episode_log, } json.dump(stats, f, indent=4) # close and delete the env env.env.close() del env.env del env if __name__ == "__main__": logging.basicConfig(level=logging.INFO) tyro.cli(eval_main)