Ronaldo-GOAT's picture
eval_stack: pi0.5 negmesh8 160-episode replay eval harness, job scripts, protocol docs
51c2c72 verified
Raw History Blame Contribute Delete
24.6 kB
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: <replay_root>/episodes/<object>/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)